fix(flows): flag noLogs jobs and lazily resolve them in log panel (#9099)

* fix(flows): flag noLogs jobs and lazily resolve them in log panel

* fix appending to flag

* fix: preserve WM_LOGS_SKIPPED sentinel on SSE/replay completion

pickMoreCompleteLogs resolved both sentinel and undefined to '', so the
SSE completion event (whose job field is fetched .without_logs()) would
clobber the sentinel placed by flagSkippedLogs. The module log panel
then saw '' instead of the sentinel, defeating the lazy-resolve path.

Also wire onLogsResolved on the OutputPickerInner inline LogViewer so a
lazy resolve writes back to flowStateStore.previewLogs, matching
ModulePreviewResultViewer and avoiding repeated fetches on remount.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
Diego Imbert
2026-05-20 16:15:11 +00:00
committed by GitHub
co-authored by Claude Opus 4.7
parent 2db1c0a1fc
commit 740a35bf7b
5 changed files with 172 additions and 44 deletions
+81 -31
View File
@@ -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<T extends Job>(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<boolean> {
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)
+53 -10
View File
@@ -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}
<Button
on:click={() => copyToClipboard(content)}
on:click={() => copyToClipboard(effectiveContent)}
color="light"
size="xs"
startIcon={{
@@ -287,8 +325,11 @@
<div>
<pre
class="bg-surface-secondary text-primary text-xs w-full p-2 whitespace-pre-wrap border rounded-md"
>{#if content}{@const len =
(content?.length ?? 0) +
>{#if resolvingSkippedLogs}<Loader2
class="animate-spin"
size={14}
/>{:else if effectiveContent}{@const len =
(effectiveContent?.length ?? 0) +
(loadedFromObjectStore?.length ?? 0)}{#if splitHtml}{@html splitHtml.before}<button
onclick={getStoreLogs}
>Show more... <Tooltip>{tooltipText(prefixInfo)}</Tooltip></button
@@ -387,9 +428,11 @@
small ? '!text-2xs' : '!text-xs',
noPadding ? '' : 'p-2'
)}
>{#if content}{@const len =
(content?.length ?? 0) + (loadedFromObjectStore?.length ?? 0)}{#if splitHtml}<span
>{@html splitHtml.before}</span
>{#if resolvingSkippedLogs}<span class="flex p-2"
><Loader2 class="animate-spin" size={14} /></span
>{:else if effectiveContent}{@const len =
(effectiveContent?.length ?? 0) +
(loadedFromObjectStore?.length ?? 0)}{#if splitHtml}<span>{@html splitHtml.before}</span
><button onclick={getStoreLogs}
>Show more... &nbsp;<Tooltip>{tooltipText(prefixInfo)}</Tooltip></button
><span>{@html splitHtml.after}</span>{:else if downloadStartUrl}<button
@@ -44,7 +44,7 @@
tagLabel = undefined
}: Props = $props()
const { stepsInputArgs } = getContext<FlowEditorContext>('FlowEditorContext')
const { stepsInputArgs, flowStateStore } = getContext<FlowEditorContext>('FlowEditorContext')
let outputPickerInner: OutputPickerInner | undefined = $state(undefined)
export function getOutputPickerInner() {
@@ -115,6 +115,18 @@
isLoading={(testIsLoading && logJob?.['running'] == false) || loadingJob}
tag={logJob?.['tag']}
{tagLabel}
onLogsResolved={(logs) => {
if (
mod.id &&
flowStateStore?.val[mod.id] &&
flowStateStore.val[mod.id].previewJobId === logJob?.id
) {
flowStateStore.val[mod.id] = {
...flowStateStore.val[mod.id],
previewLogs: logs
}
}
}}
/>
{/if}
</Pane>
@@ -367,7 +367,9 @@
)}
/>
{:else if isLoadingAndNotMock && !mock?.enabled}
<div class="flex flex-row w-fit items-center justify-between gap-2 rounded-md bg-surface-secondary p-1 px-2 min-w-16 min-h-[23px]">
<div
class="flex flex-row w-fit items-center justify-between gap-2 rounded-md bg-surface-secondary p-1 px-2 min-w-16 min-h-[23px]"
>
<Loader2 size={12} class="animate-spin text-secondary shrink-0" />
<span class="text-xs text-secondary w-[56px]">&nbsp;</span>
</div>
@@ -527,7 +529,10 @@
size="xs2"
color="light"
variant="contained"
btnClasses={twMerge('h-[27px]', showLogs ? 'bg-blue-500/10 text-blue-800 dark:text-blue-200' : 'bg-transparent')}
btnClasses={twMerge(
'h-[27px]',
showLogs ? 'bg-blue-500/10 text-blue-800 dark:text-blue-200' : 'bg-transparent'
)}
startIcon={{ icon: ScrollText }}
on:click={() => {
showLogs = !showLogs
@@ -628,6 +633,18 @@
content={selectedJob['logs']}
isLoading={false}
tag={selectedJob['tag']}
onLogsResolved={(logs) => {
if (
moduleId &&
flowStateStore?.val[moduleId] &&
flowStateStore.val[moduleId].previewJobId === selectedJob?.id
) {
flowStateStore.val[moduleId] = {
...flowStateStore.val[moduleId],
previewLogs: logs
}
}
}}
/>
{:else if isLoadingAndNotMock}
<div class="flex flex-col items-center justify-center">
+6
View File
@@ -1,5 +1,11 @@
import type { DbType } from './components/dbTypes'
// Sentinel written into a job's `logs` field when it was fetched with
// `no_logs=true` (e.g. nested flow-status viewers during "Test flow").
// It is NOT empty logs — it means logs were intentionally skipped and can
// be lazily fetched. The log panel detects this and resolves the real logs.
export const WM_LOGS_SKIPPED = '__WM_LOGS_SKIPPED__'
export const DEFAULT_WEBHOOK_TYPE: 'async' | 'sync' = 'async'
export const HOME_SHOW_HUB = true