Skip to content
Merged
Show file tree
Hide file tree
Changes from 18 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
21 changes: 21 additions & 0 deletions .informant/config.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
version = 1
timeout_minutes = 30
environment = { CI = "true" }
cache = [{ paths = ["~/.bun/install/cache"], shared = true }]
triggers = [
{ event = "commit", branch = { names = ["main"] } },
{ event = "commit", pull_request = { state = "open", draft = false } },
]

[container]
image = "docker.io/oven/bun:1.3.10"
prepareInputs = ["bun.lock", "package.json"]
prepare = """
set -eu
export DEBIAN_FRONTEND=noninteractive
apt-get update
apt-get install -y --no-install-recommends git
rm -rf /var/lib/apt/lists/*
cd "$INFORMANT_PREPARE_ROOT"
bun install --frozen-lockfile --ignore-scripts
"""
9 changes: 9 additions & 0 deletions .informant/jobs/build.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
name = "build"
needs = ["test", "typecheck"]
command = """
set -eu
bun run build
git diff --exit-code -- web/dist
"""
timeout_minutes = 15
container = { cpu = 2, memory_mb = 4096 }
7 changes: 7 additions & 0 deletions .informant/jobs/test.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
name = "test"
command = """
set -eu
bun run test
"""
timeout_minutes = 20
container = { cpu = 2, memory_mb = 4096 }
7 changes: 7 additions & 0 deletions .informant/jobs/typecheck.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
name = "typecheck"
command = """
set -eu
bun run check
"""
timeout_minutes = 15
container = { cpu = 2, memory_mb = 4096 }
17 changes: 9 additions & 8 deletions extensions/subagents.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1064,14 +1064,15 @@ class SubagentManager {
agent.completedAt = Date.now();
this.activity(agent, "terminated and released session resources");
}
if (remove) {
this.agents.delete(id);
this.archivedAgents.delete(id);
this.webTranscriptCursors.delete(id);
this.webStreamingSnapshots.delete(id);
this.footerSelected = false;
this.publishFooter();
}
this.agents.delete(id);
this.webTranscriptCursors.delete(id);
this.webStreamingSnapshots.delete(id);
if (remove) this.archivedAgents.delete(id);
else if (shouldArchiveTerminalSubagent(agent)) this.archivedAgents.set(id, agent);
this.footerSelected = false;
this.publishFooter();
if (this.webStatusPublishTimer) clearTimeout(this.webStatusPublishTimer);
this.publishWebStatus();
return agent;
}

Expand Down
17 changes: 17 additions & 0 deletions extensions/terminal-output.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
import { isBashToolResult, type ExtensionAPI } from "@earendil-works/pi-coding-agent";
import { renderTerminalOutput } from "../terminal-output.js";

export default function terminalOutputExtension(pi: ExtensionAPI) {
pi.on("tool_result", (event) => {
if (!isBashToolResult(event)) return;
let changed = false;
const content = event.content.map((part) => {
if (part.type !== "text") return part;
const text = renderTerminalOutput(part.text);
if (text === part.text) return part;
changed = true;
return { ...part, text };
});
return changed ? { content } : undefined;
});
}
130 changes: 98 additions & 32 deletions extensions/web-sessions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import type {
AgentCommand,
AgentEventMessage,
AgentHelloMessage,
AgentHistoryMessage,
AgentResponseMessage,
AgentSessionReplacedMessage,
AgentSubagentsMessage,
Expand All @@ -32,9 +33,13 @@ import type {
ServerToAgentMessage,
WebSession,
} from "../web/protocol.js";
import { WEB_STATE_VERSION } from "../web/protocol.js";
import { mergeWebSubagentUpdates, WEB_STATE_VERSION } from "../web/protocol.js";
import { expandSlashCommand, isSkillSlashCommand } from "../web/slash-commands.js";
import { formatWorktreeCreateCommandArgs } from "../web/worktree-command.js";
import { WEB_COMPACT_COMMAND } from "../web/compact-command.js";
import { boundedWebHistory } from "../web/history.js";
import { agentEndTerminalNotice } from "../web/assistant-message.js";
import { managedWorktreeFromEntries } from "../web/server/worktrees.js";
import { readWebTailscaleSetting, writeWebTailscaleSetting } from "./web-settings.js";
import {
consumeWorktreeReplacement,
Expand Down Expand Up @@ -77,6 +82,14 @@ function bridgeCommandList(pi: ExtensionAPI) {
if (!commands.some((command) => command.name === "reload")) {
commands.unshift({ name: "reload", description: "Reload extensions, skills, prompts, themes, and context files", source: "extension", location: "temporary" });
}
if (!commands.some((command) => command.name === "compact")) {
commands.unshift({
name: WEB_COMPACT_COMMAND.name,
description: WEB_COMPACT_COMMAND.description,
source: "extension",
location: "temporary",
});
}
return commands;
}

Expand All @@ -89,7 +102,7 @@ export function splitWebWorktreeCommandArgs(args: string): { token: string; work
}

/** Abort the main session and wait for subagent abort operations registered through waitUntil. */
export async function abortSessionAndSubagents(options: {
export function abortSessionAndSubagents(options: {
sessionId: string;
abortMain(): void;
emit(request: SubagentAbortRequest): void;
Expand All @@ -107,7 +120,7 @@ export async function abortSessionAndSubagents(options: {
// A broken optional listener must never prevent the main Stop request.
}
options.abortMain();
await Promise.allSettled(operations);
return Promise.allSettled(operations).then(() => undefined);
}

/** Apply a route change and roll it back if the matching settings write fails. */
Expand Down Expand Up @@ -287,6 +300,11 @@ function addWebUsage(target: NonNullable<WebSession["usage"]>, value: unknown):
}
}

function contextUsage(ctx: ExtensionContext): WebSession["contextUsage"] {
const usage = ctx.getContextUsage();
return usage ? { ...usage } : ctx.model ? { tokens: null, contextWindow: ctx.model.contextWindow, percent: null } : undefined;
}

function sessionMetrics(ctx: ExtensionContext): Pick<WebSession, "usage" | "contextUsage"> {
const usage = zeroWebUsage();
for (const entry of ctx.sessionManager.getEntries()) {
Expand All @@ -296,11 +314,14 @@ function sessionMetrics(ctx: ExtensionContext): Pick<WebSession, "usage" | "cont
addWebUsage(usage, entry.usage);
}
}
const contextUsage = ctx.getContextUsage();
return {
usage,
contextUsage: contextUsage ? { ...contextUsage } : ctx.model ? { tokens: null, contextWindow: ctx.model.contextWindow, percent: null } : undefined,
};
return { usage, contextUsage: contextUsage(ctx) };
}

function refreshIncrementalMetrics(state: BridgeState, usageValue: unknown): void {
const current = state.metrics.usage ?? zeroWebUsage();
const usage = { ...current, cost: { ...current.cost } };
addWebUsage(usage, usageValue);
state.metrics = { usage, contextUsage: contextUsage(state.ctx) };
}

function safeClone(value: unknown): Record<string, unknown> {
Expand Down Expand Up @@ -364,6 +385,12 @@ function flush(state: BridgeState): void {
for (const message of state.pending.splice(0)) state.socket.send(JSON.stringify(message));
}

function discardPendingCoveredByHello(state: BridgeState): void {
// Hello includes authoritative history and the fully merged subagent snapshot.
// Replaying an old delta after it would duplicate transcript/streaming content.
state.pending = state.pending.filter((message) => message.type === "agent.response");
}

function statusForContext(ctx: ExtensionContext): WebSession["status"] {
return ctx.isIdle() ? "idle" : "working";
}
Expand Down Expand Up @@ -394,8 +421,15 @@ async function refreshGitMetadata(pi: ExtensionAPI, state: BridgeState): Promise
if (!state.closed) updateSession(state, { branch, pullRequest });
}

function updateSession(state: BridgeState, patch: Partial<WebSession> = {}, refreshMetrics = false): void {
if (refreshMetrics) state.metrics = sessionMetrics(state.ctx);
export function applySubagentStatusToSession(session: WebSession, event: SubagentStatusEvent): WebSession {
return {
...session,
subagents: event.remove ? [] : mergeWebSubagentUpdates(session.subagents, event.agents),
subagentUsage: event.usage,
};
}

function updateSession(state: BridgeState, patch: Partial<WebSession> = {}): void {
state.session = {
...state.session,
...state.metrics,
Expand Down Expand Up @@ -509,14 +543,18 @@ async function executeAgentCommand(
respond(state, requestId, true);
return;
}
case "abort":
await abortSessionAndSubagents({
case "abort": {
// Invoke the main abort before acknowledging, then let subagent teardown
// settle in the background. Compaction can delay that settlement well
// past the browser's command bound even though Stop has taken effect.
void abortSessionAndSubagents({
sessionId: state.session.id,
abortMain: () => state.ctx.abort(),
emit: (request) => pi.events.emit(SUBAGENT_ABORT_EVENT, request),
});
respond(state, requestId, true);
respond(state, requestId, true, { accepted: true });
return;
}
case "replace_queue":
respond(state, requestId, true);
return;
Expand Down Expand Up @@ -730,13 +768,14 @@ async function connect(pi: ExtensionAPI, state: BridgeState): Promise<void> {
return;
}
state.reconnectAttempt = 0;
discardPendingCoveredByHello(state);
const hello: AgentHelloMessage = {
type: "agent.hello",
session: state.session,
// The server can read JSONL history from session.file. Sending every
// entry here makes large sessions exceed WebSocket frame limits and
// prevents the native bridge from registering after a daemon restart.
entries: [],
historyMode: "replace",
// Send only active, compaction-aware history and bound its encoded size.
// The append-only JSONL can be hundreds of MB after old context is gone.
entries: boundedWebHistory(state.ctx.sessionManager.buildContextEntries()),
};
socket.send(JSON.stringify(hello));
if (state.sourceReplacement) {
Expand All @@ -755,15 +794,26 @@ async function connect(pi: ExtensionAPI, state: BridgeState): Promise<void> {
};
}

function previewFromMessage(message: unknown): string | undefined {
if (!isRecord(message) || (message.role !== "user" && message.role !== "assistant")) return undefined;
if (typeof message.content === "string") return message.content.slice(0, 180);
if (!Array.isArray(message.content)) return undefined;
const preview = message.content
.filter((item): item is Record<string, unknown> => isRecord(item) && item.type === "text" && typeof item.text === "string")
.map((item) => item.text as string)
.join("");
return preview ? preview.slice(0, 180) : undefined;
}

function makeSession(ctx: ExtensionContext, branch: string | undefined): WebSession {
const entries = ctx.sessionManager.getEntries();
const header = ctx.sessionManager.getHeader();
const firstUser = entries.find((entry) => entry.type === "message" && entry.message.role === "user");
let preview: string | undefined;
if (firstUser?.type === "message" && firstUser.message.role === "user") {
preview = typeof firstUser.message.content === "string"
? firstUser.message.content.slice(0, 180)
: firstUser.message.content.filter((item) => item.type === "text").map((item) => item.text).join("").slice(0, 180);
for (let index = entries.length - 1; index >= 0; index -= 1) {
const entry = entries[index];
if (entry?.type !== "message" || (entry.message.role !== "user" && entry.message.role !== "assistant")) continue;
preview = previewFromMessage(entry.message);
if (preview) break;
}
return {
id: ctx.sessionManager.getSessionId(),
Expand All @@ -780,6 +830,7 @@ function makeSession(ctx: ExtensionContext, branch: string | undefined): WebSess
messageCount: entries.filter((entry) => entry.type === "message").length,
preview,
parentSession: header?.parentSession,
managedWorktree: managedWorktreeFromEntries(entries),
...sessionMetrics(ctx),
};
}
Expand All @@ -806,6 +857,9 @@ export default function webSessions(pi: ExtensionAPI): void {
pi.events.on(SUBAGENT_STATUS_EVENT, (value) => {
const event = value as SubagentStatusEvent;
if (!bridge || event.sessionId !== bridge.session.id) return;
// Retain the authoritative merged snapshot for hello. The incremental frame
// remains the normal low-bandwidth path while connected.
bridge.session = applySubagentStatusToSession(bridge.session, event);
send(bridge, {
type: "agent.subagents",
sessionId: event.sessionId,
Expand All @@ -816,17 +870,21 @@ export default function webSessions(pi: ExtensionAPI): void {

const forward = (event: unknown, ctx: ExtensionContext, status?: WebSession["status"], refreshMetrics = false): void => {
if (!bridge || bridge.closed || ctx.sessionManager.getSessionId() !== bridge.session.id) return;
if (refreshMetrics) {
const message = isRecord(event) && event.type === "message_end" ? event.message : undefined;
refreshIncrementalMetrics(bridge, isRecord(message) ? message.usage : undefined);
}
send(bridge, {
type: "agent.event",
sessionId: bridge.session.id,
event: safeClone(event),
} satisfies AgentEventMessage);
const latestPreview = isRecord(event) && event.type === "message_end" ? previewFromMessage(event.message) : undefined;
updateSession(bridge, {
status,
messageCount: refreshMetrics
? ctx.sessionManager.getEntries().filter((entry) => entry.type === "message").length
: bridge.session.messageCount,
}, refreshMetrics);
preview: latestPreview ?? bridge.session.preview,
messageCount: refreshMetrics ? bridge.session.messageCount + 1 : bridge.session.messageCount,
});
};

pi.registerCommand("web", {
Expand Down Expand Up @@ -1060,15 +1118,17 @@ export default function webSessions(pi: ExtensionAPI): void {
forward(event, ctx);
});
pi.on("agent_start", (event, ctx) => forward(event, ctx, "working"));
// The visible run is complete at agent_end. agent_settled remains the
// authoritative point for queue delivery and teardown, but the session must
// not keep presenting itself as working while those final hooks drain.
pi.on("agent_end", (event, ctx) => forward(event, ctx, "idle"));
// The visible run is complete at agent_end. Surface provider/runtime failures
// instead of making an unfinished transcript look successfully idle.
pi.on("agent_end", (event, ctx) => {
const status = agentEndTerminalNotice(event)?.kind === "error" ? "error" : "idle";
forward(event, ctx, status);
});
pi.on("agent_settled", (event, ctx) => {
if (bridge?.session.compaction) {
endBridgeCompaction(bridge, { aborted: false, willRetry: false, errorMessage: "Compaction stopped before completion" });
}
forward(event, ctx, "idle");
forward(event, ctx, bridge?.session.status === "error" ? "error" : "idle");
});
pi.on("turn_start", (event, ctx) => forward(event, ctx, "working"));
pi.on("turn_end", (event, ctx) => forward(event, ctx));
Expand All @@ -1088,8 +1148,14 @@ export default function webSessions(pi: ExtensionAPI): void {
}, { once: true });
});
pi.on("session_compact", (event, ctx) => {
if (bridge) updateSession(bridge, {}, true);
if (bridge && ctx.sessionManager.getSessionId() === bridge.session.id) {
refreshIncrementalMetrics(bridge, event.compactionEntry.usage);
updateSession(bridge);
send(bridge, {
type: "agent.history",
sessionId: bridge.session.id,
entries: boundedWebHistory(ctx.sessionManager.buildContextEntries()),
} satisfies AgentHistoryMessage);
endBridgeCompaction(bridge, {
aborted: false,
willRetry: event.willRetry,
Expand Down
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
"extensions": [
"./extensions/session-footer.ts",
"./extensions/pr-footer.ts",
"./extensions/terminal-output.ts",
"./extensions/subagents.ts",
"./extensions/worktree.ts",
"./extensions/web-sessions.ts",
Expand Down
Loading