feat(acp): multiplex N concurrent ACP sessions + bash task ownership (RFC 011)
Lifts the RFC 010 single-session-per-connection cap: the bridge now runs N concurrent sessions over one connection, each mapped to its own LoopAgent. - packages/acp: live sessions held in a Map<sessionId, SessionRecord> with an agent→sessionId reverse WeakMap so agent/* events (which carry only the Agent) demux in O(1). Every session/event and agent/status is routed strictly to its owning record — concurrent sessions never cross-settle or interleave their session/update notifications. Per-session state: one in-flight prompt each, session/cancel aborts+settles only its own agent/prompt, session/load reserves a per-id load slot (distinct ids load concurrently; re-loading a live id is rejected), and disposal drains every live session in parallel to quiescence. - packages/tool-bash: record each background task's owning agent at spawn and keep it for the executor's lifetime (NOT cleared on completion). bash_output/bash_kill reject a task owned by a different agent (a task with no owner is open; a no-agent caller can't access an owned task). Task ids are global and predictable, so this is the fence that stops one session's agent from reading/killing another session's background task. - Per-session permission ownership and a per-agent disposer seam stay deferred (depend on the deferred permission gate); the reverse map the gate will route through is in place. RFC 011 stays `proposed`. - Tests: two sessions stream concurrently without interleave; cross-session cancel isolation; per-session in-flight enforcement; dispose-all-to-quiescence; bash cross-session read/kill rejected (+ no-agent and unowned-task cases). - Docs: RFC 011 implementation-status note; acp + tool-bash READMEs; example MVP-limitations updated. 100% per-file coverage maintained.
This commit is contained in:
@@ -15,9 +15,13 @@
|
||||
* that ends in `error` rejects the RPC)
|
||||
* - `session/cancel` → `agent.abort()` + settle the in-flight prompt
|
||||
*
|
||||
* Single-session for the MVP (a 2nd `session/new` is rejected); RFC 011 lifts
|
||||
* that. The `tools/execute` permission gate is deferred — see the
|
||||
* TODO(rfc010-permission-gate) note below.
|
||||
* Multi-session (RFC 011): N concurrent sessions per connection, each mapped to
|
||||
* its own `LoopAgent`. Sessions are keyed by id in `sessions` (forward) with an
|
||||
* `agent→sessionId` reverse map for O(1) demux of `agent/*` events; every
|
||||
* `session/event` and `agent/*` event is routed strictly to its owning session
|
||||
* record, so two sessions streaming at once never interleave their
|
||||
* `session/update` notifications. The `tools/execute` permission gate is
|
||||
* deferred — see the TODO(rfc010-permission-gate) note below.
|
||||
*
|
||||
* stdout is the protocol: this plugin must run in an example that loads NO
|
||||
* stdout logger (the console logger writes to stdout and would corrupt the
|
||||
@@ -121,9 +125,8 @@ export const Config: Schema<AcpConfig> = Schema.object({
|
||||
})
|
||||
|
||||
/**
|
||||
* Per-session bridge state. Single-entry in this MVP (RFC 011 makes the maps
|
||||
* multi-entry); kept as a record from the start so RFC 011 generalizes the
|
||||
* container, not the shape.
|
||||
* Per-session bridge state. One per live ACP session; held in the `sessions`
|
||||
* map keyed by id (RFC 011 multi-session).
|
||||
*/
|
||||
interface SessionRecord {
|
||||
sessionId: string
|
||||
@@ -169,23 +172,21 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
const agentName = config.agentName ?? 'deepseek-harness-acp'
|
||||
const agentVersion = config.agentVersion ?? '0.0.1'
|
||||
|
||||
// Single live session for the MVP. RFC 011 turns this into maps keyed by
|
||||
// sessionId plus an agent→sessionId reverse map for the permission gate.
|
||||
let record: SessionRecord | undefined
|
||||
// True while a `session/load` is between reserving the single-session slot and
|
||||
// installing its `record` (resume() is async). The session guards check BOTH
|
||||
// `record` and `loading` so a pipelined load/new cannot slip past while the
|
||||
// first load's resume() is pending and leak a second live agent.
|
||||
let loading = false
|
||||
// Live sessions keyed by id (RFC 011 multi-session), plus an agent→sessionId
|
||||
// reverse map so `agent/*` events (which carry only the Agent) demux in O(1).
|
||||
// The two stay in lockstep: a record is added to `sessions` and the agent to
|
||||
// `bySession` together, and removed together.
|
||||
const sessions = new Map<string, SessionRecord>()
|
||||
const bySession = new WeakMap<Agent, string>()
|
||||
// Session ids whose `session/load` is mid-`resume()` (the slot is reserved
|
||||
// before the async resume so a pipelined load/new for the SAME id can't create
|
||||
// two agents). Distinct ids load concurrently; a given id loads once at a time.
|
||||
const loadingIds = new Set<string>()
|
||||
// Set once the bridge has torn down (disposal or client disconnect). An async
|
||||
// `session/load` that was mid-`resume()` when teardown ran must observe this
|
||||
// after its await and NOT install a `record` (which would resurrect a live
|
||||
// agent/listeners after the bridge closed). Checked after every load await.
|
||||
// `session/load` mid-`resume()` when teardown ran must observe this after its
|
||||
// await and NOT install a record (which would resurrect a live agent/listeners
|
||||
// after the bridge closed). Checked after every load await.
|
||||
let closed = false
|
||||
// Ownership marker: agents this bridge created. The deferred permission gate
|
||||
// (TODO(rfc010-permission-gate)) and RFC 011 build on this; laid down now so
|
||||
// the seam exists. A WeakMap so a disposed agent's entry is collectable.
|
||||
const owned = new WeakMap<Agent, string>()
|
||||
|
||||
// Assigned at the bottom, before any agent event can fire (a session only
|
||||
// exists after `newSession`, which the client calls after construction), so
|
||||
@@ -207,10 +208,11 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
|
||||
/** Resolve the live record for a sessionId, or throw an ACP error. */
|
||||
const requireSession = (sessionId: string): SessionRecord => {
|
||||
if (record === undefined || record.sessionId !== sessionId) {
|
||||
const rec = sessions.get(sessionId)
|
||||
if (rec === undefined) {
|
||||
throw invalidParams(`unknown session: ${sessionId}`)
|
||||
}
|
||||
return record
|
||||
return rec
|
||||
}
|
||||
|
||||
/** Push a `session/update` notification, swallowing post-close rejections. */
|
||||
@@ -252,10 +254,11 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
// already-cancelled turn whose end arrives late is ignored (see
|
||||
// SessionRecord.inflight). A turn that ends `error` REJECTS the prompt (ACP
|
||||
// has no error stop reason); other reasons resolve via the codec. Demux
|
||||
// strictly by session id.
|
||||
// strictly by session id: a `session/event` is routed to its own record, so
|
||||
// two sessions streaming at once never cross-settle or interleave updates.
|
||||
ctx.on('session/event', (session, event: SessionEvent) => {
|
||||
const rec = record
|
||||
if (rec === undefined || session.header.id !== rec.sessionId) return
|
||||
const rec = sessions.get(session.header.id)
|
||||
if (rec === undefined) return
|
||||
streamSessionEventUpdate(rec.sessionId, event, notify)
|
||||
const inflight = rec.inflight
|
||||
if (inflight === undefined) return
|
||||
@@ -333,9 +336,12 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
// On a settle to idle/disposed, reconcile any still-pending prompt from the
|
||||
// log (covers a starved `session/event` listener — see settleFromLog). A mid-
|
||||
// step disposal that never appended a clean turn/end resolves `cancelled`.
|
||||
// Demux via the agent→sessionId reverse map.
|
||||
ctx.on('agent/status', (agent, status: AgentStatus) => {
|
||||
const rec = record
|
||||
if (rec === undefined || owned.get(agent) !== rec.sessionId) return
|
||||
const sessionId = bySession.get(agent)
|
||||
if (sessionId === undefined) return
|
||||
const rec = sessions.get(sessionId)
|
||||
if (rec === undefined) return
|
||||
if (status === 'idle' || status === 'disposed') settleFromLog(rec)
|
||||
})
|
||||
|
||||
@@ -369,9 +375,6 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
|
||||
newSession(params: NewSessionRequest): Promise<NewSessionResponse> {
|
||||
assertOpen()
|
||||
if (record !== undefined || loading) {
|
||||
throw invalidParams('this agent supports a single session; a session already exists (RFC 011 will lift this)')
|
||||
}
|
||||
validateWorkspaceParams(params)
|
||||
const sessionId = randomUUID()
|
||||
const agent = ctx.agents.create({
|
||||
@@ -380,24 +383,23 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
meta: { cwd: params.cwd },
|
||||
agentOptions: agentOptions(config),
|
||||
})
|
||||
owned.set(agent, sessionId)
|
||||
record = { sessionId, agent, inflight: undefined }
|
||||
bySession.set(agent, sessionId)
|
||||
sessions.set(sessionId, { sessionId, agent, inflight: undefined })
|
||||
return Promise.resolve({ sessionId })
|
||||
},
|
||||
|
||||
async loadSession(params: LoadSessionRequest): Promise<LoadSessionResponse> {
|
||||
assertOpen()
|
||||
if (record !== undefined || loading) {
|
||||
throw invalidParams('this agent supports a single session; a session already exists (RFC 011 will lift this)')
|
||||
if (sessions.has(params.sessionId) || loadingIds.has(params.sessionId)) {
|
||||
throw invalidParams(`session ${params.sessionId} is already loaded`)
|
||||
}
|
||||
validateWorkspaceParams(params)
|
||||
// Reserve the single-session slot BEFORE the await. Without this, two
|
||||
// pipelined load/new requests could both pass the guard above while the
|
||||
// first load's resume() is pending, then both install a record and leak
|
||||
// a second live agent. `loading` claims the slot; it is cleared in
|
||||
// `finally` so a rejected load (bad id, cwd mismatch) never wedges all
|
||||
// future sessions on this connection.
|
||||
loading = true
|
||||
// Reserve THIS id's load slot BEFORE the await. Without it, two pipelined
|
||||
// loads for the same id could both pass the guard above while the first
|
||||
// resume() is pending, then both install a record and leak a second
|
||||
// agent. (Distinct ids load concurrently — the set is keyed by id.) The
|
||||
// slot is released in `finally` so a rejected load never wedges the id.
|
||||
loadingIds.add(params.sessionId)
|
||||
try {
|
||||
// Validate the PERSISTED cwd BEFORE resuming — `list()` is a
|
||||
// metadata-only read (no full-log parse) — so a mismatch rejects
|
||||
@@ -419,7 +421,7 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
agentOptions: agentOptions(config),
|
||||
})
|
||||
// The bridge may have torn down (disposal / client disconnect) while
|
||||
// resume() was pending. Its listeners are gone, so installing `record`
|
||||
// resume() was pending. Its listeners are gone, so installing a record
|
||||
// now would resurrect a live agent the bridge can no longer drive or
|
||||
// tear down. Bail: the just-resumed agent is reclaimed with the host
|
||||
// context (no per-agent disposer — TODO(rfc010-agent-disposal)).
|
||||
@@ -430,8 +432,8 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
if (closed) {
|
||||
throw invalidParams('connection closed during session/load')
|
||||
}
|
||||
owned.set(agent, params.sessionId)
|
||||
record = { sessionId: params.sessionId, agent, inflight: undefined }
|
||||
bySession.set(agent, params.sessionId)
|
||||
sessions.set(params.sessionId, { sessionId: params.sessionId, agent, inflight: undefined })
|
||||
// Replay the persisted event log to the client as session/update. Use
|
||||
// the raw event log (NOT deriveMessages, which drops assistant/chunk
|
||||
// and trace events): RFC 010's load contract reconstructs the streamed
|
||||
@@ -442,7 +444,7 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
}
|
||||
return {}
|
||||
} finally {
|
||||
loading = false
|
||||
loadingIds.delete(params.sessionId)
|
||||
}
|
||||
},
|
||||
|
||||
@@ -477,13 +479,15 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
},
|
||||
|
||||
cancel(params: CancelNotification): Promise<void> {
|
||||
const rec = record
|
||||
if (rec === undefined || rec.sessionId !== params.sessionId) return Promise.resolve()
|
||||
const rec = sessions.get(params.sessionId)
|
||||
if (rec === undefined) return Promise.resolve()
|
||||
// RFC 010: session/cancel maps to agent.abort(reason). This aborts a
|
||||
// RUNNING step (the turn ends 'aborted' → 'cancelled' via turn-end).
|
||||
// It also settles the in-flight prompt as cancelled directly, in case
|
||||
// the abort lands in the pre-step window (queued-but-not-started) where
|
||||
// abort() has no AbortController to signal — see the README
|
||||
// It aborts and settles ONLY this session's agent/prompt — a cancel in
|
||||
// one session never touches another's stream or pending prompt (RFC 011
|
||||
// isolation). It also settles the in-flight prompt as cancelled directly,
|
||||
// in case the abort lands in the pre-step window (queued-but-not-started)
|
||||
// where abort() has no AbortController to signal — see the README
|
||||
// TODO(rfc010-cancel-prestep): a not-yet-started queued turn may still
|
||||
// run to completion until a loop-level cancel lands. Best-effort abort
|
||||
// plus honest RPC/UI cancellation. A secondary consequence of that same
|
||||
@@ -515,11 +519,12 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
conn = new AgentSideConnection(makeAgent, stream)
|
||||
|
||||
/**
|
||||
* Tear the live session down to quiescence (AGENTS.md "dispose must reach
|
||||
* quiescence"): settle any pending prompt `cancelled`, abort the agent, and
|
||||
* AWAIT it draining via the interface-level `whenIdle()` signal (NOT
|
||||
* `agent/status('disposed')`, which fires before the driver exits). Idempotent
|
||||
* — clears `record` first, so a second call (close racing dispose) is a no-op.
|
||||
* Tear ALL live sessions down to quiescence (AGENTS.md "dispose must reach
|
||||
* quiescence"): for each session settle any pending prompt `cancelled`, abort
|
||||
* the agent, and AWAIT it draining via the interface-level `whenIdle()` signal
|
||||
* (NOT `agent/status('disposed')`, which fires before the driver exits). The
|
||||
* agents drain in parallel. Idempotent — clears the `sessions` map first and
|
||||
* memoizes, so a second call (close racing dispose) is a no-op.
|
||||
* Shared by Cordis disposal AND client disconnect (`conn.closed`).
|
||||
*
|
||||
* Caveat (same window as TODO(rfc010-cancel-prestep)): if teardown lands in
|
||||
@@ -528,38 +533,36 @@ export function apply(ctx: Context, config: AcpConfig): void {
|
||||
* `whenIdle()` returns immediately (status is still `idle`), so that queued
|
||||
* turn may still start and run after teardown returns. Reaching true
|
||||
* quiescence in that window needs a queue-aware loop cancel primitive (a
|
||||
* loop-level change, out of the RFC 010 MVP scope); for `newSession` agents
|
||||
* the worst case is one short queued turn, since the bridge enforces a single
|
||||
* in-flight prompt.
|
||||
* loop-level change); the single-in-flight-per-session rule bounds the worst
|
||||
* case to one short queued turn per session.
|
||||
*
|
||||
* The agent itself is NOT individually disposed/unregistered here — the
|
||||
* factory (`ctx.agents.create`/`resume`) registers it on the AgentLoop fiber
|
||||
* and returns no per-agent disposer, so the registry entry is reclaimed when
|
||||
* The agents themselves are NOT individually disposed/unregistered here — the
|
||||
* factory (`ctx.agents.create`/`resume`) registers each on the AgentLoop fiber
|
||||
* and returns no per-agent disposer, so registry entries are reclaimed when
|
||||
* the host context disposes. On a bare client disconnect (without a host
|
||||
* dispose) the idled agent therefore lingers in `ctx.agents` until shutdown;
|
||||
* since the MVP is single-session-per-connection and a reconnect spins up a
|
||||
* fresh context, this does not strand work. A per-agent disposal seam is
|
||||
* RFC 011 follow-up (TODO(rfc010-agent-disposal)).
|
||||
* dispose) the idled agents linger in `ctx.agents` until shutdown; a reconnect
|
||||
* spins up a fresh context, so this does not strand work. A per-agent disposal
|
||||
* seam is a follow-up (TODO(rfc010-agent-disposal)).
|
||||
*/
|
||||
let quiescing: Promise<void> | undefined
|
||||
const quiesce = (): Promise<void> => {
|
||||
// Memoize: disposal and client-disconnect can both fire. The first call owns
|
||||
// the teardown; later callers await the SAME promise so `fiber.dispose()`
|
||||
// never returns before an in-flight close teardown has finished (using
|
||||
// `record === undefined` as the only guard would let the second caller
|
||||
// return early while the first is still awaiting whenIdle()).
|
||||
// never returns before an in-flight close teardown has finished.
|
||||
if (quiescing !== undefined) return quiescing
|
||||
// Mark closed BEFORE the record check: a `session/load` mid-`resume()` (no
|
||||
// record installed yet) must observe this after its await and refuse to
|
||||
// install a post-teardown record. Set even when there is nothing else to do.
|
||||
// Mark closed BEFORE draining: a `session/load` mid-`resume()` (no record
|
||||
// installed yet) must observe this after its await and refuse to install a
|
||||
// post-teardown record. Set even when there are no live sessions.
|
||||
closed = true
|
||||
const rec = record
|
||||
record = undefined
|
||||
if (rec === undefined) return Promise.resolve()
|
||||
const recs = [...sessions.values()]
|
||||
sessions.clear()
|
||||
if (recs.length === 0) return Promise.resolve()
|
||||
quiescing = (async () => {
|
||||
settlePrompt(rec, 'cancelled')
|
||||
rec.agent.abort('disposed')
|
||||
await rec.agent.whenIdle()
|
||||
await Promise.all(recs.map(async (rec) => {
|
||||
settlePrompt(rec, 'cancelled')
|
||||
rec.agent.abort('disposed')
|
||||
await rec.agent.whenIdle()
|
||||
}))
|
||||
})()
|
||||
return quiescing
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user