import { describe, expect, it } from "bun:test" import type { AgentSideConnection } from "@agentclientprotocol/sdk" import type { Event, Message, OpencodeClient, Part, SessionMessageResponse } from "@opencode-ai/sdk/v2" import { Effect, ManagedRuntime } from "effect" import { ACPNextEvent } from "@/acp-next/event" import * as ACPNextService from "@/acp-next/service" import { Directory } from "@/acp-next/directory" import { ACPNextSession } from "@/acp-next/session" type SessionUpdateParams = Parameters[0] type GlobalEventEnvelope = { payload?: Event } type DeltaPartType = Extract["type"] const pollUntil = async ( check: () => boolean | Promise, message: string, opts?: { timeoutMs?: number; intervalMs?: number }, ) => { const started = Date.now() while (true) { if (await check()) return if (Date.now() - started > (opts?.timeoutMs ?? 2000)) throw new Error(message) await new Promise((resolve) => setTimeout(resolve, opts?.intervalMs ?? 5)) } } function makeSessionService() { return ManagedRuntime.make(ACPNextSession.defaultLayer).runSync( ACPNextSession.Service.use((service) => Effect.succeed(service)), ) } function createEventStream() { const queue: GlobalEventEnvelope[] = [] const waiters: Array<(value: GlobalEventEnvelope | undefined) => void> = [] const state = { closed: false } const push = (event: GlobalEventEnvelope) => { const waiter = waiters.shift() if (waiter) { waiter(event) return } queue.push(event) } const close = () => { state.closed = true for (const waiter of waiters.splice(0)) { waiter(undefined) } } const stream = async function* (signal?: AbortSignal) { while (true) { if (signal?.aborted) return const next = queue.shift() if (next) { yield next continue } if (state.closed) return const value = await new Promise((resolve) => { waiters.push(resolve) signal?.addEventListener("abort", () => resolve(undefined), { once: true }) }) if (!value) return yield value } } return { push, close, stream } } function createHarness(messages: Record = {}) { const updates: SessionUpdateParams[] = [] const calls = { eventSubscribe: 0, message: 0, } const events = createEventStream() const sdk = { global: { event: (options?: { signal?: AbortSignal }) => { calls.eventSubscribe++ return Promise.resolve({ stream: events.stream(options?.signal) }) }, }, session: { message: (input: { messageID: string }) => { calls.message++ return Promise.resolve({ data: messages[input.messageID] }) }, get: () => Promise.resolve({ data: { id: "ses_loaded" } }), messages: () => Promise.resolve({ data: [] }), }, } as unknown as OpencodeClient const connection = { sessionUpdate: (params: SessionUpdateParams) => { updates.push(params) return Promise.resolve() }, } satisfies Pick const session = makeSessionService() const subscription = new ACPNextEvent.Subscription({ sdk, connection, session }) return { calls, connection, events, sdk, session, subscription, updates } } function textDelta(sessionID: string, messageID: string, partID: string, delta: string): Event { return { id: `evt_${sessionID}_${messageID}_${partID}_${delta}`, type: "message.part.delta", properties: { sessionID, messageID, partID, field: "text", delta, }, } } function partUpdated(sessionID: string, messageID: string, partID: string, type: DeltaPartType): Event { return { id: `evt_${sessionID}_${messageID}_${partID}`, type: "message.part.updated", properties: { sessionID, time: Date.now(), part: type === "text" ? { id: partID, sessionID, messageID, type: "text", text: "", } : { id: partID, sessionID, messageID, type: "reasoning", text: "", time: { start: Date.now() }, }, }, } } function assistantMessage(sessionID: string, messageID: string, partID: string, type: DeltaPartType) { return { info: { id: messageID, sessionID, role: "assistant", time: { created: Date.now() }, parentID: "msg_parent", modelID: "model", providerID: "provider", mode: "build", agent: "build", path: { cwd: "/workspace", root: "/workspace" }, cost: 0, tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, }, parts: [ type === "text" ? { id: partID, sessionID, messageID, type: "text", text: "", } : { id: partID, sessionID, messageID, type: "reasoning", text: "", time: { start: Date.now() }, }, ], } satisfies SessionMessageResponse } async function createKnownSession( session: ACPNextSession.Interface, sessionId: string, part: { messageId: string; partId: string; partType: Part["type"]; role?: Message["role"] }, ) { await Effect.runPromise(session.create({ id: sessionId, cwd: "/workspace" })) await Effect.runPromise( session.recordPartMetadata({ sessionId, messageId: part.messageId, partId: part.partId, partType: part.partType, role: part.role ?? "assistant", }), ) } describe("acp-next event routing", () => { it("routes message.part.delta by sessionID without cross-session pollution", async () => { const harness = createHarness() await createKnownSession(harness.session, "ses_a", { messageId: "msg_a", partId: "part_a", partType: "text" }) await createKnownSession(harness.session, "ses_b", { messageId: "msg_b", partId: "part_b", partType: "text" }) await harness.subscription.handle(textDelta("ses_b", "msg_b", "part_b", "hello")) expect(harness.updates.map((update) => update.sessionId)).toEqual(["ses_b"]) expect(harness.updates[0]?.update.sessionUpdate).toBe("agent_message_chunk") }) it("keeps interleaved sessions isolated for text and reasoning deltas", async () => { const harness = createHarness() await createKnownSession(harness.session, "ses_a", { messageId: "msg_a", partId: "part_a", partType: "text" }) await createKnownSession(harness.session, "ses_b", { messageId: "msg_b", partId: "part_b", partType: "reasoning", }) await harness.subscription.handle(textDelta("ses_a", "msg_a", "part_a", "A1")) await harness.subscription.handle(textDelta("ses_b", "msg_b", "part_b", "B1")) await harness.subscription.handle(textDelta("ses_a", "msg_a", "part_a", "A2")) await harness.subscription.handle(textDelta("ses_b", "msg_b", "part_b", "B2")) expect( harness.updates.filter((update) => update.sessionId === "ses_a").map((update) => update.update.sessionUpdate), ).toEqual(["agent_message_chunk", "agent_message_chunk"]) expect( harness.updates.filter((update) => update.sessionId === "ses_b").map((update) => update.update.sessionUpdate), ).toEqual(["agent_thought_chunk", "agent_thought_chunk"]) }) it("does not create extra subscriptions on repeated loadSession", async () => { const harness = createHarness() let subscription: ACPNextEvent.Subscription | undefined const service = ACPNextService.make({ sdk: harness.sdk, connection: harness.connection, directory: { get: () => Effect.succeed( Directory.build({ directory: "/workspace", providers: {}, modes: [], defaultModeID: "build", commands: [], }), ), refresh: () => Effect.succeed( Directory.build({ directory: "/workspace", providers: {}, modes: [], defaultModeID: "build", commands: [], }), ), variants: Directory.variants, }, session: harness.session, eventSubscription: (started) => { subscription = started }, }) await pollUntil(() => harness.calls.eventSubscribe === 1, "event subscription did not start") await Effect.runPromise(service.loadSession({ cwd: "/workspace", sessionId: "ses_loaded", mcpServers: [] })) await Effect.runPromise(service.loadSession({ cwd: "/workspace", sessionId: "ses_loaded", mcpServers: [] })) await Effect.runPromise(service.loadSession({ cwd: "/workspace", sessionId: "ses_loaded", mcpServers: [] })) expect(harness.calls.eventSubscribe).toBe(1) subscription?.stop() harness.events.close() }) it("does not call sdk.session.message repeatedly when metadata is known", async () => { const harness = createHarness() await createKnownSession(harness.session, "ses_a", { messageId: "msg_a", partId: "part_a", partType: "text" }) for (const delta of ["a", "b", "c", "d", "e"]) { await harness.subscription.handle(textDelta("ses_a", "msg_a", "part_a", delta)) } expect(harness.calls.message).toBe(0) expect(harness.updates).toHaveLength(5) }) it("fetches unknown part metadata once and reuses it for later deltas", async () => { const harness = createHarness({ msg_a: assistantMessage("ses_a", "msg_a", "part_a", "text"), }) await Effect.runPromise(harness.session.create({ id: "ses_a", cwd: "/workspace" })) await harness.subscription.handle(partUpdated("ses_a", "msg_a", "part_a", "text")) await harness.subscription.handle(textDelta("ses_a", "msg_a", "part_a", "a")) await harness.subscription.handle(textDelta("ses_a", "msg_a", "part_a", "b")) expect(harness.calls.message).toBe(1) expect(harness.updates).toHaveLength(2) }) it("ignores unknown sessions and live user parts without user_message_chunk duplication", async () => { const harness = createHarness() await createKnownSession(harness.session, "ses_user", { messageId: "msg_user", partId: "part_user", partType: "text", role: "user", }) await harness.subscription.handle(textDelta("ses_missing", "msg_missing", "part_missing", "ignored")) await harness.subscription.handle(partUpdated("ses_user", "msg_user", "part_live", "text")) await harness.subscription.handle(textDelta("ses_user", "msg_user", "part_user", "hello")) expect(harness.updates).toHaveLength(0) }) })