Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
73 changes: 73 additions & 0 deletions tests/web-server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -726,6 +726,79 @@ function waitForSemanticHistory(url: string, sessionId: string): Promise<unknown
});
}

test("TUI metadata changes update every connected web catalog", async () => {
tempDir = await mkdtemp(join(tmpdir(), "pi-kit-live-tui-metadata-test-"));
const agentDir = join(tempDir, "pi-agent");
const sessionsDir = join(agentDir, "sessions", "project");
const statePath = join(tempDir, "web", "server.json");
const selectedId = `selected-${crypto.randomUUID()}`;
const tuiId = `tui-${crypto.randomUUID()}`;
await mkdir(sessionsDir, { recursive: true });
await writeFile(join(sessionsDir, `${selectedId}.jsonl`), `${JSON.stringify({ type: "session", version: 3, id: selectedId, cwd: tempDir, timestamp: new Date().toISOString() })}\n`);
child = Bun.spawn({
cmd: ["bun", "run", "web/server/index.ts"], cwd: process.cwd(),
env: { ...process.env, PI_WEB_PORT: "0", PI_WEB_ROOT: process.cwd(), PI_WEB_STATE_FILE: statePath, PI_CODING_AGENT_DIR: agentDir },
stdout: "ignore", stderr: "ignore",
});
const { port } = await waitForState(statePath);
const client = browserSocket(`ws://127.0.0.1:${port}/ws/client`);
const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent`);
const agentOpened = new Promise<void>((resolve, reject) => {
agent.onopen = () => resolve();
agent.onerror = () => reject(new Error("TUI metadata agent socket failed"));
});
const tuiSession = {
id: tuiId,
cwd: tempDir,
name: "Before rename",
status: "idle" as const,
source: "tui" as const,
createdAt: Date.now(),
updatedAt: Date.now(),
messageCount: 0,
};
let resolveSubscribed!: () => void;
let resolveRegistered!: () => void;
let resolveRenamed!: () => void;
let resolveBranch!: () => void;
const subscribed = new Promise<void>((resolve) => { resolveSubscribed = resolve; });
const registered = new Promise<void>((resolve) => { resolveRegistered = resolve; });
const renamed = new Promise<void>((resolve) => { resolveRenamed = resolve; });
const branchUpdated = new Promise<void>((resolve) => { resolveBranch = resolve; });
let timeout: ReturnType<typeof setTimeout>;
const timedOut = new Promise<never>((_resolve, reject) => {
timeout = setTimeout(() => {
client.close();
agent.close();
reject(new Error("TUI metadata did not propagate to the web catalog"));
}, 5_000);
});
client.onopen = () => client.send(JSON.stringify({ type: "client.hello" }));
client.onmessage = ({ data }) => {
const message = JSON.parse(String(data)) as { type?: string; sessionId?: string; session?: { id?: string; name?: string; branch?: string } };
if (message.type === "server.snapshot") client.send(JSON.stringify({ type: "client.subscribe", sessionId: selectedId }));
if (message.type === "server.history" && message.sessionId === selectedId) resolveSubscribed();
if (message.type !== "server.session" || message.session?.id !== tuiId) return;
if (message.session.name === "Before rename") resolveRegistered();
if (message.session.name === "Renamed in TUI") resolveRenamed();
if (message.session.branch === "feature/live-metadata") resolveBranch();
};
try {
await Promise.race([subscribed, timedOut]);
await Promise.race([agentOpened, timedOut]);
agent.send(JSON.stringify({ type: "agent.hello", session: tuiSession, entries: [] }));
await Promise.race([registered, timedOut]);
agent.send(JSON.stringify({ type: "agent.event", sessionId: tuiId, event: { type: "session_info_changed", name: "Renamed in TUI" } }));
await Promise.race([renamed, timedOut]);
agent.send(JSON.stringify({ type: "agent.update", session: { ...tuiSession, name: "Renamed in TUI", branch: "feature/live-metadata", updatedAt: Date.now() } }));
await Promise.race([branchUpdated, timedOut]);
} finally {
clearTimeout(timeout);
client.close();
agent.close();
}
}, 10_000);

test("an idle native session flushes its restored web follow-up queue on hello", async () => {
tempDir = await mkdtemp(join(tmpdir(), "pi-kit-restored-queue-test-"));
const statePath = join(tempDir, "web", "server.json");
Expand Down
58 changes: 46 additions & 12 deletions web/server/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -763,6 +763,30 @@ function broadcastToAll(message: ServerToClientMessage): void {
}
}

function broadcastSessionToAll(record: SessionRecord): void {
if (sessions.get(record.id) !== record) return;
broadcastToAll({ type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
}

function catalogSessionChanged(previous: WebSession | undefined, next: WebSession): boolean {
if (!previous) return true;
return previous.file !== next.file
|| previous.cwd !== next.cwd
|| previous.name !== next.name
|| previous.branch !== next.branch
|| previous.model !== next.model
|| previous.thinkingLevel !== next.thinkingLevel
|| previous.status !== next.status
|| previous.source !== next.source
|| previous.messageCount !== next.messageCount
|| previous.preview !== next.preview
Comment thread
ianwalter marked this conversation as resolved.
|| previous.parentSession !== next.parentSession
|| previous.pullRequest?.number !== next.pullRequest?.number
|| previous.pullRequest?.url !== next.pullRequest?.url
|| previous.compaction?.reason !== next.compaction?.reason
|| previous.compaction?.startedAt !== next.compaction?.startedAt;
}

function sessionSnapshot(): WebSession[] {
void reconcileMissingSessionFiles();
const merged = new Map<string, WebSession>();
Expand Down Expand Up @@ -891,7 +915,7 @@ async function hydrateGitMetadata(record: SessionRecord): Promise<void> {
// A branch without an open PR is expected.
}
}
broadcast(record.id, { type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
broadcastSessionToAll(record);
Comment thread
ianwalter marked this conversation as resolved.
}

class CommandRejectedError extends Error {
Expand Down Expand Up @@ -2002,11 +2026,7 @@ async function refreshManagedSession(
record.pendingWorktreeSourceDeletion = undefined;
}
identityTransition = undefined;
const message = JSON.stringify({
type: "server.session",
session: sessionToClientPayload(record),
} satisfies ServerSessionMessage);
for (const socket of record.clientSockets) socket.send(message);
broadcastSessionToAll(record);
} catch (error) {
finishQueueMigration?.();
if (identityTransition) {
Expand Down Expand Up @@ -2246,13 +2266,21 @@ async function createManagedSessionUnlocked(cwd: string, name?: string, sessionF
sessionId: record.id,
event,
} satisfies ServerEventMessage);
broadcast(record.id, { type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
const catalogChanged = event.type === "agent_start"
|| event.type === "turn_start"
|| event.type === "agent_end"
|| event.type === "agent_settled"
|| event.type === "compaction_start"
|| event.type === "compaction_end"
|| event.type === "message_end";
if (catalogChanged) broadcastSessionToAll(record);
else broadcast(record.id, { type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
},
onExit: () => {
record.status = "offline";
record.active = false;
record.managed = undefined;
broadcast(record.id, { type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
broadcastSessionToAll(record);
},
});
record.managed = managed;
Expand Down Expand Up @@ -2885,7 +2913,7 @@ async function handleAgentMessage(socket: Bun.ServerWebSocket<AgentSocketData>,
);
}
if (record.file) sessionsByFile.set(normalizePath(record.file), record);
broadcast(record.id, { type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
broadcastSessionToAll(record);
Comment thread
ianwalter marked this conversation as resolved.
// Native sessions use the same bounded semantic history as managed sessions;
// no browser connection asks Pi to paint an additional TUI viewport.
broadcast(record.id, { type: "server.history", sessionId: record.id, entries: record.history.slice(-600) } satisfies ServerHistoryMessage);
Expand Down Expand Up @@ -2948,6 +2976,10 @@ async function handleAgentMessage(socket: Bun.ServerWebSocket<AgentSocketData>,
}
const subagentsChanged = updateSubagentsFromToolEvent(record, event.event);
let sessionMetadataChanged = false;
if (event.event.type === "session_info_changed") {
record.name = typeof event.event.name === "string" && event.event.name ? event.event.name : undefined;
sessionMetadataChanged = true;
}
if (event.event.type === "message_end" && isRecord(event.event.message)) {
record.history.push({
type: "message",
Expand All @@ -2970,7 +3002,7 @@ async function handleAgentMessage(socket: Bun.ServerWebSocket<AgentSocketData>,
}
broadcast(event.sessionId, { type: "server.event", sessionId: event.sessionId, event: event.event } satisfies ServerEventMessage);
if (sessionMetadataChanged) {
broadcastToAll({ type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
broadcastSessionToAll(record);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
} else if (lifecycleChanged || subagentsChanged || event.event.type === "compaction_start" || event.event.type === "compaction_end") {
broadcast(record.id, { type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
}
Expand Down Expand Up @@ -3001,10 +3033,12 @@ async function handleAgentMessage(socket: Bun.ServerWebSocket<AgentSocketData>,
const session = existing?.agentRunning === false && update.session.status === "working"
? { ...update.session, status: "idle" as const }
: update.session;
const catalogChanged = catalogSessionChanged(existing, session);
const record = upsertSession(session, "external", existing?.history ?? []);
record.agentSockets.add(socket);
record.updatedAt = update.session.updatedAt;
broadcast(update.session.id, { type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
if (catalogChanged) broadcastSessionToAll(record);
else broadcast(update.session.id, { type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
Comment thread
ianwalter marked this conversation as resolved.
return;
}
if (message.type === "agent.response") {
Expand Down Expand Up @@ -3513,7 +3547,7 @@ function handleWebSocketClose(socket: Bun.ServerWebSocket<SocketData>): void {
if (record.agentSockets.size === 0 && record.kind === "external") {
record.status = "offline";
record.active = false;
broadcast(record.id, { type: "server.session", session: sessionToClientPayload(record) } satisfies ServerSessionMessage);
broadcastSessionToAll(record);
}
}
}
Expand Down