From 20ff26a3ba2f05045e3e60f45a9289dbeb7b1fe9 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Fri, 15 May 2026 17:39:53 +0000 Subject: [PATCH] perf: guard no-op poll re-layout; dedupe write-asset extraction - skip reactive ids/states/events reassignment when unchanged, so an idle poll tick no longer re-runs the full sugiyama layout every 3-6s - bound countedJobIds (rebuilt from eventsById in lockstep with prune) - extract shared extractWrites() helper, replacing 4 copy-pasted write-asset filter/map blocks in the pipeline page - compute activeRunnable node-id once, reuse for the active-edge set and the optimistic badge (flattened ternary); trim narrating docs Co-Authored-By: Claude Opus 4.7 (1M context) --- .../assets/AssetGraph/AssetGraphCanvas.svelte | 16 +++-- .../AssetGraph/activeRunnables.svelte.ts | 59 ++++++++++++------- frontend/src/lib/components/assets/lib.ts | 12 ++++ .../(logged)/pipeline/[folder]/+page.svelte | 46 ++++----------- 4 files changed, 73 insertions(+), 60 deletions(-) diff --git a/frontend/src/lib/components/assets/AssetGraph/AssetGraphCanvas.svelte b/frontend/src/lib/components/assets/AssetGraph/AssetGraphCanvas.svelte index 440661eccf..3eeda7dffe 100644 --- a/frontend/src/lib/components/assets/AssetGraph/AssetGraphCanvas.svelte +++ b/frontend/src/lib/components/assets/AssetGraph/AssetGraphCanvas.svelte @@ -142,6 +142,12 @@ unsaved?: boolean } + // Graph-id of the script the user just launched (zero-latency hint), + // computed once and reused by the optimistic badge + the active-edge set. + let activeRunnableNodeId = $derived( + activeRunnable ? `${activeRunnable.kind}:${activeRunnable.path}` : undefined + ) + // Lineage edges (parsed r/w usages): writer → asset, asset → reader. // Trigger edges (`// on `): asset → script or schedule → script. The // lineage subgraph is informational; the trigger subgraph is executable. @@ -243,10 +249,10 @@ // poll's known run count. Falls back to the polled state (which // carries the cascade + the final success/failure) otherwise. const polledRunState = runStates?.get(rid) - const runState = - activeRunnable && `${activeRunnable.kind}:${activeRunnable.path}` === rid - ? { status: 'running' as const, runs: polledRunState?.runs ?? 0 } - : polledRunState + let runState = polledRunState + if (activeRunnableNodeId === rid) { + runState = { status: 'running' as const, runs: polledRunState?.runs ?? 0 } + } nodes.push({ id: rid, type: 'runnable', @@ -462,7 +468,7 @@ let activeRunnableIdSet = $derived( new Set([ ...(activeRunnableIds ?? []), - ...(activeRunnable ? [`${activeRunnable.kind}:${activeRunnable.path}`] : []) + ...(activeRunnableNodeId ? [activeRunnableNodeId] : []) ]) ) let flowEdges = $derived.by(() => diff --git a/frontend/src/lib/components/assets/AssetGraph/activeRunnables.svelte.ts b/frontend/src/lib/components/assets/AssetGraph/activeRunnables.svelte.ts index 59c5445721..a085771fd9 100644 --- a/frontend/src/lib/components/assets/AssetGraph/activeRunnables.svelte.ts +++ b/frontend/src/lib/components/assets/AssetGraph/activeRunnables.svelte.ts @@ -17,22 +17,39 @@ export type PipelineEvent = { } const MAX_EVENTS = 50 +// Equality guards so an unchanged poll tick doesn't reassign the reactive +// state — a no-op `states`/`ids` reassignment would otherwise re-derive the +// whole canvas (full d3-dag Sugiyama layout + edge re-route) every 3–6s. +function setEq(a: Set, b: Set): boolean { + if (a.size !== b.size) return false + for (const v of a) if (!b.has(v)) return false + return true +} +function statesEq(a: Map, b: Map): boolean { + if (a.size !== b.size) return false + for (const [k, v] of a) { + const w = b.get(k) + if (!w || w.status !== v.status || w.runs !== v.runs) return false + } + return true +} +function eventsEq(a: PipelineEvent[], b: PipelineEvent[]): boolean { + if (a.length !== b.length) return false + for (let i = 0; i < a.length; i++) { + if (a[i].id !== b[i].id || a[i].status !== b[i].status) return false + } + return true +} + /** - * Tracks which pipeline runnables currently have an in-flight (queued/running) - * job — plus any that started AND finished since the last poll, so a fast - * cascade hop that completes between two ticks still pulses once. Returns a - * reactive `Set` of `${kind}:${path}` ids (matching the graph's - * runnable node ids) for the canvas to animate. + * Polls folder jobs and exposes, as reactive `${kind}:${path}`-keyed state: + * the in-flight set (`ids`, incl. a one-tick pulse for hops that start and + * finish between ticks), per-runnable badge state (`states`), and a capped + * activity log (`events`). * - * Cost model — deliberately cheap: - * - ZERO requests while idle. Polling only runs after `arm()` (a user- - * initiated run) and self-sustains only while jobs are in flight. - * - One folder-scoped, `perPage`-capped `listExtendedJobs` request per tick. - * - Auto-disarms after `MAX_IDLE_TICKS` consecutive empty ticks (rides the - * brief gaps between cascade hops), with a hard total-duration cap. - * - * Read-only (`listExtendedJobs` is a read endpoint, no editor dependency), so - * a future read-only / operator pipeline view can reuse this as-is. + * Zero requests while idle: polling only runs after `arm()` (a launched run) + * or while `setObserving(true)` (log open), and self-disarms when idle. + * Read-only (`listExtendedJobs`), so a future operator view can reuse it. */ export function useActiveRunnableIds( getWorkspace: () => string | undefined, @@ -174,10 +191,9 @@ export function useActiveRunnableIds( next = ids anyInFlight = ids.size > 0 } - ids = next + if (!setEq(ids, next)) ids = next // Rebuild the badge snapshot: in-flight wins (spinner), otherwise the - // last completed status. Fresh Map each tick → reactive without - // SvelteMap. completedHistory persists across `stop()`. + // last completed status. completedHistory persists across `stop()`. const snap = new Map() for (const [id, h] of completedHistory) { snap.set(id, { status: inFlightThisTick.has(id) ? 'running' : h.lastStatus, runs: h.runs }) @@ -185,14 +201,17 @@ export function useActiveRunnableIds( for (const id of inFlightThisTick) { if (!snap.has(id)) snap.set(id, { status: 'running', runs: 0 }) } - states = snap + if (!statesEq(states, snap)) states = snap // Activity-log snapshot: newest-first, capped. Prune the backing map - // so a long session doesn't grow unbounded. + // (and the dedup set in lockstep) so a long session stays bounded. const sorted = Array.from(eventsById.values()).sort((a, b) => b.at.localeCompare(a.at)) if (sorted.length > MAX_EVENTS * 4) { for (const e of sorted.slice(MAX_EVENTS * 4)) eventsById.delete(e.id) + countedJobIds.clear() + for (const id of eventsById.keys()) countedJobIds.add(id) } - events = sorted.slice(0, MAX_EVENTS) + const nextEvents = sorted.slice(0, MAX_EVENTS) + if (!eventsEq(events, nextEvents)) events = nextEvents // Look back slightly so a job finishing in the request window isn't // missed by the next tick's catch-up comparison. lastPollTs = new Date(pollStartedMs - CATCHUP_OVERLAP_MS).toISOString() diff --git a/frontend/src/lib/components/assets/lib.ts b/frontend/src/lib/components/assets/lib.ts index 4dae4c20b1..0ce80b7556 100644 --- a/frontend/src/lib/components/assets/lib.ts +++ b/frontend/src/lib/components/assets/lib.ts @@ -112,6 +112,18 @@ export function getAccessType(asset: AssetWithAltAccessType): AssetUsageAccessTy if (asset.alt_access_type) return asset.alt_access_type } +/** Write/rw assets reduced to their `{ kind, path }` identity. */ +export function extractWrites( + assets: AssetWithAltAccessType[] +): Array<{ kind: AssetKind; path: string }> { + return assets + .filter((a) => { + const at = getAccessType(a) + return at === 'w' || at === 'rw' + }) + .map((a) => ({ kind: a.kind, path: a.path })) +} + export function getFlowModuleAssets( flowModuleValue: FlowModule, additionalAssetsMap?: Record diff --git a/frontend/src/routes/(root)/(logged)/pipeline/[folder]/+page.svelte b/frontend/src/routes/(root)/(logged)/pipeline/[folder]/+page.svelte index 6ed53c4959..3feea4b9f9 100644 --- a/frontend/src/routes/(root)/(logged)/pipeline/[folder]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/pipeline/[folder]/+page.svelte @@ -9,7 +9,7 @@ import PipelineEventLog from '$lib/components/assets/AssetGraph/PipelineEventLog.svelte' import AssetGraphDetailsPane from '$lib/components/assets/AssetGraph/AssetGraphDetailsPane.svelte' import PipelinePickerModal from '$lib/components/assets/AssetGraph/PipelinePickerModal.svelte' - import type { AssetWithAltAccessType } from '$lib/components/assets/lib' + import { extractWrites, type AssetWithAltAccessType } from '$lib/components/assets/lib' import type { AssetGraphResponse, AssetGraphSelection @@ -194,12 +194,7 @@ // re-opens the draft and types something new. const liveWritesSnapshot = liveBodyAssets.scriptPath != undefined && drafts.has(liveBodyAssets.scriptPath) - ? liveBodyAssets.assets - .filter((a) => { - const at = a.access_type ?? a.alt_access_type - return at === 'w' || at === 'rw' - }) - .map((a) => ({ kind: a.kind, path: a.path })) + ? extractWrites(liveBodyAssets.assets) : undefined const liveWritesPath = liveBodyAssets.scriptPath const serialized = Array.from(drafts.entries()).map(([p, d]) => { @@ -619,12 +614,7 @@ // registers it as a dep, and writing the new Map immediately // re-fires the effect → infinite loop. if (scriptPath) { - const writes = assets - .filter((a) => { - const at = a.access_type ?? a.alt_access_type - return at === 'w' || at === 'rw' - }) - .map((a) => ({ kind: a.kind, path: a.path })) + const writes = extractWrites(assets) untrack(() => { const next = new Map(inferredWritesByPath) if (writes.length > 0) next.set(scriptPath, writes) @@ -709,10 +699,7 @@ const liveForThisDraft = liveBodyAssets.scriptPath === path const writeOuts: Array<{ kind: AssetKind; path: string }> = [] if (liveForThisDraft) { - for (const a of liveBodyAssets.assets) { - const at = a.access_type ?? a.alt_access_type - if (at === 'w' || at === 'rw') writeOuts.push({ kind: a.kind, path: a.path }) - } + writeOuts.push(...extractWrites(liveBodyAssets.assets)) } else if (d.outputAssets) { writeOuts.push(...d.outputAssets) } @@ -1083,18 +1070,12 @@ } ) - // On pipeline load, eagerly infer body assets for EVERY persisted script - // in the folder and seed the same `inferredWritesByPath` overlay the - // open-script live path uses. Without this, a script whose persisted - // asset rows are missing (e.g. object-form writeS3File the server parser - // didn't extract) has no edges until you click it — which visibly - // re-layouts the graph. Doing it up-front makes the graph complete and - // layout-stable from first paint, independent of the stale `asset` table. - // - // One-shot per (workspace, base-graph) load — deps are only those; the - // drafts/inferred maps are read via `untrack` so a keystroke or a new - // draft doesn't re-trigger the sweep. Generation token cancels a stale - // sweep if the folder/graph changes mid-flight. Pool-capped fetches. + // Eagerly infer every folder script's writes on load and seed the same + // `inferredWritesByPath` overlay the open-script path uses — otherwise a + // script whose persisted asset rows are missing only gets edges when + // clicked, which re-layouts the graph. WHY untrack: deps must stay + // (workspace, base-graph) only so a keystroke/new draft doesn't re-sweep; + // the generation token cancels an in-flight sweep on folder change. let assetPrefetchGen = 0 $effect(() => { const ws = $workspaceStore @@ -1118,12 +1099,7 @@ if (gen !== assetPrefetchGen) return const res = await inferAssets(s.language, s.content ?? '') if (gen !== assetPrefetchGen) return - const writes = ((res?.assets ?? []) as AssetWithAltAccessType[]) - .filter((a) => { - const at = a.access_type ?? a.alt_access_type - return at === 'w' || at === 'rw' - }) - .map((a) => ({ kind: a.kind, path: a.path })) + const writes = extractWrites((res?.assets ?? []) as AssetWithAltAccessType[]) if (writes.length > 0) { untrack(() => { // A live edit / prior sweep may have filled it meanwhile.