diff --git a/frontend/src/lib/components/JobLoader.svelte b/frontend/src/lib/components/JobLoader.svelte index 78429177a7..7913fd4be7 100644 --- a/frontend/src/lib/components/JobLoader.svelte +++ b/frontend/src/lib/components/JobLoader.svelte @@ -15,6 +15,7 @@ type OpenFlow } from '$lib/gen' import { workspaceStore } from '$lib/stores' + import { WM_LOGS_SKIPPED } from '$lib/consts' import { getContext, onDestroy, tick, untrack } from 'svelte' import type { SupportedLanguage } from '$lib/common' import { sendUserToast } from '$lib/toast' @@ -129,6 +130,37 @@ } }) + function isSkippedLogsValue(logs: string | undefined): boolean { + return logs === WM_LOGS_SKIPPED + } + + function getResolvedLogs(logs: string | undefined): string { + return isSkippedLogsValue(logs) ? '' : (logs ?? '') + } + + function mergeLogs(existingLogs: string | undefined, newLogs: string | undefined): string { + const existing = getResolvedLogs(existingLogs) + const incoming = newLogs ?? '' + return existing.length === 0 ? incoming : existing.concat(incoming) + } + + function pickMoreCompleteLogs( + primaryLogs: string | undefined, + fallbackLogs: string | undefined + ): string { + const primary = getResolvedLogs(primaryLogs) + const fallback = getResolvedLogs(fallbackLogs) + // When neither side has real logs but one was the skipped sentinel, keep + // the sentinel so downstream consumers can still lazily resolve logs + // instead of treating the job as having genuinely produced none. + if (primary.length === 0 && fallback.length === 0) { + return isSkippedLogsValue(primaryLogs) || isSkippedLogsValue(fallbackLogs) + ? WM_LOGS_SKIPPED + : '' + } + return primary.length >= fallback.length ? primary : fallback + } + function clearCurrentId() { if (currentId) { if (allowConcurentRequests) { @@ -251,7 +283,8 @@ function refreshLogOffset() { if (logOffset == 0) { - logOffset = job?.logs?.length ? job.logs?.length + 1 : 0 + const currentLogs = getResolvedLogs(job?.logs) + logOffset = currentLogs.length ? currentLogs.length + 1 : 0 } } export async function getLogs() { @@ -264,7 +297,7 @@ logOffset: logOffset }) - if ((job?.logs ?? '').length == 0) { + if (getResolvedLogs(job?.logs).length == 0) { job.logs = getUpdate.new_logs ?? '' logOffset = getUpdate.log_offset ?? 0 } @@ -385,11 +418,9 @@ if (event.data.completed) { const njob = (event.data as any).job as Job & { result_stream?: string } if (njob) { - // Use whichever logs are more complete (longer) - const streamedLogs = job?.logs ?? '' - const completedLogs = njob.logs ?? '' - njob.logs = - streamedLogs.length >= completedLogs.length ? streamedLogs : completedLogs + // Use whichever logs are more complete (longer), but never + // let the WM_LOGS_SKIPPED sentinel win over real logs. + njob.logs = pickMoreCompleteLogs(job?.logs, njob.logs) const streamedResult = job?.result_stream ?? '' const completedResult = njob.result_stream ?? '' njob.result_stream = @@ -451,11 +482,10 @@ } if (previewJobUpdates.new_logs) { - if (logOffset == 0) { - job.logs = previewJobUpdates.new_logs ?? '' - } else { - job.logs = (job?.logs ?? '').concat(previewJobUpdates.new_logs) - } + job.logs = + logOffset == 0 + ? (previewJobUpdates.new_logs ?? '') + : mergeLogs(job?.logs, previewJobUpdates.new_logs) } if (previewJobUpdates.new_result_stream) { @@ -500,6 +530,17 @@ callbacks?.change?.(job) } } + // When a job is fetched with no_logs=true the server omits logs entirely. + // Flag it with a sentinel so consumers (the log panel) can tell "logs were + // intentionally skipped" apart from "job genuinely produced no logs", and + // lazily resolve the real logs on demand. + function flagSkippedLogs(j: T, effectiveNoLogs: boolean): T { + if (effectiveNoLogs && !(j as Job & { logs?: string }).logs) { + ;(j as Job & { logs?: string }).logs = WM_LOGS_SKIPPED + } + return j + } + async function loadTestJob(id: string, callbacks?: Callbacks): Promise { let isCompleted = false if (isCurrentJob(id)) { @@ -519,23 +560,29 @@ }) if ((previewJobUpdates.running ?? false) || (previewJobUpdates.completed ?? false)) { - job = await JobService.getJob({ - workspace: workspace!, - id, - noCode, - noLogs: onlyResult || noLogs - }) + job = flagSkippedLogs( + await JobService.getJob({ + workspace: workspace!, + id, + noCode, + noLogs: onlyResult || noLogs + }), + onlyResult || noLogs + ) callbacks?.change?.(job) } updateJobFromProgress(previewJobUpdates, job, callbacks) } else { - job = await JobService.getJob({ - workspace: workspace!, - id, - noLogs: onlyResult || noLogs, - noCode - }) + job = flagSkippedLogs( + await JobService.getJob({ + workspace: workspace!, + id, + noLogs: onlyResult || noLogs, + noCode + }), + onlyResult || noLogs + ) } jobUpdateLastFetch = new Date() @@ -628,12 +675,15 @@ try { // First load the job to get initial state if ((!job || job.id == '') && !onlyResult) { - job = await JobService.getJob({ - workspace: workspace!, - id, - noLogs: noLogs, - noCode - }) + job = flagSkippedLogs( + await JobService.getJob({ + workspace: workspace!, + id, + noLogs: noLogs, + noCode + }), + noLogs + ) callbacks?.change?.(job) getActiveRecording()?.recordInitialJob(id, job) @@ -775,7 +825,7 @@ clearCurrentId() } else { const njob = previewJobUpdates.job as Job & { result_stream?: string } - njob.logs = job?.logs ?? '' + njob.logs = pickMoreCompleteLogs(job?.logs, njob.logs) njob.result_stream = job?.result_stream ?? '' job = njob onJobCompleted(id, job, callbacks) diff --git a/frontend/src/lib/components/LogViewer.svelte b/frontend/src/lib/components/LogViewer.svelte index 39089dec8b..ca90e3968f 100644 --- a/frontend/src/lib/components/LogViewer.svelte +++ b/frontend/src/lib/components/LogViewer.svelte @@ -21,6 +21,7 @@ import { AnsiUp } from 'ansi_up' import NoWorkerWithTagWarning from './runs/NoWorkerWithTagWarning.svelte' import { JobService } from '$lib/gen' + import { WM_LOGS_SKIPPED } from '$lib/consts' import Tooltip from './Tooltip.svelte' import { twMerge } from 'tailwind-merge' import QueuePosition from './QueuePosition.svelte' @@ -42,6 +43,10 @@ tagLabel?: string noPadding?: boolean navigationId?: string + /** Called once after we resolve a WM_LOGS_SKIPPED sentinel by fetching the + * full job. Use this to write the real logs back into the source of truth + * (e.g. flowStateStore) so subsequent mounts don't refetch. */ + onLogsResolved?: (logs: string) => void } let { @@ -60,7 +65,8 @@ customEmptyMessage = 'No logs are available yet', tagLabel = undefined, noPadding = false, - navigationId = undefined + navigationId = undefined, + onLogsResolved }: Props = $props() // @ts-ignore @@ -80,6 +86,34 @@ let loadedFromObjectStore = $state('') + // `content` is the WM_LOGS_SKIPPED sentinel when the job was fetched with + // no_logs=true. If an older in-memory value accidentally has real bytes + // concatenated after the sentinel, treat that as skipped too and refetch. + let isLogsSkipped = $derived((content ?? '').startsWith(WM_LOGS_SKIPPED)) + let resolvedSkippedLogs: string | undefined = $state(undefined) + let fetchedSkippedJobId: string | undefined = $state(undefined) + let effectiveContent = $derived(isLogsSkipped ? resolvedSkippedLogs : content) + let resolvingSkippedLogs = $derived(isLogsSkipped && !!jobId && resolvedSkippedLogs === undefined) + + $effect(() => { + if (!isLogsSkipped || !jobId || fetchedSkippedJobId === jobId) { + return + } + const id = jobId + fetchedSkippedJobId = id + untrack(() => { + JobService.getJob({ workspace: $workspaceStore ?? '', id }) + .then((j) => { + if (fetchedSkippedJobId === id) { + const logs = (j as { logs?: string })['logs'] ?? '' + resolvedSkippedLogs = logs + onLogsResolved?.(logs) + } + }) + .catch((e) => console.error('Failed to resolve skipped logs', e)) + }) + }) + function findPrefixInfo( truncateContent: string ): { prefixIndex: number; position: number } | undefined { @@ -181,7 +215,7 @@ })) as string LOG_LIMIT += Math.min(LOG_INC, res.length) loadedFromObjectStore = res + loadedFromObjectStore - let newC = truncateContent(content, loadedFromObjectStore, LOG_LIMIT) + let newC = truncateContent(effectiveContent, loadedFromObjectStore, LOG_LIMIT) LOG_LIMIT -= newC.indexOf('\n') + 1 } else { console.error('No file detected to download from') @@ -191,7 +225,7 @@ function showMoreTruncate(len: number) { scroll = false LOG_LIMIT += LOG_INC - let newC = truncateContent(content, loadedFromObjectStore, LOG_LIMIT) + let newC = truncateContent(effectiveContent, loadedFromObjectStore, LOG_LIMIT) let newlineIndex = newC.indexOf('\n') + 1 if (newlineIndex < LOG_INC / 2) { LOG_LIMIT -= newlineIndex @@ -203,12 +237,16 @@ loadedFromObjectStore = '' LOG_LIMIT = LOG_INC scroll = true + resolvedSkippedLogs = undefined + fetchedSkippedJobId = undefined } }) let logsApiPath = $derived(`/w/${$workspaceStore}/jobs_u/get_logs/${jobId}`) let downloadHref = $derived(withExternalDomain(`${base}/api${logsApiPath}`)) let downloadName = $derived(`windmill_logs_${jobId}.txt`) - let truncatedContent = $derived(truncateContent(content, loadedFromObjectStore, LOG_LIMIT)) + let truncatedContent = $derived( + truncateContent(effectiveContent, loadedFromObjectStore, LOG_LIMIT) + ) let prefixInfo = $derived(findPrefixInfo(truncatedContent)) let downloadStartUrl = $derived(findStartUrl(truncatedContent, prefixInfo)) $effect.pre(() => { @@ -274,7 +312,7 @@ {/if} {#if content}{@const len = - (content?.length ?? 0) + (loadedFromObjectStore?.length ?? 0)}{#if splitHtml}{@html splitHtml.before}{#if resolvingSkippedLogs}{:else if effectiveContent}{@const len = + (effectiveContent?.length ?? 0) + + (loadedFromObjectStore?.length ?? 0)}{#if splitHtml}{@html splitHtml.before}{@html splitHtml.after}{:else if downloadStartUrl}