Skip to content
Merged
11 changes: 11 additions & 0 deletions apps/cli/src/app.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import { HintPill } from "./components/HintPill";
import { InputBox } from "./components/InputBox";
import { MessageView } from "./components/MessageView";
import { Notification } from "./components/Notification";
import { TaskPanel } from "./components/TaskPanel";
import { WelcomeBanner } from "./components/WelcomeBanner";
import {
useCurrentThreadHasChild,
Expand All @@ -36,6 +37,7 @@ import {
useSetPendingRequestId,
useSetRun,
useSetStatus,
useSetTasks,
useStatus
} from "./hooks/use-app";
import { dispatchKey } from "./keymap/dispatch";
Expand Down Expand Up @@ -141,6 +143,7 @@ export function App({
const setIsCompacting = useSetIsCompacting();
const setCurrentThreadHasChild = useSetCurrentThreadHasChild();
const setPendingConflict = useSetPendingConflict();
const setTasks = useSetTasks();

// Per-slice subscription so App only re-renders when the conflict frame
// appears/disappears, not on every other store update.
Expand Down Expand Up @@ -186,6 +189,7 @@ export function App({
setNotice,
setIsCompacting,
setPendingConflict,
setTasks,
send: (raw) => agentRef.current?.send(raw),
switchThread
});
Expand All @@ -198,6 +202,7 @@ export function App({
setNotice,
setIsCompacting,
setPendingConflict,
setTasks,
switchThread
]
)
Expand Down Expand Up @@ -240,6 +245,8 @@ export function App({
setPendingRequestId(null);
setFocusedToolCallId(null);
setIsCompacting(false);
// Clear before reconnect; kernel's onConnect will re-broadcast.
setTasks([]);
// Optimistic: assume the freshly-selected thread is live. The /threads
// fetch below corrects this if the user is actually viewing an archive.
setCurrentThreadHasChild(false);
Expand All @@ -250,6 +257,7 @@ export function App({
setPendingRequestId,
setFocusedToolCallId,
setIsCompacting,
setTasks,
setCurrentThreadHasChild
]);

Expand Down Expand Up @@ -480,6 +488,9 @@ export function App({
onClose={() => setActivePanel(null)}
/>
) : null}

<TaskPanel />

{pendingConflict ? (
<ConflictResolver frame={pendingConflict} />
) : inApprovalMode ? (
Expand Down
44 changes: 44 additions & 0 deletions apps/cli/src/components/TaskPanel.tsx
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
import { Box, Text } from "ink";
import { memo } from "react";

import type { TaskRow } from "@agent-os/protocol";

import { useStore } from "../store";

// Status-to-glyph mapping. Keep in sync with packages/models task.status enum.
// `?` is the defensive fallback if a future status enum value reaches the TUI
// before this map is updated.
const STATUS_GLYPH: Record<TaskRow["status"], { icon: string; color: string }> =
{
pending: { icon: "☐", color: "gray" },
in_progress: { icon: "▶", color: "yellow" },
complete: { icon: "✓", color: "green" },
failed: { icon: "✗", color: "red" },
cancelled: { icon: "⊘", color: "gray" }
};

function TaskPanelImpl() {
// Selector subscribes only this component to `tasks` — App stays silent on
// task changes. Matches the LiveRow optimization pattern: per-field
// subscription prevents per-frame re-renders of the wider tree.
const tasks = useStore((s) => s.tasks);

if (tasks.length === 0) return null;

return (
<Box flexDirection="column" marginTop={1} paddingX={1}>
<Text dimColor>Tasks ({tasks.length})</Text>
{tasks.map((t) => {
const g = STATUS_GLYPH[t.status] ?? { icon: "?", color: "gray" };
return (
<Box key={t.id}>
<Text color={g.color}>{g.icon}</Text>
<Text> {t.name}</Text>
</Box>
);
})}
</Box>
);
}

export const TaskPanel = memo(TaskPanelImpl);
98 changes: 98 additions & 0 deletions apps/cli/src/components/tools/SpawnSubAgentTool.tsx
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
import { Box, Text } from "ink";
import { memo } from "react";

import type { ToolRendererProps } from "./registry";

interface SpawnSubAgentInput {
model: "opus" | "sonnet" | "haiku";
prompt: string;
task_id?: string;
}

interface SpawnSubAgentOutput {
runId: string;
finalText: string;
}

const MODEL_COLOR: Record<SpawnSubAgentInput["model"], string> = {
opus: "magenta",
sonnet: "blue",
haiku: "green"
};

function SpawnSubAgentToolImpl({ part, focused }: ToolRendererProps) {
const input = part.input as Partial<SpawnSubAgentInput> | undefined;
const output = part.output as SpawnSubAgentOutput | undefined;
const model: SpawnSubAgentInput["model"] = input?.model ?? "opus";
const modelColor = MODEL_COLOR[model];
const prompt = input?.prompt ?? "";
const taskId = input?.task_id;

const isRunning = part.state === "input-available";

// Failure state takes precedence over `state` — errorText is set on any
// terminal-failure path (abort, throw, model rejection). Also fold
// `output-error` in here: format.ts treats it as equivalent, and without
// this check an `output-error` with empty errorText would fall through
// to `null`.
if (part.errorText || part.state === "output-error") {
return (
<Box flexDirection="column">
<Text>
<Text color="red">✗</Text> Sub-agent (
<Text color={modelColor}>{model}</Text>) · failed
</Text>
<Text color="red"> {part.errorText ?? "Sub-agent failed"}</Text>
</Box>
);
}

if (part.state === "input-streaming") {
return (
<Text>
<Text color="yellow">⏳</Text> Sub-agent (
<Text color={modelColor}>{model}</Text>) …
</Text>
);
}

if (isRunning) {
return (
<Box flexDirection="column">
<Text>
<Text color="yellow">🤖</Text> Sub-agent (
<Text color={modelColor}>{model}</Text>) · running
</Text>
{prompt ? <Text dimColor> &quot;{prompt}&quot;</Text> : null}
</Box>
);
}

// state === "output-available"
if (!output) return null;
// Approximation: ~4 chars per token. Marked with `~` in the UI.
const approxTok = Math.ceil(output.finalText.length / 4);

return (
<Box flexDirection="column">
<Text>
<Text color="green">✓</Text> Sub-agent (
<Text color={modelColor}>{model}</Text>) · ~{approxTok} tok
</Text>
{prompt ? <Text dimColor> &quot;{prompt}&quot;</Text> : null}
{taskId ? <Text dimColor> → task: {taskId}</Text> : null}
{focused ? (
<Box
flexDirection="column"
borderStyle="round"
marginTop={1}
paddingX={1}
>
<Text>{output.finalText}</Text>
</Box>
) : null}
</Box>
);
}

export const SpawnSubAgentTool = memo(SpawnSubAgentToolImpl);
6 changes: 4 additions & 2 deletions apps/cli/src/components/tools/registry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,11 @@ import { FindFilesTool } from "./FindFilesTool";
import { GenericTool } from "./GenericTool";
import { ListFilesTool } from "./ListFilesTool";
import { ReadTool } from "./ReadTool";
import { SpawnSubAgentTool } from "./SpawnSubAgentTool";
import type { ToolPart } from "./types";
import { WebFetchTool } from "./WebFetchTool";
import { WebSearchTool } from "./WebSearchTool";
import { WriteTool } from "./WriteTool";
import type { ToolPart } from "./types";

export type { ToolPart };

Expand All @@ -29,7 +30,8 @@ export const renderers: Record<string, FC<ToolRendererProps>> = {
web_search: WebSearchTool,
web_fetch: WebFetchTool,
computer_bash: BashTool,
exec_code: ExecCodeTool
exec_code: ExecCodeTool,
spawn_sub_agent: SpawnSubAgentTool
};

export function pickRenderer(toolName: string): FC<ToolRendererProps> {
Expand Down
1 change: 1 addition & 0 deletions apps/cli/src/hooks/use-app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,3 +36,4 @@ export const useSetCurrentThreadHasChild = () =>
export const useSetPendingConflict = () =>
useStore((s) => s.setPendingConflict);
export const useSetSend = () => useStore((s) => s.setSend);
export const useSetTasks = () => useStore((s) => s.setTasks);
9 changes: 9 additions & 0 deletions apps/cli/src/protocol/handle-frame.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,15 @@ import {
COMPACTION_STATUS,
MERGE_CONFLICT,
RUN_STATUS_CHANGED,
TASKS_UPDATED,
THREAD_COMPACTED,
THREAD_FROZEN,
USER_MESSAGE_REJECTED,
type CompactRejectedFrame,
type CompactionStatusFrame,
type MergeConflictFrame,
type RunStatusChangedFrame,
type TasksUpdatedFrame,
type ThreadCompactedFrame,
type ThreadFrozenFrame,
type UserMessageRejectedFrame
Expand Down Expand Up @@ -40,6 +42,7 @@ export type ProtocolActions = Pick<
| "setNotice"
| "setIsCompacting"
| "setPendingConflict"
| "setTasks"
> & {
applyStreamChunk: (chunk: { type: string; [k: string]: unknown }) => void;
send: (raw: string) => void;
Expand Down Expand Up @@ -158,6 +161,12 @@ export function handleFrame(raw: unknown, actions: ProtocolActions): void {
return;
}

case TASKS_UPDATED: {
const f = frame as TasksUpdatedFrame;
actions.setTasks(f.tasks);
return;
}

case THREAD_COMPACTED: {
const f = frame as unknown as ThreadCompactedFrame;
actions.setIsCompacting(false);
Expand Down
12 changes: 11 additions & 1 deletion apps/cli/src/store/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,8 @@ import {
type Attachment,
type MergeConflictFrame,
type MergeConflictResolutionFrame,
type RunStatus
type RunStatus,
type TaskRow
} from "@agent-os/protocol";
import type { ApprovalEntry } from "../keymap/dispatch";

Expand Down Expand Up @@ -186,6 +187,12 @@ type Store = {
currentThreadHasChild: boolean;
setCurrentThreadHasChild: (v: boolean) => void;

// Per-thread task rows from the kernel. Pushed via TASKS_UPDATED frames
// whenever a tool or the kernel itself mutates the thread's task list.
// Consumed by TaskPanel (next task).
tasks: TaskRow[];
setTasks: (tasks: TaskRow[]) => void;

// Outbound WebSocket send. Wired up by App once the agent connection is
// live, so store actions (sendConflictResolution / sendConflictAbort) can
// emit frames without prop-drilling the agent down to every consumer.
Expand Down Expand Up @@ -431,6 +438,9 @@ export const useStore = create<Store>((set, get) => ({
currentThreadHasChild: false,
setCurrentThreadHasChild: (v) => set({ currentThreadHasChild: v }),

tasks: [],
setTasks: (tasks) => set({ tasks }),

// Default no-op until App wires the real agent.send. Calling before wiring
// is a programming error — the conflict UI only opens after a frame has
// arrived, by which point the socket is already up.
Expand Down
41 changes: 39 additions & 2 deletions apps/kernel/src/broadcasts.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,16 @@
import { type Connection } from "agents";
import { type OutgoingMessage } from "@cloudflare/ai-chat/types";
import { type Connection } from "agents";

import { asc, eq, schema } from "@agent-os/models";
import {
RUN_STATUS_CHANGED,
TASKS_UPDATED,
type CompactRejectedFrame,
type CompactionStatusFrame,
type RunStatus,
type RunStatusChangedFrame,
type TaskRow,
type TasksUpdatedFrame,
type ThreadCompactedFrame,
type ThreadFrozenFrame,
type UserMessageRejectedFrame
Expand All @@ -21,7 +25,8 @@ type Outbound =
| ThreadCompactedFrame
| ThreadFrozenFrame
| CompactRejectedFrame
| UserMessageRejectedFrame;
| UserMessageRejectedFrame
| TasksUpdatedFrame;

// Send a single frame to one connection. Wrapper exists for the typed
// `msg` union — without it, ad-hoc `JSON.stringify` calls drift in shape.
Expand Down Expand Up @@ -63,3 +68,35 @@ export function broadcastRunStatus(
...(messageId ? { messageId } : {})
});
}

// Query the thread's current task list and fan it out to every connection
// watching the thread. Called from three sites:
// - kernel.onConnect (initial state seed)
// - tools/tasks.ts execute() (after manage_tasks mutations)
// - tools/spawn-sub-agent.ts (after writing task.result)
//
// Full-snapshot rather than delta: simpler client, ordered WS makes it
// safe, and the row count per thread is always tiny.
export async function broadcastTasksUpdated(
kernel: Kernel,
threadId: string
): Promise<void> {
const rows = await kernel.db.query.task.findMany({
where: eq(schema.task.threadId, threadId),
orderBy: [asc(schema.task.createdAt)]
});
const tasks: TaskRow[] = rows.map((t) => ({
id: t.id,
name: t.name,
description: t.description,
status: t.status,
result: t.result,
createdAt: t.createdAt.toISOString(),
updatedAt: t.updatedAt.toISOString()
}));
broadcastToThread(kernel, threadId, {
type: TASKS_UPDATED,
threadId,
tasks
});
}
Loading
Loading