diff --git a/devlog/2026-09-16_serialize-session-init-transforms/DESIGN.md b/devlog/2026-09-16_serialize-session-init-transforms/DESIGN.md new file mode 100644 index 00000000..cac40efe --- /dev/null +++ b/devlog/2026-09-16_serialize-session-init-transforms/DESIGN.md @@ -0,0 +1,52 @@ +# DESIGN - Serialize same-session state initialization and transforms + +- Task ID: `2026-09-16_serialize-session-init-transforms` +- Home Repo: `opencode-acp` +- Created: 2026-09-16 +- Status: Accepted + +## 1. Problem Statement + +- **What problem are we solving?** Issue #404: `SessionState` is mutable per-session state shared across many async entry points (message transform, compress/decompress tools, event hook saves, system hook limit writes). Nothing serialized them across `await`s, so (a) concurrent initializations could observe partially loaded state, and (b) a stale in-flight transaction could persist over a newer committed one. +- **Why now?**: Reproduced deterministically with a slow-host mock; same failure shape as the baseline-reset class of bugs (#5.7.3) — silent cross-turn state corruption invisible to single-turn tests. + +## 2. Goals & Non-Goals + +- **Goals**: One-writer-at-a-time per session across ALL mutation paths; init coalescing so N racers pay for 1 initialization; zero steady-state overhead; no cross-session blocking; no persisted-format changes. +- **Non-Goals**: Cross-process locking; redesigning the debounced save queue (#384); v2 runtime support (#395 → billion-context#735). + +## 3. Current Architecture + +- **How it works today**: `SessionStateRegistry.states: Map`. The transform hook calls `getOrCreate` → `ensureSessionInitialized` (idempotency via `state.sessionId === sessionId`, which was assigned synchronously pre-first-await), then a long mutate pipeline, then `saveContext`/`saveSessionState`. Tools call `prepareSession` → `ensureSessionInitialized` directly on the registry state. Event/system hooks grab state by id and write+save. +- **Pain points**: Every step above can interleave with another same-session trigger at any await boundary. + +## 4. Proposed Architecture + +- **Overview**: + ``` + trigger (transform | tool | event | system) + │ + ▼ + registry.withSessionGuard(sessionId, fn) ← FIFO promise-chain mutex + │ (await previous chain entry, then run fn exclusively) + ▼ + fn: getOrCreate/prepareSession → mutate → save ← single writer + ``` +- **Key components**: + - `createSessionGuard()` (`lib/state/state.ts`): returns `{ run(sessionId, fn) }`; internal `Map` chain; each task enqueues before awaiting its predecessor; map entry deleted when the tail task completes → empty when idle. Rejections propagate to the caller but never poison the chain (predecessor awaited via `.catch(()=>{})`). + - `ensureSessionInitialized`: now a coalescing wrapper over private `runSessionInitialization`. Module-level `WeakMap>` keyed by the STATE OBJECT (not session id) so soft-cap eviction + recreation starts fresh rather than awaiting a stale promise. +- **Data flow**: unchanged apart from ordering guarantees. Snapshot/restore (`restoreCompressionState`) already preserves `compressionTiming` identity — locked in by test 7. + +## 5. Critical Subtleties (load-bearing) + +1. **Inflight check precedes fast path.** `runSessionInitialization` assigns `state.sessionId` synchronously before its first await. If the `state.sessionId === sessionId` early-return ran first, racing callers would skip coalescing and return mid-init — the original bug. Order: inflight check → fast path → start+track. +2. **Guard granularity = whole transaction, not individual mutations.** A read-modify-write must be atomic end-to-end; wrapping only `saveSessionState` would still allow stale reads during the pipeline. Hence the transform body became `runPipeline(state)` invoked inside one guard acquisition. +3. **Ephemeral transform branch unguarded.** When no user message exists, the handler builds a throwaway `createSessionState()` per request — independent objects, nothing shared, no lock needed (locking on a synthetic key would only add overhead). +4. **Tools hold the guard across `prepareSession` (incl. permission `ask`).** In OpenCode v1 a tool runs while the session turn is paused, so no same-session transform is concurrently in flight — the lock is free in practice; if a trigger did interleave, serialization is exactly the desired behavior. No deadlock possible: guards are non-reentrant and no code path acquires two locks. +5. **Event handler skips states without sessionId** inside the guarded section (defensive; such states have no persistence identity). + +## 6. Alternatives Considered + +- **Lock only around save**: insufficient — stale READS during the pipeline already corrupted the snapshot being saved. +- **Module-level lock instead of registry method**: rejected — tools receive `registry` via `ToolFactoryContext`; a module singleton would duplicate ownership and complicate test stubs. Test stubs compose the real factory (`tests/registry-stub.ts`) to avoid drift. +- **Reentrant guard / ref-counting**: rejected — no nesting exists today; reentrancy invites future misuse. diff --git a/devlog/2026-09-16_serialize-session-init-transforms/REQ.md b/devlog/2026-09-16_serialize-session-init-transforms/REQ.md new file mode 100644 index 00000000..e5ee0ef5 --- /dev/null +++ b/devlog/2026-09-16_serialize-session-init-transforms/REQ.md @@ -0,0 +1,50 @@ +# REQ - Serialize same-session state initialization and transforms + +- Task ID: `2026-09-16_serialize-session-init-transforms` +- Home Repo: `opencode-acp` +- Created: 2026-09-16 +- Status: Done +- Priority: P1 +- Owner: ranxianglei +- References: https://github.com/ranxianglei/opencode-acp/issues/404 + +## 1. Background & Problem Statement + +- **Context**: ACP keeps mutable per-session `SessionState` in a process-wide registry. Every LLM request runs the `messages.transform` hook (init + mutate + persist), and the `compress`/`decompress` tools, the `event` hook (duration attach + save) and the `system.transform` hook (model-limit write + save) also read/mutate/persist the same state. +- **Current behavior (symptom)**: Node is single-threaded but async work interleaves at every `await`. Two triggers: + 1. **Init race** — `ensureSessionInitialized` assigned `state.sessionId = sessionId` synchronously before its first await. A concurrent caller for the same session then hit the idempotency fast path and returned *partially initialized* state (persisted blocks/messageIds not yet loaded). + 2. **Stale transaction** — a transform that awaits mid-pipeline can resume after a newer transform already committed; last-write-wins persistence then stores the stale snapshot over committed state (lost updates / corrupted prune state). +- **Expected behavior**: All same-session state work is serialized; concurrent initializations coalesce into one; no caller ever observes partially initialized state; committed state is never overwritten by a stale snapshot. +- **Impact**: Intermittent data corruption in sessions under concurrency (subagent fan-out, retries, fast successive requests): lost compression blocks, wrong nudge baselines, duplicated message refs. + +## 2. Reproduction (if applicable) + +- **Environment**: Node 22/24, any OS — the race is in-process, timing-dependent. +- **Minimal reproduction steps**: + 1) Seed persisted state for session S (e.g. `modelContextLimit`). + 2) Fire two `getOrCreate(client, "S", ...)` calls concurrently with a slow host (`client.session.get` delayed ~25 ms). + 3) The second caller returns before `loadSessionState` completes → reads `modelContextLimit === undefined` instead of the persisted value. +- **Relevant configuration**: none — default config exercises the path. + +## 3. Constraints & Non-Goals + +- **Constraints**: + - Backward compatibility: persisted state format unchanged; exported API only grows (`createSessionGuard`, `SessionGuard`, `registry.withSessionGuard`); internal `dcp` naming untouched. + - Performance requirements: O(1) constant overhead on the steady-state path (one Map get/set/delete plus a microtask hop per guarded request when nothing else is in flight); different sessions must never block each other. + - Resource limits: lock map must be empty when idle (no per-session leak); init-coalescing map keyed by state object (WeakMap) so eviction + recreation starts fresh. +- **Non-Goals** (explicitly out of scope): + - Cross-process locking (out of scope for an in-process plugin). + - Changing `saveSessionState`'s existing debounced queue semantics (#384). + - OpenCode v2 runtime support (v2 moved to billion-context#735 per #395). + +## 4. Acceptance Criteria (must be testable) + +- **Correctness**: + - [x] Concurrent `getOrCreate` calls for one session coalesce into a single initialization; every waiter observes fully loaded state at return time. + - [x] Same-session tasks run FIFO through the guard; other sessions are unaffected. + - [x] Guard releases on rejection; subsequent tasks proceed. + - [x] Serialized read-modify-write: a stale transaction cannot persist over a newer committed value. + - [x] Regression test FAILS when coalescing is disabled (verified by temporarily disabling it). +- **Performance / Stability**: + - [x] Full suite green: 1270 tests, 0 failures (was 1263 before this change added 7). + - [x] `tsc --noEmit` clean; `npm run build` clean. diff --git a/devlog/2026-09-16_serialize-session-init-transforms/WORKLOG.md b/devlog/2026-09-16_serialize-session-init-transforms/WORKLOG.md new file mode 100644 index 00000000..3363ec07 --- /dev/null +++ b/devlog/2026-09-16_serialize-session-init-transforms/WORKLOG.md @@ -0,0 +1,65 @@ +# WORKLOG - Serialize same-session state initialization and transforms + +- Task ID: `2026-09-16_serialize-session-init-transforms` +- Home Repo: `opencode-acp` +- Status: Done +- Updated: 2026-09-16 (dual-review fix round) + +## 1. Summary + +- **What was done** (1–3 sentences): Added a per-session FIFO mutex (`createSessionGuard` / `registry.withSessionGuard`) and in-flight init coalescing (`inflightInits` WeakMap) to `lib/state/state.ts`; wrapped every same-session mutation path — message transform pipeline, compress/decompress tool execution, event-hook duration attach+save, system-hook model-limit write+save — in the guard. +- **Why** (1–3 sentences): Same-session async work interleaved at awaits: racing initializations handed out partially loaded state, and stale transactions could persist over newer committed state (issue #404). Serialization makes the single-writer invariant hold across awaits. +- **Behavior / compatibility changes**: No persisted-format or API breakage. New exported symbols only. Steady-state fast paths unchanged (guard map empty when idle; init fast path preserved behind the inflight check). One deliberate ordering change: the init inflight check now precedes the `sessionId` fast path (required for correctness — see DESIGN §5). +- **Risk level**: Medium (touches the core request pipeline; mitigated by full-suite green + new regression tests + FIFO guard semantics that cannot deadlock across sessions). + +## 2. Change Log + +### Commits + +| Commit | Description | +|--------|-------------| +| `571ccf0` | fix: serialize same-session state initialization and transforms (#404) | +| `see-branch` | docs: fill devlog commit SHAs (self-referential SHA omitted) | +| `see-branch` | review fixes: transform wiring regression test, event-handler cross-session parallelism, command-handler guard, prettier width wraps, REQ wording | + +### Key Files + +- `lib/state/state.ts` — `createSessionGuard()`/`SessionGuard` factory (FIFO chain mutex); `ensureSessionInitialized` split into coalescing wrapper + `runSessionInitialization`; registry exposes `withSessionGuard` (field intentionally not `readonly` so the wiring test can wrap it as an observation seam). +- `lib/hooks.ts` — system hook limit write wrapped in guard; message transform restructured into `runPipeline(state)` closure invoked inside `registry.withSessionGuard(sessionID, ...)` (session branch only; ephemeral branch unguarded by design); event handler per-state apply+save guarded per session AND run concurrently across sessions (`Promise.allSettled`) so one long-held guard cannot head-of-line-block other sessions' duration attachment; `/acp` command handler `getOrCreate`+permission-sync kept under the guard (getOrCreate can initialize + persist). +- `lib/compress/range.ts` — entire `execute` body from `prepareSession` onward wrapped in `factoryCtx.registry.withSessionGuard(toolCtx.sessionID, ...)`. +- `lib/compress/decompress.ts` — same wrap from `prepareDecompressSession` onward. +- `tests/registry-stub.ts` — both stubs compose the real `createSessionGuard()` (no drift). +- `tests/session-guard.test.ts` — 8 new tests (FIFO order, cross-session independence, release-on-reject, coalescing regression, failed-init semantics, stale read-modify-write, timing-identity preservation, **end-to-end transform wiring**: two concurrent same-session requests through the real `createChatMessageTransformHandler` must produce strictly nested guard enter/exit pairs). + +## 3. Design & Implementation Notes + +- **Entry point / key function**: `createSessionGuard()` in `lib/state/state.ts` — chain-based promise queue keyed by sessionId; entry deleted when the tail task finishes (no idle leak). +- **Key logic explanation**: See DESIGN.md. Critical subtlety: the init inflight check must run BEFORE the `state.sessionId === sessionId` fast path, because `runSessionInitialization` assigns `sessionId` synchronously before its first await — otherwise racing callers early-return mid-init (the original bug shape). + +## 4. Testing & Verification + +### Build & Test Commands + +```sh +npx tsc --noEmit # clean +npm run build # clean (tsup + d.ts) +npm run test # 1271 pass / 0 fail +node --import tsx --test tests/session-guard.test.ts # 8/8 +``` + +### Regression verification (per AGENTS.md §5.7.3) + +Temporarily disabled the inflight check in `ensureSessionInitialized` → `tests/session-guard.test.ts` test 4 ("concurrent getOrCreate coalesces…") FAILED with `bLimitAtReturn === undefined` (the pre-fix partial-snapshot observation). Re-enabled → all green. + +Wiring test (test 8): temporarily bypassed the transform-branch guard in `lib/hooks.ts` (body ran unguarded) → test 8 FAILED (`events === []`, no serialization observed). Restored → green. Reviewer 2 independently repeated the coalescing red/green check. + +### Dual-agent review (AGENTS.md §5.3 / §5.6) + +Both reviewers: **APPROVE-WITH-NITS**, zero blockers/majors; each ran tsc + full suite independently (1271/1271). Fix round applied directly to the branch: + +- MAJOR (tests): missing end-to-end wiring coverage → added test 8 (above). +- MINOR: event handler awaited each session's guard sequentially (head-of-line blocking) → `Promise.allSettled` across sessions, per-session containment preserved. +- MINOR: command handler called `getOrCreate` outside any guard → wrapped with permission-sync under one acquisition. +- MINOR: two new-code prettier width violations (`state.ts`, `hooks.ts`) → wrapped. +- NIT: REQ.md "zero overhead" claim imprecise → "O(1) constant overhead". +- Follow-ups filed separately (source-marked per §5.1.3): guard leak if a host abandons (without rejecting) tool execution while the guard spans a permission ask; sticky failed-init (a transient first-request init failure permanently suppresses persisted-state load for that session until process restart — pre-existing behavior, preserved). diff --git a/lib/compress/decompress.ts b/lib/compress/decompress.ts index 39945575..67edc53f 100644 --- a/lib/compress/decompress.ts +++ b/lib/compress/decompress.ts @@ -275,6 +275,10 @@ export function createDecompressTool(factoryCtx: ToolFactoryContext): ReturnType args: buildSchema(), async execute(args, toolCtx) { const ctx = resolveToolContext(factoryCtx, toolCtx.sessionID) + // [Issue #404] Serialize the full prepare→mutate→finalize transaction under the + // per-session guard (see compress/range.ts for rationale). The body intentionally + // keeps its original indentation inside `run` to keep this change additive-only. + const run = async () => { const { rawMessages } = await prepareDecompressSession(ctx, toolCtx) const effectiveLimitBefore = resolveEffectiveContextLimit(ctx.state, ctx.config) @@ -407,6 +411,8 @@ export function createDecompressTool(factoryCtx: ToolFactoryContext): ReturnType }) return lines.join("\n") + } + return factoryCtx.registry.withSessionGuard(toolCtx.sessionID, () => run()) }, }) } diff --git a/lib/compress/range.ts b/lib/compress/range.ts index 38709373..87b49bbc 100644 --- a/lib/compress/range.ts +++ b/lib/compress/range.ts @@ -119,6 +119,11 @@ export function createCompressRangeTool(factoryCtx: ToolFactoryContext): ReturnT ? (toolCtx as unknown as { callID: string }).callID : undefined + // [Issue #404] Serialize the full prepare→mutate→finalize transaction under the + // per-session guard so it cannot interleave with concurrent same-session work + // (message transforms, event-hook saves). The body intentionally keeps its + // original indentation inside `run` to keep this change additive-only in diff. + const run = async () => { const { rawMessages, searchContext } = await prepareSession( ctx0, toolCtx, @@ -366,6 +371,8 @@ export function createCompressRangeTool(factoryCtx: ToolFactoryContext): ReturnT ? `\n⚠️ acknowledgeRisk was ignored: no quality gate rejection was pending, so quality checks ran normally. Only pass it when retrying immediately after a quality gate rejection.\n` : "" return `Compressed ${totalCompressedMessages} messages into ${COMPRESSED_BLOCK_HEADER}.${skippedNote}${ackNote}\nIMPORTANT: This was an automatic context compression. You MUST continue your previous task exactly where you left off. Do NOT ask the user what to do next.\n💡 Tip: Use search_context('keyword') to find compressed content when you need it later.` + } + return factoryCtx.registry.withSessionGuard(toolCtx.sessionID, () => run()) }, }) } diff --git a/lib/hooks.ts b/lib/hooks.ts index e508b791..efecec12 100644 --- a/lib/hooks.ts +++ b/lib/hooks.ts @@ -124,29 +124,35 @@ export function createSystemPromptHandler( // hook is the only writer and fires AFTER messages.transform within // a request, so without this the limit is learned and lost every // message and the safety net never engages. - if (input.model?.limit?.context) { - const limit = input.model.limit.context - const providerID = input.model?.providerID - const modelID = input.model?.id - // Identity fields are only written when present: a limit without - // identity must not clobber the pair the messages hook relies on - // for staleness detection (#312). - const changed = - state.modelContextLimit !== limit || - (providerID !== undefined && state.modelProviderID !== providerID) || - (modelID !== undefined && state.modelID !== modelID) - state.modelContextLimit = limit - // [FIX #312 follow-up] Record WHICH model the limit belongs to so - // the messages hook can detect staleness on a catalog miss. - if (providerID !== undefined) { - state.modelProviderID = providerID - } - if (modelID !== undefined) { - state.modelID = modelID - } - if (changed) { - saveSessionState(state, logger).catch(() => {}) - } + if (input.model?.limit?.context && input.sessionID) { + // [Issue #404] Serialize the limit write + save against same-session + // transforms/tools: this hook fires AFTER messages.transform within a + // request, so an unsynchronized mutation+save could interleave with a + // transform's persistence of newer in-memory state. + await registry.withSessionGuard(input.sessionID, () => { + const limit = input.model.limit.context + const providerID = input.model?.providerID + const modelID = input.model?.id + // Identity fields are only written when present: a limit without + // identity must not clobber the pair the messages hook relies on + // for staleness detection (#312). + const changed = + state.modelContextLimit !== limit || + (providerID !== undefined && state.modelProviderID !== providerID) || + (modelID !== undefined && state.modelID !== modelID) + state.modelContextLimit = limit + // [FIX #312 follow-up] Record WHICH model the limit belongs to so + // the messages hook can detect staleness on a catalog miss. + if (providerID !== undefined) { + state.modelProviderID = providerID + } + if (modelID !== undefined) { + state.modelID = modelID + } + if (changed) { + saveSessionState(state, logger).catch(() => {}) + } + }) } const effectivePermission = compressPermission(state, config) @@ -197,15 +203,190 @@ export function createChatMessageTransformHandler( } const lastUserMessage = getLastUserMessage(messages) - let state: SessionState + + // [Issue #404] Shared post-resolution pipeline stages, invoked exactly + // once per request below. For real sessions it runs INSIDE the session + // guard so init + mutations + persistence serialize per session. + const runPipeline = async (state: SessionState): Promise => { + syncCompressPermissionState(state, config, hostPermissions, output.messages) + + if (state.isSubAgent && !config.allowSubAgents) { + return + } + + stripHallucinations(output.messages) + + // [#368] Drop oversized reasoning from closed-turn compress tool calls. + // compress calls are hard-exempt from compression (Bug 39), so their + // thinking otherwise rides along every request as an unreclaimable + // floor. Gated by the nested `compress.reasoning` config, resolved + // through the #344 cascade (model > provider > global) using THIS + // request's model identity (request metadata first, session state as + // fallback). Runs BEFORE token accounting / pruning so every later + // stage sees the post-drop array. + const dropReasoningModel = ( + lastUserMessage?.info as + { model?: { providerID?: string; modelID?: string } } | undefined + )?.model + const reasoningConfig = applyCompressOverrides( + config, + dropReasoningModel?.providerID ?? state.modelProviderID, + dropReasoningModel?.modelID ?? state.modelID, + ).compress.reasoning + if (reasoningConfig?.drop !== false) { + const droppedReasoning = dropCompressReasoning( + output.messages, + reasoningConfig?.threshold ?? DEFAULT_COMPRESS_REASONING.threshold, + ) + if (droppedReasoning > 0) { + logger.debug("compress.reasoning: dropped oversized reasoning parts", { + dropped: droppedReasoning, + threshold: reasoningConfig?.threshold ?? + DEFAULT_COMPRESS_REASONING.threshold, + }) + } + } + + ensureBuiltinFiltersRegistered() + const effectiveLimit = resolveEffectiveContextLimit(state, config) + applyMessageFilters(output.messages, config.messageFilters, logger, { + sessionId: state.sessionId ?? "", + isSubAgent: state.isSubAgent, + modelContextLimit: effectiveLimit?.limit, + }) + cacheSystemPromptTokens(state, output.messages) + assignMessageRefs(state, output.messages) + const activeBlockCountBefore = state.prune.messages.activeBlockIds.size // [FIX Bug 4] + const compressionStateChanged = syncCompressionBlocks(state, logger, output.messages) + if ( + compressionStateChanged || + state.prune.messages.activeBlockIds.size !== activeBlockCountBefore + ) { + // [FIX Bug 4] + saveSessionState(state, logger).catch(() => {}) // [FIX Bug 4] persist deactivations + } + syncToolCache(state, config, logger, output.messages) + buildToolIdList(state, output.messages) + const batchResult = runBatchCleanup(state, config, logger, output.messages) + if (batchResult.mergedCount > 0) { + saveSessionState(state, logger).catch(() => {}) + } + const prePruneTokens = getCurrentTokenUsage(state, output.messages) + // Keep the full post-filter projection for candidate planning. The + // nudge receives a pruned view, while range validation still needs the + // original ordering to prove tool-pair and protection parity. + // Skip the copy entirely when candidates are disabled (default). + const candidateMessages = + config.compress.candidates === true ? output.messages.slice() : undefined + prune(state, logger, config, output.messages) + hideConsumedCompressCalls(state, output.messages) + assignMessageRefs(state, output.messages) + const compressionPriorities = buildPriorityMap(config, state, output.messages) + prompts.reload() + injectCompressNudges( + state, + config, + logger, + output.messages, + prompts.getRuntimePrompts(), + compressionPriorities, + config.debug + ? (text: string) => { + // sendIgnoredMessage writes an ignored:true user msg to DB. + // opencode's runtime loop detects it as "last user" (role-only, + // ignores the flag) → phantom turn → compress → notification → + // infinite loop. Use logger.debug + toast instead. + logger.debug(`[ACP Debug] Nudge injected:\n${text}`) + client.tui + .showToast({ + body: { + title: "ACP: Nudge Injected", + message: text.slice(0, 500), + variant: "info", + duration: 5000, + }, + }) + .catch(() => {}) + } + : undefined, + prePruneTokens, + candidateMessages, + ) + // Candidate planning consumes the pre-truncation snapshot so its + // executor-parity check matches the fresh messages fetched by compress. + truncateLargeToolOutputs( + state, + config, + logger, + output.messages.filter((message) => !isSyntheticMessage(message)), + ) + // Keep candidate planning independent from the final budget guard: + // candidates are computed from the pre-truncation snapshot so their + // ranges remain valid when the compress executor fetches fresh history. + enforceContextBudget(state, config, logger, output.messages) + injectMessageIds(state, config, output.messages, compressionPriorities) + hideFailedCompressCalls(output.messages) + stripStaleMetadata(output.messages) + dropEmptyMessages(output.messages) + const postTokens = getCurrentTokenUsage(state, output.messages) + // [FIX #346] Hard guard: if the post-transform context still exceeds + // the model's real request budget (window minus system prompt + tool + // schemas + output-token reserve), the backend will reject the + // request. opencode exits 0 with zero output in that case (upstream + // behavior the plugin cannot change), so this ERROR is the only + // signal that the session has hit the length-rejection wall. + if (postTokens !== undefined && effectiveLimit) { + const budget = + effectiveLimit.limit - (state.systemPromptTokens ?? 0) - OUTPUT_RESERVE_TOKENS + if (postTokens > budget) { + logger.error( + "ACP hard guard: context exceeds model budget after in-flight reduction", + { + session: state.sessionId, + postTokens, + budget, + contextLimit: effectiveLimit.limit, + contextLimitSource: effectiveLimit.source, + hint: "request will likely be rejected; run /compact or start a new session", + }, + ) + } + } + logger.info("Chat transform complete", { + session: state.sessionId, + model: state.modelID, + messages: output.messages.length, + prePruneTokens, + postTokens, + contextLimit: effectiveLimit?.limit, + contextLimitSource: effectiveLimit?.source, + usagePct: + postTokens !== undefined && effectiveLimit + ? `${((postTokens / effectiveLimit.limit) * 100).toFixed(1)}%` + : undefined, + nudged: state.nudges.shouldInjectThisTurn, + }) + + if (state.sessionId) { + await logger.saveContext(state.sessionId, output.messages) + } + } + if (!lastUserMessage) { // Ephemeral state: no session to resolve, but keep running // state-independent stages (e.g. stripHallucinations). - state = createSessionState() - } else { + await runPipeline(createSessionState()) + return + } + + // [Issue #404] Serialize the entire same-session request path (init, + // model-limit reconciliation, mutations, persistence) so concurrent + // transforms of one session cannot interleave at awaits and persist + // stale snapshots over each other's committed state. + await registry.withSessionGuard(lastUserMessage.info.sessionID, async () => { // [FIX #33] Per-session state: each session keeps its own SessionState, // so interleaved sessions no longer reset each other's modelContextLimit. - state = await registry.getOrCreate( + const state = await registry.getOrCreate( client, lastUserMessage.info.sessionID, messages, @@ -292,169 +473,8 @@ export function createChatMessageTransformHandler( }, ) } - } - - syncCompressPermissionState(state, config, hostPermissions, output.messages) - - if (state.isSubAgent && !config.allowSubAgents) { - return - } - - stripHallucinations(output.messages) - - // [#368] Drop oversized reasoning from closed-turn compress tool calls. - // compress calls are hard-exempt from compression (Bug 39), so their - // thinking otherwise rides along every request as an unreclaimable - // floor. Gated by the nested `compress.reasoning` config, resolved - // through the #344 cascade (model > provider > global) using THIS - // request's model identity (request metadata first, session state as - // fallback). Runs BEFORE token accounting / pruning so every later - // stage sees the post-drop array. - const dropReasoningModel = ( - lastUserMessage?.info as - { model?: { providerID?: string; modelID?: string } } | undefined - )?.model - const reasoningConfig = applyCompressOverrides( - config, - dropReasoningModel?.providerID ?? state.modelProviderID, - dropReasoningModel?.modelID ?? state.modelID, - ).compress.reasoning - if (reasoningConfig?.drop !== false) { - const droppedReasoning = dropCompressReasoning( - output.messages, - reasoningConfig?.threshold ?? DEFAULT_COMPRESS_REASONING.threshold, - ) - if (droppedReasoning > 0) { - logger.debug("compress.reasoning: dropped oversized reasoning parts", { - dropped: droppedReasoning, - threshold: reasoningConfig?.threshold ?? DEFAULT_COMPRESS_REASONING.threshold, - }) - } - } - - ensureBuiltinFiltersRegistered() - const effectiveLimit = resolveEffectiveContextLimit(state, config) - applyMessageFilters(output.messages, config.messageFilters, logger, { - sessionId: state.sessionId ?? "", - isSubAgent: state.isSubAgent, - modelContextLimit: effectiveLimit?.limit, - }) - cacheSystemPromptTokens(state, output.messages) - assignMessageRefs(state, output.messages) - const activeBlockCountBefore = state.prune.messages.activeBlockIds.size // [FIX Bug 4] - const compressionStateChanged = syncCompressionBlocks(state, logger, output.messages) - if ( - compressionStateChanged || - state.prune.messages.activeBlockIds.size !== activeBlockCountBefore - ) { - // [FIX Bug 4] - saveSessionState(state, logger).catch(() => {}) // [FIX Bug 4] persist deactivations - } - syncToolCache(state, config, logger, output.messages) - buildToolIdList(state, output.messages) - const batchResult = runBatchCleanup(state, config, logger, output.messages) - if (batchResult.mergedCount > 0) { - saveSessionState(state, logger).catch(() => {}) - } - const prePruneTokens = getCurrentTokenUsage(state, output.messages) - // Keep the full post-filter projection for candidate planning. The - // nudge receives a pruned view, while range validation still needs the - // original ordering to prove tool-pair and protection parity. - // Skip the copy entirely when candidates are disabled (default). - const candidateMessages = - config.compress.candidates === true ? output.messages.slice() : undefined - prune(state, logger, config, output.messages) - hideConsumedCompressCalls(state, output.messages) - assignMessageRefs(state, output.messages) - const compressionPriorities = buildPriorityMap(config, state, output.messages) - prompts.reload() - injectCompressNudges( - state, - config, - logger, - output.messages, - prompts.getRuntimePrompts(), - compressionPriorities, - config.debug - ? (text: string) => { - // sendIgnoredMessage writes an ignored:true user msg to DB. - // opencode's runtime loop detects it as "last user" (role-only, - // ignores the flag) → phantom turn → compress → notification → - // infinite loop. Use logger.debug + toast instead. - logger.debug(`[ACP Debug] Nudge injected:\n${text}`) - client.tui - .showToast({ - body: { - title: "ACP: Nudge Injected", - message: text.slice(0, 500), - variant: "info", - duration: 5000, - }, - }) - .catch(() => {}) - } - : undefined, - prePruneTokens, - candidateMessages, - ) - // Candidate planning consumes the pre-truncation snapshot so its - // executor-parity check matches the fresh messages fetched by compress. - truncateLargeToolOutputs( - state, - config, - logger, - output.messages.filter((message) => !isSyntheticMessage(message)), - ) - // Keep candidate planning independent from the final budget guard: - // candidates are computed from the pre-truncation snapshot so their - // ranges remain valid when the compress executor fetches fresh history. - enforceContextBudget(state, config, logger, output.messages) - injectMessageIds(state, config, output.messages, compressionPriorities) - hideFailedCompressCalls(output.messages) - stripStaleMetadata(output.messages) - dropEmptyMessages(output.messages) - const postTokens = getCurrentTokenUsage(state, output.messages) - // [FIX #346] Hard guard: if the post-transform context still exceeds - // the model's real request budget (window minus system prompt + tool - // schemas + output-token reserve), the backend will reject the - // request. opencode exits 0 with zero output in that case (upstream - // behavior the plugin cannot change), so this ERROR is the only - // signal that the session has hit the length-rejection wall. - if (postTokens !== undefined && effectiveLimit) { - const budget = - effectiveLimit.limit - (state.systemPromptTokens ?? 0) - OUTPUT_RESERVE_TOKENS - if (postTokens > budget) { - logger.error( - "ACP hard guard: context exceeds model budget after in-flight reduction", - { - session: state.sessionId, - postTokens, - budget, - contextLimit: effectiveLimit.limit, - contextLimitSource: effectiveLimit.source, - hint: "request will likely be rejected; run /compact or start a new session", - }, - ) - } - } - logger.info("Chat transform complete", { - session: state.sessionId, - model: state.modelID, - messages: output.messages.length, - prePruneTokens, - postTokens, - contextLimit: effectiveLimit?.limit, - contextLimitSource: effectiveLimit?.source, - usagePct: - postTokens !== undefined && effectiveLimit - ? `${((postTokens / effectiveLimit.limit) * 100).toFixed(1)}%` - : undefined, - nudged: state.nudges.shouldInjectThisTurn, + await runPipeline(state) }) - - if (state.sessionId) { - await logger.saveContext(state.sessionId, output.messages) - } } } @@ -495,9 +515,15 @@ export function createCommandExecuteHandler( }) const messages = filterMessages(messagesResponse.data || messagesResponse) - const state = await registry.getOrCreate(client, input.sessionID, messages, config) - - syncCompressPermissionState(state, config, hostPermissions, messages) + // [Issue #404] getOrCreate can initialize + persist, and sync + // recomputes permission state in memory: keep both under the + // per-session guard so command-triggered mutations follow the same + // single-writer invariant as transforms and tools. + const state = await registry.withSessionGuard(input.sessionID, async () => { + const s = await registry.getOrCreate(client, input.sessionID, messages, config) + syncCompressPermissionState(s, config, hostPermissions, messages) + return s + }) const commandCtx = { client, @@ -603,18 +629,35 @@ export function createEventHandler(registry: SessionStateRegistry, logger: Logge durationMs, }) + // [Issue #404] Apply mutates block durations and the save persists + // them: serialize against same-session transforms so a duration + // applied here cannot clobber (or be clobbered by) a concurrent + // transform's snapshot. Sessions are independent, so acquisitions + // run concurrently across sessions — one long-held guard (e.g. an + // open permission dialog) must not delay other sessions' duration + // attachment. Per-session failures stay contained (the persistence + // layer logs save errors); one failing session must not skip the rest. + const settles: Array> = [] for (const state of registry.all()) { - const updates = applyPendingCompressionDurations(state) - if (updates > 0) { - await saveSessionState(state, logger) - logger.info("Attached compression time to blocks", { - messageID: part.messageID, - callID: part.callID, - blocks: updates, - durationMs, - }) + if (!state.sessionId) { + continue } + settles.push( + registry.withSessionGuard(state.sessionId, async () => { + const updates = applyPendingCompressionDurations(state) + if (updates > 0) { + await saveSessionState(state, logger) + logger.info("Attached compression time to blocks", { + messageID: part.messageID, + callID: part.callID, + blocks: updates, + durationMs, + }) + } + }), + ) } + await Promise.allSettled(settles) return } diff --git a/lib/state/state.ts b/lib/state/state.ts index 262c23d9..d076f733 100644 --- a/lib/state/state.ts +++ b/lib/state/state.ts @@ -60,6 +60,50 @@ export async function updatePerTurnState( // access; modelContextLimit and all persisted fields survive eviction. const REGISTRY_SOFT_CAP = 32 +/** + * [Issue #404] Per-session FIFO mutex serializing same-session mutations + * (message transforms, compress/decompress tools, event-hook saves, + * system-hook limit writes). Node is single-threaded, but async work + * interleaves at awaits: without this, a stale transform resuming after a + * newer one can persist last-write-wins garbage over committed state. + * Chain-based promise queue (not a boolean flag): waiters line up in arrival + * order; the entry is deleted once the tail task finishes, so the map stays + * empty when idle and cannot leak. + * + * Run `fn` serialized against every other guard user of the same sessionId. + * MUST NOT be nested for the same session (deadlock) — callers acquire it + * exactly once at their outermost boundary. + */ +export type SessionGuard = (sessionId: string, fn: () => Promise | T) => Promise + +export function createSessionGuard(): SessionGuard { + const locks = new Map>() + return async (sessionId, fn) => { + const previous = locks.get(sessionId) ?? Promise.resolve() + let release!: () => void + const released = new Promise((resolve) => { + release = resolve + }) + const chain = previous.then(() => released) + locks.set(sessionId, chain) + try { + await previous + } catch { + // A rejected task must not poison the queue for later tasks. + } + try { + return await fn() + } finally { + release() + // Delete only when we are still the tail — a newer waiter may have + // already chained onto us. + if (locks.get(sessionId) === chain) { + locks.delete(sessionId) + } + } + } +} + // [FIX #33] Per-session state. Replaces the single shared SessionState singleton // whose resetSessionState-on-switch wiped modelContextLimit (set only by // system.transform, which fires AFTER messages.transform) and flipped @@ -150,9 +194,18 @@ export class SessionStateRegistry { return this.states.size } + // [Issue #404] Per-session FIFO mutex serializing same-session mutations + // (message transforms, compress/decompress tools, event-hook saves, + // system-hook limit writes). Shared factory so test stubs compose the + // exact same implementation instead of drifting. + // Not readonly: the session-guard wiring test wraps this to observe + // acquisition order (Issue #404). + withSessionGuard: SessionGuard = createSessionGuard() + // Idempotent: ensureSessionInitialized returns immediately once - // state.sessionId === sessionId (assigned synchronously before any await), - // so repeat calls for the same session never re-reset. + // state.sessionId === sessionId, and coalesces concurrent inits per state + // object ([Issue #404]), so racing calls for the same session never + // re-reset or observe partially initialized state. async getOrCreate( client: any, sessionId: string, @@ -291,6 +344,15 @@ export function resetSessionState(state: SessionState): void { state.noContextLimitWarned = false } +// [Issue #404] In-flight init coalescing keyed by state object: concurrent +// callers (e.g. two same-session message transforms racing on the first +// request) share one init instead of racing past each other — the old +// synchronous `state.sessionId = sessionId` assignment let a second caller +// early-return and observe partially initialized state. Keyed by the state +// OBJECT (WeakMap) so soft-cap eviction + recreation starts a fresh init +// rather than awaiting a promise tied to an evicted instance. +const inflightInits = new WeakMap>() + export async function ensureSessionInitialized( client: any, state: SessionState, @@ -300,10 +362,43 @@ export async function ensureSessionInitialized( config?: PluginConfig, projectDir?: string, ): Promise { + // In-flight check FIRST: runSessionInitialization assigns + // state.sessionId synchronously before its first await, so the fast path + // below would otherwise let racing callers early-return mid-init. + const inflight = inflightInits.get(state) + if (inflight) { + await inflight + return + } if (state.sessionId === sessionId) { return } + const run = runSessionInitialization( + client, + state, + sessionId, + logger, + messages, + config, + projectDir, + ) + inflightInits.set(state, run) + try { + await run + } finally { + inflightInits.delete(state) + } +} +async function runSessionInitialization( + client: any, + state: SessionState, + sessionId: string, + logger: Logger, + messages: WithParts[], + config?: PluginConfig, + projectDir?: string, +): Promise { resetSessionState(state) state.sessionId = sessionId // Resolve the configured storage location once per session (transient). diff --git a/tests/registry-stub.ts b/tests/registry-stub.ts index 4793f6e9..0e713506 100644 --- a/tests/registry-stub.ts +++ b/tests/registry-stub.ts @@ -1,4 +1,4 @@ -import { createSessionState, type SessionState } from "../lib/state" +import { createSessionGuard, createSessionState, type SessionState } from "../lib/state" import { createModelLimitCatalog } from "../lib/state/model-limits" // Test helper: builds a SessionStateRegistry stub that resolves any sessionID @@ -9,6 +9,7 @@ export function singletonRegistry(state: SessionState): { all: () => SessionState[] size: number compressionTiming: SessionState["compressionTiming"] + withSessionGuard: ReturnType } { return { get: () => state, @@ -17,6 +18,7 @@ export function singletonRegistry(state: SessionState): { return 1 }, compressionTiming: state.compressionTiming, + withSessionGuard: createSessionGuard(), } } @@ -35,6 +37,7 @@ export function createTestRegistry(seedState: SessionState) { let hydratedOnce = false return { compressionTiming: sharedTiming, + withSessionGuard: createSessionGuard(), get size() { return states.size }, diff --git a/tests/session-guard.test.ts b/tests/session-guard.test.ts new file mode 100644 index 00000000..4aa140e7 --- /dev/null +++ b/tests/session-guard.test.ts @@ -0,0 +1,401 @@ +import "./test-env" +/** + * Tests for [Issue #404]: serialize same-session state initialization and + * transforms. + * + * Covers: + * - withSessionGuard FIFO ordering per session (and cross-session independence) + * - guard release on rejection + * - concurrent getOrCreate coalescing into a single initialization + * (the stale-snapshot regression from the issue) + * - failed init coalescing semantics (waiters resolve like today) + * - read-modify-write safety: a stale transaction cannot overwrite newer + * committed state while serialized + * - transform-handler wiring: concurrent same-session requests serialize + * through the guard (end-to-end, real createChatMessageTransformHandler) + * - restoreCompressionState preserves shared compressionTiming identity + */ + +import assert from "node:assert/strict" +import test, { beforeEach, afterEach } from "node:test" +import { mkdtempSync, rmSync } from "node:fs" +import { join } from "node:path" +import { tmpdir } from "node:os" +import { + SessionStateRegistry, + createSessionGuard, + createSessionState, + saveSessionState, + type WithParts, +} from "../lib/state" +import { snapshotCompressionState, restoreCompressionState } from "../lib/compress/pipeline" +import { createChatMessageTransformHandler } from "../lib/hooks" +import type { PluginConfig } from "../lib/config" +import { Logger } from "../lib/logger" + +const MESSAGES: WithParts[] = [] + +let tempDir: string +let prevData: string | undefined +let prevConfig: string | undefined + +beforeEach(() => { + tempDir = mkdtempSync(join(tmpdir(), "acp-guard-")) + prevData = process.env.XDG_DATA_HOME + prevConfig = process.env.XDG_CONFIG_HOME + process.env.XDG_DATA_HOME = tempDir + process.env.XDG_CONFIG_HOME = tempDir +}) + +afterEach(() => { + if (prevData === undefined) delete process.env.XDG_DATA_HOME + else process.env.XDG_DATA_HOME = prevData + if (prevConfig === undefined) delete process.env.XDG_CONFIG_HOME + else process.env.XDG_CONFIG_HOME = prevConfig + rmSync(tempDir, { recursive: true, force: true }) +}) + +test("withSessionGuard serializes same-session tasks in FIFO order", async () => { + const guard = createSessionGuard() + const order: string[] = [] + const run = (id: string, holdMs: number) => + guard("s1", async () => { + order.push(`start-${id}`) + await new Promise((r) => setTimeout(r, holdMs)) + order.push(`end-${id}`) + }) + await Promise.all([run("A", 30), run("B", 5), run("C", 1)]) + assert.deepEqual(order, ["start-A", "end-A", "start-B", "end-B", "start-C", "end-C"]) +}) + +test("withSessionGuard does not block other sessions", async () => { + const guard = createSessionGuard() + let s1Finished = false + const t1 = guard("s1", async () => { + await new Promise((r) => setTimeout(r, 60)) + s1Finished = true + }) + await new Promise((r) => setTimeout(r, 10)) + // s2 must not wait behind s1's long task. + const t2 = guard("s2", async () => "ok") + assert.equal(await t2, "ok") + assert.equal(s1Finished, false) + await t1 + assert.equal(s1Finished, true) +}) + +test("withSessionGuard releases the lock when fn rejects and passes the value through", async () => { + const guard = createSessionGuard() + await assert.rejects(guard("s", async () => { + throw new Error("boom") + }), /boom/) + assert.equal(await guard("s", async () => 42), 42) + assert.equal(await guard("s", () => 7), 7) +}) + +test("concurrent getOrCreate coalesces into one initialization (stale-snapshot regression)", async () => { + // Seed persisted state so partial initialization is observable: a caller + // that early-returns before loadSessionState completes sees no limit. + const logger = new Logger(false) + const seed = createSessionState() + seed.sessionId = "session-x" + seed.modelContextLimit = 200000 + await saveSessionState(seed, logger) + + let sessionGets = 0 + const client = { + session: { + get: async () => { + sessionGets++ + await new Promise((r) => setTimeout(r, 25)) + return { data: { parentID: null } } + }, + }, + } + + const registry = new SessionStateRegistry(new Logger(false)) + // Capture what the second caller observes AT RETURN time — the pre-fix + // early-return handed out the state before loadSessionState finished, so + // this read saw the pre-load snapshot (undefined limit). + let bLimitAtReturn: number | undefined + const pA = registry.getOrCreate(client, "session-x", MESSAGES) + const pB = registry.getOrCreate(client, "session-x", MESSAGES).then((s) => { + bLimitAtReturn = s.modelContextLimit + return s + }) + const pC = registry.getOrCreate(client, "session-x", MESSAGES) + const [a, b, c] = await Promise.all([pA, pB, pC]) + + assert.equal(a, b) + assert.equal(b, c) + assert.equal(sessionGets, 1) + // Every waiter must observe the fully initialized state, not a snapshot + // taken before the persisted data was loaded. + assert.equal(bLimitAtReturn, 200000) + assert.equal(a.modelContextLimit, 200000) + assert.equal(b.modelContextLimit, 200000) + assert.equal(c.modelContextLimit, 200000) +}) + +test("failed init coalesces: all waiters resolve with the same partially-initialized state", async () => { + const client = { + session: { + get: async () => { + throw new Error("host down") + }, + }, + } + const registry = new SessionStateRegistry(new Logger(false)) + const results = await Promise.allSettled([ + registry.getOrCreate(client, "s-fail", MESSAGES), + registry.getOrCreate(client, "s-fail", MESSAGES), + ]) + assert.equal(results[0].status, "fulfilled") + assert.equal(results[1].status, "fulfilled") + const a = (results[0] as PromiseFulfilledResult>).value + const b = (results[1] as PromiseFulfilledResult>).value + assert.equal(a, b) + // sessionId is assigned synchronously at init start (pre-existing semantics). + assert.equal(a.sessionId, "s-fail") +}) + +test("serialized read-modify-write: stale transaction cannot clobber newer committed state", async () => { + const guard = createSessionGuard() + const state = { value: 0 } + const persistLog: number[] = [] + const persist = () => { + persistLog.push(state.value) + } + + // T_old snapshots, does slow work, then commits snapshot+1. + const tOld = guard("s", async () => { + const snap = state.value + await new Promise((r) => setTimeout(r, 40)) + state.value = snap + 1 + persist() + }) + // T_new queues behind T_old and does its own read-modify-write. + const tNew = guard("s", async () => { + state.value = state.value * 10 + 5 + persist() + }) + await Promise.all([tOld, tNew]) + + // Serialized order: old commits 0+1=1, then new reads 1 → 15. + // An unsynchronized interleaving could persist [5, 1] (stale last write). + assert.deepEqual(persistLog, [1, 15]) + assert.equal(state.value, 15) +}) + +test("restoreCompressionState preserves the shared compressionTiming object identity", async () => { + const state = createSessionState() + const timingBefore = state.compressionTiming + assert.ok(timingBefore !== null) + + const snapshot = snapshotCompressionState(state) + state.prune.messages.activeBlockIds.add(99) + restoreCompressionState(state, snapshot) + + assert.equal(state.prune.messages.activeBlockIds.size, 0) + assert.equal(state.compressionTiming, timingBefore) +}) + +test("concurrent transforms on one session serialize through the guard (wiring regression)", async () => { + // End-to-end wiring check for [Issue #404]: two chat-message-transform + // requests for the SAME session must not interleave their + // init→mutate→persist sections. Pre-fix, the second request's getOrCreate + // hit the sessionId fast path mid-init and ran its pipeline concurrently + // on partially initialized state (stale-snapshot corruption). We observe + // guard acquisition order through the registry seam: serialized handlers + // produce strictly nested enter/exit pairs; unsynchronized ones interleave + // (or, before the guard existed, never acquire at all). + const logger = new Logger(false) + const seed = createSessionState() + seed.sessionId = "session-wire" + seed.modelContextLimit = 200000 + // [FIX #312] A persisted limit only survives the transform's reconciliation + // when its model identity matches the request; seed the pair like a real + // session so the loaded value is what we assert on. + seed.modelProviderID = "test-provider" + seed.modelID = "test-model" + await saveSessionState(seed, logger) + + const client = { + session: { + get: async () => { + // Hold the init window open so concurrent requests genuinely + // overlap (pre-fix interleaving requires this window). + await new Promise((r) => setTimeout(r, 25)) + return { data: { parentID: null } } + }, + }, + } + + const config: PluginConfig = { + enabled: true, + autoUpdate: true, + debug: false, + pruneNotification: "off", + pruneNotificationType: "chat", + commands: { enabled: true, protectedTools: [] }, + experimental: { allowSubAgents: false, customPrompts: false }, + protectedFilePatterns: [], + compress: { + mode: "message", + permission: "allow", + showCompression: false, + summaryBuffer: true, + candidates: true, + maxContextLimit: 150000, + minContextLimit: 50000, + nudgeFrequency: 5, + iterationNudgeThreshold: 15, + nudgeForce: "soft", + protectedTools: ["task"], + protectTags: false, + protectUserMessages: false, + reasoning: { drop: true, threshold: 2048 }, + }, + gc: { + algorithm: "truncate", + promotionThreshold: 5, + maxBlockAge: 15, + maxOldGenSummaryLength: 3000, + majorGcThresholdPercent: "100%", + batchCleanup: { lowThreshold: "60%", highThreshold: "75%", forceThreshold: "90%" }, + }, + } + const prompts = { + reload() {}, + getRuntimePrompts() { + return { + system: "ACP system", + compressRange: "compress range", + compressMessage: "compress message", + contextLimitNudge: "nudge", + turnNudge: "turn nudge", + iterationNudge: "iteration nudge", + manualExtension: "", + subagentExtension: "", + } + }, + } + const hostPermissions = { global: undefined, agents: {} } + + const registry = new SessionStateRegistry(logger) + const events: string[] = [] + const originalGuard = registry.withSessionGuard + // Markers are pushed INSIDE the locked fn: "enter" means the task actually + // started running under the lock (not merely arrived at the FIFO queue), so + // strictly nested pairs prove serialization while interleave proves overlap. + registry.withSessionGuard = (sessionId: string, fn: () => Promise | T): Promise => + originalGuard(sessionId, async (): Promise => { + events.push(`enter:${sessionId}`) + try { + return await fn() + } finally { + events.push(`exit:${sessionId}`) + } + }) + + const handler = createChatMessageTransformHandler( + client, + registry, + logger, + config, + prompts, + hostPermissions, + ) + + const mkMessages = (): WithParts[] => [ + { + info: { + id: "u1", + sessionID: "session-wire", + role: "user", + agent: "assistant", + model: { providerID: "test-provider", modelID: "test-model" }, + time: { created: Date.now() }, + } as WithParts["info"], + parts: [ + { + type: "text", + text: "hello", + id: "u1-p1", + sessionID: "session-wire", + messageID: "u1", + }, + ], + }, + { + info: { + id: "a1", + sessionID: "session-wire", + role: "assistant", + agent: "assistant", + parentID: "parent-placeholder", + modelID: "test-model", + providerID: "test-provider", + mode: "normal", + path: { cwd: "/", root: "/" }, + summary: false, + cost: 0, + tokens: { input: 100, output: 50, reasoning: 0, cache: { read: 0, write: 0 } }, + time: { created: Date.now() }, + } as WithParts["info"], + parts: [ + { + type: "step-start", + id: "a1-ss", + sessionID: "session-wire", + messageID: "a1", + }, + { + type: "text", + text: "hi there", + id: "a1-p1", + sessionID: "session-wire", + messageID: "a1", + }, + ], + }, + { + info: { + id: "u2", + sessionID: "session-wire", + role: "user", + agent: "assistant", + model: { providerID: "test-provider", modelID: "test-model" }, + time: { created: Date.now() }, + } as WithParts["info"], + parts: [ + { + type: "text", + text: "again", + id: "u2-p1", + sessionID: "session-wire", + messageID: "u2", + }, + ], + }, + ] + + await Promise.all([handler({}, { messages: mkMessages() }), handler({}, { messages: mkMessages() })]) + + // Same-session work ran strictly serially (nested pairs, no overlap). An + // unsynchronized pre-fix interleave would look like enter,enter,exit,exit — + // or produce no events at all before the guard existed. + assert.deepEqual(events, [ + "enter:session-wire", + "exit:session-wire", + "enter:session-wire", + "exit:session-wire", + ]) + // Both requests observed fully initialized state (persisted limit loaded), + // and refs were assigned for the shared input. + const state = registry.get("session-wire") + assert.ok(state) + assert.equal(state.modelContextLimit, 200000) + assert.equal(state.messageIds.byRef.get("m00001"), "u1") + assert.equal(state.messageIds.byRef.get("m00003"), "u2") +})