From a944acd994cf9783e023f3b5807533ff152c11b3 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Tue, 29 Sep 2026 22:33:30 -0700 Subject: [PATCH 1/2] feat(web): workers show working from send until reply, in the sidebar and pill (CL-9575) --- apps/web/src/api.ts | 33 +++--- apps/web/src/chat-path.ts | 2 + apps/web/src/pages/workbench-page.tsx | 4 +- apps/web/src/shell/workbench-list.tsx | 59 +++++++--- apps/web/src/worker-status.ts | 161 +++++++++++++++++++++----- 5 files changed, 200 insertions(+), 59 deletions(-) diff --git a/apps/web/src/api.ts b/apps/web/src/api.ts index 5a506a4cb..d38e91311 100644 --- a/apps/web/src/api.ts +++ b/apps/web/src/api.ts @@ -86,20 +86,12 @@ export type PrincipalsPage = Paginated; /** An arktype schema, seen as the validating call every `Type` provides. */ type Validator = (data: unknown) => T | ArkErrors; -// Pass a module-level schema so identity stays stable. Empty paths are -// disabled and never fetch, so a call site with an unresolved tenant -// can't hit the network with a broken URL. -export function useAPIQuery( - path: string, - schema: Validator, - options?: { readonly refetchInterval?: number | false }, -): APIQuery { - const enabled = path !== ""; - const result = useQuery({ +/** The key and fetcher `useAPIQuery` runs, for callers that batch reads with + * `useQueries` and must share the same cache entry. */ +export function apiQueryOptions(path: string, schema: Validator) { + return { queryKey: pathToQueryKey(path), - enabled, - ...(options?.refetchInterval === undefined ? {} : { refetchInterval: options.refetchInterval }), - queryFn: async () => { + queryFn: async (): Promise => { const response = await fetch(path, { headers: { accept: "application/json" }, }); @@ -115,6 +107,21 @@ export function useAPIQuery( } return parsed; }, + }; +} + +// Pass a module-level schema so identity stays stable. Empty paths are +// disabled and never fetch, so a call site with an unresolved tenant +// can't hit the network with a broken URL. +export function useAPIQuery( + path: string, + schema: Validator, + options?: { readonly refetchInterval?: number | false }, +): APIQuery { + const result = useQuery({ + ...apiQueryOptions(path, schema), + enabled: path !== "", + ...(options?.refetchInterval === undefined ? {} : { refetchInterval: options.refetchInterval }), }); return toAPIQuery(result); } diff --git a/apps/web/src/chat-path.ts b/apps/web/src/chat-path.ts index 7c3a1fcb5..312205752 100644 --- a/apps/web/src/chat-path.ts +++ b/apps/web/src/chat-path.ts @@ -7,6 +7,8 @@ export const workbenchKeys = { scope: (tenantId: string) => ["workbench", tenantId] as const, tenant: (tenantId: string) => ["workbench", tenantId, "tenant"] as const, participants: (tenantId: string) => ["workbench", tenantId, "participants"] as const, + /** Client-only: {sentAt} from a send until the worker replies. */ + pendingTurn: (tenantId: string) => ["workbench", tenantId, "pendingTurn"] as const, timeline: (tenantId: string) => ["workbench", tenantId, "timeline"] as const, /** The bench's own child tenants (its workbenches), read over the stock * tenant routes — previously `chatKeys.childTenants`, re-homed here when diff --git a/apps/web/src/pages/workbench-page.tsx b/apps/web/src/pages/workbench-page.tsx index b8ee05eef..bb8721e0d 100644 --- a/apps/web/src/pages/workbench-page.tsx +++ b/apps/web/src/pages/workbench-page.tsx @@ -25,7 +25,7 @@ import { } from "@/chat/threads-api"; import { BenchDrawer } from "../bench/bench-drawer"; import { BenchPill } from "../bench/bench-pill"; -import { useWorkerStatus } from "../worker-status"; +import { markTurnPending, useWorkerStatus } from "../worker-status"; import { GrantsTab } from "../bench/grants-tab"; import { InformationTab } from "../bench/information-tab"; import { InsightsTab } from "../bench/insights-tab"; @@ -196,10 +196,12 @@ function Workbench({ workbenchTenantId }: { readonly workbenchTenantId: string } content, ...(inReplyTo !== undefined ? { inReplyTo } : {}), }), + onMutate: () => markTurnPending(queryClient, workbenchTenantId), onSuccess: () => queryClient.invalidateQueries({ queryKey: workbenchKeys.scope(workbenchTenantId), }), + onError: () => queryClient.setQueryData(workbenchKeys.pendingTurn(workbenchTenantId), null), }); const messages = timeline.data ?? []; diff --git a/apps/web/src/shell/workbench-list.tsx b/apps/web/src/shell/workbench-list.tsx index 35b2d4cfb..f2082caec 100644 --- a/apps/web/src/shell/workbench-list.tsx +++ b/apps/web/src/shell/workbench-list.tsx @@ -8,6 +8,7 @@ import { listChatAgents } from "@/chat/threads-api"; import { Hash } from "@/lib/icons"; import { useBench } from "../bench-context"; +import { useBenchWorkerStatus } from "../worker-status"; import { tenantKeys } from "../query-client"; import { workbenchIdFromPath, workbenchPath } from "../workbench-path"; import type { HubTenant } from "../needs-converge"; @@ -98,19 +99,50 @@ export function WorkbenchList({ /> ))} - + ); } -// No client-side in-flight-turn state is queryable here, so every worker -// renders idle. +function WorkerRow({ + agent, + benches, + active, + onSelect, +}: { + readonly agent: { readonly name: string }; + readonly benches: readonly HubTenant[]; + readonly active: boolean; + readonly onSelect: () => void; +}) { + // Sidebar workers are the workspace's agents; a bench's copy shares the name. + const status = useBenchWorkerStatus(benches, (p) => p.name === agent.name); + return ( + + ); +} + function WorkerGroup({ path, onNavigate, + benches, }: { readonly path: string; readonly onNavigate: (to: string) => void; + readonly benches: readonly HubTenant[]; }) { const { selectedTenantId } = useBench(); const agents = useQuery({ @@ -125,23 +157,14 @@ function WorkerGroup({ Workers {workers.map((agent) => { const to = `/workers/${agent.id}`; - const active = path === to; return ( - + agent={agent} + benches={benches} + active={path === to} + onSelect={() => onNavigate(to)} + /> ); })} diff --git a/apps/web/src/worker-status.ts b/apps/web/src/worker-status.ts index 800a52b07..a64d60e9b 100644 --- a/apps/web/src/worker-status.ts +++ b/apps/web/src/worker-status.ts @@ -1,13 +1,22 @@ -// The one status source for a worker. The inbox stream (see -// `subscribeToInbox`) writes into the participants and approvals caches this -// reads; nothing here fetches on its own. +// The one status source for a worker. Stock Interchange has no live turn +// stream, so `working` is client-derived: a send records `pendingTurn` in the +// query cache, and it ends when a worker message newer than it lands, an +// approval is pending, or the safety timeout passes. -import { useQuery } from "@tanstack/react-query"; +import { useEffect } from "react"; +import { useQueries, useQuery, useQueryClient, type QueryClient } from "@tanstack/react-query"; import type { AvatarStatus } from "@/chat/avatar"; -import { sameAddress, type WorkbenchParticipant } from "@/chat/threads-api"; +import { + listWorkbenchParticipants, + readWorkbench, + sameAddress, + type WorkbenchMessage, + type WorkbenchParticipant, +} from "@/chat/threads-api"; import { workbenchKeys } from "./chat-path"; -import { TenantApprovalsSchema, useAPIQuery } from "./api"; +import { TenantApprovalsSchema, apiQueryOptions } from "./api"; +import { createFetchStockHub } from "./needs-converge"; import { pendingApprovalsPath } from "./pending-approvals"; export type WorkerStatus = { @@ -15,30 +24,128 @@ export type WorkerStatus = { readonly text: string; }; -// `working` has no stock signal: the inbox stream carries no turn events, and -// a live agent's run stays `running` whether or not a turn is in flight. +type PendingTurn = { readonly sentAt: number }; +type Bench = { readonly id: string; readonly domain: string }; + +const PENDING_TURN_TIMEOUT_MS = 10 * 60 * 1000; +const PENDING_POLL_MS = 5000; + +/** Call from a send mutation's `onMutate`. */ +export function markTurnPending(queryClient: QueryClient, benchId: string): void { + queryClient.setQueryData(workbenchKeys.pendingTurn(benchId), { + sentAt: Date.now(), + }); +} + +function turnAnswered(timeline: readonly WorkbenchMessage[] | undefined, sentAt: number): boolean { + return (timeline ?? []).some((m) => m.author === "other" && Date.parse(m.at) > sentAt); +} + +const RANK: Record = { "Needs you": 3, Working: 2, Live: 1, Starting: 0 }; + +/** The status of the worker `isWorker` matches, across the given benches: + * needs-you beats working beats live. */ +export function useBenchWorkerStatus( + benches: readonly Bench[], + isWorker: (participant: WorkbenchParticipant) => boolean, +): WorkerStatus { + const queryClient = useQueryClient(); + const participants = useQueries({ + queries: benches.map((bench) => ({ + queryKey: workbenchKeys.participants(bench.id), + queryFn: () => listWorkbenchParticipants(bench.id, bench.domain), + enabled: bench.domain !== "", + })), + }); + const approvals = useQueries({ + queries: benches.map((bench) => ({ + ...apiQueryOptions(pendingApprovalsPath(bench.id), TenantApprovalsSchema), + })), + }); + const pending = useQueries({ + queries: benches.map((bench) => ({ + queryKey: workbenchKeys.pendingTurn(bench.id), + queryFn: (): PendingTurn | null => null, + enabled: false, + gcTime: Infinity, + })), + }); + const timelines = useQueries({ + queries: benches.map((bench, index) => ({ + queryKey: workbenchKeys.timeline(bench.id), + queryFn: () => readWorkbench(bench.id), + enabled: pending[index]?.data != null, + refetchInterval: PENDING_POLL_MS, + })), + }); + + let best: WorkerStatus | undefined; + const offer = (status: WorkerStatus) => { + if (best === undefined || (RANK[status.text] ?? 0) > (RANK[best.text] ?? 0)) best = status; + }; + const now = Date.now(); + const expiries: { id: string; at: number }[] = []; + + benches.forEach((bench, index) => { + const agent = participants[index]?.data?.find((p) => p.kind === "agent" && isWorker(p)); + if (agent === undefined) return; + if (agent.address === "") return offer({ tone: "idle", text: "Starting" }); + const needsYou = (approvals[index]?.data?.data ?? []).some( + (row) => row.status === "pending" && sameAddress(row.agentAddress, agent.address), + ); + if (needsYou) return offer({ tone: "ready", text: "Needs you" }); + const sentAt = pending[index]?.data?.sentAt; + if (sentAt !== undefined) { + expiries.push({ id: bench.id, at: sentAt + PENDING_TURN_TIMEOUT_MS }); + if (now < sentAt + PENDING_TURN_TIMEOUT_MS && !turnAnswered(timelines[index]?.data, sentAt)) { + return offer({ tone: "working", text: "Working…" }); + } + } + offer({ tone: "idle", text: "Live" }); + }); + + // Ends a turn by dropping its entry; a timer covers the timeout when no + // other signal arrives. + const settled = benches + .map((bench, index) => { + const sentAt = pending[index]?.data?.sentAt; + const approved = (approvals[index]?.data?.data ?? []).some((row) => row.status === "pending"); + return sentAt !== undefined && (approved || turnAnswered(timelines[index]?.data, sentAt)) + ? bench.id + : null; + }) + .filter((id): id is string => id !== null); + const settledKey = settled.join(","); + const expiryKey = expiries.map((e) => `${e.id}:${e.at}`).join(","); + useEffect(() => { + for (const id of settledKey.split(",").filter(Boolean)) { + queryClient.setQueryData(workbenchKeys.pendingTurn(id), null); + } + }, [settledKey, queryClient]); + useEffect(() => { + const timers = expiries.map(({ id, at }) => + setTimeout( + () => queryClient.setQueryData(workbenchKeys.pendingTurn(id), null), + Math.max(0, at - Date.now()), + ), + ); + return () => timers.forEach(clearTimeout); + // `expiryKey` is the serialized `expiries`. + // eslint-disable-next-line react-hooks/exhaustive-deps + }, [expiryKey, queryClient]); + + return best ?? { tone: "idle", text: "Live" }; +} + export function useWorkerStatus( tenantId: string | null, agentId: string | undefined, ): WorkerStatus { - const participants = useQuery({ - queryKey: workbenchKeys.participants(tenantId ?? ""), - queryFn: (): readonly WorkbenchParticipant[] => [], - enabled: false, + const tenant = useQuery({ + queryKey: workbenchKeys.tenant(tenantId ?? ""), + queryFn: () => createFetchStockHub().getTenant(tenantId ?? ""), + enabled: tenantId !== null, }); - const approvals = useAPIQuery( - tenantId === null ? "" : pendingApprovalsPath(tenantId), - TenantApprovalsSchema, - ); - - const agent = participants.data?.find((p) => p.kind === "agent" && p.id === agentId); - if (agent === undefined) return { tone: "idle", text: "Live" }; - if (agent.address === "") return { tone: "idle", text: "Starting" }; - - const needsYou = - approvals.kind === "ready" && - approvals.data.data.some( - (row) => row.status === "pending" && sameAddress(row.agentAddress, agent.address), - ); - return needsYou ? { tone: "ready", text: "Needs you" } : { tone: "idle", text: "Live" }; + const benches = tenantId === null ? [] : [{ id: tenantId, domain: tenant.data?.domain ?? "" }]; + return useBenchWorkerStatus(benches, (p) => p.id === agentId); } From e7cadc7f391f70d51e5bd365944ce69e71ad869e Mon Sep 17 00:00:00 2001 From: Sawyer Date: Tue, 29 Sep 2026 22:36:18 -0700 Subject: [PATCH 2/2] chore(ci): retrigger checks