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
7 changes: 7 additions & 0 deletions CONTEXT.md
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,9 @@ The counter distinguishing multiple ephemeral deltas that share the same anchor
**Chat chunk**:
A live conversation event delivered to watching peers as a turn runs.

**Shell execution**:
The conversation entry persisted for a **Shell input** run: command, stream-tagged lines (stdout/stderr), exit code, and start/complete times. A dedicated entry kind alongside message/thought/tool/error, not a synthetic tool call. Identified by its own execution id — distinct from a **Turn** id — since a shell execution runs independently of any turn and is cancelled independently (its own Stop control, not the composer's turn-scoped Stop button).

**Ephemeral delta**:
A streamed fragment (token or thought), ordered by its anchor seq plus sub — broadcast live but never persisted, and never held in the replay buffer; the completed message is persisted once, in full.
_Avoid_: partial message
Expand Down Expand Up @@ -151,6 +154,10 @@ The health check that spawns an enabled agent and verifies the ACP handshake com
**Composer**:
The chat input area: prompt box, agent/model picker, mode/effort/persona controls, prompt queue, and the panel where approvals and elicitations are answered.

**Shell input**:
Composer state entered when `!` is the first character typed in a thread's composer; the text is sent as a shell command instead of a prompt. Thread-scoped only — requires an existing thread. Distinct from an agent's **mode** (`modeId`).
_Avoid_: shell mode (collides with modeId)

**Prompt queue**:
Messages queued per thread while a turn is active, sent sequentially once the turn completes.

Expand Down
2 changes: 2 additions & 0 deletions apps/cli/src/handlers/controller/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import { fsHandlers } from "./fs";
import { gitHandlers } from "./git";
import { interactiveHandlers } from "./interactive";
import { projectsHandlers } from "./projects";
import { shellHandlers } from "./shell";
import { threadsHandlers } from "./threads";

export function createControllerRouter(runtime: WorkerRuntime) {
Expand All @@ -29,6 +30,7 @@ export function createControllerRouter(runtime: WorkerRuntime) {
...personaHandlers(deps),
...contextHandlers(deps),
...chatHandlers(deps),
...shellHandlers(deps),
...interactiveHandlers(deps),
};
}
49 changes: 49 additions & 0 deletions apps/cli/src/handlers/controller/shell.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
import { appendConversation } from "@cyrus/database/repositories/conversations";
import { resolveThreadGitCwd } from "@cyrus/database/repositories/git";
import { orpcOk, throwOrpc } from "@cyrus/errors/orpc";
import { randomId } from "@cyrus/utils/identity";
import { log } from "evlog";
import { cancelActiveShellExecution, runShellExecution } from "@/shell/run";
import type { ControllerDeps } from "./deps";

export function shellHandlers({ os }: ControllerDeps) {
Comment thread
soorya-u marked this conversation as resolved.
return {
executeShellInput: os.executeShellInput.handler(
async ({ input, context }) => {
const { threadId, command } = input;
const cwd = orpcOk(await resolveThreadGitCwd(threadId));
const shellExecutionId = randomId();

const started = await appendConversation(threadId, {
threadId,
shellExecutionId,
event: { type: "shell_execution_start", command },
});
if (started.isErr()) throwOrpc(started.error);
context.eventBus.publish(started.value.chunk);

runShellExecution({
command,
context,
cwd,
shellExecutionId,
threadId,
}).catch((error) => {
log.error({
kind: "shell_execution_failed",
error,
shellExecutionId,
threadId,
});
});

return { shellExecutionId };
}
),

cancelShellExecution: os.cancelShellExecution.handler(({ input }) => {
cancelActiveShellExecution(input.threadId, input.shellExecutionId);
return {};
}),
};
}
5 changes: 5 additions & 0 deletions apps/cli/src/lib/env.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,11 @@ export const env = createEnv({
.int()
.positive()
.default(30 * 60 * 1000),
CYRUS_SHELL_INPUT_TIMEOUT_MS: z.coerce
.number()
.int()
.positive()
.default(60_000),
},
runtimeEnv: process.env,
emptyStringAsUndefined: true,
Expand Down
45 changes: 45 additions & 0 deletions apps/cli/src/queue/bus.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,18 @@ function tokenChunk(turnId = "turn-1"): ChatChunk {
};
}

function shellPersistedChunk(
seq: number,
shellExecutionId = "shell-1"
): ChatChunk {
return {
threadId: "thread-1",
shellExecutionId,
seq,
event: { type: "shell_execution_start", command: "ls" },
};
}

function createReader(iterator: AsyncGenerator<ChatChunk>) {
let pending = iterator.next();
return {
Expand Down Expand Up @@ -116,4 +128,37 @@ describe("thread event bus", () => {
expect(await second.next()).toMatchObject({ seq: 3 });
expect(await second.next()).toMatchObject({ seq: 4 });
});

test("fans out shell execution chunks live, same as turn chunks", async () => {
const bus = createThreadEventBus();
bus.watch("peer-1", "thread-1");
const reader = createReader(bus.subscribe("peer-1"));

bus.publish(persistedChunk(1));
bus.publish(shellPersistedChunk(2));

expect(await reader.next()).toMatchObject({ seq: 1, turnId: "turn-1" });
expect(await reader.next()).toMatchObject({
seq: 2,
shellExecutionId: "shell-1",
});
});

test("does not replay shell execution chunks to a late-joining watcher", async () => {
const bus = createThreadEventBus();
bus.watch("peer-1", "thread-1");
const first = createReader(bus.subscribe("peer-1"));

bus.publish(persistedChunk(1));
bus.publish(shellPersistedChunk(2));
await first.next();
await first.next();

bus.watch("peer-2", "thread-1");
const second = createReader(bus.subscribe("peer-2"));
expect(await second.next()).toMatchObject({ seq: 1 });

bus.publish(tokenChunk());
expect(await second.next()).toMatchObject({ event: { type: "token" } });
});
});
8 changes: 4 additions & 4 deletions apps/cli/src/queue/bus.ts
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,7 @@ export function createThreadEventBus(
while (log.length > maxChunksPerTurn) log.shift();
}

function appendToTurnLog(chunk: ChatChunk): void {
function appendToTurnLog(chunk: ChatChunk & { turnId: string }): void {
const { turnId } = chunk;
let log = activeTurnLogs.get(turnId);
if (!log) {
Expand Down Expand Up @@ -137,11 +137,11 @@ export function createThreadEventBus(
const stamped = stampChunk(chunk);
const terminal = isTerminalEvent(stamped.event);

if (!terminal && stamped.sub === undefined) {
appendToTurnLog(stamped);
if (!terminal && stamped.sub === undefined && stamped.turnId) {
appendToTurnLog(stamped as ChatChunk & { turnId: string });
}
fanOut(stamped);
if (terminal) {
if (terminal && stamped.turnId) {
evictTurnLog(stamped.turnId);
}
},
Expand Down
66 changes: 66 additions & 0 deletions apps/cli/src/shell/run.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
import { describe, expect, test } from "bun:test";
import type { ShellExecutionLine } from "@cyrus/schemas/rtc/chat";
import { pumpLines } from "./run";

function streamFromChunks(chunks: Uint8Array[]): ReadableStream<Uint8Array> {
return new ReadableStream({
start(controller) {
for (const chunk of chunks) controller.enqueue(chunk);
controller.close();
},
});
}

function encode(text: string): Uint8Array {
return new TextEncoder().encode(text);
}

async function collectLines(
stream: ReadableStream<Uint8Array> | null
): Promise<ShellExecutionLine[]> {
const lines: ShellExecutionLine[] = [];
await pumpLines(stream, "stdout", (line) => lines.push(line));
return lines;
}

describe("pumpLines", () => {
test("splits a single chunk into its newline-terminated lines", async () => {
const lines = await collectLines(streamFromChunks([encode("a\nb\nc\n")]));
expect(lines).toEqual([
{ stream: "stdout", text: "a" },
{ stream: "stdout", text: "b" },
{ stream: "stdout", text: "c" },
]);
});

test("joins a line split across two read chunks", async () => {
const lines = await collectLines(
streamFromChunks([encode("hello wor"), encode("ld\n")])
);
expect(lines).toEqual([{ stream: "stdout", text: "hello world" }]);
});

test("emits trailing text with no final newline as its own line", async () => {
const lines = await collectLines(streamFromChunks([encode("no newline")]));
expect(lines).toEqual([{ stream: "stdout", text: "no newline" }]);
});

test("decodes a multi-byte UTF-8 character split across chunk boundaries", async () => {
// "🎉" is 4 bytes in UTF-8 — split it in the middle of the codepoint.
const bytes = encode("party 🎉\n");
const lines = await collectLines(
streamFromChunks([bytes.slice(0, 7), bytes.slice(7)])
);
expect(lines).toEqual([{ stream: "stdout", text: "party 🎉" }]);
});

test("emits nothing for an empty stream", async () => {
const lines = await collectLines(streamFromChunks([]));
expect(lines).toEqual([]);
});

test("resolves immediately for a null stream (stdout/stderr absent)", async () => {
const lines = await collectLines(null);
expect(lines).toEqual([]);
});
});
Loading