From af9548d6632446c6ccba59224ab2c61b4d32a3d3 Mon Sep 17 00:00:00 2001 From: Fabro Date: Mon, 4 May 2026 19:39:56 +0000 Subject: [PATCH] fabro(01KQT1TWWJYWZGDT8F05E29H9D): simplify_gpt (succeeded) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fabro-Run: 01KQT1TWWJYWZGDT8F05E29H9D Fabro-Completed: 7 Fabro-Checkpoint: 44aca9aa7d9b8658bfded82b69ba58f1e9609cac ⚒️ Generated with [Fabro](https://fabro.sh) --- apps/fabro-web/app/hooks/use-run-toasts.ts | 67 +++++++++++++++++++ apps/fabro-web/app/lib/run-events.test.tsx | 28 ++++++++ apps/fabro-web/app/lib/run-events.ts | 16 ++++- apps/fabro-web/app/lib/sse.ts | 29 ++++++-- apps/fabro-web/app/routes/run-detail.tsx | 24 ++++++- .../fabro-workflow/src/operations/start.rs | 9 ++- lib/crates/fabro-workflow/src/steering_hub.rs | 25 ++++--- 7 files changed, 178 insertions(+), 20 deletions(-) create mode 100644 apps/fabro-web/app/hooks/use-run-toasts.ts diff --git a/apps/fabro-web/app/hooks/use-run-toasts.ts b/apps/fabro-web/app/hooks/use-run-toasts.ts new file mode 100644 index 000000000..087b1fb65 --- /dev/null +++ b/apps/fabro-web/app/hooks/use-run-toasts.ts @@ -0,0 +1,67 @@ +import { useEffect, useRef } from "react"; + +import { useToast } from "../components/toast"; +import { subscribeToRunEvents, type RunEventPayload } from "../lib/run-events"; +import type { MutateFn } from "../lib/sse"; + +const NOOP_MUTATE = (() => undefined) as MutateFn; + +export function useRunToasts(runId: string | undefined) { + const { push } = useToast(); + const seenEventIdsRef = useRef(new Set()); + + useEffect(() => { + if (!runId) return; + + seenEventIdsRef.current.clear(); + return subscribeToRunEvents(runId, NOOP_MUTATE, undefined, { + onEvent: (payload) => { + const dedupeId = eventDedupeId(payload); + if (dedupeId) { + if (seenEventIdsRef.current.has(dedupeId)) return; + seenEventIdsRef.current.add(dedupeId); + } + + const message = steeringToastMessage(payload); + if (message) { + push({ message }); + } + }, + }); + }, [push, runId]); +} + +function eventDedupeId(payload: RunEventPayload): string | null { + if (typeof payload.id === "string") return payload.id; + if (typeof payload.seq === "number") return `seq:${payload.seq}`; + return null; +} + +function steeringToastMessage(payload: RunEventPayload): string | null { + const props = payload.properties ?? {}; + + switch (payload.event) { + case "agent.steering.injected": { + const kind = props.kind; + if (kind === "append") return "Steer delivered."; + if (kind === "interrupt") { + return "Agent interrupted — your message is the next turn."; + } + return null; + } + case "agent.steer.buffered": + return "Steer queued — will apply when an agent stage runs."; + case "agent.steer.dropped": { + const reason = props.reason; + if (reason === "queue_full") { + return "Steer rate limit reached; oldest queued steer dropped."; + } + if (reason === "run_ended") { + return "Run ended before queued steer(s) could apply."; + } + return null; + } + default: + return null; + } +} diff --git a/apps/fabro-web/app/lib/run-events.test.tsx b/apps/fabro-web/app/lib/run-events.test.tsx index 14acf4f23..19b0675b4 100644 --- a/apps/fabro-web/app/lib/run-events.test.tsx +++ b/apps/fabro-web/app/lib/run-events.test.tsx @@ -68,6 +68,34 @@ describe("subscribeToRunEvents", () => { expect(source.closed).toBe(true); }); + test("runs payload callbacks for later subscribers on a shared source", () => { + const source = new FakeEventSource(); + const seen: string[] = []; + const keys: string[] = []; + const mutate = (key: string) => { + keys.push(key); + return Promise.resolve(); + }; + + const firstCleanup = subscribeToRunEvents("run-shared-payload", mutate, () => source, { debounceMs: 0 }); + const secondCleanup = subscribeToRunEvents("run-shared-payload", mutate, () => { + throw new Error("source should be reused"); + }, { + debounceMs: 0, + onEvent: (payload) => { + if (payload.event) seen.push(payload.event); + }, + }); + + source.emit({ id: "evt-1", event: "agent.steer.buffered", properties: { kind: "append" } }); + + expect(seen).toEqual(["agent.steer.buffered"]); + expect(keys).toEqual([queryKeys.runs.events("run-shared-payload", 1000)]); + + firstCleanup(); + secondCleanup(); + }); + test("terminal events close the source after invalidating keys", () => { const source = new FakeEventSource(); const keys: string[] = []; diff --git a/apps/fabro-web/app/lib/run-events.ts b/apps/fabro-web/app/lib/run-events.ts index bcae6c8de..948b31acf 100644 --- a/apps/fabro-web/app/lib/run-events.ts +++ b/apps/fabro-web/app/lib/run-events.ts @@ -11,7 +11,9 @@ import { type SharedEventSubscription, } from "./sse"; -interface RunEventPayload extends EventPayload { +export interface RunEventPayload extends EventPayload { + id?: string; + seq?: number; event?: string; node_id?: string; properties?: Record; @@ -119,7 +121,13 @@ export function subscribeToRunEvents( runId: string, mutate: MutateFn, eventSourceFactory: (url: string) => EventSourceLike = createBrowserEventSource, - { debounceMs = 300 }: { debounceMs?: number } = {}, + { + debounceMs = 300, + onEvent, + }: { + debounceMs?: number; + onEvent?: (payload: RunEventPayload) => void; + } = {}, ): () => void { return subscribeToSharedEventSource({ subscriptions, @@ -129,6 +137,8 @@ export function subscribeToRunEvents( eventSourceFactory, debounceMs, resolveInvalidation: (payload) => { + onEvent?.(payload); + const event = payload.event; if (!event) return { keys: [] }; @@ -157,4 +167,4 @@ export function useRunEvents(runId: string | undefined) { if (!runId) return; return subscribeToRunEvents(runId, mutate as MutateFn); }, [mutate, runId]); -} \ No newline at end of file +} diff --git a/apps/fabro-web/app/lib/sse.ts b/apps/fabro-web/app/lib/sse.ts index 408ff3c99..29597f873 100644 --- a/apps/fabro-web/app/lib/sse.ts +++ b/apps/fabro-web/app/lib/sse.ts @@ -18,10 +18,13 @@ export interface EventInvalidation { immediate?: boolean; } +type EventResolver = (payload: EventPayload) => EventInvalidation; + export interface SharedEventSubscription { source: EventSourceLike; refcount: number; mutators: Map; + resolvers: Map; pendingKeys: Set; debounceTimer: ReturnType | null; } @@ -54,6 +57,7 @@ export function subscribeToSharedEventSource({ source, refcount: 0, mutators: new Map(), + resolvers: new Map(), pendingKeys: new Set(), debounceTimer: null, }; @@ -70,18 +74,29 @@ export function subscribeToSharedEventSource({ return; } - const invalidation = resolveInvalidation(payload); - queueInvalidations(current, invalidation.keys, { - debounceMs, - immediate: invalidation.immediate, - }); + const keys = new Set(); + let close = false; + let immediate = false; + for (const resolver of current.resolvers.values()) { + const invalidation = resolver(payload); + for (const key of invalidation.keys) keys.add(key); + close ||= Boolean(invalidation.close); + immediate ||= Boolean(invalidation.immediate); + } - if (invalidation.close) { + queueInvalidations(current, [...keys], { debounceMs, immediate }); + + if (close) { closeSharedEventSource(subscriptions, subscriptionKey, { flushPending: true }); } }; } + const resolverId = Symbol(subscriptionKey); + subscription.resolvers.set( + resolverId, + resolveInvalidation as EventResolver, + ); subscription.refcount += 1; subscription.mutators.set(mutate, (subscription.mutators.get(mutate) ?? 0) + 1); @@ -89,6 +104,8 @@ export function subscribeToSharedEventSource({ const current = subscriptions.get(subscriptionKey); if (!current) return; + current.resolvers.delete(resolverId); + const mutateCount = current.mutators.get(mutate) ?? 0; if (mutateCount <= 1) { current.mutators.delete(mutate); diff --git a/apps/fabro-web/app/routes/run-detail.tsx b/apps/fabro-web/app/routes/run-detail.tsx index 2732fdf03..0b9493da4 100644 --- a/apps/fabro-web/app/routes/run-detail.tsx +++ b/apps/fabro-web/app/routes/run-detail.tsx @@ -1,8 +1,9 @@ -import { useEffect, useRef } from "react"; +import { useEffect, useRef, useState } from "react"; import { ArrowPathIcon, ChevronRightIcon } from "@heroicons/react/20/solid"; import { Link, Outlet, useLocation } from "react-router"; import { InterviewDock } from "../components/interview-dock"; +import { SteerComposer } from "../components/steer-composer"; import { ErrorState } from "../components/state"; import { useToast } from "../components/toast"; import { PRIMARY_BUTTON_CLASS, SECONDARY_BUTTON_CLASS } from "../components/ui"; @@ -22,6 +23,7 @@ import { type PreviewMutationResult, } from "../lib/mutations"; import { useRunEvents } from "../lib/run-events"; +import { useRunToasts } from "../hooks/use-run-toasts"; import { useRun, useRunQuestions } from "../lib/queries"; import { canArchive, @@ -115,8 +117,10 @@ export default function RunDetail({ params }: { params: { id: string } }) { const { push, dismiss } = useToast(); const tabs = allTabs.filter((t) => !t.demoOnly || demoMode); const lifecycleToastStateRef = useRef(INITIAL_LIFECYCLE_TOAST_STATE); + const [steerOpen, setSteerOpen] = useState(false); useRunEvents(params.id); + useRunToasts(params.id); useEffect(() => { if (previewMutation.data?.intent === "preview") { @@ -204,6 +208,18 @@ export default function RunDetail({ params }: { params: { id: string } }) {
+ {statusKind === "running" && ( +
+ +
+ )} + {visibility.showPrimaryCancel && (
+ setSteerOpen(false)} + /> + {isBlocked && pendingQuestions.length > 0 && ( <>