Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 20 additions & 13 deletions apps/web/src/api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -86,20 +86,12 @@ export type PrincipalsPage = Paginated<Principal>;
/** An arktype schema, seen as the validating call every `Type` provides. */
type Validator<T> = (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<T>(
path: string,
schema: Validator<T>,
options?: { readonly refetchInterval?: number | false },
): APIQuery<T> {
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<T>(path: string, schema: Validator<T>) {
return {
queryKey: pathToQueryKey(path),
enabled,
...(options?.refetchInterval === undefined ? {} : { refetchInterval: options.refetchInterval }),
queryFn: async () => {
queryFn: async (): Promise<T> => {
const response = await fetch(path, {
headers: { accept: "application/json" },
});
Expand All @@ -115,6 +107,21 @@ export function useAPIQuery<T>(
}
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<T>(
path: string,
schema: Validator<T>,
options?: { readonly refetchInterval?: number | false },
): APIQuery<T> {
const result = useQuery({
...apiQueryOptions(path, schema),
enabled: path !== "",
...(options?.refetchInterval === undefined ? {} : { refetchInterval: options.refetchInterval }),
});
return toAPIQuery(result);
}
Expand Down
2 changes: 2 additions & 0 deletions apps/web/src/chat-path.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 3 additions & 1 deletion apps/web/src/pages/workbench-page.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ import { VoiceOverlay } from "../voice/voice-overlay";
import { BenchDrawer } from "../bench/bench-drawer";
import { BenchPill } from "../bench/bench-pill";
import { ArtifactsTab } from "../bench/artifacts-tab";
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";
Expand Down Expand Up @@ -228,10 +228,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 ?? [];
Expand Down
59 changes: 41 additions & 18 deletions apps/web/src/shell/workbench-list.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -98,19 +99,50 @@ export function WorkbenchList({
/>
))}
</div>
<WorkerGroup path={path} onNavigate={onNavigate} />
<WorkerGroup path={path} onNavigate={onNavigate} benches={workbenches} />
</div>
);
}

// 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 (
<button
type="button"
className="shell-ch-row"
aria-current={active ? "true" : undefined}
data-active={active ? "true" : undefined}
onClick={onSelect}
>
<WorkbenchAvatar kind="worker" name={agent.name} size="sm" status={status.tone} />
<span className="shell-ch-meta">
<span className="shell-ch-name-row">
<span className="shell-ch-name">{agent.name}</span>
</span>
</span>
</button>
);
}

function WorkerGroup({
path,
onNavigate,
benches,
}: {
readonly path: string;
readonly onNavigate: (to: string) => void;
readonly benches: readonly HubTenant[];
}) {
const { selectedTenantId } = useBench();
const agents = useQuery({
Expand All @@ -125,23 +157,14 @@ function WorkerGroup({
<SectionLabel>Workers</SectionLabel>
{workers.map((agent) => {
const to = `/workers/${agent.id}`;
const active = path === to;
return (
<button
<WorkerRow
key={agent.id}
type="button"
className="shell-ch-row"
aria-current={active ? "true" : undefined}
data-active={active ? "true" : undefined}
onClick={() => onNavigate(to)}
>
<WorkbenchAvatar kind="worker" name={agent.name} size="sm" status="idle" />
<span className="shell-ch-meta">
<span className="shell-ch-name-row">
<span className="shell-ch-name">{agent.name}</span>
</span>
</span>
</button>
agent={agent}
benches={benches}
active={path === to}
onSelect={() => onNavigate(to)}
/>
);
})}
</div>
Expand Down
161 changes: 134 additions & 27 deletions apps/web/src/worker-status.ts
Original file line number Diff line number Diff line change
@@ -1,44 +1,151 @@
// 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 = {
readonly tone: AvatarStatus;
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<PendingTurn>(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<string, number> = { "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);
}
Loading