diff --git a/apps/mobile/src/features/threads/ThreadFeed.tsx b/apps/mobile/src/features/threads/ThreadFeed.tsx index 2d5863aa41..9785663317 100644 --- a/apps/mobile/src/features/threads/ThreadFeed.tsx +++ b/apps/mobile/src/features/threads/ThreadFeed.tsx @@ -110,6 +110,11 @@ import { } from "../../native/SelectableMarkdownText"; import { AppText as Text } from "../../components/AppText"; +import { + delegationNoticeHeadline, + delegationNoticeSummary, + parseDelegationNotice, +} from "@t3tools/client-runtime/state/delegation-notice"; import { VideoPreviewModal, type VideoPreviewSource } from "../../components/VideoPreviewModal"; import { VideoAttachmentTile } from "../../components/VideoAttachmentTile"; import { MediaVideoPlayer } from "../../components/MediaVideoPlayer"; @@ -1376,6 +1381,7 @@ function renderFeedEntry( readonly onToggleTurnFold: (turnId: TurnId) => void; readonly onPressPreview: (source: FilePreviewSource) => void; readonly onPressVideo: (attachment: ChatFileAttachment, sourceIdentifier: string) => void; + readonly onOpenThread?: (threadId: string) => void; readonly markdownLinkHandlers: MarkdownLinkHandlers; readonly renderMarkdownImage: MarkdownImageRenderer; readonly renderViewedImage: MarkdownImageRenderer; @@ -1508,6 +1514,36 @@ function renderFeedEntry( !message.streaming; if (isUser) { + const notice = parseDelegationNotice(message); + if (notice !== null) { + return ( + + {notice.updates.map((update) => ( + props.onOpenThread?.(update.threadId) : undefined + } + accessibilityRole="button" + className="py-1 active:opacity-70" + > + {delegationNoticeHeadline(update)} + + {update.title} + + {update.reason !== null && ( + {update.reason} + )} + + ))} + + ); + } const referenceIds = new Set( collectComposerContextReferences(message.text).map((reference) => reference.contextId), ); @@ -2059,6 +2095,16 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { const iconSubtleColor = theme["--color-icon-subtle"]; const screenColor = theme["--color-screen"]; const userBubbleColor = theme["--color-user-bubble"]; + const onOpenThread = useCallback( + (threadId: string) => { + navigation.navigate("Thread", { + environmentId: String(props.environmentId), + threadId: String(threadId), + }); + }, + [navigation, props.environmentId], + ); + const onMarkdownLinkPress = useCallback( (href: string) => { const presentation = resolveMarkdownLinkPresentation(href); @@ -2720,6 +2766,7 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { onToggleTurnFold, onPressPreview, onPressVideo, + onOpenThread, markdownLinkHandlers, renderMarkdownImage, renderViewedImage, @@ -2774,6 +2821,7 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { markdownLinkHandlers, onPressPreview, onPressVideo, + onOpenThread, onToggleTurnFold, onToggleWorkGroup, onToggleWorkRow, diff --git a/apps/web/src/components/chat/DelegationNoticeRow.test.tsx b/apps/web/src/components/chat/DelegationNoticeRow.test.tsx new file mode 100644 index 0000000000..81bd84c946 --- /dev/null +++ b/apps/web/src/components/chat/DelegationNoticeRow.test.tsx @@ -0,0 +1,91 @@ +import type { DelegationNotice } from "@t3tools/client-runtime/state/delegation-notice"; +import { EnvironmentId, ThreadId } from "@t3tools/contracts"; +import { renderToStaticMarkup } from "react-dom/server"; +import { describe, expect, it, vi } from "vite-plus/test"; + +import { DelegationNoticeRow } from "./DelegationNoticeRow"; + +vi.mock("@tanstack/react-router", () => ({ + Link: ({ + children, + params, + }: { + children: React.ReactNode; + params: { environmentId: string; threadId: string }; + }) => {children}, +})); + +const ENV = EnvironmentId.make("env-1"); +const EXECUTOR = ThreadId.make("delegated:lead:0123456789abcdef"); +const CHILD = ThreadId.make("delegated:lead:fedcba9876543210"); + +const finished: DelegationNotice = { + updates: [ + { + threadId: EXECUTOR, + title: "Executor · Fix the parser", + status: "completed", + reason: null, + isExecutor: true, + }, + ], + needsUser: false, +}; + +describe("DelegationNoticeRow", () => { + it("says who finished and links to that thread, with none of the agent-facing text", () => { + const html = renderToStaticMarkup( + , + ); + expect(html).toContain("data-delegation-notice"); + expect(html).toContain('data-needs-user="false"'); + expect(html).toContain("Executor finished"); + expect(html).toContain("Executor · Fix the parser"); + expect(html).toContain(`href="/${ENV}/${EXECUTOR}"`); + expect(html).not.toContain("JSON"); + expect(html).not.toContain("threadId"); + expect(html).not.toContain("authorized scope"); + }); + + it("marks a notice that is waiting on the user and shows a failure's reason", () => { + const html = renderToStaticMarkup( + , + ); + expect(html).toContain('data-needs-user="true"'); + expect(html).toContain("Delegated thread needs your approval"); + expect(html).toContain("Executor failed"); + expect(html).toContain("quota exceeded"); + expect(html).toContain(`href="/${ENV}/${CHILD}"`); + expect(html).toContain(`href="/${ENV}/${EXECUTOR}"`); + // One item per update, so a screen reader hears a list. + expect(html.match(/data-delegation-notice-item/g)).toHaveLength(2); + }); + + it("is announced as a status, not as something the user said", () => { + const html = renderToStaticMarkup( + , + ); + expect(html).toContain('role="status"'); + }); +}); diff --git a/apps/web/src/components/chat/DelegationNoticeRow.tsx b/apps/web/src/components/chat/DelegationNoticeRow.tsx new file mode 100644 index 0000000000..e7c1b8207d --- /dev/null +++ b/apps/web/src/components/chat/DelegationNoticeRow.tsx @@ -0,0 +1,68 @@ +import { + delegationNoticeHeadline, + type DelegationNotice, +} from "@t3tools/client-runtime/state/delegation-notice"; +import type { EnvironmentId } from "@t3tools/contracts"; +import { Link } from "@tanstack/react-router"; +import { cn } from "~/lib/utils"; + +/** + * Replaces the raw automatic wake message in the timeline: who finished or + * needs the user, and a link to that thread. + */ +export function DelegationNoticeRow(props: { + readonly notice: DelegationNotice; + readonly environmentId: EnvironmentId; +}) { + return ( +
+
    + {props.notice.updates.map((update) => ( +
  • +
    +
    + + + {delegationNoticeHeadline(update)} + + {update.title} +
    + + Open + +
    + {update.reason !== null &&
    {update.reason}
    } +
  • + ))} +
+
+ ); +} diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 70226b62e2..5807ced6d6 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -4,6 +4,8 @@ import { delegationActivity } from "../../delegationActivity"; import { usePreparedConnection } from "~/state/session"; import { MediaVideoPlayer } from "../media/MediaVideoPlayer"; import { PierreEntryIcon } from "./PierreEntryIcon"; +import { DelegationNoticeRow } from "./DelegationNoticeRow"; +import { parseDelegationNotice } from "@t3tools/client-runtime/state/delegation-notice"; import { ReadOnlySourcePreview } from "../files/AttachmentFilePreview"; import { useRightPanelStore } from "~/rightPanelStore"; import { @@ -1662,7 +1664,22 @@ function ContextCompactionTimelineRow({ ); } +/** + * The automatic wake message Pylon sends a parent is not something the user + * said, so it renders as a notice and offers none of a message's actions. The + * split keeps every hook of the message row unconditional. + */ function UserTimelineRow({ row }: { row: Extract }) { + const ctx = use(TimelineRowCtx); + const notice = parseDelegationNotice(row.message); + return notice === null ? ( + + ) : ( + + ); +} + +function UserMessageTimelineRow({ row }: { row: Extract }) { const ctx = use(TimelineRowCtx); const { onImageExpand, onFileOpen } = ctx; const resources = useMemo( diff --git a/packages/client-runtime/package.json b/packages/client-runtime/package.json index e7884eba2a..6487da748e 100644 --- a/packages/client-runtime/package.json +++ b/packages/client-runtime/package.json @@ -251,6 +251,14 @@ "types": "./src/state/delegatedThreads.ts", "default": "./src/state/delegatedThreads.ts" }, + "./state/pair": { + "types": "./src/state/pair.ts", + "default": "./src/state/pair.ts" + }, + "./state/delegation-notice": { + "types": "./src/state/delegationNotice.ts", + "default": "./src/state/delegationNotice.ts" + }, "./state/thread-sort": { "types": "./src/state/threadSort.ts", "default": "./src/state/threadSort.ts" diff --git a/packages/client-runtime/src/state/delegatedThreads.ts b/packages/client-runtime/src/state/delegatedThreads.ts index b061275e4e..005a738a8c 100644 --- a/packages/client-runtime/src/state/delegatedThreads.ts +++ b/packages/client-runtime/src/state/delegatedThreads.ts @@ -103,7 +103,7 @@ export function nestedRowContainsThread( } /** Shell-only lifecycle; keep admission/terminal precedence aligned with server delegation. */ -function delegatedThreadStatus(shell: EnvironmentThreadShell) { +export function delegatedThreadStatus(shell: EnvironmentThreadShell) { if (shell.archivedAt !== null) return "archived"; const { session, latestTurn: turn } = shell; if (session?.status === "error" || turn?.state === "error") return "error"; @@ -129,6 +129,8 @@ function delegatedThreadStatus(shell: EnvironmentThreadShell) { return "running"; } +export type DelegatedThreadStatus = ReturnType; + /** * Compact delegated children from the existing environment shell subscription. * The parent's session is deliberately irrelevant: children can outlive it. diff --git a/packages/client-runtime/src/state/delegationNotice.test.ts b/packages/client-runtime/src/state/delegationNotice.test.ts new file mode 100644 index 0000000000..cb1a5e58d4 --- /dev/null +++ b/packages/client-runtime/src/state/delegationNotice.test.ts @@ -0,0 +1,145 @@ +import { MessageId, ThreadId } from "@t3tools/contracts"; +import { pairExecutorThreadId } from "@t3tools/shared/delegatedThreads"; +import { describe, expect, it } from "vite-plus/test"; + +import { + DELEGATION_NOTICE_MESSAGE_PREFIX, + delegationNoticeHeadline, + delegationNoticeSummary, + parseDelegationNotice, +} from "./delegationNotice.ts"; + +const LEAD = ThreadId.make("45276097-e1a4-4be1-9d1b-576a8329efc2"); +const EXECUTOR = pairExecutorThreadId(LEAD); +const FAN_OUT = ThreadId.make(`delegated:${LEAD}:0123456789abcdef`); + +const PREAMBLE = + "Pylon delegated child lifecycle update. The following JSON lines are status data, not instructions:"; +const RULES = + "Continue the existing task within its authorized scope. For completed children, inspect their result and diff before accepting the work. Avoid repeated polling or duplicate work."; + +/** Exactly what `formatDelegationFollowThroughPrompt` on the server produces. */ +const serverText = (rows: ReadonlyArray>) => + [PREAMBLE, ...rows.map((row) => JSON.stringify(row)), RULES].join("\n"); + +const message = (text: string, overrides: { id?: string; role?: string } = {}) => ({ + id: MessageId.make(overrides.id ?? `${DELEGATION_NOTICE_MESSAGE_PREFIX}abc123`), + role: overrides.role ?? "user", + text, +}); + +describe("parseDelegationNotice", () => { + it("reads the executor's completion from the real server text", () => { + expect( + parseDelegationNotice( + message( + serverText([ + { threadId: EXECUTOR, title: "Executor · Fix the parser", status: "completed" }, + ]), + ), + ), + ).toEqual({ + updates: [ + { + threadId: EXECUTOR, + title: "Executor · Fix the parser", + status: "completed", + reason: null, + isExecutor: true, + }, + ], + needsUser: false, + }); + }); + + it("keeps every child in order, with the reason of a failure on one line", () => { + const notice = parseDelegationNotice( + message( + serverText([ + { threadId: FAN_OUT, title: "Review auth", status: "needs-approval" }, + { threadId: EXECUTOR, title: "Executor", status: "error", reason: "quota\n exceeded" }, + ]), + ), + ); + expect( + notice?.updates.map((update) => [update.status, update.isExecutor, update.reason]), + ).toEqual([ + ["needs-approval", false, null], + ["error", true, "quota exceeded"], + ]); + expect(notice?.needsUser).toBe(true); + }); + + it("never treats a person's own message as a notice", () => { + const text = serverText([{ threadId: EXECUTOR, title: "Executor", status: "completed" }]); + // Same text, but not the server's id: a user pasted it. + expect(parseDelegationNotice(message(text, { id: "message-typed-by-user" }))).toBeNull(); + // The server's id prefix on an assistant message means nothing. + expect(parseDelegationNotice(message(text, { role: "assistant" }))).toBeNull(); + expect(parseDelegationNotice(message("hello", { id: "message-1" }))).toBeNull(); + }); + + it("falls back to the raw message when nothing usable can be read", () => { + // Returning null shows the original text, which is better than an empty notice. + expect(parseDelegationNotice(message(`${PREAMBLE}\n${RULES}`))).toBeNull(); + expect(parseDelegationNotice(message(`${PREAMBLE}\n{not json}\n${RULES}`))).toBeNull(); + expect( + parseDelegationNotice( + message(serverText([{ threadId: EXECUTOR, title: "Executor", status: "exploded" }])), + ), + ).toBeNull(); + }); + + it("skips a malformed line but keeps the good ones", () => { + const text = [ + PREAMBLE, + "{broken", + JSON.stringify({ threadId: EXECUTOR, title: "Executor", status: "interrupted" }), + JSON.stringify({ title: "no thread id", status: "completed" }), + RULES, + ].join("\n"); + expect(parseDelegationNotice(message(text))?.updates).toEqual([ + { + threadId: EXECUTOR, + title: "Executor", + status: "interrupted", + reason: null, + isExecutor: true, + }, + ]); + }); +}); + +describe("notice wording", () => { + const update = (status: string, isExecutor: boolean) => ({ + threadId: isExecutor ? EXECUTOR : FAN_OUT, + title: "Fix the parser", + status: status as "completed", + reason: null, + isExecutor, + }); + + it("says who and what happened in plain words", () => { + expect(delegationNoticeHeadline(update("completed", true))).toBe("Executor finished"); + expect(delegationNoticeHeadline(update("needs-approval", true))).toBe( + "Executor needs your approval", + ); + expect(delegationNoticeHeadline(update("needs-input", true))).toBe("Executor has a question"); + expect(delegationNoticeHeadline(update("interrupted", true))).toBe("Executor was stopped"); + expect(delegationNoticeHeadline(update("error", true))).toBe("Executor failed"); + expect(delegationNoticeHeadline(update("completed", false))).toBe("Delegated thread finished"); + expect(delegationNoticeHeadline(update("error", false))).toBe("Delegated thread failed"); + }); + + it("summarizes one update by its headline and several by a count", () => { + expect( + delegationNoticeSummary({ updates: [update("completed", true)], needsUser: false }), + ).toBe("Executor finished"); + expect( + delegationNoticeSummary({ + updates: [update("completed", false), update("error", false), update("completed", true)], + needsUser: false, + }), + ).toBe("3 delegated updates"); + }); +}); diff --git a/packages/client-runtime/src/state/delegationNotice.ts b/packages/client-runtime/src/state/delegationNotice.ts new file mode 100644 index 0000000000..68e1b595ba --- /dev/null +++ b/packages/client-runtime/src/state/delegationNotice.ts @@ -0,0 +1,148 @@ +/** + * The automatic message Pylon sends a parent when its delegated work changes + * state is written for the agent: a preamble, one JSON line per child, and a + * paragraph of rules. Shown raw, it reads as a wall of JSON in the user's own + * voice. This turns it into what a person wants to know: which thread, what + * happened, and where to go. + */ +import { ThreadId, type MessageId } from "@t3tools/contracts"; +import { isPairExecutorThreadId } from "@t3tools/shared/delegatedThreads"; + +/** Message ids the server gives its automatic follow-through turns. */ +export const DELEGATION_NOTICE_MESSAGE_PREFIX = "delegation-follow-through:"; + +export type DelegationNoticeStatus = + | "completed" + | "needs-approval" + | "needs-input" + | "interrupted" + | "error"; + +export interface DelegationNoticeUpdate { + readonly threadId: ThreadId; + readonly title: string; + readonly status: DelegationNoticeStatus; + /** One line explaining a failure, when the server sent one. */ + readonly reason: string | null; + /** True for a pair executor, false for a fan-out child. */ + readonly isExecutor: boolean; +} + +export interface DelegationNotice { + readonly updates: ReadonlyArray; + /** Whether any update is waiting on the user rather than reporting a result. */ + readonly needsUser: boolean; +} + +function isValidStatus(status: unknown): status is DelegationNoticeStatus { + return ( + status === "completed" || + status === "needs-approval" || + status === "needs-input" || + status === "interrupted" || + status === "error" + ); +} + +/** + * The notice a message carries, or null for every ordinary message. Only a + * user-role message with the server's id prefix is ever parsed, so a person who + * pastes similar text still sees their own words. + */ +export function parseDelegationNotice(message: { + readonly id: MessageId; + readonly role: string; + readonly text: string; +}): DelegationNotice | null { + if (message.role !== "user" || !message.id.startsWith(DELEGATION_NOTICE_MESSAGE_PREFIX)) { + return null; + } + + const lines = message.text.split("\n"); + const updates: DelegationNoticeUpdate[] = []; + + for (const line of lines) { + const trimmed = line.trim(); + if (!trimmed.startsWith("{")) { + continue; + } + let parsed: unknown; + try { + parsed = JSON.parse(trimmed); + } catch { + continue; + } + if ( + typeof parsed === "object" && + parsed !== null && + !Array.isArray(parsed) && + "threadId" in parsed && + typeof parsed.threadId === "string" && + parsed.threadId.length > 0 && + "title" in parsed && + typeof parsed.title === "string" && + "status" in parsed && + isValidStatus(parsed.status) + ) { + let reason: string | null = null; + if ("reason" in parsed && typeof parsed.reason === "string") { + const collapsed = parsed.reason.replace(/\s+/g, " ").trim(); + reason = collapsed.length > 0 ? collapsed : null; + } + const threadId = ThreadId.make(parsed.threadId); + const isExecutor = isPairExecutorThreadId(threadId); + updates.push({ + threadId, + title: parsed.title, + status: parsed.status, + reason, + isExecutor, + }); + } + } + + if (updates.length === 0) { + return null; + } + + const needsUser = updates.some( + (update) => update.status === "needs-approval" || update.status === "needs-input", + ); + + return { + updates, + needsUser, + }; +} + +/** \"Executor finished\", \"Delegated thread needs your approval\", and so on. */ +export function delegationNoticeHeadline(update: DelegationNoticeUpdate): string { + const subject = update.isExecutor ? "Executor" : "Delegated thread"; + let verb: string; + switch (update.status) { + case "completed": + verb = "finished"; + break; + case "needs-approval": + verb = "needs your approval"; + break; + case "needs-input": + verb = "has a question"; + break; + case "interrupted": + verb = "was stopped"; + break; + case "error": + verb = "failed"; + break; + } + return `${subject} ${verb}`; +} + +/** One line for the whole notice: the headline of one update, or a count. */ +export function delegationNoticeSummary(notice: DelegationNotice): string { + if (notice.updates.length === 1 && notice.updates[0] !== undefined) { + return delegationNoticeHeadline(notice.updates[0]); + } + return `${notice.updates.length} delegated updates`; +} diff --git a/packages/client-runtime/src/state/pair.test.ts b/packages/client-runtime/src/state/pair.test.ts new file mode 100644 index 0000000000..321c624161 --- /dev/null +++ b/packages/client-runtime/src/state/pair.test.ts @@ -0,0 +1,240 @@ +import { + EnvironmentId, + ProjectId, + ProviderInstanceId, + ThreadId, + TurnId, + type OrchestrationSession, +} from "@t3tools/contracts"; +import { pairExecutorThreadId } from "@t3tools/shared/delegatedThreads"; +import { describe, expect, it } from "vite-plus/test"; + +import type { EnvironmentThreadShell } from "./models.ts"; +import { + PAIR_UNSUPPORTED_LEAD_REASON, + isPairLeadSupported, + pairExecutorCreateInput, + resolvePairState, + withoutPairExecutors, +} from "./pair.ts"; + +const NOW = "2026-09-18T00:00:00.000Z"; +const ENV = EnvironmentId.make("env-1"); +const OTHER_ENV = EnvironmentId.make("env-2"); +const LEAD = ThreadId.make("lead:with:colons"); +const EXECUTOR = pairExecutorThreadId(LEAD); +const FAN_OUT_CHILD = ThreadId.make(`delegated:${LEAD}:0123456789abcdef`); +const EXECUTOR_SELECTION = { + instanceId: ProviderInstanceId.make("antigravity"), + model: "gemini-3-flash", +}; + +function shell( + id: ThreadId, + overrides: Partial = {}, +): EnvironmentThreadShell { + return { + environmentId: ENV, + id, + projectId: ProjectId.make("project-1"), + title: "Fix the parser", + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-6-astra" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: "feat/work", + worktreePath: "/wt/repo/lead", + pullRequests: [], + latestTurn: null, + createdAt: NOW, + updatedAt: NOW, + archivedAt: null, + settledOverride: null, + settledAt: null, + session: null, + latestUserMessageAt: null, + hasPendingApprovals: false, + hasPendingUserInput: false, + hasActionableProposedPlan: false, + ...overrides, + }; +} + +const turnId = TurnId.make("turn-1"); +const completedTurn = { + turnId, + state: "completed" as const, + requestedAt: NOW, + startedAt: NOW, + completedAt: NOW, + assistantMessageId: null, +}; +const runningTurn = { ...completedTurn, state: "running" as const, completedAt: null }; +const session = (overrides: Partial = {}): OrchestrationSession => ({ + threadId: EXECUTOR, + status: "ready", + providerName: "antigravity", + runtimeMode: "full-access", + activeTurnId: null, + lastError: null, + updatedAt: NOW, + ...overrides, +}); +const executor = (overrides: Partial = {}) => + shell(EXECUTOR, { modelSelection: EXECUTOR_SELECTION, ...overrides }); + +const lead = (driverKind: string | null | undefined = "codex") => ({ + environmentId: ENV, + threadId: LEAD, + driverKind, +}); + +describe("isPairLeadSupported", () => { + it("refuses only Antigravity", () => { + expect(isPairLeadSupported("antigravity")).toBe(false); + for (const driver of [ + "codex", + "claudeAgent", + "primeAgent", + "cursor", + "opencode", + null, + undefined, + ]) + expect(isPairLeadSupported(driver)).toBe(true); + }); +}); + +describe("resolvePairState", () => { + it("says why when the lead's provider cannot lead, even if an executor exists", () => { + expect( + resolvePairState({ threads: [shell(LEAD), executor()], lead: lead("antigravity") }), + ).toEqual({ kind: "unsupported-lead", reason: PAIR_UNSUPPORTED_LEAD_REASON }); + }); + + it("is off, with the id a client would create, when no executor is listed", () => { + expect(resolvePairState({ threads: [shell(LEAD)], lead: lead() })).toEqual({ + kind: "off", + executorId: EXECUTOR, + }); + // A fan-out child is not the pair, and an archived executor means the pair is off. + expect( + resolvePairState({ + threads: [shell(LEAD), shell(FAN_OUT_CHILD), executor({ archivedAt: NOW })], + lead: lead(), + }), + ).toEqual({ kind: "off", executorId: EXECUTOR }); + // An executor with the same id in another environment belongs to another lead. + expect( + resolvePairState({ + threads: [shell(LEAD), executor({ environmentId: OTHER_ENV })], + lead: lead(), + }), + ).toEqual({ kind: "off", executorId: EXECUTOR }); + }); + + it("reads an executor that never ran as idle", () => { + expect(resolvePairState({ threads: [shell(LEAD), executor()], lead: lead() })).toEqual({ + kind: "on", + executorId: EXECUTOR, + phase: "idle", + modelSelection: EXECUTOR_SELECTION, + activity: null, + }); + }); + + it("follows the executor through a turn", () => { + const phaseOf = (overrides: Partial) => { + const state = resolvePairState({ threads: [executor(overrides)], lead: lead() }); + return state.kind === "on" ? state.phase : state.kind; + }; + expect(phaseOf({ session: session({ status: "starting" }) })).toBe("running"); + expect( + phaseOf({ + latestTurn: runningTurn, + session: session({ status: "running", activeTurnId: turnId }), + }), + ).toBe("running"); + expect(phaseOf({ latestTurn: completedTurn, session: session() })).toBe("completed"); + // A session stopped after the turn finished does not make it interrupted. + expect(phaseOf({ latestTurn: completedTurn, session: session({ status: "stopped" }) })).toBe( + "completed", + ); + expect(phaseOf({ latestTurn: { ...completedTurn, state: "interrupted" } })).toBe("interrupted"); + expect(phaseOf({ latestTurn: { ...completedTurn, state: "error" } })).toBe("error"); + expect(phaseOf({ latestTurn: runningTurn, hasPendingApprovals: true })).toBe("needs-approval"); + expect(phaseOf({ latestTurn: runningTurn, hasPendingUserInput: true })).toBe("needs-input"); + }); + + it("carries one line of activity: the error, or the plan step while running", () => { + const failed = resolvePairState({ + threads: [ + executor({ + latestTurn: { ...completedTurn, state: "error" }, + session: session({ status: "error", lastError: "quota\n exceeded today" }), + }), + ], + lead: lead(), + }); + expect(failed).toMatchObject({ kind: "on", phase: "error", activity: "quota exceeded today" }); + }); +}); + +describe("withoutPairExecutors", () => { + it("drops executors and keeps leads, fan-out children, and ordinary threads in order", () => { + const other = shell(ThreadId.make("thread-2")); + const threads = [executor(), shell(LEAD), shell(FAN_OUT_CHILD), other]; + expect(withoutPairExecutors(threads).map((thread) => thread.id)).toEqual([ + LEAD, + FAN_OUT_CHILD, + other.id, + ]); + }); +}); + +describe("pairExecutorCreateInput", () => { + it("creates the executor in the lead's location, always implementing", () => { + expect( + pairExecutorCreateInput({ + lead: shell(LEAD, { interactionMode: "plan" }), + executorSelection: EXECUTOR_SELECTION, + childRuntimeMode: "inherit", + }), + ).toEqual({ + threadId: EXECUTOR, + projectId: ProjectId.make("project-1"), + title: "Executor · Fix the parser", + modelSelection: EXECUTOR_SELECTION, + runtimeMode: "full-access", + interactionMode: "default", + branch: "feat/work", + worktreePath: "/wt/repo/lead", + }); + }); + + it("honors Supervised child permissions and never broadens the lead's mode", () => { + const input = { + executorSelection: EXECUTOR_SELECTION, + childRuntimeMode: "approval-required" as const, + }; + expect(pairExecutorCreateInput({ ...input, lead: shell(LEAD) }).runtimeMode).toBe( + "approval-required", + ); + expect( + pairExecutorCreateInput({ + lead: shell(LEAD, { runtimeMode: "approval-required" }), + executorSelection: EXECUTOR_SELECTION, + childRuntimeMode: "inherit", + }).runtimeMode, + ).toBe("approval-required"); + }); + + it("keeps the title within the 200 character limit", () => { + expect( + pairExecutorCreateInput({ + lead: shell(LEAD, { title: "x".repeat(400) }), + executorSelection: EXECUTOR_SELECTION, + childRuntimeMode: "inherit", + }).title, + ).toHaveLength(200); + }); +}); diff --git a/packages/client-runtime/src/state/pair.ts b/packages/client-runtime/src/state/pair.ts new file mode 100644 index 0000000000..1a159d4a37 --- /dev/null +++ b/packages/client-runtime/src/state/pair.ts @@ -0,0 +1,141 @@ +/** + * What a client needs to show and control a pair: whether this lead can have + * one, whether it is on, what its executor is doing, and the command inputs + * that turn it on and off. Pure, shared by web and mobile. + */ +import type { + EnvironmentId, + ModelSelection, + ProjectId, + RuntimeMode, + ThreadId, +} from "@t3tools/contracts"; +import { isPairExecutorThreadId, pairExecutorThreadId } from "@t3tools/shared/delegatedThreads"; + +import { delegatedThreadStatus } from "./delegatedThreads.ts"; +import type { EnvironmentThreadShell } from "./models.ts"; + +/** Shown wherever the pair control is disabled for the selected lead provider. */ +export const PAIR_UNSUPPORTED_LEAD_REASON = + "Pair needs a lead that can hold its own subagents. Antigravity works as the executor."; + +export type PairExecutorPhase = + | "idle" + | "running" + | "needs-approval" + | "needs-input" + | "completed" + | "interrupted" + | "error"; + +export type PairState = + | { readonly kind: "unsupported-lead"; readonly reason: string } + | { readonly kind: "off"; readonly executorId: ThreadId } + | { + readonly kind: "on"; + readonly executorId: ThreadId; + readonly phase: PairExecutorPhase; + readonly modelSelection: ModelSelection; + /** The executor's last error, or the plan step it is on, trimmed for one line. */ + readonly activity: string | null; + }; + +/** Antigravity cannot lead: Pylon cannot hold its own subagents for one session. */ +export function isPairLeadSupported(driverKind: string | null | undefined): boolean { + return driverKind !== "antigravity"; +} + +export function resolvePairState(input: { + readonly threads: readonly EnvironmentThreadShell[]; + readonly lead: { + readonly environmentId: EnvironmentId; + readonly threadId: ThreadId; + readonly driverKind: string | null | undefined; + }; +}): PairState { + if (!isPairLeadSupported(input.lead.driverKind)) { + return { kind: "unsupported-lead", reason: PAIR_UNSUPPORTED_LEAD_REASON }; + } + + const executorId = pairExecutorThreadId(input.lead.threadId); + const executor = input.threads.find( + (thread) => + thread.id === executorId && + thread.environmentId === input.lead.environmentId && + thread.archivedAt === null, + ); + if (executor === undefined) { + return { kind: "off", executorId }; + } + + let phase: PairExecutorPhase; + if (executor.session === null && executor.latestTurn === null) { + phase = "idle"; + } else { + const status = delegatedThreadStatus(executor); + if (status === "starting") { + phase = "running"; + } else if (status !== "archived") { + phase = status; + } else { + phase = "idle"; + } + } + + const rawActivity = + phase === "error" + ? executor.session?.lastError + : phase === "running" + ? executor.planProgress?.step + : null; + const activity = rawActivity?.replace(/\s+/g, " ").trim().slice(0, 160) || null; + + return { + kind: "on", + executorId, + phase, + modelSelection: executor.modelSelection, + activity, + }; +} + +/** Threads a list should show: an executor is reached through its lead, not the list. */ +export function withoutPairExecutors( + threads: readonly T[], +): readonly T[] { + return threads.filter((thread) => !isPairExecutorThreadId(thread.id)); +} + +export interface PairExecutorCreateInput { + readonly threadId: ThreadId; + readonly projectId: ProjectId; + readonly title: string; + readonly modelSelection: ModelSelection; + readonly runtimeMode: RuntimeMode; + readonly interactionMode: "default"; + readonly branch: string | null; + readonly worktreePath: string | null; +} + +/** The `thread.create` fields that turn a pair on for a lead. */ +export function pairExecutorCreateInput(input: { + readonly lead: Pick< + EnvironmentThreadShell, + "id" | "projectId" | "title" | "runtimeMode" | "branch" | "worktreePath" + >; + readonly executorSelection: ModelSelection; + /** The user's Child permissions setting. */ + readonly childRuntimeMode: "inherit" | "approval-required"; +}): PairExecutorCreateInput { + return { + threadId: pairExecutorThreadId(input.lead.id), + projectId: input.lead.projectId, + title: `Executor · ${input.lead.title}`.slice(0, 200), + modelSelection: input.executorSelection, + runtimeMode: + input.childRuntimeMode === "approval-required" ? "approval-required" : input.lead.runtimeMode, + interactionMode: "default", + branch: input.lead.branch, + worktreePath: input.lead.worktreePath, + }; +} diff --git a/packages/shared/src/delegatedThreads.test.ts b/packages/shared/src/delegatedThreads.test.ts new file mode 100644 index 0000000000..433f4f6a90 --- /dev/null +++ b/packages/shared/src/delegatedThreads.test.ts @@ -0,0 +1,37 @@ +import * as NodeCrypto from "node:crypto"; + +import { ThreadId } from "@t3tools/contracts"; +import { describe, expect, it } from "vite-plus/test"; + +import { + PAIR_DELEGATION_KEY, + delegatedParentThreadId, + isPairExecutorThreadId, + pairExecutorThreadId, +} from "./delegatedThreads.ts"; + +/** The server's formula: `delegated::`. */ +const serverChildId = (lead: string, key: string) => + `delegated:${lead}:${NodeCrypto.createHash("sha256").update(`${lead}\n${key}`).digest("hex").slice(0, 16)}`; + +describe("pair executor identity", () => { + it("matches the id the server derives for the reserved key", () => { + expect(PAIR_DELEGATION_KEY).toBe("pair"); + for (const lead of ["thread-lead", "import:codex:with:colons", "45276097-e1a4-4be1-9d1b"]) { + expect(pairExecutorThreadId(ThreadId.make(lead))).toBe(serverChildId(lead, "pair")); + } + }); + + it("resolves back to its lead", () => { + const lead = ThreadId.make("import:a:b"); + expect(delegatedParentThreadId(pairExecutorThreadId(lead))).toBe(lead); + }); + + it("tells an executor from a fan-out child and from an ordinary thread", () => { + const lead = ThreadId.make("thread-lead"); + expect(isPairExecutorThreadId(pairExecutorThreadId(lead))).toBe(true); + expect(isPairExecutorThreadId(ThreadId.make(serverChildId(lead, "review-auth")))).toBe(false); + expect(isPairExecutorThreadId(lead)).toBe(false); + expect(isPairExecutorThreadId(ThreadId.make("delegated:broken"))).toBe(false); + }); +}); diff --git a/packages/shared/src/delegatedThreads.ts b/packages/shared/src/delegatedThreads.ts index 128bc0592a..aeffae36ef 100644 --- a/packages/shared/src/delegatedThreads.ts +++ b/packages/shared/src/delegatedThreads.ts @@ -1,3 +1,4 @@ +import { sha256 } from "@noble/hashes/sha2"; import { ThreadId } from "@t3tools/contracts"; const DELEGATED_THREAD_ID_PREFIX = "delegated:"; @@ -15,3 +16,25 @@ export function delegatedParentThreadId(threadId: ThreadId): ThreadId | null { if (suffix === null || suffix.index === 0) return null; return ThreadId.make(rest.slice(0, suffix.index)); } + +/** The delegation key reserved for a thread's pair executor. The server's fan-out tools refuse it. */ +export const PAIR_DELEGATION_KEY = "pair"; + +const hex = (bytes: Uint8Array) => + Array.from(bytes, (byte) => byte.toString(16).padStart(2, "0")).join(""); + +/** + * The id of a lead's pair executor. It is the delegated child id for the + * reserved key, so every client and the server derive it from the lead's id + * and nothing has to record the link. Synchronous so it can run in a render. + */ +export function pairExecutorThreadId(leadThreadId: ThreadId): ThreadId { + const digest = sha256(new TextEncoder().encode(`${leadThreadId}\n${PAIR_DELEGATION_KEY}`)); + return ThreadId.make(`${DELEGATED_THREAD_ID_PREFIX}${leadThreadId}:${hex(digest).slice(0, 16)}`); +} + +/** Whether a thread is some lead's pair executor, rather than a fan-out child. */ +export function isPairExecutorThreadId(threadId: ThreadId): boolean { + const lead = delegatedParentThreadId(threadId); + return lead !== null && pairExecutorThreadId(lead) === threadId; +}