Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
98 changes: 51 additions & 47 deletions packages/codingcode/src/agent/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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 };
});
Expand All @@ -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<FrameBody>;
emit: (body: FrameBody) => void;
onEnd: () => void;
}): AsyncGenerator<FrameBody> {
const q = Effect.runSync(Queue.unbounded<FrameBody>());

const program = agentLoopInternal(opts, q);
const program = agentLoopInternal(opts, out.emit);

return (async function* () {
const fiber = Effect.runFork(opts.toolEnv.provide(program));
Expand All @@ -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<FrameBody>) {
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<FrameBody>) {
yield body;
}
} finally {
out.onEnd();
}
})();
}
Expand All @@ -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<FrameBody>): Effect.Effect<Result<string, AgentError>, AgentError> {
}, emit: (body: FrameBody) => void): Effect.Effect<Result<string, AgentError>, AgentError> {
const { state, model, profile, abortSignal, catalog, rulesText, sid, projectPath, permissionMode } = opts;
const { tools, lookup: toolLookup } = catalog;

Expand All @@ -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* () {
Expand All @@ -209,13 +222,13 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () {
let lastResult: Result<string, AgentError> | 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 = {
Expand All @@ -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[] = [];
Expand All @@ -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 },
}));
});
}
}
},
Expand All @@ -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') {
Expand All @@ -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;
}

Expand All @@ -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[] = [];
Expand Down Expand Up @@ -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,
Expand Down
17 changes: 4 additions & 13 deletions packages/codingcode/src/approval/approval.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

Expand Down Expand Up @@ -89,11 +90,9 @@ function recordAuditAndReturn(
export function runPipeline(
request: ToolCallRequest,
opts: PipelineOptions
): Effect.Effect<ApprovalDecision, never, HookService | ApprovalWaitService> {
): Effect.Effect<ApprovalDecision, never, HookService | EventSinkService | ApprovalWaitService> {
return Effect.gen(function* () {
const hooks = yield* HookService;
const approvalWait = yield* ApprovalWaitService;
const asyncConfirm = yield* approvalWait.hasEmitter(opts.sessionId);
const layers: string[] = [];

// Layer 1: Rule Engine
Expand Down Expand Up @@ -163,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,
Expand Down Expand Up @@ -206,6 +195,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);
Expand All @@ -232,6 +222,7 @@ export const ApprovalLayer = Layer.effect(ApprovalService, Effect.gen(function*
}
).pipe(
Effect.provideService(HookService, hooks),
Effect.provideService(EventSinkService, sink),
Effect.provideService(ApprovalWaitService, approvalWait)
),
};
Expand Down
6 changes: 4 additions & 2 deletions packages/codingcode/src/approval/confirmation.ts
Original file line number Diff line number Diff line change
@@ -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' }
Expand All @@ -13,12 +14,13 @@ export function userConfirmAsync(
args: Record<string, unknown>,
sessionId: string,
callId: string
): Effect.Effect<ConfirmResult, never, ApprovalWaitService> {
): Effect.Effect<ConfirmResult, never, ApprovalWaitService | EventSinkService> {
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);
});
Expand Down
7 changes: 2 additions & 5 deletions packages/codingcode/src/approval/wait-port.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,8 @@ import type { ConfirmResult } from './confirmation.js';
export interface ApprovalWaitShape {
waitForConfirm(id: string, sessionId: string): Effect.Effect<ConfirmResult>;
resolveConfirm(id: string, sessionId: string, result: ConfirmResult): Effect.Effect<boolean>;
emitApprovalRequest(sessionId: string, id: string, tool: string, args: Record<string, unknown>): Effect.Effect<void>;
registerEmitter(sessionId: string, fn: (id: string, tool: string, args: Record<string, unknown>) => void): Effect.Effect<void>;
delegateEmitter(childSessionId: string, parentSessionId: string): Effect.Effect<void>;
unregisterEmitter(sessionId: string): Effect.Effect<void>;
hasEmitter(sessionId: string): Effect.Effect<boolean>;
/** 会话结束时按 sessionId 清掉待决审批(fail-closed 成 deny),返回清理条数 */
cancelPendingFor(sessionId: string): Effect.Effect<number>;
}

export class ApprovalWaitService extends Context.Tag('ApprovalWait')<ApprovalWaitService, ApprovalWaitShape>() {}
44 changes: 9 additions & 35 deletions packages/codingcode/src/approval/wait.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, PendingEntry>();
const approvalEmitters = new Map<
string,
(id: string, tool: string, args: Record<string, unknown>) => void
>();

return {
waitForConfirm: (id: string, sessionId: string): Effect.Effect<ConfirmResult> =>
Expand All @@ -35,38 +31,16 @@ export const ApprovalWaitLayer = Layer.effect(ApprovalWaitService, Effect.gen(fu
return true;
}),

emitApprovalRequest: (
sessionId: string,
id: string,
tool: string,
args: Record<string, unknown>
): Effect.Effect<void> =>
cancelPendingFor: (sessionId: string): Effect.Effect<number> =>
Effect.sync(() => {
approvalEmitters.get(sessionId)?.(id, tool, args);
}),

registerEmitter: (
sessionId: string,
fn: (id: string, tool: string, args: Record<string, unknown>) => void
): Effect.Effect<void> =>
Effect.sync(() => {
approvalEmitters.set(sessionId, fn);
}),

delegateEmitter: (childSessionId: string, parentSessionId: string): Effect.Effect<void> =>
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<void> =>
Effect.sync(() => {
approvalEmitters.delete(sessionId);
}),

hasEmitter: (sessionId: string): Effect.Effect<boolean> =>
Effect.sync(() => approvalEmitters.has(sessionId)),
};
}));
Loading
Loading