From d90a117c7216dc2c3d6b689be0ee326e92d4654e Mon Sep 17 00:00:00 2001 From: phantom5099 <1011668688@qq.com> Date: Sun, 4 Oct 2026 19:47:09 +0800 Subject: [PATCH 1/2] =?UTF-8?q?=E5=A2=9E=E5=8A=A0eventsink=E9=80=9A?= =?UTF-8?q?=E9=81=93=EF=BC=8C=E5=9B=9E=E5=BD=92=E4=B8=8A=E4=B8=8B=E6=96=87?= =?UTF-8?q?=E5=86=85=E5=AD=98=E8=A1=8C=E4=B8=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- packages/codingcode/src/agent/agent.ts | 98 +++++----- packages/codingcode/src/approval/approval.ts | 9 +- .../codingcode/src/approval/confirmation.ts | 6 +- packages/codingcode/src/approval/wait-port.ts | 7 +- packages/codingcode/src/approval/wait.ts | 44 +---- packages/codingcode/src/context/context.ts | 170 +++++++++--------- packages/codingcode/src/context/port.ts | 13 +- packages/codingcode/src/layer.ts | 12 +- packages/codingcode/src/server/handler.ts | 10 +- .../codingcode/src/server/routes/sessions.ts | 4 +- packages/codingcode/src/sink/port.ts | 16 ++ packages/codingcode/src/sink/sink.ts | 26 +++ .../test/agent/context-compressed.test.ts | 4 +- .../test/approval/async-confirm.test.ts | 69 +++---- .../codingcode/test/approval/pipeline.test.ts | 17 +- .../test/context/budget-integration.test.ts | 12 +- .../test/context/compressor/behavior.test.ts | 19 +- .../test/context/memory-buffer.test.ts | 147 +++++++++++++++ .../codingcode/test/helpers/agent-harness.ts | 42 +++-- .../test/plan/gate-pipeline.test.ts | 24 ++- .../security/plan-profile-restart.test.ts | 9 +- .../test/server/compact-route.test.ts | 12 +- packages/codingcode/test/server/index.test.ts | 11 +- .../messages-fork-permission-mode.test.ts | 6 +- .../test/server/plan-file-route.test.ts | 10 +- packages/codingcode/test/sink/sink.test.ts | 87 +++++++++ .../test/subagent/dispatch-end-to-end.test.ts | 8 +- .../subagent/dispatch-production-path.test.ts | 8 +- .../test/subagent/runner-wiring.test.ts | 8 +- 29 files changed, 605 insertions(+), 303 deletions(-) create mode 100644 packages/codingcode/src/sink/port.ts create mode 100644 packages/codingcode/src/sink/sink.ts create mode 100644 packages/codingcode/test/context/memory-buffer.test.ts create mode 100644 packages/codingcode/test/sink/sink.test.ts diff --git a/packages/codingcode/src/agent/agent.ts b/packages/codingcode/src/agent/agent.ts index 21d64ae5..b26a8905 100644 --- a/packages/codingcode/src/agent/agent.ts +++ b/packages/codingcode/src/agent/agent.ts @@ -6,6 +6,7 @@ import type { RunTurnOptions, ToolEnv } from './port.js'; import { ApprovalService } from '../approval/port.js'; import { CheckpointService } from '../checkpoint/port.js'; import { ContextService } from '../context/port.js'; +import { EventSinkService } from '../sink/port.js'; import { HookService } from '../hooks/port.js'; import { LLMService } from '../llm/port.js'; import { McpService } from '../mcp/port.js'; @@ -49,6 +50,7 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { const skills = yield* SkillService; const mcp = yield* McpService; const context = yield* ContextService; + const sink = yield* EventSinkService; const memory = yield* MemoryService; const llm = yield* LLMService; const rules = yield* RulesService; @@ -129,13 +131,18 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { // get rules text const rulesText = yield* rules.getAllRules(state.cwd); - // run agent loop - const stream = runAgentLoop({ - state, model, profile, catalog, systemPrompt: opts.systemPrompt, - toolEnv, - abortSignal: opts.signal, rulesText, - sid: sessionId, projectPath: state.cwd, permissionMode: effectivePerm, - }); + // run agent loop:出站队列挂在 sink 上,本回合是它的唯一读者 + const q = yield* sink.attach(sessionId); + const emit = (body: FrameBody) => Effect.runSync(sink.emit(sessionId, body)); + const stream = runAgentLoop( + { + state, model, profile, catalog, systemPrompt: opts.systemPrompt, + toolEnv, + abortSignal: opts.signal, rulesText, + sid: sessionId, projectPath: state.cwd, permissionMode: effectivePerm, + }, + { q, emit, onEnd: () => Effect.runSync(sink.detach(sessionId)) } + ); return { stream, sessionId }; }); @@ -148,10 +155,12 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { rulesText: string; sid: string; projectPath: string; permissionMode: PermissionMode; systemPrompt?: string; + }, out: { + q: Queue.Queue; + emit: (body: FrameBody) => void; + onEnd: () => void; }): AsyncGenerator { - const q = Effect.runSync(Queue.unbounded()); - - const program = agentLoopInternal(opts, q); + const program = agentLoopInternal(opts, out.emit); return (async function* () { const fiber = Effect.runFork(opts.toolEnv.provide(program)); @@ -162,11 +171,15 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { if (opts.abortSignal.aborted) Effect.runFork(Fiber.interrupt(fiber)); } - const stream = Stream.fromQueue(q).pipe( - Stream.takeUntil((body: FrameBody) => isTurnEnd(body)) - ); - for await (const body of Stream.toAsyncIterable(stream) as AsyncIterable) { - yield body; + try { + const stream = Stream.fromQueue(out.q).pipe( + Stream.takeUntil((body: FrameBody) => isTurnEnd(body)) + ); + for await (const body of Stream.toAsyncIterable(stream) as AsyncIterable) { + yield body; + } + } finally { + out.onEnd(); } })(); } @@ -178,7 +191,7 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { rulesText: string; sid: string; projectPath: string; permissionMode: PermissionMode; systemPrompt?: string; - }, q: Queue.Queue): Effect.Effect, AgentError> { + }, emit: (body: FrameBody) => void): Effect.Effect, AgentError> { const { state, model, profile, abortSignal, catalog, rulesText, sid, projectPath, permissionMode } = opts; const { tools, lookup: toolLookup } = catalog; @@ -187,7 +200,7 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { Effect.sync(() => { if (ended) return; ended = true; - Effect.runSync(q.offer({ family: 'transition', transition })); + emit({ family: 'transition', transition }); }); return Effect.gen(function* () { @@ -209,13 +222,13 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { let lastResult: Result | null = null; yield* hooks.emit('agent.turn.start', { sessionId: sid, projectPath }); - yield* q.offer({ family: 'transition', transition: { to: 'start', turnId: state.currentTurnId } }); + emit({ family: 'transition', transition: { to: 'start', turnId: state.currentTurnId } }); for (let step = 0; step < maxSteps; step++) { yield* hooks.emitDecision('agent.step.before', { sessionId: sid, step: step + 1, projectPath }); if (step === 0) { - yield* q.offer({ family: 'transition', transition: { to: 'executing' } }); + emit({ family: 'transition', transition: { to: 'executing' } }); } const sessionRef: SessionRef = { @@ -225,27 +238,14 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { currentTurnId: state.currentTurnId, }; - const willCompact = yield* Effect.either(context.willCompact(sessionRef, model)); - if (Either.isLeft(willCompact)) { - yield* offerEnd({ to: 'end', reason: 'error', error: toFrameError(willCompact.left) }); + const history = yield* Effect.either(context.getHistory(sessionRef, model)); + if (Either.isLeft(history)) { + yield* offerEnd({ to: 'end', reason: 'error', error: toFrameError(history.left) }); yield* hooks.emit('agent.turn.end', { sessionId: sid, turnId: state.currentTurnId, status: 'error', projectPath }); - return Result.err(willCompact.left); - } - if (willCompact.right) { - yield* q.offer({ family: 'transition', transition: { to: 'compress' } }); - } - - const assembled = yield* Effect.either(context.assemblePayload(sessionRef, model)); - if (Either.isLeft(assembled)) { - yield* offerEnd({ to: 'end', reason: 'error', error: toFrameError(assembled.left) }); - yield* hooks.emit('agent.turn.end', { sessionId: sid, turnId: state.currentTurnId, status: 'error', projectPath }); - return Result.err(assembled.left); - } - if (willCompact.right) { - yield* q.offer({ family: 'transition', transition: { to: 'executing' } }); + return Result.err(history.left); } - const llmMessages = [...assembled.right]; + const llmMessages = [...history.right]; let content = ''; const toolCalls: ToolCall[] = []; @@ -257,19 +257,19 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { if (abortSignal?.aborted) break; if (part.type === 'text') { content += part.text; - Effect.runSync(q.offer({ family: 'event', event: { type: 'text_delta', text: part.text } })); + emit({ family: 'event', event: { type: 'text_delta', text: part.text } }); } else if (part.type === 'tool_call') { toolCalls.push({ id: part.id, name: part.name, arguments: part.arguments }); - Effect.runSync(q.offer({ + emit({ family: 'event', event: { type: 'tool_call', id: part.id, name: part.name, args: part.arguments }, - })); + }); } else { responded = part.usage ? { usage: part.usage } : {}; - Effect.runSync(q.offer({ + emit({ family: 'transition', transition: { to: 'executing', responded }, - })); + }); } } }, @@ -282,7 +282,8 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { } if (toolCalls.length === 0) { - yield* session.recordAssistant(state, content, [], responded.usage); + const assistantEv = yield* session.recordAssistant(state, content, [], responded.usage); + yield* context.absorb(sessionRef, [assistantEv]); const stopDecision = yield* hooks.emitDecision('agent.turn.stop', { sessionId: sid, content, turnId: state.currentTurnId, projectPath }); if (stopDecision && stopDecision.decision === 'continue') { @@ -295,7 +296,8 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { } stopContinuations++; const injection = stopDecision.injection ?? '(continue)'; - yield* session.recordSystem(state, injection); + const systemEv = yield* session.recordSystem(state, injection); + yield* context.absorb(sessionRef, [systemEv]); continue; } @@ -305,7 +307,8 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { break; } - yield* session.recordAssistant(state, content, toolCalls, responded.usage); + const assistantToolEv = yield* session.recordAssistant(state, content, toolCalls, responded.usage); + yield* context.absorb(sessionRef, [assistantToolEv]); const approvedCalls: any[] = []; const deniedResults: any[] = []; @@ -342,13 +345,14 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { let todoPrinted = false; for (const r of allResults) { const resultOut = r.status === 'denied' ? '' : r.output; - yield* session.recordToolResult(state, r.name, r.id, resultOut); + const toolEv = yield* session.recordToolResult(state, r.name, r.id, resultOut); + yield* context.absorb(sessionRef, [toolEv]); const outcome = toolOutcomeOf(r); const todos = !todoPrinted && r.status === 'ok' && r.name === 'todo_write' ? todo.read(sid) : undefined; if (todos) todoPrinted = true; - yield* q.offer({ + emit({ family: 'event', event: { type: 'tool_result', id: r.id, name: r.name, outcome, diff --git a/packages/codingcode/src/approval/approval.ts b/packages/codingcode/src/approval/approval.ts index 01eeccc1..b07e537a 100644 --- a/packages/codingcode/src/approval/approval.ts +++ b/packages/codingcode/src/approval/approval.ts @@ -8,6 +8,7 @@ import { PLAN_ALLOWED_TOOLS } from '../contracts/permission.js'; import { createRuleEngine, type RuleEngine } from './rule-engine.js'; import { userConfirmAsync } from './confirmation.js'; import { ApprovalWaitService } from './wait-port.js'; +import { EventSinkService } from '../sink/port.js'; import { ApprovalService } from './port.js'; import type { ApprovalRequest } from './port.js'; @@ -89,11 +90,11 @@ function recordAuditAndReturn( export function runPipeline( request: ToolCallRequest, opts: PipelineOptions -): Effect.Effect { +): Effect.Effect { return Effect.gen(function* () { const hooks = yield* HookService; - const approvalWait = yield* ApprovalWaitService; - const asyncConfirm = yield* approvalWait.hasEmitter(opts.sessionId); + const sink = yield* EventSinkService; + const asyncConfirm = yield* sink.has(opts.sessionId); const layers: string[] = []; // Layer 1: Rule Engine @@ -206,6 +207,7 @@ export function runPipeline( export const ApprovalLayer = Layer.effect(ApprovalService, Effect.gen(function* () { const hooks = yield* HookService; + const sink = yield* EventSinkService; const approvalWait = yield* ApprovalWaitService; const ruleEngine: RuleEngine = createRuleEngine(); const destructiveTools = new Set(DANGEROUS_TOOL_NAMES); @@ -232,6 +234,7 @@ export const ApprovalLayer = Layer.effect(ApprovalService, Effect.gen(function* } ).pipe( Effect.provideService(HookService, hooks), + Effect.provideService(EventSinkService, sink), Effect.provideService(ApprovalWaitService, approvalWait) ), }; diff --git a/packages/codingcode/src/approval/confirmation.ts b/packages/codingcode/src/approval/confirmation.ts index 6c3f6045..b8dc9732 100644 --- a/packages/codingcode/src/approval/confirmation.ts +++ b/packages/codingcode/src/approval/confirmation.ts @@ -1,6 +1,7 @@ import { Effect } from 'effect'; import type { PermissionRule } from './types.js'; import { ApprovalWaitService } from './wait-port.js'; +import { EventSinkService } from '../sink/port.js'; export type ConfirmResult = | { type: 'allow' } @@ -13,12 +14,13 @@ export function userConfirmAsync( args: Record, sessionId: string, callId: string -): Effect.Effect { +): Effect.Effect { return Effect.gen(function* () { const waitSvc = yield* ApprovalWaitService; + const sink = yield* EventSinkService; const id = callId; - yield* waitSvc.emitApprovalRequest(sessionId, id, tool, args); + yield* sink.emit(sessionId, { family: 'event', event: { type: 'approval_request', id, tool, args } }); return yield* waitSvc.waitForConfirm(id, sessionId); }); diff --git a/packages/codingcode/src/approval/wait-port.ts b/packages/codingcode/src/approval/wait-port.ts index 8df55c24..cf2e21ff 100644 --- a/packages/codingcode/src/approval/wait-port.ts +++ b/packages/codingcode/src/approval/wait-port.ts @@ -5,11 +5,8 @@ import type { ConfirmResult } from './confirmation.js'; export interface ApprovalWaitShape { waitForConfirm(id: string, sessionId: string): Effect.Effect; resolveConfirm(id: string, sessionId: string, result: ConfirmResult): Effect.Effect; - emitApprovalRequest(sessionId: string, id: string, tool: string, args: Record): Effect.Effect; - registerEmitter(sessionId: string, fn: (id: string, tool: string, args: Record) => void): Effect.Effect; - delegateEmitter(childSessionId: string, parentSessionId: string): Effect.Effect; - unregisterEmitter(sessionId: string): Effect.Effect; - hasEmitter(sessionId: string): Effect.Effect; + /** 会话结束时按 sessionId 清掉待决审批(fail-closed 成 deny),返回清理条数 */ + cancelPendingFor(sessionId: string): Effect.Effect; } export class ApprovalWaitService extends Context.Tag('ApprovalWait')() {} diff --git a/packages/codingcode/src/approval/wait.ts b/packages/codingcode/src/approval/wait.ts index d76aa5e7..f2845fe6 100644 --- a/packages/codingcode/src/approval/wait.ts +++ b/packages/codingcode/src/approval/wait.ts @@ -7,12 +7,8 @@ interface PendingEntry { sessionId: string; } -export const ApprovalWaitLayer = Layer.effect(ApprovalWaitService, Effect.gen(function* () { +export const ApprovalWaitLayer = Layer.effect(ApprovalWaitService, Effect.sync(() => { const pendingConfirmations = new Map(); - const approvalEmitters = new Map< - string, - (id: string, tool: string, args: Record) => void - >(); return { waitForConfirm: (id: string, sessionId: string): Effect.Effect => @@ -35,38 +31,16 @@ export const ApprovalWaitLayer = Layer.effect(ApprovalWaitService, Effect.gen(fu return true; }), - emitApprovalRequest: ( - sessionId: string, - id: string, - tool: string, - args: Record - ): Effect.Effect => + cancelPendingFor: (sessionId: string): Effect.Effect => Effect.sync(() => { - approvalEmitters.get(sessionId)?.(id, tool, args); - }), - - registerEmitter: ( - sessionId: string, - fn: (id: string, tool: string, args: Record) => void - ): Effect.Effect => - Effect.sync(() => { - approvalEmitters.set(sessionId, fn); - }), - - delegateEmitter: (childSessionId: string, parentSessionId: string): Effect.Effect => - Effect.sync(() => { - const parentFn = approvalEmitters.get(parentSessionId); - if (parentFn) { - approvalEmitters.set(childSessionId, parentFn); + let cleared = 0; + for (const [id, entry] of pendingConfirmations) { + if (entry.sessionId !== sessionId) continue; + pendingConfirmations.delete(id); + Deferred.unsafeDone(entry.deferred, Effect.succeed({ type: 'deny' } as ConfirmResult)); + cleared++; } + return cleared; }), - - unregisterEmitter: (sessionId: string): Effect.Effect => - Effect.sync(() => { - approvalEmitters.delete(sessionId); - }), - - hasEmitter: (sessionId: string): Effect.Effect => - Effect.sync(() => approvalEmitters.has(sessionId)), }; })); diff --git a/packages/codingcode/src/context/context.ts b/packages/codingcode/src/context/context.ts index c6777fc7..9af0219c 100644 --- a/packages/codingcode/src/context/context.ts +++ b/packages/codingcode/src/context/context.ts @@ -19,6 +19,7 @@ import type { SessionEvent, AssistantEvent, ToolResultEvent, CompactEvent, Summa import { AgentError } from '../core/error.js'; import { ContextService } from './port.js'; import type { CompressResult } from './port.js'; +import { EventSinkService } from '../sink/port.js'; export function transcriptPathFor(ref: SessionRef): string { const sessionsDir = join( @@ -86,15 +87,17 @@ function applyVisibilityEvents(events: SessionEvent[]): { return { hiddenTurnIds, hiddenOpUuids, compactedTurnIds }; } +export function passesContextFilter(ev: SessionEvent): boolean { + return ev.type !== 'session_meta' && ev.type !== 'rollback' && ev.type !== 'compact'; +} + export function filterForContext(events: SessionEvent[]): { visible: SessionEvent[]; compactedTurnIds: Set; } { const { hiddenTurnIds, hiddenOpUuids, compactedTurnIds } = applyVisibilityEvents(events); const visible = events.filter((ev) => { - if (ev.type === 'session_meta') return false; - if (ev.type === 'rollback') return false; - if (ev.type === 'compact') return false; + if (!passesContextFilter(ev)) return false; if (ev.type === 'summary' && hiddenOpUuids.has(ev.uuid)) return false; if ('turnId' in ev && hiddenTurnIds.has(ev.turnId)) return false; return true; @@ -194,50 +197,48 @@ export function estimatePromptTokensFrom(events: SessionEvent[]): number { return estimateTokens(buildContextMessages(visible, compactedTurnIds)); } -interface PayloadState { +interface ContextBuffer { + turnId: number; jsonlPath: string; - currentTurnId: number; - visible: SessionEvent[]; + events: SessionEvent[]; compactedTurnIds: Set; } export const ContextLayer = Layer.effect(ContextService, Effect.gen(function* () { const session = yield* SessionService; const llm = yield* LLMService; + const sink = yield* EventSinkService; - const readState = ( - transcriptPath: string, - currentTurnId: number - ): Effect.Effect => + // 回合内内存态:键 = sessionId,仅在换回合(turnId 变化)时重建 + const buffers = new Map(); + + const buildBuffer = (ref: SessionRef): Effect.Effect => Effect.gen(function* () { - const events = yield* session.readEvents(transcriptPath); + const jsonlPath = transcriptPathFor(ref); + const events = yield* session.readEvents(jsonlPath); // 唯一读盘点 const { visible, compactedTurnIds } = filterForContext(events); - return { jsonlPath: transcriptPath, currentTurnId, visible, compactedTurnIds }; + return { turnId: ref.currentTurnId, jsonlPath, events: visible, compactedTurnIds }; }); - const estimateFor = (s: PayloadState): number => - estimateTokens(buildContextMessages(s.visible, s.compactedTurnIds)); - - const applyOldTurnCompact = ( - events: SessionEvent[], - currentTurnId: number, - jsonlPath: string - ): Effect.Effect => + const ensureBuffer = (ref: SessionRef): Effect.Effect => Effect.gen(function* () { - const compactedTurnIds = new Set(); - for (const ev of events) { - if (ev.type === 'compact') { - for (let t = ev.startTurnId; t <= ev.endTurnId; t++) { - compactedTurnIds.add(t); - } - } - } + const cached = buffers.get(ref.sessionId); + if (cached && cached.turnId === ref.currentTurnId) return cached; // 回合内复用 + const buf = yield* buildBuffer(ref); + buffers.set(ref.sessionId, buf); + return buf; + }); + + const estimateFor = (buf: ContextBuffer): number => + estimateTokens(buildContextMessages(buf.events, buf.compactedTurnIds)); + const applyOldTurnCompact = (buf: ContextBuffer): Effect.Effect => + Effect.gen(function* () { const oldResults: ToolResultEvent[] = []; - for (const ev of events) { + for (const ev of buf.events) { if (ev.type !== 'tool_result') continue; - if (ev.turnId >= currentTurnId - 1) continue; - if (compactedTurnIds.has(ev.turnId)) continue; + if (ev.turnId >= buf.turnId - 1) continue; + if (buf.compactedTurnIds.has(ev.turnId)) continue; if (!COMPACTABLE_TOOLS.has(ev.toolName.toLowerCase())) continue; if (ev.output.length <= MICRO_COMPACT_MIN_CHARS) continue; oldResults.push(ev); @@ -255,32 +256,29 @@ export const ContextLayer = Layer.effect(ContextService, Effect.gen(function* () startTurnId, endTurnId, }; - yield* session.appendEvent(jsonlPath, compactEvent); + yield* session.appendEvent(buf.jsonlPath, compactEvent); + for (let t = startTurnId; t <= endTurnId; t++) buf.compactedTurnIds.add(t); return true; }); const runMicroCompact = ( - s: PayloadState, + buf: ContextBuffer, contextWindow: number - ): Effect.Effect => + ): Effect.Effect => Effect.gen(function* () { - if (estimateFor(s) <= contextWindow * MICRO_COMPACT_THRESHOLD) return s; - const applied = yield* applyOldTurnCompact(s.visible, s.currentTurnId, s.jsonlPath); - if (applied) { - return yield* readState(s.jsonlPath, s.currentTurnId); - } - return s; + if (estimateFor(buf) <= contextWindow * MICRO_COMPACT_THRESHOLD) return; + yield* applyOldTurnCompact(buf); }); const tryCompaction = ( - s: PayloadState, + buf: ContextBuffer, model: string ): Effect.Effect => Effect.gen(function* () { - const endTurn = s.currentTurnId - KEEP_RECENT_TURNS - 1; + const endTurn = buf.turnId - KEEP_RECENT_TURNS - 1; if (endTurn < 1) return 0; - const inRange = s.visible.filter((ev) => { + const inRange = buf.events.filter((ev) => { if (ev.type === 'session_meta') return false; if ('turnId' in ev && (ev as any).turnId >= 1 && (ev as any).turnId <= endTurn) return true; return false; @@ -290,7 +288,7 @@ export const ContextLayer = Layer.effect(ContextService, Effect.gen(function* () const targetEvents = getIncrementalEvents(inRange); if (targetEvents.length === 0) return 0; - const msgs = buildContextMessages(targetEvents, s.compactedTurnIds); + const msgs = buildContextMessages(targetEvents, buf.compactedTurnIds); const totalTokens = estimateTokens(msgs); const configured = loadConfig().context.compactionModel?.trim(); @@ -315,31 +313,32 @@ export const ContextLayer = Layer.effect(ContextService, Effect.gen(function* () endTurnId, summaryText: summary, }; - yield* session.appendEvent(s.jsonlPath, summaryEvent); + yield* session.appendEvent(buf.jsonlPath, summaryEvent); + + // 就地更新:被摘要的 turn 移出可见集,摘要追加到末尾(与重读盘后的顺序一致) + buf.events = buf.events.filter( + (ev) => !('turnId' in ev && (ev as any).turnId >= startTurnId && (ev as any).turnId <= endTurnId) + ); + buf.events.push(summaryEvent); const summaryMsg: Message = { role: 'system', name: 'compacted_history', content: summary }; return Math.max(0, totalTokens - estimateMessageTokens(summaryMsg)); }); - const needsCompaction = (s: PayloadState, contextWindow: number): boolean => - estimateFor(s) > contextWindow * COMPACTION_THRESHOLD; + const needsCompaction = (buf: ContextBuffer, contextWindow: number): boolean => + estimateFor(buf) > contextWindow * COMPACTION_THRESHOLD; const summarizeToFit = ( - s: PayloadState, + buf: ContextBuffer, contextWindow: number, model: string - ): Effect.Effect<{ state: PayloadState; released: number }, AgentError> => + ): Effect.Effect => Effect.gen(function* () { - let cur = s; - let releasedTotal = 0; for (let i = 0; i < MAX_AUTO_COMPACT_PASSES; i++) { - if (!needsCompaction(cur, contextWindow)) break; - const released = yield* tryCompaction(cur, model); + if (!needsCompaction(buf, contextWindow)) break; + const released = yield* tryCompaction(buf, model); if (released <= 0) break; - releasedTotal += released; - cur = yield* readState(cur.jsonlPath, cur.currentTurnId); } - return { state: cur, released: releasedTotal }; }); function getIncrementalEvents(inRange: SessionEvent[]): SessionEvent[] { @@ -383,51 +382,54 @@ export const ContextLayer = Layer.effect(ContextService, Effect.gen(function* () return raw.trim(); } - const willCompact = ( - ref: SessionRef, - model: string - ): Effect.Effect => + const getHistory = (ref: SessionRef, model: string): Effect.Effect => Effect.gen(function* () { - const transcriptPath = transcriptPathFor(ref); + const buf = yield* ensureBuffer(ref); const contextWindow = contextWindowOf(model); - const s = yield* runMicroCompact(yield* readState(transcriptPath, ref.currentTurnId), contextWindow); - return needsCompaction(s, contextWindow); + yield* runMicroCompact(buf, contextWindow); + if (needsCompaction(buf, contextWindow)) { + // 压缩判定与压缩帧都归 context,agent 不参与 + yield* sink.emit(ref.sessionId, { family: 'transition', transition: { to: 'compress' } }); + yield* summarizeToFit(buf, contextWindow, model); + yield* sink.emit(ref.sessionId, { family: 'transition', transition: { to: 'executing' } }); + } + return buildContextMessages(buf.events, buf.compactedTurnIds); }); - const assemblePayload = ( - ref: SessionRef, - model: string - ): Effect.Effect => - Effect.gen(function* () { - const transcriptPath = transcriptPathFor(ref); - const contextWindow = contextWindowOf(model); - let s = yield* readState(transcriptPath, ref.currentTurnId); - s = yield* runMicroCompact(s, contextWindow); - const { state } = yield* summarizeToFit(s, contextWindow, model); - return buildContextMessages(state.visible, state.compactedTurnIds); + const absorb = (ref: SessionRef, events: readonly SessionEvent[]): Effect.Effect => + Effect.sync(() => { + const buf = buffers.get(ref.sessionId); + // 未建(事件已在盘上,重建时会读到)或已换回合 ⇒ no-op + if (!buf || buf.turnId !== ref.currentTurnId) return; + for (const ev of events) if (passesContextFilter(ev)) buf.events.push(ev); }); - const compactWithLLM = ( + const compact = ( ref: SessionRef, model: string, usage?: number ): Effect.Effect => Effect.gen(function* () { - const transcriptPath = transcriptPathFor(ref); + const buf = yield* ensureBuffer(ref); const contextWindow = contextWindowOf(model); - let s = yield* runMicroCompact(yield* readState(transcriptPath, ref.currentTurnId), contextWindow); - const preEstimate = usage ?? estimateFor(s); - const released = yield* tryCompaction(s, model); + yield* runMicroCompact(buf, contextWindow); + const preEstimate = usage ?? estimateFor(buf); + const released = yield* tryCompaction(buf, model); if (released <= 0) { return { didCompress: false, released: 0, promptEstimate: preEstimate }; } - s = yield* readState(transcriptPath, ref.currentTurnId); - return { didCompress: true, released, promptEstimate: estimateFor(s) }; + return { didCompress: true, released, promptEstimate: estimateFor(buf) }; + }); + + const dispose = (sessionId: string): Effect.Effect => + Effect.sync(() => { + buffers.delete(sessionId); }); return { - willCompact, - assemblePayload, - compactWithLLM, + getHistory, + absorb, + compact, + dispose, }; })); diff --git a/packages/codingcode/src/context/port.ts b/packages/codingcode/src/context/port.ts index 8cb49d27..4a80a2fa 100644 --- a/packages/codingcode/src/context/port.ts +++ b/packages/codingcode/src/context/port.ts @@ -1,7 +1,7 @@ import { Context } from 'effect'; import type { Effect } from 'effect'; import type { Message } from '../contracts/types.js'; -import type { SessionRef } from '../contracts/session.js'; +import type { SessionRef, SessionEvent } from '../contracts/session.js'; import type { AgentError } from '../core/error.js'; export interface CompressResult { @@ -11,9 +11,14 @@ export interface CompressResult { } export interface ContextShape { - willCompact(ref: SessionRef, model: string): Effect.Effect; - assemblePayload(ref: SessionRef, model: string): Effect.Effect; - compactWithLLM(ref: SessionRef, model: string, usage?: number): Effect.Effect; + /** 取该会话当前给模型的 history。回合内首次调用读一次盘,之后走内存;换回合自动重建 */ + getHistory(ref: SessionRef, model: string): Effect.Effect; + /** 把本回合自己写进 transcript 的事件并入内存态(零 IO) */ + absorb(ref: SessionRef, events: readonly SessionEvent[]): Effect.Effect; + /** 手动压缩入口(HTTP /compact),不受阈值限制 */ + compact(ref: SessionRef, model: string, usage?: number): Effect.Effect; + /** 会话删除时丢弃缓存 */ + dispose(sessionId: string): Effect.Effect; } export class ContextService extends Context.Tag('Context')() {} diff --git a/packages/codingcode/src/layer.ts b/packages/codingcode/src/layer.ts index eea42d81..4c41f063 100644 --- a/packages/codingcode/src/layer.ts +++ b/packages/codingcode/src/layer.ts @@ -7,6 +7,7 @@ import { McpLayer } from './mcp/mcp.js'; import { CheckpointLayer } from './checkpoint/checkpoint.js'; import { ApprovalLayer } from './approval/approval.js'; import { ApprovalWaitLayer } from './approval/wait.js'; +import { EventSinkLayer } from './sink/sink.js'; import { TodoLayer } from './todo/todo.js'; import { SessionLayer } from './session/session.js'; import { ToolExecutorLayer } from './tools/tools.js'; @@ -19,14 +20,18 @@ import { SchedulerLayer } from './scheduler/scheduler.js'; // base layers const InfraLayer = Layer.mergeAll( - HookLayer, RulesLayer, SkillLayer, McpLayer, ApprovalWaitLayer, TodoLayer, + HookLayer, RulesLayer, SkillLayer, McpLayer, EventSinkLayer, ApprovalWaitLayer, TodoLayer, ); -const ApprovalWithDeps = ApprovalLayer.pipe(Layer.provide(Layer.mergeAll(HookLayer, ApprovalWaitLayer))); +const ApprovalWithDeps = ApprovalLayer.pipe( + Layer.provide(Layer.mergeAll(HookLayer, EventSinkLayer, ApprovalWaitLayer)) +); const ToolExecutorWithDeps = ToolExecutorLayer.pipe( Layer.provide(Layer.mergeAll(HookLayer, ApprovalWithDeps)) ); -const ContextWithDeps = ContextLayer.pipe(Layer.provide(Layer.mergeAll(SessionLayer, LlmLayer))); +const ContextWithDeps = ContextLayer.pipe( + Layer.provide(Layer.mergeAll(SessionLayer, LlmLayer, EventSinkLayer)) +); const MemoryWithDeps = MemoryLayer.pipe(Layer.provide(LlmLayer)); // agent 直接消费的宽服务集合 @@ -55,6 +60,7 @@ export const AppLayer = Layer.mergeAll( AgentWithDeps, SubagentWithDeps, SchedulerLayer, + EventSinkLayer, ); export const createAppRuntime = () => ManagedRuntime.make(AppLayer); diff --git a/packages/codingcode/src/server/handler.ts b/packages/codingcode/src/server/handler.ts index 7eb52828..c6857b98 100644 --- a/packages/codingcode/src/server/handler.ts +++ b/packages/codingcode/src/server/handler.ts @@ -27,14 +27,6 @@ export function createSseHandler(rt: ManagedRt) { return yield* ApprovalWaitService; }) ); - Effect.runSync( - waitService.registerEmitter( - sessionId, - (id: string, tool: string, args: Record) => { - emit({ family: 'event', event: { type: 'approval_request', id, tool, args } }); - } - ) - ); try { const generator = createGenerator(); @@ -51,7 +43,7 @@ export function createSseHandler(rt: ManagedRt) { }, }); } finally { - Effect.runSync(waitService.unregisterEmitter(sessionId)); + Effect.runSync(waitService.cancelPendingFor(sessionId)); opts?.onDone?.(); } controller.close(); diff --git a/packages/codingcode/src/server/routes/sessions.ts b/packages/codingcode/src/server/routes/sessions.ts index baa49952..99b01594 100644 --- a/packages/codingcode/src/server/routes/sessions.ts +++ b/packages/codingcode/src/server/routes/sessions.ts @@ -117,7 +117,7 @@ export function registerSessionsRoutes(router: Hono, rt: ManagedRt): void { const context = yield* ContextService; const session = yield* SessionService; const state = yield* session.load(normalizedCwd, sessionId); - return yield* context.compactWithLLM( + return yield* context.compact( { cwd: state.cwd, sessionId: state.sessionId, @@ -142,7 +142,9 @@ export function registerSessionsRoutes(router: Hono, rt: ManagedRt): void { await runWithLayer( Effect.gen(function* () { const session = yield* SessionService; + const context = yield* ContextService; yield* session.deleteSession(sessionId, cwd); + yield* context.dispose(sessionId); }) as any ); return c.json({ ok: true }); diff --git a/packages/codingcode/src/sink/port.ts b/packages/codingcode/src/sink/port.ts new file mode 100644 index 00000000..7b747547 --- /dev/null +++ b/packages/codingcode/src/sink/port.ts @@ -0,0 +1,16 @@ +import { Context } from 'effect'; +import type { Effect, Queue } from 'effect'; +import type { FrameBody } from '../contracts/frame.js'; + +export interface EventSinkShape { + /** 建队列并挂载为该会话的出站队列;调用者即这条帧流的唯一读者。覆盖式:每次调用都换新队列 */ + attach(sessionId: string): Effect.Effect>; + /** 消费者退出时摘掉挂载;此后 emit 静默丢弃 */ + detach(sessionId: string): Effect.Effect; + /** 任何模块投帧,插在同一队尾,与回合自己的帧严格全序 */ + emit(sessionId: string, body: FrameBody): Effect.Effect; + /** 该会话当前是否有消费者(approval 用它判「有没有 UI」) */ + has(sessionId: string): Effect.Effect; +} + +export class EventSinkService extends Context.Tag('EventSink')() {} diff --git a/packages/codingcode/src/sink/sink.ts b/packages/codingcode/src/sink/sink.ts new file mode 100644 index 00000000..fb57d250 --- /dev/null +++ b/packages/codingcode/src/sink/sink.ts @@ -0,0 +1,26 @@ +import { Effect, Layer, Queue } from 'effect'; +import type { FrameBody } from '../contracts/frame.js'; +import { EventSinkService } from './port.js'; + +export const EventSinkLayer = Layer.effect( + EventSinkService, + Effect.sync(() => { + const queues = new Map>(); + + return { + attach: (sessionId: string) => + Effect.sync(() => { + const q = Effect.runSync(Queue.unbounded()); + queues.set(sessionId, q); + return q; + }), + detach: (sessionId: string) => Effect.sync(() => void queues.delete(sessionId)), + emit: (sessionId: string, body: FrameBody) => + Effect.sync(() => { + const q = queues.get(sessionId); + if (q) Effect.runSync(Queue.offer(q, body)); + }), + has: (sessionId: string) => Effect.sync(() => queues.has(sessionId)), + }; + }) +); diff --git a/packages/codingcode/test/agent/context-compressed.test.ts b/packages/codingcode/test/agent/context-compressed.test.ts index de676058..12fd0166 100644 --- a/packages/codingcode/test/agent/context-compressed.test.ts +++ b/packages/codingcode/test/agent/context-compressed.test.ts @@ -13,7 +13,7 @@ function phaseOrder(events: readonly unknown[]): string[] { } describe('compaction transition', () => { - it('emits the compress signal when willCompact is true', async () => { + it('emits the compress signal (emitted by context via sink) when compaction triggers', async () => { const { events } = await runAgentTurn( { llm: makePlainLlm(), @@ -40,7 +40,7 @@ describe('compaction transition', () => { expect(tos[tos.indexOf('compress') + 1]).toBe('executing'); }); - it('emits no compress signal when willCompact is false', async () => { + it('emits no compress signal when compaction does not trigger', async () => { const { events } = await runAgentTurn( { llm: makePlainLlm(), diff --git a/packages/codingcode/test/approval/async-confirm.test.ts b/packages/codingcode/test/approval/async-confirm.test.ts index cf038f63..4a0ea72e 100644 --- a/packages/codingcode/test/approval/async-confirm.test.ts +++ b/packages/codingcode/test/approval/async-confirm.test.ts @@ -1,5 +1,5 @@ import { describe, it, expect } from 'vitest'; -import { Effect, Layer } from 'effect'; +import { Effect, Fiber, Layer } from 'effect'; import { ApprovalWaitService } from '../../src/approval/wait-port.js'; import type { ConfirmResult } from '../../src/approval/confirmation.js'; import { ApprovalWaitLayer } from '../../src/approval/wait.js'; @@ -60,61 +60,42 @@ describe('ApprovalWaitService', () => { }); }); -describe('delegateEmitter', () => { - it('delegates parent emitter to child session', async () => { - const parentSid = 'parent-' + Math.random().toString(36).slice(2); - const childSid = 'child-' + Math.random().toString(36).slice(2); - const calls: Array<[string, string, Record]> = []; +describe('cancelPendingFor', () => { + it('fails pending approvals of that session closed as deny and returns the count', async () => { + const sid = 'sess-' + Math.random().toString(36).slice(2); + const other = 'other-' + Math.random().toString(36).slice(2); - await run( + const results = await run( Effect.gen(function* () { const svc = yield* ApprovalWaitService; - yield* svc.registerEmitter( - parentSid, - (id: string, tool: string, args: Record) => calls.push([id, tool, args]) - ); - - expect(yield* svc.hasEmitter(parentSid)).toBe(true); - expect(yield* svc.hasEmitter(childSid)).toBe(false); - - yield* svc.delegateEmitter(childSid, parentSid); - - expect(yield* svc.hasEmitter(childSid)).toBe(true); + const mine = yield* Effect.fork(svc.waitForConfirm('a1', sid)); + const alsoMine = yield* Effect.fork(svc.waitForConfirm('a2', sid)); + const theirs = yield* Effect.fork(svc.waitForConfirm('b1', other)); + yield* Effect.sleep('5 millis'); - yield* svc.unregisterEmitter(childSid); - yield* svc.unregisterEmitter(parentSid); + const cleared = yield* svc.cancelPendingFor(sid); + const mineRes = yield* Fiber.join(mine); + const alsoMineRes = yield* Fiber.join(alsoMine); + // 其它会话的待决项不受影响 + const stillPending = yield* svc.resolveConfirm('b1', other, { type: 'allow' }); + yield* Fiber.join(theirs); + return { cleared, mineRes, alsoMineRes, stillPending }; }) ); - }); - - it('child emitter fires the same callback as parent', async () => { - const parentSid = 'parent-cb-' + Math.random().toString(36).slice(2); - const childSid = 'child-cb-' + Math.random().toString(36).slice(2); - const received: string[] = []; - await run( - Effect.gen(function* () { - const svc = yield* ApprovalWaitService; - yield* svc.registerEmitter(parentSid, (id: string) => received.push(id)); - yield* svc.delegateEmitter(childSid, parentSid); - - // Since we can't directly access the private map, we verify via hasEmitter - expect(yield* svc.hasEmitter(childSid)).toBe(true); - - yield* svc.unregisterEmitter(childSid); - yield* svc.unregisterEmitter(parentSid); - }) - ); + expect(results.cleared).toBe(2); + expect(results.mineRes).toEqual({ type: 'deny' }); + expect(results.alsoMineRes).toEqual({ type: 'deny' }); + expect(results.stillPending).toBe(true); }); - it('delegateEmitter is a no-op when parent has no emitter', async () => { - const childSid = 'child-noop-' + Math.random().toString(36).slice(2); - await run( + it('returns 0 when the session has no pending approvals', async () => { + const cleared = await run( Effect.gen(function* () { const svc = yield* ApprovalWaitService; - yield* svc.delegateEmitter(childSid, 'nonexistent-parent'); - expect(yield* svc.hasEmitter(childSid)).toBe(false); + return yield* svc.cancelPendingFor('nobody'); }) ); + expect(cleared).toBe(0); }); }); diff --git a/packages/codingcode/test/approval/pipeline.test.ts b/packages/codingcode/test/approval/pipeline.test.ts index 75777fbb..709b25c6 100644 --- a/packages/codingcode/test/approval/pipeline.test.ts +++ b/packages/codingcode/test/approval/pipeline.test.ts @@ -4,6 +4,7 @@ import { runPipeline } from '../../src/approval/approval.js'; import { createRuleEngine } from '../../src/approval/rule-engine.js'; import type { PermissionRule } from '../../src/approval/types.js'; import { ApprovalWaitService } from '../../src/approval/wait-port.js'; +import { EventSinkService } from '../../src/sink/port.js'; import { HookService } from '../../src/hooks/port.js'; const mockHookService = { @@ -15,16 +16,20 @@ const mockHookService = { const mockApprovalWaitService = { waitForConfirm: () => Effect.dieMessage('not implemented'), resolveConfirm: () => Effect.succeed(false), - emitApprovalRequest: () => Effect.succeed(undefined), - registerEmitter: () => Effect.succeed(undefined), - delegateEmitter: () => Effect.succeed(undefined), - unregisterEmitter: () => Effect.succeed(undefined), - hasEmitter: () => Effect.succeed(false), + cancelPendingFor: () => Effect.succeed(0), +}; + +const mockEventSink = { + attach: () => Effect.succeed({} as any), + detach: () => Effect.void, + emit: () => Effect.void, + has: () => Effect.succeed(false), }; const HookTestLayer = Layer.succeed(HookService, mockHookService); const WaitTestLayer = Layer.succeed(ApprovalWaitService, mockApprovalWaitService); -const TestLayer = Layer.mergeAll(HookTestLayer, WaitTestLayer); +const SinkTestLayer = Layer.succeed(EventSinkService, mockEventSink as any); +const TestLayer = Layer.mergeAll(HookTestLayer, WaitTestLayer, SinkTestLayer); function runWithLayer(eff: Effect.Effect): Promise { return Effect.runPromise(eff.pipe(Effect.provide(TestLayer))); diff --git a/packages/codingcode/test/context/budget-integration.test.ts b/packages/codingcode/test/context/budget-integration.test.ts index a04e56ca..34137b1c 100644 --- a/packages/codingcode/test/context/budget-integration.test.ts +++ b/packages/codingcode/test/context/budget-integration.test.ts @@ -10,15 +10,17 @@ import { LLMService } from '../../src/llm/port.js'; import type { SessionRef } from '../../src/contracts/session.js'; import { useTempProjectBase } from '../helpers/project-base.js'; import { ContextLayer, transcriptPathFor } from '../../src/context/context.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; useTempProjectBase(); -const TestLayer = Layer.merge( +const TestLayer = Layer.mergeAll( SessionLayer, Layer.succeed(LLMService, { complete: () => Effect.fail(new Error('no llm')), completeStream: () => (async function* () {})(), - } as any) + } as any), + EventSinkLayer ); async function getCtxService(): Promise { @@ -31,7 +33,7 @@ async function getCtxService(): Promise { const CWD = '/tmp/test'; -describe('assemblePayload integration', () => { +describe('getHistory integration', () => { let ref: SessionRef; let transcriptPath: string; @@ -86,7 +88,7 @@ describe('assemblePayload integration', () => { it('returns messages assembled from the transcript', async () => { const ctx = await getCtxService(); - const messages = await Effect.runPromise(ctx.assemblePayload(ref, 'test-model')); + const messages = await Effect.runPromise(ctx.getHistory(ref, 'test-model')); expect(messages.length).toBeGreaterThan(0); }); @@ -96,7 +98,7 @@ describe('assemblePayload integration', () => { const emptyPath = transcriptPathFor(emptyRef); writeFileSync(emptyPath, '', 'utf8'); const ctx = await getCtxService(); - const messages = await Effect.runPromise(ctx.assemblePayload(emptyRef, 'test-model')); + const messages = await Effect.runPromise(ctx.getHistory(emptyRef, 'test-model')); expect(messages).toEqual([]); }); }); diff --git a/packages/codingcode/test/context/compressor/behavior.test.ts b/packages/codingcode/test/context/compressor/behavior.test.ts index dd6ca4eb..38d587d2 100644 --- a/packages/codingcode/test/context/compressor/behavior.test.ts +++ b/packages/codingcode/test/context/compressor/behavior.test.ts @@ -14,6 +14,7 @@ import { readHistory } from '../../../src/session/file-ops.js'; import { estimateTokens } from '../../../src/context/tokens.js'; import { useTempProjectBase } from '../../helpers/project-base.js'; import { ContextLayer } from '../../../src/context/context.js'; +import { EventSinkLayer } from '../../../src/sink/sink.js'; // 上下文窗口现在由 catalog 按模型值现取,测试里钉死成一个可控值 const windowState = vi.hoisted(() => ({ value: 128000 })); @@ -107,7 +108,7 @@ const FailingLLM = { } as any; function makeTestLayer(llm: unknown) { - return Layer.merge(SessionLayer, Layer.succeed(LLMService, llm as any)); + return Layer.mergeAll(SessionLayer, Layer.succeed(LLMService, llm as any), EventSinkLayer); } async function getCtxService(llm: unknown): Promise { @@ -130,7 +131,7 @@ describe('compressor behavior', () => { '## Compacted History\n\n### Goal\nfix bug\n\n### Instructions\nbe careful\n\n### Discoveries\nrace condition\n\n### Accomplished\npatched\n\n### Relevant Files\nsrc/x.ts'; windowState.value = 1000; const ctx = await getCtxService(makeMockLLM(summary)); - await run(ctx.compactWithLLM(fx.ref, 'test-model')); + await run(ctx.compact(fx.ref, 'test-model')); const summaries = readSummaryEvents(fx.transcriptPath); expect(summaries.length).toBe(1); expect(summaries[0]!.summaryText).toContain('### Goal'); @@ -148,7 +149,7 @@ describe('compressor behavior', () => { try { windowState.value = 1000; const ctx = await getCtxService(FailingLLM); - const result = await run(ctx.compactWithLLM(fx.ref, 'test-model')); + const result = await run(ctx.compact(fx.ref, 'test-model')); expect(result.didCompress).toBe(false); const summaries = readSummaryEvents(fx.transcriptPath); expect(summaries).toHaveLength(0); @@ -168,7 +169,7 @@ describe('compressor behavior', () => { '## Compacted History\n\n### Goal\na\n\n### Instructions\nb\n\n### Discoveries\nc\n\n### Accomplished\nd\n\n### Relevant Files\ne' ) ); - await run(ctx.compactWithLLM(fx.ref, 'test-model')); + await run(ctx.compact(fx.ref, 'test-model')); const summaries = readSummaryEvents(fx.transcriptPath); expect(summaries).toHaveLength(1); @@ -180,7 +181,7 @@ describe('compressor behavior', () => { }); }); - describe('compactWithLLM result', () => { + describe('compact result', () => { it('returns promptEstimate after compression', async () => { const fx = makeFixture({ numTurns: 5 }); try { @@ -194,7 +195,7 @@ describe('compressor behavior', () => { '## Compacted History\n\n### Goal\na\n\n### Instructions\nb\n\n### Discoveries\nc\n\n### Accomplished\nd\n\n### Relevant Files\ne' ) ); - const result = await run(ctx.compactWithLLM(fx.ref, 'test-model')); + const result = await run(ctx.compact(fx.ref, 'test-model')); expect(result.didCompress).toBe(true); expect(result.promptEstimate).toBeGreaterThan(0); expect(result.promptEstimate).toBeLessThan(before); @@ -205,7 +206,7 @@ describe('compressor behavior', () => { }); }); - describe('assemblePayload compaction', () => { + describe('getHistory compaction', () => { const SUMMARY = '## Compacted History\n\n### Goal\na\n\n### Instructions\nb\n\n### Discoveries\nc\n\n### Accomplished\nd\n\n### Relevant Files\ne'; @@ -214,7 +215,7 @@ describe('compressor behavior', () => { try { windowState.value = 1000; const ctx = await getCtxService(makeMockLLM(SUMMARY)); - const messages = await run(ctx.assemblePayload(fx.ref, 'test-model')); + const messages = await run(ctx.getHistory(fx.ref, 'test-model')); expect(messages.length).toBeGreaterThan(0); expect(messages.some((m) => m.name === 'compacted_history')).toBe(true); } finally { @@ -227,7 +228,7 @@ describe('compressor behavior', () => { try { windowState.value = 2_000_000; const ctx = await getCtxService(makeMockLLM(SUMMARY)); - const messages = await run(ctx.assemblePayload(fx.ref, 'test-model')); + const messages = await run(ctx.getHistory(fx.ref, 'test-model')); expect(messages.some((m) => m.name === 'compacted_history')).toBe(false); } finally { cleanup(fx.dir); diff --git a/packages/codingcode/test/context/memory-buffer.test.ts b/packages/codingcode/test/context/memory-buffer.test.ts new file mode 100644 index 00000000..495e54b2 --- /dev/null +++ b/packages/codingcode/test/context/memory-buffer.test.ts @@ -0,0 +1,147 @@ +import { describe, it, expect, vi } from 'vitest'; +import { Effect, Layer } from 'effect'; +import { ContextService } from '../../src/context/port.js'; +import type { ContextShape } from '../../src/context/port.js'; +import { ContextLayer } from '../../src/context/context.js'; +import { SessionService } from '../../src/session/port.js'; +import { LLMService } from '../../src/llm/port.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; +import type { SessionEvent, SessionRef } from '../../src/contracts/session.js'; + +// 上下文窗口钉死成可控值(与其余 context 用例一致) +const windowState = vi.hoisted(() => ({ value: 128000 })); +vi.mock('../../src/infra/models.js', () => ({ + contextWindowOf: () => windowState.value, +})); + +const SUMMARY = '## Compacted History\n\n### Goal\nx\n\n### Instructions\ny\n\n### Discoveries\nz\n\n### Accomplished\nw\n\n### Relevant Files\nf'; + +/** 计数 + 内存中的假 transcript:把「读盘」变成可断言的数字 */ +function makeCountingSession(seed: SessionEvent[]) { + const reads = { count: 0 }; + const appended: SessionEvent[] = []; + const svc = { + readEvents: () => + Effect.sync(() => { + reads.count++; + return [...seed, ...appended]; + }), + appendEvent: (_path: string, ev: SessionEvent) => + Effect.sync(() => { + appended.push(ev); + }), + }; + return { svc, reads, appended }; +} + +function makeLayer(counting: ReturnType, summary = SUMMARY) { + const sessionLayer = Layer.succeed(SessionService, counting.svc as any); + const llmLayer = Layer.succeed(LLMService, { + complete: () => Effect.succeed({ content: summary }), + completeStream: () => (async function* () {})(), + } as any); + return ContextLayer.pipe(Layer.provide(Layer.mergeAll(sessionLayer, llmLayer, EventSinkLayer))); +} + +async function getCtx(counting: ReturnType): Promise { + return Effect.runPromise( + Effect.gen(function* () { + return yield* ContextService; + }).pipe(Effect.provide(makeLayer(counting)) as any) + ); +} + +const REF: SessionRef = { cwd: '/tmp', sessionId: 's1', currentTurnId: 1 }; + +function seedEvents(): SessionEvent[] { + return [ + { + type: 'session_meta', + sessionId: 's1', + cwd: '/tmp', + createdAt: new Date().toISOString(), + model: 'test-model', + title: 't', + activeProfile: 'build', + permissionMode: 'ask', + }, + { type: 'user', turnId: 1, content: 'q1' }, + ]; +} + +describe('context memory buffer', () => { + it('回合内多次 getHistory 只读一次盘', async () => { + const counting = makeCountingSession(seedEvents()); + const ctx = await getCtx(counting); + + for (let i = 0; i < 3; i++) { + await Effect.runPromise(ctx.getHistory(REF, 'test-model')); + } + + expect(counting.reads.count).toBe(1); + }); + + it('换回合(turnId 变化)后重读一次盘', async () => { + const counting = makeCountingSession(seedEvents()); + const ctx = await getCtx(counting); + + await Effect.runPromise(ctx.getHistory(REF, 'test-model')); + await Effect.runPromise(ctx.getHistory({ ...REF, currentTurnId: 2 }, 'test-model')); + + expect(counting.reads.count).toBe(2); + }); + + it('absorb 的事件在下一步立即可见,且不触发回盘', async () => { + const counting = makeCountingSession(seedEvents()); + const ctx = await getCtx(counting); + + await Effect.runPromise(ctx.getHistory(REF, 'test-model')); + await Effect.runPromise(ctx.absorb(REF, [{ type: 'assistant', turnId: 1, content: 'r1', toolCalls: [] }])); + const messages = await Effect.runPromise(ctx.getHistory(REF, 'test-model')); + + expect(messages.some((m) => m.content === 'r1')).toBe(true); + expect(counting.reads.count).toBe(1); + }); + + it('absorb 早于首次 getHistory 时是 no-op(事件本就在盘上)', async () => { + const counting = makeCountingSession(seedEvents()); + const ctx = await getCtx(counting); + + await Effect.runPromise(ctx.absorb(REF, [{ type: 'user', turnId: 1, content: 'early' }])); + const messages = await Effect.runPromise(ctx.getHistory(REF, 'test-model')); + + expect(counting.reads.count).toBe(1); + expect(messages.some((m) => m.content === 'q1')).toBe(true); + }); + + it('压缩轮不再额外读盘', async () => { + windowState.value = 1000; + try { + const big = 'X'.repeat(8000); + const counting = makeCountingSession([ + ...seedEvents(), + { type: 'assistant', turnId: 1, content: 'r1', toolCalls: [{ id: 'tc1', name: 'bash', arguments: {} }] }, + { type: 'tool_result', turnId: 1, toolName: 'bash', toolCallId: 'tc1', output: big }, + ]); + const ctx = await getCtx(counting); + + const messages = await Effect.runPromise(ctx.getHistory({ ...REF, currentTurnId: 3 }, 'test-model')); + + expect(messages.some((m) => m.name === 'compacted_history')).toBe(true); + expect(counting.reads.count).toBe(1); + } finally { + windowState.value = 128000; + } + }); + + it('dispose 后缓存被清,下次 getHistory 重新读盘', async () => { + const counting = makeCountingSession(seedEvents()); + const ctx = await getCtx(counting); + + await Effect.runPromise(ctx.getHistory(REF, 'test-model')); + await Effect.runPromise(ctx.dispose(REF.sessionId)); + await Effect.runPromise(ctx.getHistory(REF, 'test-model')); + + expect(counting.reads.count).toBe(2); + }); +}); diff --git a/packages/codingcode/test/helpers/agent-harness.ts b/packages/codingcode/test/helpers/agent-harness.ts index dc0495a9..cdc537d9 100644 --- a/packages/codingcode/test/helpers/agent-harness.ts +++ b/packages/codingcode/test/helpers/agent-harness.ts @@ -12,6 +12,8 @@ import { McpService } from '../../src/mcp/port.js'; import { MemoryService } from '../../src/memory/port.js'; import { RulesService } from '../../src/rules/port.js'; import { SessionService } from '../../src/session/port.js'; +import { EventSinkService } from '../../src/sink/port.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; import { SkillService } from '../../src/skills/port.js'; import { SubagentRunnerService } from '../../src/subagent/port.js'; import { TodoService } from '../../src/todo/port.js'; @@ -19,7 +21,7 @@ import { ToolExecutorService } from '../../src/tools/port.js'; import type { FrameBody, RuntimeEvent, Transition } from '../../src/contracts/frame.js'; import type { TokenUsage } from '../../src/contracts/types.js'; import type { LLMStreamPart } from '../../src/contracts/provider.js'; -import type { SessionStoreState } from '../../src/contracts/session.js'; +import type { SessionStoreState, SessionRef } from '../../src/contracts/session.js'; // ---- LLM 部件构造器 ---- @@ -130,9 +132,9 @@ export interface HarnessMocks { }; todo?: Map>; memorySnapshot?: string; - /** 可选:覆盖 ContextService.assemblePayload 的返回(默认一条 user 消息)。 */ + /** 可选:覆盖 ContextService.getHistory 的返回(默认一条 user 消息)。 */ contextAssemble?: () => Promise>; - /** 可选:覆盖 ContextService.willCompact(默认 false)。 */ + /** 可选:模拟压缩触发(默认 false)——为 true 时 getHistory 经 sink 发 compress/executing 帧。 */ contextWillCompact?: () => Promise; /** 可选:覆盖 SessionService 的个别方法(默认实现见 makeAgentLayer)。 */ session?: Partial<{ @@ -239,14 +241,29 @@ export function makeAgentLayer(mocks: HarnessMocks): Layer.Layer { const skills = { extractSkill: (_cwd: string, query: string) => Effect.succeed([undefined, query]), }; - const context = { - willCompact: () => - Effect.promise(async () => (mocks.contextWillCompact ? mocks.contextWillCompact() : false)), - assemblePayload: () => - Effect.promise(async () => - mocks.contextAssemble ? mocks.contextAssemble() : [{ role: 'user' as const, content: 'hi' }] - ), - }; + // 压缩帧由 context 自己经 sink 发出(agent 不再参与),mock 同样遵守这个归属 + const ContextMockLayer = Layer.effect( + ContextService, + Effect.gen(function* () { + const sink = yield* EventSinkService; + return { + getHistory: (ref: SessionRef) => + Effect.gen(function* () { + const shouldCompact = mocks.contextWillCompact ? yield* Effect.promise(mocks.contextWillCompact) : false; + if (shouldCompact) { + yield* sink.emit(ref.sessionId, { family: 'transition', transition: { to: 'compress' } }); + yield* sink.emit(ref.sessionId, { family: 'transition', transition: { to: 'executing' } }); + } + return mocks.contextAssemble + ? yield* Effect.promise(mocks.contextAssemble) + : [{ role: 'user' as const, content: 'hi' }]; + }), + absorb: () => Effect.void, + compact: () => Effect.succeed({ didCompress: false, released: 0, promptEstimate: 0 }), + dispose: () => Effect.void, + } as any; + }) + ).pipe(Layer.provide(EventSinkLayer)); const memory = { loadMemoryForPrompt: () => Effect.succeed(mocks.memorySnapshot ?? ''), flushSessionToMemory: () => Effect.succeed({ written: false, bytes: 0 }), @@ -268,7 +285,8 @@ export function makeAgentLayer(mocks: HarnessMocks): Layer.Layer { evaluate: () => Effect.succeed({ type: 'allow', source: 'test' }), } as any), Layer.succeed(SkillService, skills as any), - Layer.succeed(ContextService, context as any), + ContextMockLayer, + EventSinkLayer, Layer.succeed(MemoryService, memory as any), Layer.succeed(LLMService, { complete: () => Effect.fail(new Error('complete not implemented in harness')), diff --git a/packages/codingcode/test/plan/gate-pipeline.test.ts b/packages/codingcode/test/plan/gate-pipeline.test.ts index 5818d178..79914d58 100644 --- a/packages/codingcode/test/plan/gate-pipeline.test.ts +++ b/packages/codingcode/test/plan/gate-pipeline.test.ts @@ -7,6 +7,7 @@ import { runPipeline } from '../../src/approval/approval.js'; import { createRuleEngine } from '../../src/approval/rule-engine.js'; import { HookService } from '../../src/hooks/port.js'; import { ApprovalWaitService } from '../../src/approval/wait-port.js'; +import { EventSinkService } from '../../src/sink/port.js'; import type { ProfileName } from '../../src/contracts/types.js'; import { useTempProjectBase } from '../helpers/project-base.js'; @@ -24,14 +25,22 @@ function makeMockApprovalWait() { return { waitForConfirm: () => Effect.succeed({ type: 'deny' }) as any, resolveConfirm: () => Effect.succeed(false), - emitApprovalRequest: (sessionId: string, id: string, tool: string, args: any) => + cancelPendingFor: () => Effect.succeed(0), + }; +} + +// 审批请求现在经 EventSink 出站,捕获点从 wait 的 emitter 迁到 sink.emit +function makeMockEventSink() { + return { + attach: () => Effect.succeed({} as any), + detach: () => Effect.void, + emit: (sessionId: string, body: any) => Effect.sync(() => { - capturedApproval = { sessionId, id, tool, args }; + if (body?.family === 'event' && body.event?.type === 'approval_request') { + capturedApproval = { sessionId, id: body.event.id, tool: body.event.tool, args: body.event.args }; + } }), - registerEmitter: () => Effect.succeed(undefined), - delegateEmitter: () => Effect.succeed(undefined), - unregisterEmitter: () => Effect.succeed(undefined), - hasEmitter: () => Effect.succeed(true), + has: () => Effect.succeed(true), }; } @@ -47,7 +56,8 @@ function runPipelineWithMock(opts: { const mockWait = makeMockApprovalWait(); const HookTestLayer = Layer.succeed(HookService, mockHookService as any); const WaitTestLayer = Layer.succeed(ApprovalWaitService, mockWait as any); - const TestLayer = Layer.mergeAll(HookTestLayer, WaitTestLayer); + const SinkTestLayer = Layer.succeed(EventSinkService, makeMockEventSink() as any); + const TestLayer = Layer.mergeAll(HookTestLayer, WaitTestLayer, SinkTestLayer); return Effect.runPromise( runPipeline( { tool: opts.tool, input: opts.input }, diff --git a/packages/codingcode/test/security/plan-profile-restart.test.ts b/packages/codingcode/test/security/plan-profile-restart.test.ts index c71026b3..90fef911 100644 --- a/packages/codingcode/test/security/plan-profile-restart.test.ts +++ b/packages/codingcode/test/security/plan-profile-restart.test.ts @@ -8,6 +8,7 @@ import { SessionLayer } from '../../src/session/session.js'; import { HookService } from '../../src/hooks/port.js'; import { ApprovalService } from '../../src/approval/port.js'; import { ApprovalWaitService } from '../../src/approval/wait-port.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; import type { ProfileName } from '../../src/contracts/types.js'; import { useTempProjectBase } from '../helpers/project-base.js'; import { ApprovalLayer } from '../../src/approval/approval.js'; @@ -23,11 +24,7 @@ const mockHookService = { const mockApprovalWaitService = { waitForConfirm: () => Effect.dieMessage('not implemented'), resolveConfirm: () => Effect.succeed(false), - emitApprovalRequest: () => Effect.succeed(undefined), - registerEmitter: () => Effect.succeed(undefined), - delegateEmitter: () => Effect.succeed(undefined), - unregisterEmitter: () => Effect.succeed(undefined), - hasEmitter: () => Effect.succeed(false), + cancelPendingFor: () => Effect.succeed(0), }; function makeLayer() { @@ -36,6 +33,7 @@ function makeLayer() { Layer.provide( Layer.mergeAll( HookTestLayer, + EventSinkLayer, Layer.succeed(ApprovalWaitService, mockApprovalWaitService as any) ) ) @@ -43,6 +41,7 @@ function makeLayer() { return Layer.mergeAll( SessionLayer, HookTestLayer, + EventSinkLayer, ApprovalTestLayer, Layer.succeed(ApprovalWaitService, mockApprovalWaitService as any) ); diff --git a/packages/codingcode/test/server/compact-route.test.ts b/packages/codingcode/test/server/compact-route.test.ts index 8adc1563..18dde267 100644 --- a/packages/codingcode/test/server/compact-route.test.ts +++ b/packages/codingcode/test/server/compact-route.test.ts @@ -14,6 +14,7 @@ import { CheckpointService } from '../../src/checkpoint/port.js'; import { HookLayer } from '../../src/hooks/hooks.js'; import { ApprovalWaitLayer } from '../../src/approval/wait.js'; import { ApprovalLayer } from '../../src/approval/approval.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; const mockCompactWithLLM = vi.fn(); @@ -51,7 +52,7 @@ const MockSessionLayer = Layer.succeed(SessionService, { } as any); const MockApprovalLayer = ApprovalLayer.pipe( - Layer.provide(Layer.mergeAll(HookLayer, ApprovalWaitLayer)) + Layer.provide(Layer.mergeAll(HookLayer, EventSinkLayer, ApprovalWaitLayer)) ); const MockSkillLayer = Layer.succeed(SkillService, { @@ -82,8 +83,10 @@ const MockSchedulerLayer = Layer.succeed(SchedulerService, { } as any); const MockContextLayer = Layer.succeed(ContextService, { - assemblePayload: () => Effect.succeed([]), - compactWithLLM: mockCompactWithLLM, + getHistory: () => Effect.succeed([]), + absorb: () => Effect.void, + compact: mockCompactWithLLM, + dispose: () => Effect.void, } as any); const MockCheckpointLayer = Layer.succeed(CheckpointService, { @@ -112,6 +115,7 @@ const TestLayer = Layer.mergeAll( MockSessionLayer, MockApprovalLayer, HookLayer, + EventSinkLayer, ApprovalWaitLayer, MockSkillLayer, MockMcpLayer, @@ -135,7 +139,7 @@ describe('POST /api/sessions/:id/compact (manual compact)', () => { ); }); - it('should pass the requested model through to compactWithLLM', async () => { + it('should pass the requested model through to compact', async () => { const app = await createServer(rt); const res = await app.request('/api/sessions/test-sid/compact', { method: 'POST', diff --git a/packages/codingcode/test/server/index.test.ts b/packages/codingcode/test/server/index.test.ts index bd3b9af6..7dc12342 100644 --- a/packages/codingcode/test/server/index.test.ts +++ b/packages/codingcode/test/server/index.test.ts @@ -15,6 +15,7 @@ import { CheckpointService } from '../../src/checkpoint/port.js'; import { HookLayer } from '../../src/hooks/hooks.js'; import { ApprovalWaitLayer } from '../../src/approval/wait.js'; import { ApprovalLayer } from '../../src/approval/approval.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; const MockSessionLayer = Layer.succeed(SessionService, { create: () => Effect.succeed({ sessionId: 'test', cwd: '/tmp/test' }), @@ -41,7 +42,7 @@ const MockLLMFactoryLayer = Layer.succeed(LLMFactoryService, { } as any); const MockApprovalLayer = ApprovalLayer.pipe( - Layer.provide(Layer.mergeAll(HookLayer, ApprovalWaitLayer)) + Layer.provide(Layer.mergeAll(HookLayer, EventSinkLayer, ApprovalWaitLayer)) ); const MockSkillLayer = Layer.succeed(SkillService, { @@ -71,7 +72,12 @@ const MockSchedulerLayer = Layer.succeed(SchedulerService, { runOnce: () => Promise.resolve('session-id'), } as any); -const MockContextLayer = Layer.succeed(ContextService, {} as any); +const MockContextLayer = Layer.succeed(ContextService, { + getHistory: () => Effect.succeed([]), + absorb: () => Effect.void, + compact: () => Effect.succeed({ didCompress: false, released: 0, promptEstimate: 0 }), + dispose: () => Effect.void, +} as any); const MockCheckpointLayer = Layer.succeed(CheckpointService, { _tag: 'Checkpoint' as const, @@ -100,6 +106,7 @@ const TestLayer = Layer.mergeAll( MockLLMFactoryLayer, MockApprovalLayer, HookLayer, + EventSinkLayer, ApprovalWaitLayer, MockSkillLayer, MockMcpLayer, diff --git a/packages/codingcode/test/server/messages-fork-permission-mode.test.ts b/packages/codingcode/test/server/messages-fork-permission-mode.test.ts index dbe987de..3d769120 100644 --- a/packages/codingcode/test/server/messages-fork-permission-mode.test.ts +++ b/packages/codingcode/test/server/messages-fork-permission-mode.test.ts @@ -23,11 +23,7 @@ const mockHookService = { const mockApprovalWaitService = { waitForConfirm: () => Effect.dieMessage('not implemented'), resolveConfirm: () => Effect.succeed(false), - emitApprovalRequest: () => Effect.succeed(undefined), - registerEmitter: () => Effect.succeed(undefined), - delegateEmitter: () => Effect.succeed(undefined), - unregisterEmitter: () => Effect.succeed(undefined), - hasEmitter: () => Effect.succeed(false), + cancelPendingFor: () => Effect.succeed(0), }; // The message-send path now lives in AgentService.runTurn. A real runTurn loads diff --git a/packages/codingcode/test/server/plan-file-route.test.ts b/packages/codingcode/test/server/plan-file-route.test.ts index dfdcaf9a..608c9693 100644 --- a/packages/codingcode/test/server/plan-file-route.test.ts +++ b/packages/codingcode/test/server/plan-file-route.test.ts @@ -24,6 +24,7 @@ import { projectBaseDir } from '../helpers/project-base.js'; import { HookLayer } from '../../src/hooks/hooks.js'; import { ApprovalWaitLayer } from '../../src/approval/wait.js'; import { ApprovalLayer } from '../../src/approval/approval.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; const MockSessionLayer = Layer.succeed(SessionService, { create: () => @@ -96,7 +97,7 @@ const MockLLMFactoryLayer = Layer.succeed(LLMFactoryService, { } as any); const MockApprovalLayer = ApprovalLayer.pipe( - Layer.provide(Layer.mergeAll(HookLayer, ApprovalWaitLayer)) + Layer.provide(Layer.mergeAll(HookLayer, EventSinkLayer, ApprovalWaitLayer)) ); const MockSkillLayer = Layer.succeed(SkillService, { @@ -127,8 +128,10 @@ const MockSchedulerLayer = Layer.succeed(SchedulerService, { } as any); const MockContextLayer = Layer.succeed(ContextService, { - assemblePayload: () => Effect.succeed([]), - compactWithLLM: () => Effect.succeed({ didCompress: false, released: 0, promptEstimate: 0 }), + getHistory: () => Effect.succeed([]), + absorb: () => Effect.void, + compact: () => Effect.succeed({ didCompress: false, released: 0, promptEstimate: 0 }), + dispose: () => Effect.void, } as any); const MockCheckpointLayer = Layer.succeed(CheckpointService, { @@ -158,6 +161,7 @@ const TestLayer = Layer.mergeAll( MockLLMFactoryLayer, MockApprovalLayer, HookLayer, + EventSinkLayer, ApprovalWaitLayer, MockSkillLayer, MockMcpLayer, diff --git a/packages/codingcode/test/sink/sink.test.ts b/packages/codingcode/test/sink/sink.test.ts new file mode 100644 index 00000000..4a4cac8d --- /dev/null +++ b/packages/codingcode/test/sink/sink.test.ts @@ -0,0 +1,87 @@ +import { describe, it, expect } from 'vitest'; +import { Effect, Queue } from 'effect'; +import { EventSinkService } from '../../src/sink/port.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; +import type { FrameBody } from '../../src/contracts/frame.js'; + +const run = (eff: Effect.Effect): Promise => + Effect.runPromise(eff.pipe(Effect.provide(EventSinkLayer))); + +const textDelta = (t: string): FrameBody => ({ family: 'event', event: { type: 'text_delta', text: t } }); +const toolCall = (id: string): FrameBody => ({ + family: 'event', + event: { type: 'tool_call', id, name: 'read_file', args: {} }, +}); + +function takeN(q: Queue.Queue, n: number): Effect.Effect { + return Effect.gen(function* () { + const out: FrameBody[] = []; + for (let i = 0; i < n; i++) out.push(yield* Queue.take(q)); + return out; + }); +} + +describe('EventSink', () => { + it('attach 覆盖式:旧队列被替换,emit 只进新队列', async () => { + const result = await run( + Effect.gen(function* () { + const sink = yield* EventSinkService; + const first = yield* sink.attach('s1'); + const second = yield* sink.attach('s1'); + + yield* sink.emit('s1', textDelta('a')); + + const fromSecond = yield* Queue.poll(second); + const fromFirst = yield* Queue.poll(first); + return { fromSecond, fromFirst }; + }) + ); + + expect(result.fromSecond._tag).toBe('Some'); + expect(result.fromFirst._tag).toBe('None'); + }); + + it('detach 后 has=false 且 emit 静默(不抛错)', async () => { + const result = await run( + Effect.gen(function* () { + const sink = yield* EventSinkService; + yield* sink.attach('s2'); + const before = yield* sink.has('s2'); + yield* sink.detach('s2'); + const after = yield* sink.has('s2'); + yield* sink.emit('s2', textDelta('dropped')); + return { before, after }; + }) + ); + + expect(result.before).toBe(true); + expect(result.after).toBe(false); + }); + + it('未挂载的会话 emit 是静默丢弃', async () => { + await expect(run(Effect.gen(function* () { + const sink = yield* EventSinkService; + yield* sink.emit('never-attached', textDelta('x')); + return yield* sink.has('never-attached'); + }))).resolves.toBe(false); + }); + + it('同一队列内先入先出:投递顺序 == 取出顺序', async () => { + const frames = await run( + Effect.gen(function* () { + const sink = yield* EventSinkService; + const q = yield* sink.attach('s3'); + yield* sink.emit('s3', textDelta('first')); + yield* sink.emit('s3', toolCall('t1')); + yield* sink.emit('s3', textDelta('third')); + return yield* takeN(q, 3); + }) + ); + + expect(frames.map((f) => (f.family === 'event' ? f.event.type : f.family))).toEqual([ + 'text_delta', + 'tool_call', + 'text_delta', + ]); + }); +}); diff --git a/packages/codingcode/test/subagent/dispatch-end-to-end.test.ts b/packages/codingcode/test/subagent/dispatch-end-to-end.test.ts index 3a2cdd75..3081a377 100644 --- a/packages/codingcode/test/subagent/dispatch-end-to-end.test.ts +++ b/packages/codingcode/test/subagent/dispatch-end-to-end.test.ts @@ -10,6 +10,7 @@ import { AgentService } from '../../src/agent/port.js'; import { ApprovalService } from '../../src/approval/port.js'; import { CheckpointService } from '../../src/checkpoint/port.js'; import { ContextService } from '../../src/context/port.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; import { LLMService } from '../../src/llm/port.js'; import { MemoryService } from '../../src/memory/port.js'; import { RulesService } from '../../src/rules/port.js'; @@ -98,9 +99,12 @@ const AgentDeps = Layer.mergeAll( extractSkill: (_cwd: string, query: string) => Effect.succeed([undefined, query]), } as any), Layer.succeed(ContextService, { - willCompact: () => Effect.succeed(false), - assemblePayload: (ref: SessionRef) => Effect.sync(() => readMessages(transcriptPathFor(ref))), + getHistory: (ref: SessionRef) => Effect.sync(() => readMessages(transcriptPathFor(ref))), + absorb: () => Effect.void, + compact: () => Effect.succeed({ didCompress: false, released: 0, promptEstimate: 0 }), + dispose: () => Effect.void, } as any), + EventSinkLayer, Layer.succeed(MemoryService, { loadMemoryForPrompt: () => Effect.succeed(''), flushSessionToMemory: () => Effect.succeed({ written: false, bytes: 0 }), diff --git a/packages/codingcode/test/subagent/dispatch-production-path.test.ts b/packages/codingcode/test/subagent/dispatch-production-path.test.ts index 2f911a9d..1e95b6a7 100644 --- a/packages/codingcode/test/subagent/dispatch-production-path.test.ts +++ b/packages/codingcode/test/subagent/dispatch-production-path.test.ts @@ -11,6 +11,7 @@ import { SubagentRunnerLayer } from '../../src/subagent/subagent.js'; import { ApprovalService } from '../../src/approval/port.js'; import { CheckpointService } from '../../src/checkpoint/port.js'; import { ContextService } from '../../src/context/port.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; import { LLMService } from '../../src/llm/port.js'; import { MemoryService } from '../../src/memory/port.js'; import { RulesService } from '../../src/rules/port.js'; @@ -127,9 +128,12 @@ const AgentDeps = Layer.mergeAll( extractSkill: (_cwd: string, query: string) => Effect.succeed([undefined, query]), } as any), Layer.succeed(ContextService, { - willCompact: () => Effect.succeed(false), - assemblePayload: (ref: SessionRef) => Effect.sync(() => readMessages(transcriptPathFor(ref))), + getHistory: (ref: SessionRef) => Effect.sync(() => readMessages(transcriptPathFor(ref))), + absorb: () => Effect.void, + compact: () => Effect.succeed({ didCompress: false, released: 0, promptEstimate: 0 }), + dispose: () => Effect.void, } as any), + EventSinkLayer, Layer.succeed(MemoryService, { loadMemoryForPrompt: () => Effect.succeed(''), flushSessionToMemory: () => Effect.succeed({ written: false, bytes: 0 }), diff --git a/packages/codingcode/test/subagent/runner-wiring.test.ts b/packages/codingcode/test/subagent/runner-wiring.test.ts index 3f4aa5bc..80857dbc 100644 --- a/packages/codingcode/test/subagent/runner-wiring.test.ts +++ b/packages/codingcode/test/subagent/runner-wiring.test.ts @@ -11,6 +11,7 @@ import { SubagentRunnerService } from '../../src/subagent/port.js'; import { ApprovalService } from '../../src/approval/port.js'; import { CheckpointService } from '../../src/checkpoint/port.js'; import { ContextService } from '../../src/context/port.js'; +import { EventSinkLayer } from '../../src/sink/sink.js'; import { LLMService } from '../../src/llm/port.js'; import { AgentError } from '../../src/core/error.js'; import { MemoryService } from '../../src/memory/port.js'; @@ -132,9 +133,12 @@ const AgentDeps = Layer.mergeAll( extractSkill: (_cwd: string, query: string) => Effect.succeed([undefined, query]), } as any), Layer.succeed(ContextService, { - willCompact: () => Effect.succeed(false), - assemblePayload: (ref: SessionRef) => Effect.sync(() => readMessages(transcriptPathFor(ref))), + getHistory: (ref: SessionRef) => Effect.sync(() => readMessages(transcriptPathFor(ref))), + absorb: () => Effect.void, + compact: () => Effect.succeed({ didCompress: false, released: 0, promptEstimate: 0 }), + dispose: () => Effect.void, } as any), + EventSinkLayer, Layer.succeed(MemoryService, { loadMemoryForPrompt: () => Effect.succeed(''), flushSessionToMemory: () => Effect.succeed({ written: false, bytes: 0 }), From 05be66025dc84ee4b0e4524ce956fe059c0ea99e Mon Sep 17 00:00:00 2001 From: phantom5099 <1011668688@qq.com> Date: Sun, 4 Oct 2026 23:20:47 +0800 Subject: [PATCH 2/2] =?UTF-8?q?=E5=88=A0=E6=8E=89=E5=86=97=E4=BD=99?= =?UTF-8?q?=E6=A0=A1=E9=AA=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- packages/codingcode/src/approval/approval.ts | 12 ------------ packages/codingcode/src/sink/port.ts | 2 -- packages/codingcode/src/sink/sink.ts | 1 - packages/codingcode/test/approval/pipeline.test.ts | 8 +++----- .../codingcode/test/plan/gate-pipeline.test.ts | 1 - .../test/security/plan-profile-restart.test.ts | 2 +- packages/codingcode/test/sink/sink.test.ts | 14 +++++--------- 7 files changed, 9 insertions(+), 31 deletions(-) diff --git a/packages/codingcode/src/approval/approval.ts b/packages/codingcode/src/approval/approval.ts index b07e537a..cfb4eecf 100644 --- a/packages/codingcode/src/approval/approval.ts +++ b/packages/codingcode/src/approval/approval.ts @@ -93,8 +93,6 @@ export function runPipeline( ): Effect.Effect { return Effect.gen(function* () { const hooks = yield* HookService; - const sink = yield* EventSinkService; - const asyncConfirm = yield* sink.has(opts.sessionId); const layers: string[] = []; // Layer 1: Rule Engine @@ -164,16 +162,6 @@ export function runPipeline( { layers.push(LAYER_NAMES[3]); - if (!asyncConfirm) { - const result: ApprovalDecision = { - type: 'deny', - reason: 'Approval required but no UI available', - source: 'system', - }; - const final = yield* recordAuditAndReturn(hooks, request, result, layers, opts.projectPath); - return final; - } - const confirmResult = yield* userConfirmAsync( request.tool, request.input, diff --git a/packages/codingcode/src/sink/port.ts b/packages/codingcode/src/sink/port.ts index 7b747547..7cb235a3 100644 --- a/packages/codingcode/src/sink/port.ts +++ b/packages/codingcode/src/sink/port.ts @@ -9,8 +9,6 @@ export interface EventSinkShape { detach(sessionId: string): Effect.Effect; /** 任何模块投帧,插在同一队尾,与回合自己的帧严格全序 */ emit(sessionId: string, body: FrameBody): Effect.Effect; - /** 该会话当前是否有消费者(approval 用它判「有没有 UI」) */ - has(sessionId: string): Effect.Effect; } export class EventSinkService extends Context.Tag('EventSink')() {} diff --git a/packages/codingcode/src/sink/sink.ts b/packages/codingcode/src/sink/sink.ts index fb57d250..c6afcffd 100644 --- a/packages/codingcode/src/sink/sink.ts +++ b/packages/codingcode/src/sink/sink.ts @@ -20,7 +20,6 @@ export const EventSinkLayer = Layer.effect( const q = queues.get(sessionId); if (q) Effect.runSync(Queue.offer(q, body)); }), - has: (sessionId: string) => Effect.sync(() => queues.has(sessionId)), }; }) ); diff --git a/packages/codingcode/test/approval/pipeline.test.ts b/packages/codingcode/test/approval/pipeline.test.ts index 709b25c6..f16a09c5 100644 --- a/packages/codingcode/test/approval/pipeline.test.ts +++ b/packages/codingcode/test/approval/pipeline.test.ts @@ -14,7 +14,7 @@ const mockHookService = { }; const mockApprovalWaitService = { - waitForConfirm: () => Effect.dieMessage('not implemented'), + waitForConfirm: () => Effect.succeed({ type: 'deny' } as const), resolveConfirm: () => Effect.succeed(false), cancelPendingFor: () => Effect.succeed(0), }; @@ -23,7 +23,6 @@ const mockEventSink = { attach: () => Effect.succeed({} as any), detach: () => Effect.void, emit: () => Effect.void, - has: () => Effect.succeed(false), }; const HookTestLayer = Layer.succeed(HookService, mockHookService); @@ -55,7 +54,7 @@ describe('Approval Pipeline — PermissionMode auto-allow (merged from ReadonlyW expect((decision as any).source).toContain('rule:'); }); - it('ask mode does NOT auto-allow read-only tools (no UI → system deny)', async () => { + it('ask mode routes read-only tools to user confirmation', async () => { const decision = await runWithLayer( runPipeline( { tool: 'read_file', input: { path: '/safe/file.txt' } }, @@ -68,8 +67,7 @@ describe('Approval Pipeline — PermissionMode auto-allow (merged from ReadonlyW ) ); expect((decision as any).type).toBe('deny'); - expect((decision as any).source).toBe('system'); - expect((decision as any).reason).toBe('Approval required but no UI available'); + expect((decision as any).source).toBe('user-confirm'); }); it('acceptEdits mode auto-allows read-only tools (read-only merged into non-destructive)', async () => { diff --git a/packages/codingcode/test/plan/gate-pipeline.test.ts b/packages/codingcode/test/plan/gate-pipeline.test.ts index 79914d58..d2107b69 100644 --- a/packages/codingcode/test/plan/gate-pipeline.test.ts +++ b/packages/codingcode/test/plan/gate-pipeline.test.ts @@ -40,7 +40,6 @@ function makeMockEventSink() { capturedApproval = { sessionId, id: body.event.id, tool: body.event.tool, args: body.event.args }; } }), - has: () => Effect.succeed(true), }; } diff --git a/packages/codingcode/test/security/plan-profile-restart.test.ts b/packages/codingcode/test/security/plan-profile-restart.test.ts index 90fef911..08af6fec 100644 --- a/packages/codingcode/test/security/plan-profile-restart.test.ts +++ b/packages/codingcode/test/security/plan-profile-restart.test.ts @@ -22,7 +22,7 @@ const mockHookService = { }; const mockApprovalWaitService = { - waitForConfirm: () => Effect.dieMessage('not implemented'), + waitForConfirm: () => Effect.succeed({ type: 'deny' } as const), resolveConfirm: () => Effect.succeed(false), cancelPendingFor: () => Effect.succeed(0), }; diff --git a/packages/codingcode/test/sink/sink.test.ts b/packages/codingcode/test/sink/sink.test.ts index 4a4cac8d..f2b43091 100644 --- a/packages/codingcode/test/sink/sink.test.ts +++ b/packages/codingcode/test/sink/sink.test.ts @@ -41,29 +41,25 @@ describe('EventSink', () => { expect(result.fromFirst._tag).toBe('None'); }); - it('detach 后 has=false 且 emit 静默(不抛错)', async () => { + it('detach 后 emit 静默丢弃(不抛错,旧队列收不到)', async () => { const result = await run( Effect.gen(function* () { const sink = yield* EventSinkService; - yield* sink.attach('s2'); - const before = yield* sink.has('s2'); + const q = yield* sink.attach('s2'); yield* sink.detach('s2'); - const after = yield* sink.has('s2'); yield* sink.emit('s2', textDelta('dropped')); - return { before, after }; + return yield* Queue.poll(q); }) ); - expect(result.before).toBe(true); - expect(result.after).toBe(false); + expect(result._tag).toBe('None'); }); it('未挂载的会话 emit 是静默丢弃', async () => { await expect(run(Effect.gen(function* () { const sink = yield* EventSinkService; yield* sink.emit('never-attached', textDelta('x')); - return yield* sink.has('never-attached'); - }))).resolves.toBe(false); + }))).resolves.toBeUndefined(); }); it('同一队列内先入先出:投递顺序 == 取出顺序', async () => {