feat(acp): implemented multi-session support with forking and lifecycle management

- Added multi-session support to ACP mode enabling session forking, resumption, and closure operations.
- Added session model state reporting and unstable_setSessionModel RPC command for direct model configuration.
- Added turn-level usage tracking and stable message ID tracking for assistant chunks in ACP responses.
- Added Settings.cloneForCwd() method and ExtensionRunner.getFlagValues() for configuration isolation across sessions.
- Fixed ACP session cleanup to properly cancel in-flight prompts and dispose resources on disconnect.
- Refactored AcpAgent to manage multiple concurrent sessions with per-session lifecycle instead of single session.
This commit is contained in:
can1357
2026-04-08 06:54:33 +02:00
parent 58c0d52679
commit 53c8d33758
9 changed files with 1083 additions and 167 deletions
+16
View File
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
### Breaking Changes
- Simplified chunk edit operations: removed `append_child`, `prepend_child`, `append_sibling`, `prepend_sibling`, and `replace_body` ops in favor of unified `replace`, `before`, `after`, `prepend`, and `append` with region targeting (`@container`, `@prologue`, `@body`, `@epilogue`)
@@ -9,6 +10,16 @@
### Added
- Multi-session support in ACP mode: agents can now manage multiple concurrent sessions with independent state, models, and configurations
- Session forking in ACP mode: `unstable_forkSession` creates a new session from an existing one's history
- Session resumption in ACP mode: `unstable_resumeSession` reloads a previously saved session
- Session closure in ACP mode: `unstable_closeSession` cleanly shuts down a session and releases resources
- Model state reporting in ACP mode: `SessionModelState` with available models and current selection in session responses
- Direct model setting in ACP mode: `unstable_setSessionModel` RPC command for changing the active model
- Turn-level usage tracking in ACP mode: prompt responses now include `usage` with input/output/cached token counts
- Message ID tracking in ACP mode: stable message IDs for assistant chunks enabling client-side message correlation
- Settings cloning: `Settings.cloneForCwd()` method to create isolated settings instances for different working directories
- Extension flag value retrieval: `ExtensionRunner.getFlagValues()` to inspect current flag state
- Exported autoresearch module and submodules via `./autoresearch` and `./autoresearch/*` package paths
- Exported autoresearch tools via `./autoresearch/tools/*` package path
- Exported CLI commands via `./cli/commands/*` package path
@@ -55,6 +66,10 @@
### Changed
- ACP agent now manages multiple sessions instead of a single session; session lifecycle and configuration are now per-session
- ACP session creation now uses a factory function to support creating new sessions for different working directories
- ACP event mapping now accepts optional `getMessageId` callback for stable message ID assignment to assistant chunks
- ACP session initialization now registers connection cleanup handlers to dispose all sessions on disconnect
- Reorganized package.json exports: moved `./edit` exports before `./plan-mode` for better logical grouping
- Notebook conversion logic now checks for raw read mode or non-chunk mode before converting via markit, allowing chunk-mode reads of `.ipynb` files to use chunk parsing instead of conversion
- Go receiver methods now render as top-level siblings instead of nested under their receiver type in chunk read output
@@ -109,6 +124,7 @@
### Fixed
- ACP session cleanup now properly cancels in-flight prompts and disposes resources when sessions are closed or connection aborts
- Removed unused `_createErrorToolResult` helper function from RPC host-tools module
- Fixed Go receiver method indentation in append operations to preserve relative indentation from the anchor chunk
- Fixed Go type chunk line counts to report only the type body lines instead of including grouped receiver methods
@@ -269,6 +269,21 @@ export class Settings {
}
}
async cloneForCwd(cwd: string): Promise<Settings> {
const cloned = new Settings({
cwd,
agentDir: this.#agentDir,
inMemory: !this.#persist,
});
cloned.#storage = this.#storage;
cloned.#global = structuredClone(this.#global);
cloned.#project = this.#persist ? await cloned.#loadProjectSettings() : structuredClone(this.#project);
cloned.#overrides = structuredClone(this.#overrides);
cloned.#rebuildMerged();
cloned.#fireAllHooks();
return cloned;
}
// ─────────────────────────────────────────────────────────────────────────
// Accessors
// ─────────────────────────────────────────────────────────────────────────
@@ -262,6 +262,10 @@ export class ExtensionRunner {
return allFlags;
}
getFlagValues(): Map<string, boolean | string> {
return new Map(this.runtime.flagValues);
}
setFlagValue(name: string, value: boolean | string): void {
this.runtime.flagValues.set(name, value);
}
+23 -1
View File
@@ -822,10 +822,32 @@ export async function runRootCommand(parsed: Args, rawArgs: string[]): Promise<v
process.exit(1);
}
const extensionFlagValues = session.extensionRunner?.getFlagValues() ?? new Map<string, boolean | string>();
const createAcpSession = async (cwd: string) => {
const nextSettings = await session.settings.cloneForCwd(cwd);
const nextSessionManager = SessionManager.create(cwd, parsedArgs.sessionDir);
const { session: nextSession } = await createAgentSession({
...sessionOptions,
cwd,
sessionManager: nextSessionManager,
settings: nextSettings,
authStorage,
modelRegistry,
searchDb: session.searchDb,
hasUI: false,
});
if (nextSession.extensionRunner) {
for (const [flagName, value] of extensionFlagValues) {
nextSession.extensionRunner.setFlagValue(flagName, value);
}
}
return nextSession;
};
if (mode === "rpc") {
await runRpcMode(session);
} else if (mode === "acp") {
await runAcpMode(session);
await runAcpMode(session, createAcpSession);
} else if (isInteractive) {
const versionCheckPromise = checkForNewVersion(VERSION).catch(() => undefined);
const changelogMarkdown = await getChangelogForDisplay(parsedArgs);
File diff suppressed because it is too large Load Diff
@@ -8,6 +8,10 @@ import type {
import type { AgentSessionEvent } from "../../session/agent-session";
import type { TodoStatus } from "../../tools/todo-write";
interface AcpEventMapperOptions {
getMessageId?: (message: unknown) => string | undefined;
}
interface ContentArrayContainer {
content?: unknown;
}
@@ -118,10 +122,11 @@ export function mapToolKind(toolName: string): ToolKind {
export function mapAgentSessionEventToAcpSessionUpdates(
event: AgentSessionEvent,
sessionId: string,
options: AcpEventMapperOptions = {},
): SessionNotification[] {
switch (event.type) {
case "message_update":
return mapAssistantMessageUpdate(event, sessionId);
return mapAssistantMessageUpdate(event, sessionId, options);
case "tool_execution_start": {
const update: SessionUpdate = {
sessionUpdate: "tool_call",
@@ -181,6 +186,7 @@ export function mapAgentSessionEventToAcpSessionUpdates(
function mapAssistantMessageUpdate(
event: Extract<AgentSessionEvent, { type: "message_update" }>,
sessionId: string,
options: AcpEventMapperOptions,
): SessionNotification[] {
if (!isAssistantMessage(event.message)) {
return [];
@@ -208,10 +214,12 @@ function mapAssistantMessageUpdate(
return [];
}
const messageId = options.getMessageId?.(event.message);
return [
toSessionNotification(sessionId, {
sessionUpdate,
content: { type: "text", text },
messageId,
}),
];
}
@@ -3,11 +3,13 @@ import { AgentSideConnection, ndJsonStream } from "@agentclientprotocol/sdk";
import type { AgentSession } from "../../session/agent-session";
import { AcpAgent } from "./acp-agent";
export async function runAcpMode(session: AgentSession): Promise<never> {
export type AcpSessionFactory = (cwd: string) => Promise<AgentSession>;
export async function runAcpMode(session: AgentSession, createSession: AcpSessionFactory): Promise<never> {
const input = stream.Writable.toWeb(process.stdout);
const output = stream.Readable.toWeb(process.stdin);
const transport = ndJsonStream(input, output);
const connection = new AgentSideConnection(conn => new AcpAgent(conn, session), transport);
const connection = new AgentSideConnection(conn => new AcpAgent(conn, session, createSession), transport);
await connection.closed;
process.exit(0);
}
@@ -0,0 +1,382 @@
import { afterEach, describe, expect, it } from "bun:test";
import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import type { AgentSideConnection, PromptRequest, SessionNotification } from "@agentclientprotocol/sdk";
import type { Model } from "@oh-my-pi/pi-ai";
import { getConfigRootDir, setAgentDir } from "@oh-my-pi/pi-utils";
import { AcpAgent } from "../src/modes/acp/acp-agent";
import type { AgentSession, AgentSessionEvent } from "../src/session/agent-session";
import { SessionManager } from "../src/session/session-manager";
const TEST_MODELS: Model[] = [
{
id: "claude-sonnet-4-20250514",
name: "Claude Sonnet",
api: "anthropic-messages",
provider: "anthropic",
baseUrl: "https://example.invalid",
reasoning: true,
input: ["text", "image"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 200_000,
maxTokens: 8_192,
},
{
id: "gpt-5.4",
name: "GPT-5.4",
api: "openai-responses",
provider: "openai",
baseUrl: "https://example.invalid",
reasoning: true,
input: ["text", "image"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 200_000,
maxTokens: 8_192,
},
];
function makeAssistantMessage(text: string, thinking?: string) {
const content: Array<{ type: "text"; text: string } | { type: "thinking"; thinking: string }> = [
{ type: "text", text },
];
if (thinking) {
content.push({ type: "thinking" as const, thinking });
}
return {
role: "assistant" as const,
content,
api: "anthropic-messages" as const,
provider: "anthropic" as const,
model: TEST_MODELS[0].id,
usage: {
input: 10,
output: 5,
cacheRead: 2,
cacheWrite: 1,
totalTokens: 18,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
stopReason: "stop" as const,
timestamp: Date.now(),
};
}
class FakeAgentSession {
sessionManager: SessionManager;
sessionId: string;
agent: { sessionId: string; waitForIdle: () => Promise<void> };
model: Model | undefined;
thinkingLevel: string | undefined;
customCommands: [] = [];
extensionRunner = undefined;
searchDb = undefined;
isStreaming = false;
queuedMessageCount = 0;
systemPrompt = "system";
disposed = false;
#listeners = new Set<(event: AgentSessionEvent) => void>();
constructor(
cwd: string,
private readonly models: Model[] = TEST_MODELS,
) {
this.sessionManager = SessionManager.create(cwd);
this.sessionId = this.sessionManager.getSessionId();
this.agent = {
sessionId: this.sessionId,
waitForIdle: async () => {},
};
this.model = models[0];
}
get sessionName(): string {
return this.sessionManager.getHeader()?.title ?? `Session ${this.sessionId}`;
}
get modelRegistry(): { getApiKey: (model: Model) => Promise<string> } {
return {
getApiKey: async (_model: Model) => "test-key",
};
}
getAvailableModels(): Model[] {
return this.models;
}
getAvailableThinkingLevels(): ReadonlyArray<string> {
return ["low", "medium", "high"];
}
setThinkingLevel(level: string | undefined): void {
this.thinkingLevel = level;
}
async setModel(model: Model): Promise<void> {
this.model = model;
}
subscribe(listener: (event: AgentSessionEvent) => void): () => void {
this.#listeners.add(listener);
return () => {
this.#listeners.delete(listener);
};
}
async prompt(text: string): Promise<void> {
this.isStreaming = true;
this.sessionManager.appendMessage({ role: "user", content: text, timestamp: Date.now() });
const assistantMessage = makeAssistantMessage("pong");
for (const listener of this.#listeners) {
listener({
type: "message_update",
message: assistantMessage,
assistantMessageEvent: { type: "text_delta", delta: "pong" },
} as AgentSessionEvent);
}
this.sessionManager.appendMessage(assistantMessage);
for (const listener of this.#listeners) {
listener({
type: "agent_end",
messages: [assistantMessage],
} as AgentSessionEvent);
}
this.isStreaming = false;
}
async abort(): Promise<void> {
this.isStreaming = false;
}
async refreshMCPTools(_tools: unknown[]): Promise<void> {}
getContextUsage(): undefined {
return undefined;
}
async switchSession(sessionPath: string): Promise<boolean> {
await this.sessionManager.setSessionFile(sessionPath);
this.sessionId = this.sessionManager.getSessionId();
this.agent.sessionId = this.sessionId;
return true;
}
async dispose(): Promise<void> {
this.disposed = true;
await this.sessionManager.close();
}
async reload(): Promise<void> {}
async newSession(): Promise<boolean> {
await this.sessionManager.newSession();
this.sessionId = this.sessionManager.getSessionId();
this.agent.sessionId = this.sessionId;
return true;
}
async branch(_entryId: string): Promise<{ cancelled: boolean }> {
return { cancelled: false };
}
async navigateTree(_targetId: string): Promise<{ cancelled: boolean }> {
return { cancelled: false };
}
getActiveToolNames(): string[] {
return [];
}
getAllToolNames(): string[] {
return [];
}
setActiveToolsByName(_toolNames: string[]): void {}
async sendCustomMessage(_message: string, _options?: unknown): Promise<void> {}
async sendUserMessage(_content: string, _options?: unknown): Promise<void> {}
async compact(_instructions?: string, _options?: unknown): Promise<void> {}
}
interface AgentHarness {
agent: AcpAgent;
updates: SessionNotification[];
abortController: AbortController;
sessions: FakeAgentSession[];
cwdA: string;
cwdB: string;
findSession(sessionId: string): FakeAgentSession | undefined;
}
function getChunkMessageId(notification: SessionNotification): string | undefined {
const update = notification.update as { messageId?: string | null };
return typeof update.messageId === "string" ? update.messageId : undefined;
}
const cleanupRoots: string[] = [];
const originalAgentDir = process.env.PI_CODING_AGENT_DIR;
const fallbackAgentDir = path.join(getConfigRootDir(), "agent");
afterEach(async () => {
if (originalAgentDir) {
setAgentDir(originalAgentDir);
} else {
setAgentDir(fallbackAgentDir);
delete process.env.PI_CODING_AGENT_DIR;
}
for (const root of cleanupRoots.splice(0)) {
await fs.promises.rm(root, { recursive: true, force: true });
}
});
async function createHarness(): Promise<AgentHarness> {
const root = await fs.promises.mkdtemp(path.join(os.tmpdir(), "omp-acp-test-"));
cleanupRoots.push(root);
const agentDir = path.join(root, "agent");
const cwdA = path.join(root, "cwd-a");
const cwdB = path.join(root, "cwd-b");
await fs.promises.mkdir(agentDir, { recursive: true });
await fs.promises.mkdir(cwdA, { recursive: true });
await fs.promises.mkdir(cwdB, { recursive: true });
setAgentDir(agentDir);
const updates: SessionNotification[] = [];
const abortController = new AbortController();
const sessions: FakeAgentSession[] = [];
const connection = {
sessionUpdate: async (notification: SessionNotification) => {
updates.push(notification);
},
signal: abortController.signal,
closed: Promise.withResolvers<void>().promise,
} as unknown as AgentSideConnection;
const initialSession = new FakeAgentSession(cwdA);
sessions.push(initialSession);
const factory = async (cwd: string): Promise<AgentSession> => {
const session = new FakeAgentSession(cwd);
sessions.push(session);
return session as unknown as AgentSession;
};
return {
agent: new AcpAgent(connection, initialSession as unknown as AgentSession, factory),
updates,
abortController,
sessions,
cwdA,
cwdB,
findSession: (sessionId: string) => sessions.find(session => session.sessionId === sessionId),
};
}
describe("ACP agent", () => {
it("supports multiple live ACP sessions with model and lifecycle handlers", async () => {
const harness = await createHarness();
const first = await harness.agent.newSession({ cwd: harness.cwdA, mcpServers: [] });
const second = await harness.agent.newSession({ cwd: harness.cwdB, mcpServers: [] });
expect(first.models?.availableModels.map(model => model.modelId)).toEqual(
TEST_MODELS.map(model => `${model.provider}/${model.id}`),
);
await harness.agent.unstable_setSessionModel({
sessionId: first.sessionId,
modelId: `${TEST_MODELS[1]!.provider}/${TEST_MODELS[1]!.id}`,
});
await harness.agent.setSessionConfigOption({
sessionId: first.sessionId,
configId: "thinking",
value: "high",
});
const firstSession = harness.findSession(first.sessionId);
const secondSession = harness.findSession(second.sessionId);
expect(firstSession?.model?.id).toBe(TEST_MODELS[1]!.id);
expect(firstSession?.thinkingLevel).toBe("high");
expect(secondSession?.model?.id).toBe(TEST_MODELS[0]!.id);
expect(secondSession?.thinkingLevel).toBeUndefined();
firstSession?.sessionManager.appendMessage({ role: "user", content: "fork me", timestamp: Date.now() });
await firstSession?.sessionManager.flush();
const forked = await harness.agent.unstable_forkSession({
sessionId: first.sessionId,
cwd: harness.cwdA,
mcpServers: [],
});
const forkedSession = harness.findSession(forked.sessionId);
const forkedMessages = forkedSession?.sessionManager.buildSessionContext().messages ?? [];
expect(forked.sessionId).not.toBe(first.sessionId);
expect(forkedMessages.some(message => message.role === "user" && message.content === "fork me")).toBe(true);
await harness.agent.unstable_closeSession({ sessionId: forked.sessionId });
await expect(harness.agent.setSessionMode({ sessionId: forked.sessionId, modeId: "default" })).rejects.toThrow(
"Unsupported ACP session",
);
harness.abortController.abort();
await Bun.sleep(0);
});
it("replays messageIds and returns turn usage for prompts", async () => {
const harness = await createHarness();
const stored = new FakeAgentSession(harness.cwdA);
harness.sessions.push(stored);
stored.sessionManager.appendMessage({ role: "user", content: "hello", timestamp: Date.now() });
stored.sessionManager.appendMessage(makeAssistantMessage("reply", "reasoning"));
await stored.sessionManager.ensureOnDisk();
await stored.sessionManager.flush();
await harness.agent.loadSession({ sessionId: stored.sessionId, cwd: harness.cwdA, mcpServers: [] });
const replayChunks = harness.updates.filter(
update =>
update.sessionId === stored.sessionId &&
(update.update.sessionUpdate === "user_message_chunk" ||
update.update.sessionUpdate === "agent_message_chunk" ||
update.update.sessionUpdate === "agent_thought_chunk"),
);
const replayAssistantChunks = replayChunks.filter(
update =>
update.update.sessionUpdate === "agent_message_chunk" ||
update.update.sessionUpdate === "agent_thought_chunk",
);
expect(
replayChunks.every(
update => typeof getChunkMessageId(update) === "string" && getChunkMessageId(update)!.length > 0,
),
).toBe(true);
expect(new Set(replayAssistantChunks.map(update => getChunkMessageId(update))).size).toBe(1);
const live = await harness.agent.newSession({ cwd: harness.cwdB, mcpServers: [] });
const response = await harness.agent.prompt({
sessionId: live.sessionId,
messageId: "05b17a6f-b310-4be7-b767-6b4f3a84eb63",
prompt: [{ type: "text", text: "ping" }],
} as PromptRequest);
const liveChunks = harness.updates.filter(
update => update.sessionId === live.sessionId && update.update.sessionUpdate === "agent_message_chunk",
);
expect(response.userMessageId).toBe("05b17a6f-b310-4be7-b767-6b4f3a84eb63");
expect(response.usage).toEqual({
inputTokens: 10,
outputTokens: 5,
cachedReadTokens: 2,
cachedWriteTokens: 1,
totalTokens: 18,
});
expect(
liveChunks.some(
update => typeof getChunkMessageId(update) === "string" && getChunkMessageId(update)!.length > 0,
),
).toBe(true);
harness.abortController.abort();
await Bun.sleep(0);
});
});
@@ -0,0 +1,64 @@
import { describe, expect, it } from "bun:test";
import { mapAgentSessionEventToAcpSessionUpdates } from "../src/modes/acp/acp-event-mapper";
import type { AgentSessionEvent } from "../src/session/agent-session";
function makeAssistantMessage(text: string) {
return {
role: "assistant" as const,
content: [{ type: "text" as const, text }],
api: "anthropic-messages" as const,
provider: "anthropic" as const,
model: "claude-sonnet-4-20250514",
usage: {
input: 10,
output: 5,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 15,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
stopReason: "stop" as const,
timestamp: Date.now(),
};
}
function getChunkMessageId(event: { update: object }): string | undefined {
const update = event.update as { messageId?: string | null };
return typeof update.messageId === "string" ? update.messageId : undefined;
}
describe("ACP event mapper", () => {
it("attaches a stable messageId to live assistant chunks", () => {
const assistantMessage = makeAssistantMessage("chunk");
const getMessageId = (message: unknown): string | undefined =>
message === assistantMessage ? "a80f1ff7-4f0a-4e6b-9f09-c94857b62a4a" : undefined;
const textUpdates = mapAgentSessionEventToAcpSessionUpdates(
{
type: "message_update",
message: assistantMessage,
assistantMessageEvent: { type: "text_delta", delta: "chunk" },
} as AgentSessionEvent,
"session-1",
{ getMessageId },
);
const thoughtUpdates = mapAgentSessionEventToAcpSessionUpdates(
{
type: "message_update",
message: assistantMessage,
assistantMessageEvent: { type: "thinking_delta", delta: "plan" },
} as AgentSessionEvent,
"session-1",
{ getMessageId },
);
expect(textUpdates).toHaveLength(1);
expect(thoughtUpdates).toHaveLength(1);
expect(textUpdates[0] ? getChunkMessageId(textUpdates[0]) : undefined).toBe(
"a80f1ff7-4f0a-4e6b-9f09-c94857b62a4a",
);
expect(thoughtUpdates[0] ? getChunkMessageId(thoughtUpdates[0]) : undefined).toBe(
"a80f1ff7-4f0a-4e6b-9f09-c94857b62a4a",
);
});
});