diff --git a/packages/agent-core-v2/docs/di-scope-domains.puml b/packages/agent-core-v2/docs/di-scope-domains.puml
index bea8498ad..c456a503e 100644
--- a/packages/agent-core-v2/docs/di-scope-domains.puml
+++ b/packages/agent-core-v2/docs/di-scope-domains.puml
@@ -63,7 +63,7 @@ package "Session scope (per session)" #EAFAF1 {
package "Agent scope (per agent)" #FDF5E6 {
rectangle "wireRecord\nAgent\n IAgentWireRecordService (event hub)" as wireRecord #FDEBD0
- rectangle "eventSink\nAgent\n IAgentEventSinkService" as eventSink #FDEBD0
+ rectangle "record\nAgent\n IAgentRecordService" as record #FDEBD0
rectangle "blobStore\nAgent\n IAgentBlobStoreService" as blobStore #FDEBD0
rectangle "contextMemory\nAgent\n IAgentContextMemoryService" as contextMemory #FDEBD0
rectangle "contextProjector\nAgent\n IAgentContextProjectorService" as contextProjector #FDEBD0
@@ -109,7 +109,7 @@ package "Agent scope (per agent)" #FDF5E6 {
' ---- DI injection (solid) ----
gateway --> session_lifecycle #34495E
gateway --> log #34495E
-gateway --> eventSink #34495E
+gateway --> record #34495E
sessionIndex --> bootstrap #34495E
sessionIndex --> storage #34495E
session_lifecycle --> bootstrap #34495E
@@ -165,19 +165,19 @@ wireRecord --> bootstrap #34495E
wireRecord --> storage #34495E
blobStore --> storage #34495E
filestore --> storage #34495E
-eventSink --> wireRecord #34495E
+record --> wireRecord #34495E
contextMemory --> wireRecord #34495E
contextMemory --> replayBuilder #34495E
contextInjector --> contextMemory #34495E
contextInjector --> turn #34495E
contextInjector --> systemReminder #34495E
contextSize --> contextMemory #34495E
-contextSize --> eventSink #34495E
+contextSize --> record #34495E
contextSize --> wireRecord #34495E
systemReminder --> contextMemory #34495E
replayBuilder --> wireRecord #34495E
profile --> wireRecord #34495E
-profile --> eventSink #34495E
+profile --> record #34495E
profile --> replayBuilder #34495E
profile --> telemetry #34495E
profile --> config #34495E
@@ -186,9 +186,9 @@ profile --> chatProvider #34495E
prompt --> contextMemory #34495E
prompt --> turn #34495E
prompt --> wireRecord #34495E
-prompt --> eventSink #34495E
+prompt --> record #34495E
turn --> loop #34495E
-turn --> eventSink #34495E
+turn --> record #34495E
turn --> wireRecord #34495E
turn --> contextMemory #34495E
turn --> externalHooks #34495E
@@ -197,7 +197,7 @@ loop --> contextMemory #34495E
loop --> contextProjector #34495E
loop --> contextSize #34495E
loop --> llmRequester #34495E
-loop --> eventSink #34495E
+loop --> record #34495E
loop --> toolRegistry #34495E
loop --> toolExecutor #34495E
loop --> profile #34495E
@@ -226,7 +226,7 @@ permissionGate --> externalHooks #34495E
permissionGate --> telemetry #34495E
permissionGate --> toolExecutor #34495E
permissionMode --> wireRecord #34495E
-permissionMode --> eventSink #34495E
+permissionMode --> record #34495E
permissionMode --> replayBuilder #34495E
permissionMode --> contextInjector #34495E
permissionRules --> wireRecord #34495E
@@ -234,7 +234,7 @@ permissionRules --> replayBuilder #34495E
permissionRules --> config #34495E
plan --> contextMemory #34495E
plan --> wireRecord #34495E
-plan --> eventSink #34495E
+plan --> record #34495E
plan --> agentFs #34495E
plan --> profile #34495E
plan --> replayBuilder #34495E
@@ -242,13 +242,13 @@ plan --> toolRegistry #34495E
plan --> contextInjector #34495E
plan --> telemetry #34495E
goal --> wireRecord #34495E
-goal --> eventSink #34495E
+goal --> record #34495E
goal --> systemReminder #34495E
goal --> replayBuilder #34495E
goal --> telemetry #34495E
goal --> contextInjector #34495E
skill --> prompt #34495E
-skill --> eventSink #34495E
+skill --> record #34495E
skill --> wireRecord #34495E
skill --> telemetry #34495E
skill --> sessionSkillCatalog #34495E
@@ -259,7 +259,7 @@ questionTools --> telemetry #34495E
userTool --> toolRegistry #34495E
userTool --> profile #34495E
userTool --> wireRecord #34495E
-background --> eventSink #34495E
+background --> record #34495E
background --> wireRecord #34495E
background --> telemetry #34495E
background --> prompt #34495E
@@ -269,7 +269,7 @@ background --> config #34495E
background --> storage #34495E
background --> session_context #34495E
cron --> prompt #34495E
-cron --> eventSink #34495E
+cron --> record #34495E
cron --> wireRecord #34495E
cron --> turn #34495E
cron --> telemetry #34495E
@@ -277,11 +277,11 @@ cron --> toolRegistry #34495E
cron --> config #34495E
cron --> storage #34495E
swarm --> wireRecord #34495E
-swarm --> eventSink #34495E
+swarm --> record #34495E
swarm --> systemReminder #34495E
swarm --> agentTool #34495E
mcp --> toolRegistry #34495E
-mcp --> eventSink #34495E
+mcp --> record #34495E
mcp --> toolExecutor #34495E
fullCompaction --> contextMemory #34495E
fullCompaction --> contextProjector #34495E
@@ -291,7 +291,7 @@ fullCompaction --> profile #34495E
fullCompaction --> toolStore #34495E
fullCompaction --> telemetry #34495E
fullCompaction --> wireRecord #34495E
-fullCompaction --> eventSink #34495E
+fullCompaction --> record #34495E
fullCompaction --> replayBuilder #34495E
fullCompaction --> externalHooks #34495E
fullCompaction --> turn #34495E
@@ -310,7 +310,7 @@ todoList --> toolStore #34495E
todoList --> toolRegistry #34495E
todoList --> contextInjector #34495E
usage --> wireRecord #34495E
-usage --> eventSink #34495E
+usage --> record #34495E
rpc --> prompt #34495E
rpc --> turn #34495E
rpc --> profile #34495E
diff --git a/packages/agent-core-v2/examples/config.example.ts b/packages/agent-core-v2/examples/config.example.ts
index 97d553426..505e30e76 100644
--- a/packages/agent-core-v2/examples/config.example.ts
+++ b/packages/agent-core-v2/examples/config.example.ts
@@ -56,7 +56,7 @@ import { IAgentContextMemoryService } from '#/agent/contextMemory';
import { IAgentContextProjectorService } from '#/agent/contextProjector';
import { IAgentContextSizeService } from '#/agent/contextSize';
import { AgentCronService } from '#/agent/cron';
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import { AgentExternalHooksService, IAgentExternalHooksService } from '#/agent/externalHooks';
import { IFlagService } from '#/app/flag';
import { IAgentLLMRequesterService } from '#/agent/llmRequester';
@@ -186,7 +186,7 @@ describe('config slice (every section owner against one shared registry)', () =>
// Collaborators declared but not touched during construction — empty
// stubs keep the container strict-clean (no "unknown service" warnings).
- reg.definePartialInstance(IAgentEventSinkService, {});
+ reg.definePartialInstance(IAgentRecordService, {});
reg.definePartialInstance(ITelemetryService, { track: () => {} });
reg.definePartialInstance(IAgentPromptService, {});
reg.definePartialInstance(IAtomicDocumentStore, {});
diff --git a/packages/agent-core-v2/scripts/check-domain-layers.mjs b/packages/agent-core-v2/scripts/check-domain-layers.mjs
index e72179ce3..d0edc5d44 100644
--- a/packages/agent-core-v2/scripts/check-domain-layers.mjs
+++ b/packages/agent-core-v2/scripts/check-domain-layers.mjs
@@ -81,7 +81,6 @@ const DOMAIN_LAYER = new Map([
['model', 2],
['session-index', 2],
['sessionStore', 2],
- ['eventSink', 2],
// L3 — registries & capabilities
['tool', 3],
['skill', 3],
@@ -136,6 +135,9 @@ const DOMAIN_LAYER = new Map([
['mcp', 5],
['cron', 5],
['agentTool', 5],
+ // `btw` forks a single side-question sub-agent via `agent-lifecycle`, mirroring
+ // the `agentTool` shape (Agent-scope, spawns one child) — same layer.
+ ['btw', 5],
// L6 — coordination
['agent-lifecycle', 6],
['session-lifecycle', 6],
@@ -251,6 +253,7 @@ const ALLOWED_EXCEPTIONS = new Set([
'skill>prompt',
'swarm>agentTool',
'swarm>session-metadata',
+ 'btw>agent-lifecycle',
'agentTool>agent-lifecycle',
'agentTool>session-metadata',
'toolExecutor>loop',
diff --git a/packages/agent-core-v2/scripts/dep-graph/analyzer/analyze.ts b/packages/agent-core-v2/scripts/dep-graph/analyzer/analyze.ts
index 46bb658fc..a2a93ab5a 100644
--- a/packages/agent-core-v2/scripts/dep-graph/analyzer/analyze.ts
+++ b/packages/agent-core-v2/scripts/dep-graph/analyzer/analyze.ts
@@ -7,7 +7,7 @@
* - `ctor` — constructor DI (`@IToken` param decorators)
* - `accessor` — runtime lookups (`.get(IToken)`)
* - `publish`/`subscribe` — `IEventService` usage from a class field
- * - `emit`/`on` — `IAgentEventSinkService` usage from a class field
+ * - `signal`/`append`/`on` — `IAgentRecordService` usage from a class field
*
* Deliberately parse-only (no type checker) so the whole tree runs in ~1s.
* We rely on the codebase convention that constructor DI params carry an
@@ -40,12 +40,13 @@ export const REPO_ROOT = resolve(PKG_ROOT, '..', '..');
export const SRC_ROOT = join(PKG_ROOT, 'src');
export const SNAPSHOT_PATH = join(PKG_ROOT, '.local', 'dep-graph.json');
-const EVENT_BUS_TOKENS = new Set(['IEventService', 'IAgentEventSinkService']);
+const EVENT_BUS_TOKENS = new Set(['IEventService', 'IAgentRecordService']);
const EVENT_METHOD_KIND: Record = {
publish: 'publish',
subscribe: 'subscribe',
- emit: 'emit',
+ append: 'emit',
+ signal: 'emit',
on: 'on',
};
@@ -160,9 +161,36 @@ function pushEdge(
token: string,
kind: EdgeKind,
ref: EdgeRef,
+ /**
+ * When set, resolve `token` from this scope instead of the source service's
+ * scope. Used for `.accessor.get(IX)` where the handle is statically
+ * known to belong to an inner scope (e.g. an `IAgentScopeHandle`), so the
+ * lookup resolves against that inner scope rather than the source's.
+ */
+ overrideScope?: ServiceScope,
): void {
- const target = resolveFromScope(acc.bindings, token, source.scope);
- const toId = target ? target.id : `unresolved::${token}`;
+ const target = resolveFromScope(acc.bindings, token, overrideScope ?? source.scope);
+
+ // Classify a miss: the token is either unknown everywhere (genuinely
+ // unresolved) or registered at a scope the resolution scope can't see (a
+ // scope mismatch). The two render differently in the viewer.
+ let toId: string;
+ let extra: Pick;
+ if (target) {
+ toId = target.id;
+ extra = {};
+ } else {
+ const scopeMap = acc.bindings.get(token);
+ const actualScope = scopeMap ? innermostScope(scopeMap) : undefined;
+ if (actualScope !== undefined) {
+ toId = `scopeMismatch::${token}`;
+ extra = { scopeMismatch: true as const, actualScope };
+ } else {
+ toId = `unresolved::${token}`;
+ extra = { unresolved: true as const };
+ }
+ }
+
const key = edgeKey(fromId, toId, kind);
const existing = acc.edges.get(key);
if (existing) {
@@ -177,10 +205,23 @@ function pushEdge(
token,
kind,
refs: [ref],
- ...(target ? {} : { unresolved: true as const }),
+ ...extra,
};
acc.edges.set(key, edge);
- if (!target) acc.unknownRefs.add(token);
+ if (extra.unresolved) acc.unknownRefs.add(token);
+}
+
+function innermostScope(scopeMap: Map): ServiceScope | undefined {
+ let best: ServiceScope | undefined;
+ let bestLevel = -1;
+ for (const s of scopeMap.keys()) {
+ const lvl = SCOPE_LEVEL[s];
+ if (lvl > bestLevel) {
+ bestLevel = lvl;
+ best = s;
+ }
+ }
+ return best;
}
function sameRef(a: EdgeRef, b: EdgeRef): boolean {
@@ -438,6 +479,251 @@ function chainedMethodName(getCall: CallExpression): string | undefined {
return pae.getName();
}
+/**
+ * Map scope-typed handle aliases (and the generic `IScopeHandle`
+ * form) to their scope. Used to resolve `.accessor.get(IX)` against the
+ * handle's real scope rather than the source service's scope.
+ */
+const HANDLE_ALIAS_SCOPE: Record = {
+ IAppScopeHandle: 'App',
+ ISessionScopeHandle: 'Session',
+ IAgentScopeHandle: 'Agent',
+};
+
+const FUNCTION_LIKE_KINDS = new Set([
+ SyntaxKind.MethodDeclaration,
+ SyntaxKind.FunctionDeclaration,
+ SyntaxKind.ArrowFunction,
+ SyntaxKind.FunctionExpression,
+ SyntaxKind.Constructor,
+ SyntaxKind.GetAccessor,
+ SyntaxKind.SetAccessor,
+]);
+
+/** Strip `Promise<...>`, `| undefined` / `| null`, array brackets, `readonly`. */
+function stripTypeWrappers(text: string): string {
+ let t = text.trim();
+ t = t.replace(/\s*\|\s*(undefined|null)\s*/g, '').trim();
+ const promise = /^Promise\s*<\s*(.+?)\s*>$/.exec(t);
+ if (promise) t = promise[1].trim();
+ t = t.replace(/\[\]\s*$/, '').trim();
+ t = t.replace(/^readonly\s+/, '').trim();
+ return t;
+}
+
+function handleScopeFromTypeText(text: string | undefined): ServiceScope | undefined {
+ if (text === undefined) return undefined;
+ const t = stripTypeWrappers(text);
+ const alias = HANDLE_ALIAS_SCOPE[t];
+ if (alias !== undefined) return alias;
+ const generic = /^IScopeHandle\s*<\s*LifecycleScope\.(App|Session|Agent)\s*>$/.exec(t);
+ if (generic) return generic[1] as ServiceScope;
+ return undefined;
+}
+
+/** Nearest function-like ancestor (method / function / arrow / ctor / accessor). */
+function enclosingFunction(node: Node): Node | undefined {
+ let cur: Node | undefined = node.getParent();
+ while (cur) {
+ if (FUNCTION_LIKE_KINDS.has(cur.getKind())) return cur;
+ cur = cur.getParent();
+ }
+ return undefined;
+}
+
+/** Parameters of a function-like node (all such nodes carry `getParameters`). */
+function getParams(fn: Node): ParameterDeclaration[] {
+ return (fn as unknown as { getParameters(): ParameterDeclaration[] }).getParameters();
+}
+
+function isAccessorReceiver(node: Node): boolean {
+ if (node.getKind() !== SyntaxKind.PropertyAccessExpression) return false;
+ return node.asKindOrThrow(SyntaxKind.PropertyAccessExpression).getName() === 'accessor';
+}
+
+/**
+ * Per-interface method → declared return-type text. Lets the analyzer resolve
+ * what `agents.getHandle(...)` returns once it knows `agents: IAgentLifecycleService`.
+ */
+function collectInterfaceMethodReturns(
+ interfacesByName: Map,
+): Map> {
+ const out = new Map>();
+ for (const [name, iface] of interfacesByName) {
+ const methods = new Map();
+ for (const member of iface.getMembers()) {
+ if (member.getKind() === SyntaxKind.MethodSignature) {
+ const m = member.asKindOrThrow(SyntaxKind.MethodSignature);
+ const rt = m.getReturnTypeNode()?.getText();
+ if (rt) methods.set(m.getName(), rt);
+ }
+ }
+ out.set(name, methods);
+ }
+ return out;
+}
+
+/**
+ * Best-effort, parse-only type-text inference for an expression within a
+ * function. Handles just enough to follow scope-typed handles:
+ * - parameter / variable annotations,
+ * - `.accessor.get(IToken)` → the token (DI accessor returns its type),
+ * - `this.method(...)` → the class method's declared return type,
+ * - `.method(...)` → the interface method's declared return type,
+ * - `await X` and single-step identifier aliases.
+ * Returns `undefined` when it can't tell — callers fall back to source scope.
+ * `depth` bounds identifier-chasing so we don't loop on aliased locals.
+ */
+function inferExprTypeText(
+ expr: Node,
+ cls: ClassDeclaration,
+ ifaceMethods: Map>,
+ fn: Node,
+ depth = 0,
+): string | undefined {
+ if (depth > 6) return undefined;
+ const kind = expr.getKind();
+
+ if (kind === SyntaxKind.AwaitExpression) {
+ const inner = (expr as unknown as { getExpression(): Node }).getExpression();
+ return inferExprTypeText(inner, cls, ifaceMethods, fn, depth + 1);
+ }
+
+ if (kind === SyntaxKind.AsExpression || kind === SyntaxKind.NonNullExpression) {
+ const inner = (expr as unknown as { getExpression(): Node }).getExpression();
+ return inferExprTypeText(inner, cls, ifaceMethods, fn, depth + 1);
+ }
+
+ if (kind === SyntaxKind.CallExpression) {
+ const call = expr.asKindOrThrow(SyntaxKind.CallExpression);
+ const callee = call.getExpression();
+ if (callee.getKind() !== SyntaxKind.PropertyAccessExpression) return undefined;
+ const pae = callee.asKindOrThrow(SyntaxKind.PropertyAccessExpression);
+ const methodName = pae.getName();
+ const base = pae.getExpression();
+
+ // `.accessor.get(IToken)` → the DI accessor returns the token's type.
+ if (methodName === 'get' && isAccessorReceiver(base)) {
+ const first = call.getArguments()[0];
+ if (first && first.getKind() === SyntaxKind.Identifier) return first.getText();
+ return undefined;
+ }
+
+ // `this.method(...)` → class method's declared return type.
+ if (base.getKind() === SyntaxKind.ThisKeyword) {
+ return cls.getMethod(methodName)?.getReturnTypeNode()?.getText();
+ }
+
+ // `.method(...)` → resolve base to an interface, look up the method.
+ const baseType = inferExprTypeText(base, cls, ifaceMethods, fn, depth + 1);
+ if (baseType === undefined) return undefined;
+ return ifaceMethods.get(stripTypeWrappers(baseType))?.get(methodName);
+ }
+
+ if (kind === SyntaxKind.Identifier) {
+ return resolveIdentifierTypeText(expr, cls, ifaceMethods, fn, depth + 1);
+ }
+
+ // `this.` → the field's declared type (ctor parameter property or
+ // class property annotation).
+ if (kind === SyntaxKind.PropertyAccessExpression) {
+ const pae = expr.asKindOrThrow(SyntaxKind.PropertyAccessExpression);
+ if (pae.getExpression().getKind() === SyntaxKind.ThisKeyword) {
+ return thisFieldTypeText(cls, pae.getName());
+ }
+ return undefined;
+ }
+
+ // `a ?? b` → either branch's type.
+ if (kind === SyntaxKind.BinaryExpression) {
+ const bin = expr.asKindOrThrow(SyntaxKind.BinaryExpression);
+ if (bin.getOperatorToken().getKind() === SyntaxKind.QuestionQuestionToken) {
+ return (
+ inferExprTypeText(bin.getLeft(), cls, ifaceMethods, fn, depth + 1) ??
+ inferExprTypeText(bin.getRight(), cls, ifaceMethods, fn, depth + 1)
+ );
+ }
+ return undefined;
+ }
+
+ // `cond ? a : b` → either branch's type.
+ if (kind === SyntaxKind.ConditionalExpression) {
+ const cond = expr.asKindOrThrow(SyntaxKind.ConditionalExpression);
+ return (
+ inferExprTypeText(cond.getWhenTrue(), cls, ifaceMethods, fn, depth + 1) ??
+ inferExprTypeText(cond.getWhenFalse(), cls, ifaceMethods, fn, depth + 1)
+ );
+ }
+
+ return undefined;
+}
+
+function thisFieldTypeText(cls: ClassDeclaration, fieldName: string): string | undefined {
+ // Constructor parameter property: `constructor(@IX private readonly foo: IX)`.
+ const ctor = cls.getConstructors()[0];
+ if (ctor) {
+ for (const p of ctor.getParameters()) {
+ if (p.getName() !== fieldName) continue;
+ const t = p.getTypeNode()?.getText();
+ if (t) return t;
+ }
+ }
+ // Class property with an explicit type annotation.
+ return cls.getProperty(fieldName)?.getTypeNode()?.getText();
+}
+
+function resolveIdentifierTypeText(
+ id: Node,
+ cls: ClassDeclaration,
+ ifaceMethods: Map>,
+ fn: Node,
+ depth: number,
+): string | undefined {
+ const name = id.getText();
+
+ for (const p of getParams(fn)) {
+ if (p.getName() === name) {
+ const t = p.getTypeNode()?.getText();
+ if (t) return t;
+ }
+ }
+
+ const decls = fn.getDescendantsOfKind(SyntaxKind.VariableDeclaration);
+ for (const decl of decls) {
+ if (decl.getName() !== name) continue;
+ if (decl.getStart() > id.getStart()) continue;
+ const annotated = decl.getTypeNode()?.getText();
+ if (annotated) return annotated;
+ const init = decl.getInitializer();
+ if (init) {
+ const inferred = inferExprTypeText(init, cls, ifaceMethods, fn, depth + 1);
+ if (inferred) return inferred;
+ }
+ }
+ return undefined;
+}
+
+/**
+ * For a `.accessor.get(IX)` call, return the scope of the handle ``
+ * when it can be inferred from a scope-typed handle alias. Returns `undefined`
+ * for the scope-agnostic cases (base `IScopeHandle`, unions, injected
+ * accessors) so the caller keeps the default source-scope resolution.
+ */
+function inferAccessorScope(
+ getCall: CallExpression,
+ cls: ClassDeclaration,
+ ifaceMethods: Map>,
+): ServiceScope | undefined {
+ const getExpr = getCall.getExpression();
+ if (getExpr.getKind() !== SyntaxKind.PropertyAccessExpression) return undefined;
+ const receiver = getExpr.asKindOrThrow(SyntaxKind.PropertyAccessExpression).getExpression();
+ if (!isAccessorReceiver(receiver)) return undefined;
+ const obj = receiver.asKindOrThrow(SyntaxKind.PropertyAccessExpression).getExpression();
+ const fn = enclosingFunction(getCall);
+ if (fn === undefined) return undefined;
+ return handleScopeFromTypeText(inferExprTypeText(obj, cls, ifaceMethods, fn));
+}
+
/**
* Pass 2 — for a given impl class, walk method bodies and detect:
* - `.get(IToken)[.method(...)]` → accessor edge (with optional `toMethod`)
@@ -450,6 +736,7 @@ function collectRuntimeEdges(
source: ServiceNode,
injectedFields: Map,
acc: EdgeAccumulator,
+ ifaceMethods: Map>,
): void {
const filePath = relFromRepo(cls.getSourceFile().getFilePath());
@@ -476,7 +763,11 @@ function collectRuntimeEdges(
const toMethod = chainedMethodName(call);
const ref: EdgeRef = { ...baseRef };
if (toMethod !== undefined) ref.toMethod = toMethod;
- pushEdge(acc, source.id, source, tokenName, 'accessor', ref);
+ // If this is `.accessor.get(IX)` and the handle's scope is
+ // statically known (e.g. `agent: IAgentScopeHandle`), resolve IX against
+ // that scope instead of the source service's scope.
+ const accessorScope = inferAccessorScope(call, cls, ifaceMethods);
+ pushEdge(acc, source.id, source, tokenName, 'accessor', ref, accessorScope);
continue;
}
@@ -542,6 +833,7 @@ export function analyze(options: { srcRoot?: string; generatedAt?: string } = {}
const { services, implClasses, bindings } = collectServices(sourceFiles);
const interfacesByName = collectInterfaces(sourceFiles);
+ const ifaceMethods = collectInterfaceMethodReturns(interfacesByName);
// Seed the framework tokens as synthetic nodes so edges to them resolve
// like any other registered service. They are marked domain=`framework`
@@ -630,7 +922,7 @@ export function analyze(options: { srcRoot?: string; generatedAt?: string } = {}
if (dep.token === svc.token) continue;
pushEdge(acc, svc.id, svc, dep.token, 'ctor', { file: filePath, line: dep.line });
}
- collectRuntimeEdges(cls, svc, injectedFields, acc);
+ collectRuntimeEdges(cls, svc, injectedFields, acc, ifaceMethods);
}
// Synthesise interface-only nodes for tokens referenced by edges but with no
@@ -680,6 +972,35 @@ export function analyze(options: { srcRoot?: string; generatedAt?: string } = {}
services.push(node);
}
+ // Synthesise scope-mismatch nodes: tokens that ARE registered but referenced
+ // from a scope that can't see them. Placed at the token's real registered
+ // scope (from `actualScope`) and flagged so the viewer styles them apart
+ // from genuinely missing implementations.
+ const mismatchTokens = new Map();
+ for (const edge of acc.edges.values()) {
+ if (!edge.scopeMismatch || edge.actualScope === undefined) continue;
+ if (!mismatchTokens.has(edge.token)) mismatchTokens.set(edge.token, edge.actualScope);
+ }
+ for (const [token, scope] of mismatchTokens) {
+ const registered = acc.bindings.get(token)?.get(scope);
+ const node: ServiceNode = {
+ id: `scopeMismatch::${token}`,
+ token,
+ impl: token,
+ scope,
+ domain: registered?.domain ?? 'unknown',
+ file: '',
+ line: 0,
+ scopeMismatch: true,
+ };
+ const iface = interfacesByName.get(token);
+ if (iface) {
+ const members = collectInterfaceMembers(iface);
+ if (members.length > 0) node.publicMembers = members;
+ }
+ services.push(node);
+ }
+
return {
generatedAt: options.generatedAt ?? new Date(0).toISOString(),
services: services.sort(
diff --git a/packages/agent-core-v2/scripts/dep-graph/analyzer/types.ts b/packages/agent-core-v2/scripts/dep-graph/analyzer/types.ts
index d11de159f..df621905f 100644
--- a/packages/agent-core-v2/scripts/dep-graph/analyzer/types.ts
+++ b/packages/agent-core-v2/scripts/dep-graph/analyzer/types.ts
@@ -15,9 +15,9 @@ export type EdgeKind =
| 'publish'
/** `.subscribe(...)` — subscribes to `IEventService`. */
| 'subscribe'
- /** `.emit(...)` — emits on `IAgentEventSinkService`. */
+ /** `.signal(...)` / `.append(...)` — emits on `IAgentRecordService`. */
| 'emit'
- /** `.on(...)` — listens on `IAgentEventSinkService`. */
+ /** `.on(...)` — listens on `IAgentRecordService`. */
| 'on';
export interface ServiceNode {
@@ -57,6 +57,13 @@ export interface ServiceNode {
* services rather than being dropped as dangling edges.
*/
unresolved?: true;
+ /**
+ * True for synthesized scope-mismatch nodes: the token IS registered, but at
+ * a scope invisible to the edge's source. Rendered distinctly (and placed at
+ * the token's real registered scope) so a cross-scope reach reads differently
+ * from a genuinely missing implementation.
+ */
+ scopeMismatch?: true;
}
export interface EdgeRef {
@@ -90,18 +97,29 @@ export interface Edge {
* Resolved target `ServiceNode.id` — the concrete registration that the
* DI container would actually pick when the source is instantiated. For
* `unresolved: true` edges this is the token that couldn't be resolved,
- * prefixed with `unresolved::` for visualisation purposes.
+ * prefixed with `unresolved::`; for `scopeMismatch: true` edges it is
+ * prefixed with `scopeMismatch::`.
*/
to: string;
/** The interface/decorator name that appears at the source site. */
token: string;
kind: EdgeKind;
/**
- * True when there is no impl registered for `token` at a scope visible
- * from the source (walking source scope up to root). A `ctor` edge in
- * this state would crash the container at instantiation time.
+ * True when there is no impl registered for `token` at ANY scope — the
+ * token is simply unknown to the container. A `ctor` edge in this state
+ * would crash the container at instantiation time.
*/
unresolved?: true;
+ /**
+ * True when the token IS registered, but only at a scope that is not
+ * visible from the source (e.g. an App-scope service reaching for an
+ * Agent-scope token through an accessor whose scope the analyzer couldn't
+ * pin down). Distinct from `unresolved`: an implementation exists, the
+ * edge just can't be satisfied from where it is requested.
+ */
+ scopeMismatch?: true;
+ /** When `scopeMismatch`, the innermost scope where `token` is registered. */
+ actualScope?: ServiceScope;
/** One or more locations that produced this edge (deduped). */
refs: EdgeRef[];
}
diff --git a/packages/agent-core-v2/src/agent/agentTool/runChildAgent.ts b/packages/agent-core-v2/src/agent/agentTool/runChildAgent.ts
index 29a91455f..4df2b41ba 100644
--- a/packages/agent-core-v2/src/agent/agentTool/runChildAgent.ts
+++ b/packages/agent-core-v2/src/agent/agentTool/runChildAgent.ts
@@ -21,14 +21,14 @@ import {
import { linkAbortSignal, userCancellationReason } from '#/_base/utils/abort';
import { IAgentLifecycleService } from '#/session/agent-lifecycle';
-import type { IScopeHandle } from '#/_base/di/scope';
+import type { IAgentScopeHandle } from '#/_base/di/scope';
import {
IAgentContextMemoryService,
type ContextMessage,
type PromptOrigin,
} from '#/agent/contextMemory';
import { ErrorCodes, toKimiErrorPayload, type KimiErrorPayload } from '#/errors';
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import { IAgentExternalHooksService } from '#/agent/externalHooks';
import { isAbortError } from '#/agent/loop/errors';
import { IAgentProfileService } from '#/agent/profile';
@@ -183,7 +183,7 @@ export function cancelAllChildren(
async function ensureParent(
lifecycle: IAgentLifecycleService,
parentAgentId: string,
-): Promise {
+): Promise {
const existing = lifecycle.getHandle(parentAgentId);
if (existing !== undefined) return existing;
throw new Error(`Parent agent "${parentAgentId}" does not exist`);
@@ -194,7 +194,7 @@ async function requireChild(
parentAgentId: string,
metadata: ISessionMetadata | undefined,
agentId: string,
-): Promise {
+): Promise {
if (metadata !== undefined) {
const meta = (await metadata.read()).agents?.[agentId];
if (meta === undefined) throw new Error(`Agent instance "${agentId}" does not exist`);
@@ -212,13 +212,13 @@ async function requireChild(
return child;
}
-function ensureAgentTool(child: IScopeHandle): void {
+function ensureAgentTool(child: IAgentScopeHandle): void {
// Force-instantiate the child agent's `Agent` tool registrar so its `Agent`
// tool is registered before the child's first turn builds its tool list.
child.accessor.get(IAgentToolService);
}
-function configureChild(parent: IScopeHandle, child: IScopeHandle, profileName: string): void {
+function configureChild(parent: IAgentScopeHandle, child: IAgentScopeHandle, profileName: string): void {
const parentProfile = parent.accessor.get(IAgentProfileService);
const childProfile = child.accessor.get(IAgentProfileService);
const parentData = parentProfile.data();
@@ -241,13 +241,13 @@ function configureChild(parent: IScopeHandle, child: IScopeHandle, profileName:
}
function emitSpawned(
- parent: IScopeHandle,
+ parent: IAgentScopeHandle,
parentAgentId: string,
subagentId: string,
profileName: string,
options: RunSubagentOptions,
): void {
- parent.accessor.get(IAgentEventSinkService)?.emit({
+ parent.accessor.get(IAgentRecordService)?.signal({
type: 'subagent.spawned',
subagentId,
subagentName: profileName,
@@ -264,17 +264,17 @@ function emitSpawned(
});
}
-function emitStarted(parent: IScopeHandle, subagentId: string): void {
- parent.accessor.get(IAgentEventSinkService)?.emit({ type: 'subagent.started', subagentId });
+function emitStarted(parent: IAgentScopeHandle, subagentId: string): void {
+ parent.accessor.get(IAgentRecordService)?.signal({ type: 'subagent.started', subagentId });
}
function emitCompleted(
- parent: IScopeHandle,
+ parent: IAgentScopeHandle,
subagentId: string,
resultSummary: string,
usage?: TokenUsage,
): void {
- parent.accessor.get(IAgentEventSinkService)?.emit({
+ parent.accessor.get(IAgentRecordService)?.signal({
type: 'subagent.completed',
subagentId,
resultSummary,
@@ -283,14 +283,14 @@ function emitCompleted(
}
function emitFailed(
- parent: IScopeHandle,
+ parent: IAgentScopeHandle,
subagentId: string,
error: unknown,
options: RunSubagentOptions,
): void {
if (isAbortError(error)) return;
if (shouldSuppressQueuedAttemptFailureEvent(options, error)) return;
- parent.accessor.get(IAgentEventSinkService)?.emit({
+ parent.accessor.get(IAgentRecordService)?.signal({
type: 'subagent.failed',
subagentId,
error: errorMessage(error),
@@ -299,7 +299,7 @@ function emitFailed(
async function triggerSubagentStart(
- parent: IScopeHandle,
+ parent: IAgentScopeHandle,
profileName: string,
prompt: string,
signal: AbortSignal,
@@ -313,7 +313,7 @@ async function triggerSubagentStart(
);
}
-function triggerSubagentStop(parent: IScopeHandle, profileName: string, result: string): void {
+function triggerSubagentStop(parent: IAgentScopeHandle, profileName: string, result: string): void {
parent.accessor.get(IAgentExternalHooksService)?.triggerSubagentStop({
agentName: profileName,
response: result.slice(0, HOOK_TEXT_PREVIEW_LENGTH),
@@ -327,9 +327,9 @@ function observeFirstRequest(turn: Turn, options: RunSubagentOptions): void {
async function runWithActiveChild(
parentAgentId: string,
- child: IScopeHandle,
+ child: IAgentScopeHandle,
options: RunSubagentOptions,
- parent: IScopeHandle,
+ parent: IAgentScopeHandle,
profileName: string,
run: (
turn: { current?: Turn },
@@ -359,8 +359,8 @@ async function runWithActiveChild(
}
async function runPromptTurn(
- child: IScopeHandle,
- parent: IScopeHandle,
+ child: IAgentScopeHandle,
+ parent: IAgentScopeHandle,
options: RunSubagentOptions,
profileName: string,
turnRef: { current?: Turn },
@@ -389,8 +389,8 @@ async function runPromptTurn(
}
async function runRetryTurn(
- child: IScopeHandle,
- parent: IScopeHandle,
+ child: IAgentScopeHandle,
+ parent: IAgentScopeHandle,
options: RunSubagentOptions,
profileName: string,
turnRef: { current?: Turn },
@@ -429,7 +429,7 @@ async function awaitTurn(
}
async function completeSummary(
- child: IScopeHandle,
+ child: IAgentScopeHandle,
controller: AbortController,
turnRef: { current?: Turn },
): Promise {
diff --git a/packages/agent-core-v2/src/agent/background/backgroundService.ts b/packages/agent-core-v2/src/agent/background/backgroundService.ts
index b7d4d2eeb..29242ffdb 100644
--- a/packages/agent-core-v2/src/agent/background/backgroundService.ts
+++ b/packages/agent-core-v2/src/agent/background/backgroundService.ts
@@ -4,9 +4,9 @@
* Owns the agent's registry of running and restored background tasks:
* registers and drives tasks to completion, retains a bounded output ring,
* persists task state and output through `background` persistence, reads
- * limits through `config`, records lifecycle through `wireRecord`, delivers
- * terminal notifications through `contextMemory`, and broadcasts through
- * `eventSink`. Bound at Agent scope.
+ * limits through `config`, records lifecycle and broadcasts through `record`,
+ * and delivers terminal notifications through `contextMemory`. Bound at Agent
+ * scope.
*/
import { randomBytes } from 'node:crypto';
@@ -27,15 +27,13 @@ import {
import { IAgentContextMemoryService } from '#/agent/contextMemory';
import { IConfigService } from '#/app/config';
-import { IAgentEventSinkService } from '#/agent/eventSink';
import { IAgentExternalHooksService } from '#/agent/externalHooks';
import { IAgentPromptService } from '#/agent/prompt';
import { ISessionContext } from '#/session/session-context';
import { IAtomicDocumentStore, IStorageService } from '#/app/storage';
import { ITelemetryService } from '#/app/telemetry';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
-import type { WireRecord } from '#/agent/wireRecord';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentRecordService, type AgentRecord } from '#/agent/record';
import {
IAgentBackgroundService,
type BackgroundLoadOptions,
@@ -133,8 +131,7 @@ export class AgentBackgroundService extends Disposable implements IAgentBackgrou
private readonly persistence: BackgroundTaskPersistence;
constructor(
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@ITelemetryService private readonly telemetry: ITelemetryService,
@IAgentPromptService private readonly prompt: IAgentPromptService,
@IAgentExternalHooksService private readonly externalHooks: IAgentExternalHooksService,
@@ -153,17 +150,21 @@ export class AgentBackgroundService extends Disposable implements IAgentBackgrou
byteStore,
);
this._register(
- wireRecord.register('background.task.started', (record) => {
- this.applyRestoredTask(record);
+ record.define('background.task.started', {
+ resume: (r) => {
+ this.applyRestoredTask(r);
+ },
}),
);
this._register(
- wireRecord.register('background.task.terminated', (record) => {
- this.applyRestoredTask(record);
+ record.define('background.task.terminated', {
+ resume: (r) => {
+ this.applyRestoredTask(r);
+ },
}),
);
this._register(
- wireRecord.hooks.onResumeEnded.register(
+ record.hooks.onResumeEnded.register(
'background-lifecycle-resume',
async (_ctx, next) => {
await this.loadFromDisk({ replace: false });
@@ -561,7 +562,7 @@ export class AgentBackgroundService extends Disposable implements IAgentBackgrou
}
private applyRestoredTask(
- record: WireRecord<'background.task.started' | 'background.task.terminated'>,
+ record: AgentRecord<'background.task.started' | 'background.task.terminated'>,
): void {
const info = record.info;
if (this.tasks.has(info.taskId)) return;
@@ -685,16 +686,16 @@ export class AgentBackgroundService extends Disposable implements IAgentBackgrou
}
private recordTaskStarted(info: BackgroundTaskInfo): void {
- this.wireRecord.append({ type: 'background.task.started', info });
- this.events.emit({ type: 'background.task.started', info });
+ this.record.append({ type: 'background.task.started', info });
+ this.record.signal({ type: 'background.task.started', info });
this.telemetry.track('background_task_created', {
kind: info.kind === 'process' ? 'bash' : info.kind,
});
}
private recordTaskTerminated(info: BackgroundTaskInfo): void {
- this.wireRecord.append({ type: 'background.task.terminated', info });
- this.events.emit({ type: 'background.task.terminated', info });
+ this.record.append({ type: 'background.task.terminated', info });
+ this.record.signal({ type: 'background.task.terminated', info });
this.telemetry.track('background_task_completed', {
kind: info.kind,
duration: info.endedAt !== null ? info.endedAt - info.startedAt : null,
diff --git a/packages/agent-core-v2/src/agent/contextMemory/contextMemoryService.ts b/packages/agent-core-v2/src/agent/contextMemory/contextMemoryService.ts
index 00666e1b0..043b0452d 100644
--- a/packages/agent-core-v2/src/agent/contextMemory/contextMemoryService.ts
+++ b/packages/agent-core-v2/src/agent/contextMemory/contextMemoryService.ts
@@ -3,7 +3,7 @@ import {
} from "#/_base/di";
import { OrderedHookSlot } from '#/hooks';
import { IAgentReplayBuilderService } from '#/agent/replayBuilder';
-import { IAgentWireRecordService, type WireRecord } from '#/agent/wireRecord';
+import { IAgentRecordService, type AgentRecord } from '#/agent/record';
import { IAgentContextMemoryService } from './contextMemory';
import { ensureMessageId } from './messageId';
import type { ContextMessage } from './types';
@@ -35,18 +35,17 @@ export class AgentContextMemoryService extends Disposable implements IAgentConte
};
constructor(
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentReplayBuilderService private readonly replayBuilder: IAgentReplayBuilderService,
) {
super();
this._register(
- wireRecord.register(
- 'context.splice',
- (record) => {
- this.applySplice(record);
+ record.define('context.splice', {
+ resume: (r) => {
+ this.applySplice(r);
},
- {
- blobs: (record) => record.messages.map((message, index) => ({
+ blobs: (r) =>
+ r.messages.map((message, index) => ({
parts: message.content,
replace: (current, content) => ({
...current,
@@ -55,8 +54,7 @@ export class AgentContextMemoryService extends Disposable implements IAgentConte
),
}),
})),
- },
- ),
+ }),
);
}
@@ -71,18 +69,18 @@ export class AgentContextMemoryService extends Disposable implements IAgentConte
tokens?: number,
): void {
const stamped = messages.map(ensureMessageId);
- const record: WireRecord<'context.splice'> = {
+ const record: AgentRecord<'context.splice'> = {
type: 'context.splice',
start,
deleteCount,
messages: stamped,
tokens,
};
- this.wireRecord.append(record);
+ this.record.append(record);
this.applySplice(record);
}
- private applySplice(record: WireRecord<'context.splice'>): void {
+ private applySplice(record: AgentRecord<'context.splice'>): void {
const removedMessages =
record.deleteCount > 0 && record.start > 0
? this.history.slice(record.start, record.start + record.deleteCount)
diff --git a/packages/agent-core-v2/src/agent/contextSize/contextSizeService.ts b/packages/agent-core-v2/src/agent/contextSize/contextSizeService.ts
index 86082265b..ca7fc5f94 100644
--- a/packages/agent-core-v2/src/agent/contextSize/contextSizeService.ts
+++ b/packages/agent-core-v2/src/agent/contextSize/contextSizeService.ts
@@ -6,8 +6,7 @@ import {
} from "#/_base/utils/tokens";
import type { ContextMessage } from '#/agent/contextMemory';
import { IAgentContextMemoryService } from '#/agent/contextMemory';
-import { IAgentEventSinkService } from '#/agent/eventSink';
-import { IAgentWireRecordService, type WireRecord } from '#/agent/wireRecord';
+import { IAgentRecordService, type AgentRecord } from '#/agent/record';
import {
IAgentContextSizeService,
type ContextSizeStatus,
@@ -39,8 +38,7 @@ export class AgentContextSizeService
constructor(
@IAgentContextMemoryService private readonly context: IAgentContextMemoryService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
+ @IAgentRecordService private readonly records: IAgentRecordService,
) {
super();
this._register(
@@ -50,8 +48,10 @@ export class AgentContextSizeService
}),
);
this._register(
- wireRecord.register('context_size.measured', (record) => {
- this.applyMeasurement(record);
+ records.define('context_size.measured', {
+ resume: (r) => {
+ this.applyMeasurement(r);
+ },
}),
);
}
@@ -66,12 +66,12 @@ export class AgentContextSizeService
}
measured(length: number, tokens: number): void {
- const record: WireRecord<'context_size.measured'> = {
+ const record: AgentRecord<'context_size.measured'> = {
type: 'context_size.measured',
length,
tokens,
};
- this.wireRecord.append(record);
+ this.records.append(record);
this.applyMeasurement(record);
}
@@ -102,7 +102,7 @@ export class AgentContextSizeService
this.emitIfChanged();
}
- private applyMeasurement(record: WireRecord<'context_size.measured'>): void {
+ private applyMeasurement(record: AgentRecord<'context_size.measured'>): void {
const length = clampMeasuredLength(record.length, this.estimates.length);
const tokens = Math.max(0, record.tokens);
this.measuredPrefixTokens[length] = tokens;
@@ -125,7 +125,7 @@ export class AgentContextSizeService
return;
}
this.lastEmitted = status;
- this.events.emit({
+ this.records.signal({
type: 'agent.status.updated',
contextTokens: status.contextTokens,
});
diff --git a/packages/agent-core-v2/src/agent/cron/cronService.ts b/packages/agent-core-v2/src/agent/cron/cronService.ts
index 0e5d55ff3..3f3b73483 100644
--- a/packages/agent-core-v2/src/agent/cron/cronService.ts
+++ b/packages/agent-core-v2/src/agent/cron/cronService.ts
@@ -22,12 +22,11 @@ import { IConfigService } from '#/app/config';
import { IAtomicDocumentStore } from '#/app/storage';
import { ITelemetryService } from '#/app/telemetry';
import type { ContextMessage } from '#/agent/contextMemory';
-import { IAgentEventSinkService } from '#/agent/eventSink';
import { IAgentPromptService } from '#/agent/prompt';
+import { IAgentRecordService } from '#/agent/record';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
import type { Turn } from '#/agent/turn';
import { IAgentTurnService } from '#/agent/turn';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
import {
type CronConfig,
@@ -145,8 +144,7 @@ export class AgentCronService extends Disposable implements IAgentCronService {
constructor(
options: CronOptions = {},
@IAgentPromptService private readonly prompt: IAgentPromptService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentTurnService private readonly turnService: IAgentTurnService,
@ITelemetryService private readonly telemetry: ITelemetryService,
@IAgentToolRegistryService private readonly toolRegistry: IAgentToolRegistryService,
@@ -167,22 +165,28 @@ export class AgentCronService extends Disposable implements IAgentCronService {
resolveClockSources(this.cronConfig.clock, this.cronConfig.debug) ?? SYSTEM_CLOCKS;
this._register(
- wireRecord.register('cron.add', (record) => {
- if (this.enabled) this.adopt(record.task);
+ record.define('cron.add', {
+ resume: (r) => {
+ if (this.enabled) this.adopt(r.task);
+ },
}),
);
this._register(
- wireRecord.register('cron.delete', (record) => {
- if (this.enabled) this.removeByIds(record.ids);
+ record.define('cron.delete', {
+ resume: (r) => {
+ if (this.enabled) this.removeByIds(r.ids);
+ },
}),
);
this._register(
- wireRecord.register('cron.cursor', (record) => {
- if (this.enabled) this.markFired(record.id, record.lastFiredAt);
+ record.define('cron.cursor', {
+ resume: (r) => {
+ if (this.enabled) this.markFired(r.id, r.lastFiredAt);
+ },
}),
);
this._register(
- wireRecord.hooks.onResumeEnded.register('cron-lifecycle-resume', async (_ctx, next) => {
+ record.hooks.onResumeEnded.register('cron-lifecycle-resume', async (_ctx, next) => {
await this.loadFromDisk({ replace: false });
this.start();
await next();
@@ -221,7 +225,7 @@ export class AgentCronService extends Disposable implements IAgentCronService {
createdAt: this.clocks.wallNow(),
};
this.tasks.set(task.id, task);
- this.wireRecord.append({ type: 'cron.add', task });
+ this.record.append({ type: 'cron.add', task });
this.persistEnqueue(task.id, () => this.atomicDocs.set(CRON_SCOPE, cronKey(task.id), task));
return task;
}
@@ -230,7 +234,7 @@ export class AgentCronService extends Disposable implements IAgentCronService {
const removed = this.removeByIds(ids);
if (removed.length === 0) return removed;
- this.wireRecord.append({ type: 'cron.delete', ids: removed });
+ this.record.append({ type: 'cron.delete', ids: removed });
for (const id of removed) {
this.persistEnqueue(id, () => this.atomicDocs.delete(CRON_SCOPE, cronKey(id)));
}
@@ -457,7 +461,7 @@ export class AgentCronService extends Disposable implements IAgentCronService {
toolCalls: [],
origin,
};
- this.events.emit({ type: 'cron.fired', origin, prompt: task.prompt });
+ this.record.signal({ type: 'cron.fired', origin, prompt: task.prompt });
const turn = this.prompt.steer(message);
this.telemetry.track(CRON_FIRED, {
recurring: task.recurring !== false,
@@ -472,7 +476,7 @@ export class AgentCronService extends Disposable implements IAgentCronService {
const updated = this.markFired(id, lastFiredAt);
if (updated === undefined) return;
- this.wireRecord.append({ type: 'cron.cursor', id, lastFiredAt });
+ this.record.append({ type: 'cron.cursor', id, lastFiredAt });
this.persistEnqueue(id, () => this.atomicDocs.set(CRON_SCOPE, cronKey(id), updated));
}
diff --git a/packages/agent-core-v2/src/agent/eventSink/eventSink.ts b/packages/agent-core-v2/src/agent/eventSink/eventSink.ts
deleted file mode 100644
index f1e72c0e2..000000000
--- a/packages/agent-core-v2/src/agent/eventSink/eventSink.ts
+++ /dev/null
@@ -1,11 +0,0 @@
-import type { IDisposable } from "#/_base/di";
-import { createDecorator } from "#/_base/di";
-import type { AgentEvent } from '@moonshot-ai/protocol';
-
-export interface IAgentEventSinkService {
- readonly _serviceBrand: undefined;
- emit(event: AgentEvent): void;
- on(handler: (event: AgentEvent) => void): IDisposable;
-}
-
-export const IAgentEventSinkService = createDecorator('agentEventSinkService');
diff --git a/packages/agent-core-v2/src/agent/eventSink/eventSinkService.ts b/packages/agent-core-v2/src/agent/eventSink/eventSinkService.ts
deleted file mode 100644
index 48bd0d5cb..000000000
--- a/packages/agent-core-v2/src/agent/eventSink/eventSinkService.ts
+++ /dev/null
@@ -1,37 +0,0 @@
-import { InstantiationType } from '#/_base/di/extensions';
-import { LifecycleScope, registerScopedService } from '#/_base/di/scope';
-
-import {
- Disposable,
-} from "#/_base/di";
-import { Emitter } from "#/_base/event";
-
-import { IAgentWireRecordService } from '#/agent/wireRecord';
-import type { AgentEvent } from '@moonshot-ai/protocol';
-import { IAgentEventSinkService } from './eventSink';
-
-export class AgentEventSinkService extends Disposable implements IAgentEventSinkService {
- declare readonly _serviceBrand: undefined;
- private readonly onDidEmitEmitter = this._register(new Emitter());
-
- constructor(@IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService) {
- super();
- }
-
- emit(event: AgentEvent): void {
- if (this.wireRecord.restoring) return;
- this.onDidEmitEmitter.fire(event);
- }
-
- on(handler: (event: AgentEvent) => void) {
- return this.onDidEmitEmitter.event(handler);
- }
-}
-
-registerScopedService(
- LifecycleScope.Agent,
- IAgentEventSinkService,
- AgentEventSinkService,
- InstantiationType.Delayed,
- 'eventSink',
-);
diff --git a/packages/agent-core-v2/src/agent/eventSink/index.ts b/packages/agent-core-v2/src/agent/eventSink/index.ts
deleted file mode 100644
index d3f121ae5..000000000
--- a/packages/agent-core-v2/src/agent/eventSink/index.ts
+++ /dev/null
@@ -1,6 +0,0 @@
-/**
- * `eventBus` domain barrel - re-exports the eventBus service contract and implementation.
- */
-
-export * from './eventSink';
-export * from './eventSinkService';
diff --git a/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts b/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts
index 0aef4791b..3b02e5d2d 100644
--- a/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts
+++ b/packages/agent-core-v2/src/agent/fullCompaction/fullCompactionService.ts
@@ -18,7 +18,6 @@ import { estimateTokens, estimateTokensForMessages } from "#/_base/utils/tokens"
import { IAgentContextMemoryService } from '#/agent/contextMemory';
import { IAgentContextProjectorService } from '#/agent/contextProjector';
import { IAgentContextSizeService } from '#/agent/contextSize';
-import { IAgentEventSinkService } from '#/agent/eventSink';
import { IAgentExternalHooksService } from '#/agent/externalHooks';
import { IAgentLLMRequesterService, type LLMEvent } from '#/agent/llmRequester';
import { isAbortError } from '#/agent/loop/errors';
@@ -29,7 +28,7 @@ import { ITelemetryService } from '#/app/telemetry';
import { IAgentToolStoreService } from '#/agent/toolStore';
import { IAgentTurnService, type TurnContextOverflowContext } from '#/agent/turn';
import type { ContextMessage } from '#/agent/contextMemory';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentRecordService } from '#/agent/record';
import {
TODO_STORE_KEY,
renderTodoList,
@@ -104,8 +103,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull
@IAgentProfileService private readonly profile: IAgentProfileService,
@IAgentToolStoreService private readonly toolStore: IAgentToolStoreService,
@ITelemetryService private readonly telemetry: ITelemetryService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentReplayBuilderService private readonly replayBuilder: IAgentReplayBuilderService,
@IAgentExternalHooksService private readonly externalHooks: IAgentExternalHooksService,
@IAgentTurnService turnService: IAgentTurnService,
@@ -138,32 +136,38 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull
}),
);
this._register(
- wireRecord.register('full_compaction.begin', (record) => {
- this.replayBuilder.push({
- type: 'compaction',
- instruction: record.instruction,
- });
+ record.define('full_compaction.begin', {
+ resume: (r) => {
+ this.replayBuilder.push({
+ type: 'compaction',
+ instruction: r.instruction,
+ });
+ },
}),
);
this._register(
- wireRecord.register('full_compaction.cancel', () => {
- this.replayBuilder.patchLast('compaction', { result: 'cancelled' });
+ record.define('full_compaction.cancel', {
+ resume: () => {
+ this.replayBuilder.patchLast('compaction', { result: 'cancelled' });
+ },
}),
);
this._register(
- wireRecord.register('full_compaction.complete', (record) => {
- const message = compactionSummaryMessage(this.context.get());
- if (message === undefined) return;
- const summary = contextMessageText(message);
- this.replayBuilder.removeLastMessages(new Set([message]));
- this.replayBuilder.patchLast('compaction', {
- result: {
- summary,
- compactedCount: record.compactedCount,
- tokensBefore: record.tokensBefore,
- tokensAfter: record.tokensAfter,
- },
- });
+ record.define('full_compaction.complete', {
+ resume: (r) => {
+ const message = compactionSummaryMessage(this.context.get());
+ if (message === undefined) return;
+ const summary = contextMessageText(message);
+ this.replayBuilder.removeLastMessages(new Set([message]));
+ this.replayBuilder.patchLast('compaction', {
+ result: {
+ summary,
+ compactedCount: r.compactedCount,
+ tokensBefore: r.tokensBefore,
+ tokensAfter: r.tokensAfter,
+ },
+ });
+ },
}),
);
}
@@ -188,8 +192,8 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull
throw new KimiError(ErrorCodes.COMPACTION_UNABLE, 'No prefix that can be compacted in current history.');
}
- this.wireRecord.append({ type: 'full_compaction.begin', ...data });
- this.events.emit({
+ this.record.append({ type: 'full_compaction.begin', ...data });
+ this.record.signal({
type: 'compaction.started',
trigger: data.source,
instruction: data.instruction,
@@ -208,15 +212,15 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull
cancel(): void {
const active = this.compacting;
if (active === null) return;
- this.wireRecord.append({ type: 'full_compaction.cancel' });
+ this.record.append({ type: 'full_compaction.cancel' });
active.abortController.abort();
this.compacting = null;
- this.events.emit({ type: 'compaction.cancelled' });
+ this.record.signal({ type: 'compaction.cancelled' });
}
private markCompleted(result: FullCompactionCompleteData): void {
if (this.compacting === null) return;
- this.wireRecord.append({ type: 'full_compaction.complete', ...result });
+ this.record.append({ type: 'full_compaction.complete', ...result });
this.compacting = null;
}
@@ -281,7 +285,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull
}
}, { once: true });
}
- this.events.emit({ type: 'compaction.blocked', turnId });
+ this.record.signal({ type: 'compaction.blocked', turnId });
await active.promise;
}
@@ -318,7 +322,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull
if (this.compacting !== active) return;
this.markCompleted(completeData(finalResult));
- this.events.emit({ type: 'compaction.completed', result: finalResult });
+ this.record.signal({ type: 'compaction.completed', result: finalResult });
this.externalHooks.triggerPostCompact({
trigger: data.source,
estimatedTokenCount: finalResult.tokensAfter,
@@ -332,7 +336,7 @@ export class AgentFullCompactionService extends Disposable implements IAgentFull
if (blockedByTurn) {
throw error;
}
- this.events.emit({
+ this.record.signal({
type: 'error',
...toKimiErrorPayload(error),
});
diff --git a/packages/agent-core-v2/src/agent/goal/goalService.ts b/packages/agent-core-v2/src/agent/goal/goalService.ts
index 2f18aefc7..53ecb9d3d 100644
--- a/packages/agent-core-v2/src/agent/goal/goalService.ts
+++ b/packages/agent-core-v2/src/agent/goal/goalService.ts
@@ -8,15 +8,13 @@ import {
} from "#/_base/di";
import { ErrorCodes, KimiError } from "#/errors";
import { IAgentContextInjectorService } from '#/agent/contextInjector';
-import { IAgentEventSinkService } from '#/agent/eventSink';
import { IAgentPermissionModeService } from '#/agent/permissionMode';
import { IAgentReplayBuilderService } from '#/agent/replayBuilder';
import { IAgentSystemReminderService } from '#/agent/systemReminder';
import type { TelemetryProperties } from '#/app/telemetry';
import { ITelemetryService } from '#/app/telemetry';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
-import type { WireRecord } from '#/agent/wireRecord';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentRecordService, type AgentRecord } from '#/agent/record';
import {
IAgentGoalService,
type GoalReasonInput,
@@ -101,8 +99,7 @@ export class AgentGoalService extends Disposable implements IAgentGoalService {
constructor(
private readonly options: GoalServiceOptions = {},
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentSystemReminderService private readonly reminders: IAgentSystemReminderService,
@IAgentReplayBuilderService private readonly replayBuilder: IAgentReplayBuilderService,
@ITelemetryService private readonly telemetry: ITelemetryService,
@@ -121,27 +118,35 @@ export class AgentGoalService extends Disposable implements IAgentGoalService {
),
);
this._register(
- wireRecord.register('forked', (record) => {
- this.restoreForked(record);
+ record.define('forked', {
+ resume: (r) => {
+ this.restoreForked(r);
+ },
}),
);
this._register(
- wireRecord.register('goal.create', (record) => {
- this.restoreCreate(record);
+ record.define('goal.create', {
+ resume: (r) => {
+ this.restoreCreate(r);
+ },
}),
);
this._register(
- wireRecord.register('goal.update', (record) => {
- this.restoreUpdate(record);
+ record.define('goal.update', {
+ resume: (r) => {
+ this.restoreUpdate(r);
+ },
}),
);
this._register(
- wireRecord.register('goal.clear', () => {
- this.restoreClear();
+ record.define('goal.clear', {
+ resume: () => {
+ this.restoreClear();
+ },
}),
);
this._register(
- wireRecord.hooks.onResumeEnded.register('goal-normalize-after-replay', async (_ctx, next) => {
+ record.hooks.onResumeEnded.register('goal-normalize-after-replay', async (_ctx, next) => {
await next();
this.normalizeAfterReplay();
}),
@@ -207,7 +212,7 @@ export class AgentGoalService extends Disposable implements IAgentGoalService {
};
this.persistState(state);
- this.wireRecord.append({
+ this.record.append({
type: 'goal.create',
goalId: state.goalId,
objective: state.objective,
@@ -378,7 +383,7 @@ export class AgentGoalService extends Disposable implements IAgentGoalService {
this.appendStatusUpdate(state, 'runtime', reason);
}
- private restoreCreate(record: WireRecord<'goal.create'>): void {
+ private restoreCreate(record: AgentRecord<'goal.create'>): void {
const state: GoalState = {
goalId: record.goalId,
objective: record.objective,
@@ -397,7 +402,7 @@ export class AgentGoalService extends Disposable implements IAgentGoalService {
});
}
- private restoreUpdate(record: WireRecord<'goal.update'>): void {
+ private restoreUpdate(record: AgentRecord<'goal.update'>): void {
const state = this.state;
if (state === undefined) return;
@@ -441,7 +446,7 @@ export class AgentGoalService extends Disposable implements IAgentGoalService {
this.state = undefined;
}
- private restoreForked(_record: WireRecord<'forked'>): void {
+ private restoreForked(_record: AgentRecord<'forked'>): void {
const hadGoal = this.state !== undefined;
this.state = undefined;
if (!hadGoal) return;
@@ -457,7 +462,7 @@ export class AgentGoalService extends Disposable implements IAgentGoalService {
): void {
if (this.state === undefined) return;
this.persistState(undefined, { silent: opts.emit === false });
- this.wireRecord.append({ type: 'goal.clear' });
+ this.record.append({ type: 'goal.clear' });
if (opts.track !== false) this.telemetry.track('goal_cleared', { actor });
}
@@ -479,9 +484,9 @@ export class AgentGoalService extends Disposable implements IAgentGoalService {
}
private appendGoalUpdate(
- update: Omit, 'type' | 'time'>,
+ update: Omit, 'type' | 'time'>,
): void {
- this.wireRecord.append({
+ this.record.append({
type: 'goal.update',
...update,
});
@@ -518,7 +523,7 @@ export class AgentGoalService extends Disposable implements IAgentGoalService {
}
private emitGoalUpdated(snapshot: GoalSnapshot | null, change?: GoalChange): void {
- this.events.emit({ type: 'goal.updated', snapshot, change });
+ this.record.signal({ type: 'goal.updated', snapshot, change });
}
private statsOf(state: GoalState): GoalChangeStats {
diff --git a/packages/agent-core-v2/src/agent/loop/events.ts b/packages/agent-core-v2/src/agent/loop/events.ts
index 4f268077c..02801d36b 100644
--- a/packages/agent-core-v2/src/agent/loop/events.ts
+++ b/packages/agent-core-v2/src/agent/loop/events.ts
@@ -136,66 +136,8 @@ export type LoopLiveOnlyEvent =
| LoopToolProgressEvent;
export type LoopEvent = LoopRecordedEvent | LoopLiveOnlyEvent;
-export type LoopLiveEventEmitter = (event: LoopEvent) => void;
export type LoopEventDispatcher = {
(event: LoopRecordedEvent): Promise;
(event: LoopLiveOnlyEvent): void;
};
-
-export interface CreateLoopEventDispatcherInput {
- readonly appendTranscriptRecord: (record: LoopRecordedEvent) => Promise;
- readonly emitLiveEvent?: LoopLiveEventEmitter | undefined;
-}
-
-export function createLoopEventDispatcher(
- input: CreateLoopEventDispatcherInput,
-): LoopEventDispatcher {
- function dispatchEvent(event: LoopRecordedEvent): Promise;
- function dispatchEvent(event: LoopLiveOnlyEvent): void;
- function dispatchEvent(event: LoopEvent): Promise | void {
- if (isRecordedEvent(event)) {
- return recordEvent(input, event);
- }
- safeEmitLive(input.emitLiveEvent, event);
- }
- return dispatchEvent;
-}
-
-function isRecordedEvent(event: LoopEvent): event is LoopRecordedEvent {
- return (
- event.type === 'step.begin' ||
- event.type === 'step.end' ||
- event.type === 'content.part' ||
- event.type === 'tool.call' ||
- event.type === 'tool.result'
- );
-}
-
-async function recordEvent(
- input: CreateLoopEventDispatcherInput,
- event: LoopRecordedEvent,
-): Promise {
- await input.appendTranscriptRecord(event);
- safeEmitLive(input.emitLiveEvent, event);
-}
-
-function safeEmitLive(emit: LoopLiveEventEmitter | undefined, event: LoopEvent): void {
- if (emit === undefined) return;
- let maybePromise: unknown;
- try {
- maybePromise = (emit as (event: LoopEvent) => unknown)(event);
- } catch {
- return;
- }
- if (
- maybePromise !== undefined &&
- maybePromise !== null &&
- typeof (maybePromise as { then?: unknown }).then === 'function' &&
- typeof (maybePromise as { catch?: unknown }).catch === 'function'
- ) {
- (maybePromise as Promise).catch(() => {
- // Live listeners are best-effort; their failures must not affect the turn.
- });
- }
-}
diff --git a/packages/agent-core-v2/src/agent/loop/index.ts b/packages/agent-core-v2/src/agent/loop/index.ts
index 5af36c3c5..fe820ef6f 100644
--- a/packages/agent-core-v2/src/agent/loop/index.ts
+++ b/packages/agent-core-v2/src/agent/loop/index.ts
@@ -27,7 +27,6 @@ export type {
} from './types';
export type {
- CreateLoopEventDispatcherInput,
LoopContentPartEvent,
LoopRecordedEvent,
LoopStepBeginEvent,
@@ -36,7 +35,6 @@ export type {
LoopLiveOnlyEvent,
LoopEvent,
LoopInterruptReason,
- LoopLiveEventEmitter,
LoopEventDispatcher,
LoopTextDeltaEvent,
LoopThinkingDeltaEvent,
@@ -46,7 +44,6 @@ export type {
LoopToolResultEvent,
LoopTurnInterruptedEvent,
} from './events';
-export { createLoopEventDispatcher } from './events';
export type {
LLM,
diff --git a/packages/agent-core-v2/src/agent/loop/loopService.ts b/packages/agent-core-v2/src/agent/loop/loopService.ts
index 4eae62a7c..cc6ecbcbf 100644
--- a/packages/agent-core-v2/src/agent/loop/loopService.ts
+++ b/packages/agent-core-v2/src/agent/loop/loopService.ts
@@ -40,7 +40,7 @@ import { canonicalTelemetryArgs } from '#/_base/utils/canonical-args';
import { IAgentContextMemoryService, newMessageId, type ContextMessage } from '#/agent/contextMemory';
import { IAgentContextProjectorService } from '#/agent/contextProjector';
import { IAgentContextSizeService } from '#/agent/contextSize';
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import { IAgentExternalHooksService } from '#/agent/externalHooks';
import { IAgentLLMRequesterService } from '#/agent/llmRequester';
import { ILogService } from '#/app/log';
@@ -50,7 +50,6 @@ import { ITelemetryService } from '#/app/telemetry';
import { IAgentToolExecutorService } from '#/agent/toolExecutor';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
import type { Turn, TurnResult } from '#/agent/turn';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
import type {
LoopEvent,
LoopEventDispatcher,
@@ -93,12 +92,11 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
@IAgentContextProjectorService private readonly projector: IAgentContextProjectorService,
@IAgentContextSizeService private readonly contextSize: IAgentContextSizeService,
@IAgentLLMRequesterService private readonly llmRequester: IAgentLLMRequesterService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentToolRegistryService private readonly toolRegistry: IAgentToolRegistryService,
@IAgentToolExecutorService private readonly toolExecutor: IAgentToolExecutorService,
@IAgentProfileService private readonly profile: IAgentProfileService,
@ITelemetryService private readonly telemetry: ITelemetryService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
@IAgentExternalHooksService private readonly externalHooks: IAgentExternalHooksService,
@IConfigService private readonly config: IConfigService,
@ILogService private readonly log: ILogService,
@@ -110,7 +108,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
}
await next();
});
- this.wireRecord.hooks.onResumeEnded.register(
+ this.record.hooks.onResumeEnded.register(
'loop-service-finish-resume',
async (_event, next) => {
this.finishResume();
@@ -248,7 +246,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
switch (event.type) {
case 'step.begin':
this.beginTrackedStep(event.step);
- this.events.emit({
+ this.record.signal({
type: 'turn.step.started',
turnId: Number(event.turnId),
step: event.step,
@@ -256,7 +254,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
});
return;
case 'step.end':
- this.events.emit({
+ this.record.signal({
type: 'turn.step.completed',
turnId: Number(event.turnId),
step: event.step,
@@ -274,7 +272,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
});
return;
case 'step.retrying':
- this.events.emit({
+ this.record.signal({
type: 'turn.step.retrying',
turnId: Number(event.turnId),
step: event.step,
@@ -295,7 +293,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
}
}
if (this.protocolTurnId === undefined || event.activeStep === undefined) return;
- this.events.emit({
+ this.record.signal({
type: 'turn.step.interrupted',
turnId: this.protocolTurnId,
step: event.activeStep,
@@ -305,7 +303,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
return;
case 'text.delta':
if (this.protocolTurnId === undefined) return;
- this.events.emit({
+ this.record.signal({
type: 'assistant.delta',
turnId: this.protocolTurnId,
delta: event.delta,
@@ -313,7 +311,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
return;
case 'thinking.delta':
if (this.protocolTurnId === undefined) return;
- this.events.emit({
+ this.record.signal({
type: 'thinking.delta',
turnId: this.protocolTurnId,
delta: event.delta,
@@ -321,7 +319,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
return;
case 'tool.call.delta':
if (this.protocolTurnId === undefined) return;
- this.events.emit({
+ this.record.signal({
type: 'tool.call.delta',
turnId: this.protocolTurnId,
toolCallId: event.toolCallId,
@@ -331,7 +329,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
return;
case 'tool.call':
this.trackToolCallStarted(event);
- this.events.emit({
+ this.record.signal({
type: 'tool.call.started',
turnId: Number(event.turnId),
toolCallId: event.toolCallId,
@@ -343,7 +341,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
return;
case 'tool.progress':
if (this.protocolTurnId === undefined) return;
- this.events.emit({
+ this.record.signal({
type: 'tool.progress',
turnId: this.protocolTurnId,
toolCallId: event.toolCallId,
@@ -353,7 +351,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
case 'tool.result':
if (this.protocolTurnId === undefined) return;
this.trackToolCallResult(event);
- this.events.emit({
+ this.record.signal({
type: 'tool.result',
turnId: this.protocolTurnId,
toolCallId: event.toolCallId,
diff --git a/packages/agent-core-v2/src/agent/mcp/mcpService.ts b/packages/agent-core-v2/src/agent/mcp/mcpService.ts
index 9c58c3904..58b2c6c10 100644
--- a/packages/agent-core-v2/src/agent/mcp/mcpService.ts
+++ b/packages/agent-core-v2/src/agent/mcp/mcpService.ts
@@ -7,7 +7,7 @@ import {
type IDisposable,
} from "#/_base/di";
import { ErrorCodes, makeErrorPayload } from "#/errors";
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import { IAgentToolExecutorService } from '#/agent/toolExecutor';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
import { createMcpAuthTool } from '#/agent/mcp/tools/auth';
@@ -38,7 +38,7 @@ export class AgentMcpService extends Disposable implements IAgentMcpService {
constructor(
private readonly options: McpServiceOptions = {},
@IAgentToolRegistryService private readonly registry: IAgentToolRegistryService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentToolExecutorService toolExecutor: IAgentToolExecutorService,
) {
super();
@@ -103,7 +103,7 @@ export class AgentMcpService extends Disposable implements IAgentMcpService {
}
private handleMcpServerStatusChange(entry: McpServerEntry): void {
- this.events.emit({
+ this.record.signal({
type: 'mcp.server.status',
server: {
name: entry.name,
@@ -123,7 +123,7 @@ export class AgentMcpService extends Disposable implements IAgentMcpService {
}
if (entry.status === 'failed') {
this.unregisterMcpServer(entry.name);
- this.events.emit({
+ this.record.signal({
type: 'tool.list.updated',
reason: 'mcp.failed',
serverName: entry.name,
@@ -133,7 +133,7 @@ export class AgentMcpService extends Disposable implements IAgentMcpService {
if (entry.status === 'disabled' || entry.status === 'pending') {
const removed = this.unregisterMcpServer(entry.name);
if (removed) {
- this.events.emit({
+ this.record.signal({
type: 'tool.list.updated',
reason: 'mcp.disconnected',
serverName: entry.name,
@@ -152,7 +152,7 @@ export class AgentMcpService extends Disposable implements IAgentMcpService {
resolved.enabledNames,
);
this.emitMcpToolCollisions(entry.name, result.collisions);
- this.events.emit({
+ this.record.signal({
type: 'tool.list.updated',
reason: 'mcp.connected',
serverName: entry.name,
@@ -173,7 +173,7 @@ export class AgentMcpService extends Disposable implements IAgentMcpService {
const disposable = this._register(this.registry.register(tool, { source: 'mcp' }));
this.mcpTools.set(tool.name, { disposable, serverName: entry.name });
this.mcpToolsByServer.set(entry.name, [tool.name]);
- this.events.emit({
+ this.record.signal({
type: 'tool.list.updated',
reason: 'mcp.connected',
serverName: entry.name,
@@ -251,7 +251,7 @@ export class AgentMcpService extends Disposable implements IAgentMcpService {
: `"${collision.toolName}" -> ${collision.qualified} (collides with server "${collision.collidesWith.serverName}")`,
)
.join('; ');
- this.events.emit({
+ this.record.signal({
type: 'error',
...makeErrorPayload(
ErrorCodes.MCP_TOOL_NAME_COLLISION,
diff --git a/packages/agent-core-v2/src/agent/microCompaction/microCompactionService.ts b/packages/agent-core-v2/src/agent/microCompaction/microCompactionService.ts
index cfc5c291e..1ad25ad99 100644
--- a/packages/agent-core-v2/src/agent/microCompaction/microCompactionService.ts
+++ b/packages/agent-core-v2/src/agent/microCompaction/microCompactionService.ts
@@ -27,7 +27,7 @@ import { IAgentProfileService } from '#/agent/profile';
import { ITelemetryService } from '#/app/telemetry';
import { IAgentTurnService } from '#/agent/turn';
import type { ContextMessage } from '#/agent/contextMemory';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentRecordService } from '#/agent/record';
import {
IAgentMicroCompactionService,
type MicroCompactionConfig,
@@ -66,7 +66,7 @@ export class AgentMicroCompactionService
constructor(
@IAgentContextMemoryService private readonly context: IAgentContextMemoryService,
@IAgentContextSizeService private readonly contextSize: IAgentContextSizeService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IFlagService private readonly flags: IFlagService,
@IAgentProfileService private readonly profile: IAgentProfileService,
@ITelemetryService private readonly telemetry: ITelemetryService,
@@ -93,13 +93,17 @@ export class AgentMicroCompactionService
),
);
this._register(
- this.wireRecord.register('micro_compaction.apply', (record) => {
- this.apply(record.cutoff);
+ this.record.define('micro_compaction.apply', {
+ resume: (r) => {
+ this.apply(r.cutoff);
+ },
}),
);
this._register(
- this.wireRecord.register('full_compaction.complete', () => {
- this.reset();
+ this.record.define('full_compaction.complete', {
+ resume: () => {
+ this.reset();
+ },
}),
);
this._register(
@@ -119,7 +123,7 @@ export class AgentMicroCompactionService
}
private apply(cutoff: number): void {
- this.wireRecord.append({
+ this.record.append({
type: 'micro_compaction.apply',
cutoff,
});
@@ -210,7 +214,7 @@ export class AgentMicroCompactionService
}
if (context.messages.some(isAssistantCacheAnchor)) {
- this._lastAssistantAt = this.wireRecord.restoring?.time ?? Date.now();
+ this._lastAssistantAt = this.record.restoring?.time ?? Date.now();
}
}
diff --git a/packages/agent-core-v2/src/agent/permissionMode/permissionModeService.ts b/packages/agent-core-v2/src/agent/permissionMode/permissionModeService.ts
index 31c54ee5b..e5c241187 100644
--- a/packages/agent-core-v2/src/agent/permissionMode/permissionModeService.ts
+++ b/packages/agent-core-v2/src/agent/permissionMode/permissionModeService.ts
@@ -6,11 +6,8 @@ import { InstantiationType } from '#/_base/di/extensions';
import { LifecycleScope, registerScopedService } from '#/_base/di/scope';
import { IAgentContextInjectorService } from '#/agent/contextInjector';
-import { IAgentEventSinkService } from '#/agent/eventSink';
import { OrderedHookSlot } from '#/hooks';
-import { IAgentReplayBuilderService } from '#/agent/replayBuilder';
-import type { WireRecord } from '#/agent/wireRecord';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentRecordService } from '#/agent/record';
import { registerPermissionModeInjection } from '#/agent/permissionMode/injection/permissionModeInjection';
import { IAgentPermissionModeService } from './permissionMode';
@@ -35,15 +32,17 @@ export class AgentPermissionModeService extends Disposable implements IAgentPerm
};
constructor(
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
- @IAgentReplayBuilderService private readonly replayBuilder: IAgentReplayBuilderService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentContextInjectorService dynamicInjector: IAgentContextInjectorService,
) {
super();
this._register(
- wireRecord.register('permission.set_mode', (record) => {
- this.applyMode(record);
+ record.define('permission.set_mode', {
+ resume: (r) => {
+ this.applyMode(r.mode);
+ },
+ toLive: (r) => ({ type: 'agent.status.updated', permission: r.mode }),
+ toReplay: (r) => ({ type: 'permission_updated', mode: r.mode }),
}),
);
this._register(
@@ -56,18 +55,13 @@ export class AgentPermissionModeService extends Disposable implements IAgentPerm
}
setMode(mode: PermissionMode): void {
- this.wireRecord.append({ type: 'permission.set_mode', mode });
- this.applyMode({ type: 'permission.set_mode', mode });
+ this.record.append({ type: 'permission.set_mode', mode });
+ this.applyMode(mode);
}
- private applyMode(record: WireRecord<'permission.set_mode'>): void {
- this.replayBuilder.push({ type: 'permission_updated', mode: record.mode });
+ private applyMode(mode: PermissionMode): void {
const previousMode = this.currentMode;
- this.currentMode = record.mode;
- this.events.emit({
- type: 'agent.status.updated',
- permission: this.currentMode,
- });
+ this.currentMode = mode;
void this.hooks.onChanged.run({ mode: this.currentMode, previousMode });
}
}
diff --git a/packages/agent-core-v2/src/agent/plan/planService.ts b/packages/agent-core-v2/src/agent/plan/planService.ts
index e5ba79c07..5dd646d14 100644
--- a/packages/agent-core-v2/src/agent/plan/planService.ts
+++ b/packages/agent-core-v2/src/agent/plan/planService.ts
@@ -15,13 +15,11 @@ import { toInputJsonSchema } from '#/_base/tools/support/input-schema';
import { generateHeroSlug } from "#/_base/utils/hero-slug";
import { IAgentContextMemoryService, type ContextMessage } from '#/agent/contextMemory';
import { IAgentContextInjectorService } from '#/agent/contextInjector';
-import { IAgentEventSinkService } from '#/agent/eventSink';
import { ISessionAgentFileSystem } from '#/session/agentFs';
import { IAgentProfileService } from '#/agent/profile';
-import { IAgentReplayBuilderService } from '#/agent/replayBuilder';
import { ITelemetryService } from '#/app/telemetry';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentRecordService } from '#/agent/record';
import type { ToolInputDisplay } from '@moonshot-ai/protocol';
import type { ExecutableToolResult } from '#/agent/tool';
import { EnterPlanModeInputSchema } from '#/agent/plan/tools/enter-plan-mode';
@@ -65,31 +63,36 @@ export class AgentPlanService extends Disposable implements IAgentPlanService {
constructor(
@IAgentContextMemoryService private readonly context: IAgentContextMemoryService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@ISessionAgentFileSystem private readonly agentFs: ISessionAgentFileSystem,
@IAgentProfileService private readonly profile: IAgentProfileService,
- @IAgentReplayBuilderService private readonly replayBuilder: IAgentReplayBuilderService,
@IAgentToolRegistryService toolRegistry: IAgentToolRegistryService,
@IAgentContextInjectorService dynamicInjector: IAgentContextInjectorService,
@ITelemetryService private readonly telemetry: ITelemetryService,
) {
super();
this._register(
- wireRecord.register('plan_mode.enter', ({ id }) => {
- this.restoreEnter({ id });
+ record.define('plan_mode.enter', {
+ resume: ({ id }) => {
+ this.restoreEnter({ id });
+ },
+ toReplay: () => ({ type: 'plan_updated', enabled: true }),
}),
);
this._register(
- wireRecord.register('plan_mode.cancel', () => {
- this.replayBuilder.push({ type: 'plan_updated', enabled: false });
- this.applyInactive();
+ record.define('plan_mode.cancel', {
+ resume: () => {
+ this.applyInactive();
+ },
+ toReplay: () => ({ type: 'plan_updated', enabled: false }),
}),
);
this._register(
- wireRecord.register('plan_mode.exit', () => {
- this.replayBuilder.push({ type: 'plan_updated', enabled: false });
- this.applyInactive();
+ record.define('plan_mode.exit', {
+ resume: () => {
+ this.applyInactive();
+ },
+ toReplay: () => ({ type: 'plan_updated', enabled: false }),
}),
);
@@ -169,7 +172,7 @@ export class AgentPlanService extends Disposable implements IAgentPlanService {
const planFilePath = this.planFilePathFor(id);
this._planFilePath = planFilePath;
await this.ensurePlanDirectory(planFilePath);
- this.wireRecord.append({ type: 'plan_mode.enter', id });
+ this.record.append({ type: 'plan_mode.enter', id });
enterRecorded = true;
if (createFile) {
await this.writeEmptyPlanFile(planFilePath);
@@ -187,15 +190,13 @@ export class AgentPlanService extends Disposable implements IAgentPlanService {
}
private restoreEnter({ id }: { readonly id: string }): void {
- this.replayBuilder.push({ type: 'plan_updated', enabled: true });
this._active = true;
this.planId = id;
this._planFilePath = this.planFilePathFor(id);
}
cancel(id?: string): void {
- this.wireRecord.append({ type: 'plan_mode.cancel', id });
- this.replayBuilder.push({ type: 'plan_updated', enabled: false });
+ this.record.append({ type: 'plan_mode.cancel', id });
this.applyInactive();
this.emitChanged();
}
@@ -206,8 +207,7 @@ export class AgentPlanService extends Disposable implements IAgentPlanService {
}
exit(id?: string): void {
- this.wireRecord.append({ type: 'plan_mode.exit', id });
- this.replayBuilder.push({ type: 'plan_updated', enabled: false });
+ this.record.append({ type: 'plan_mode.exit', id });
this.applyInactive();
this.emitChanged();
}
@@ -343,7 +343,7 @@ export class AgentPlanService extends Disposable implements IAgentPlanService {
}
private emitChanged(): void {
- this.events.emit({ type: 'agent.status.updated', planMode: this._active });
+ this.record.signal({ type: 'agent.status.updated', planMode: this._active });
}
private trackTelemetry(
diff --git a/packages/agent-core-v2/src/agent/profile/profileService.ts b/packages/agent-core-v2/src/agent/profile/profileService.ts
index 7e3168ec8..d9fef7781 100644
--- a/packages/agent-core-v2/src/agent/profile/profileService.ts
+++ b/packages/agent-core-v2/src/agent/profile/profileService.ts
@@ -34,7 +34,7 @@ import { ISessionModelResolver, type ResolvedModel } from '#/session/modelRuntim
import { ISessionWorkspaceContext } from '#/session/workspaceContext';
import type { ResolvedAgentProfile, SystemPromptContext } from '#/agent/profile';
-import { IAgentRecordService } from '#/agent/record';
+import { IAgentRecordService, type AgentRecord } from '#/agent/record';
import { ITelemetryService } from '#/app/telemetry';
import type { ToolSource } from '#/agent/tool';
import { prepareSystemPromptContext } from './context';
@@ -119,7 +119,7 @@ export class AgentProfileService implements IAgentProfileService {
update(changed: ProfileUpdateData): void {
const { activeToolNames, ...configChanged } = changed;
if (Object.keys(configChanged).length > 0) {
- this.wireRecord.append({ type: 'config.update', ...configChanged });
+ this.record.append({ type: 'config.update', ...configChanged });
this.apply(configChanged);
}
if (activeToolNames !== undefined) {
@@ -173,7 +173,7 @@ export class AgentProfileService implements IAgentProfileService {
const { agentsMdWarning } = context;
this.agentsMdWarning = agentsMdWarning;
if (agentsMdWarning !== undefined) {
- this.events.emit({
+ this.record.signal({
type: 'warning',
message: agentsMdWarning,
code: 'agents-md-oversized',
@@ -274,7 +274,6 @@ export class AgentProfileService implements IAgentProfileService {
}
private apply(changed: ProfileUpdateData): void {
- this.replayBuilder.push({ type: 'config_updated', config: changed });
if (changed.cwd !== undefined) {
this.cwdValue = changed.cwd;
void this.optionsValue.chdir?.(changed.cwd);
@@ -298,7 +297,7 @@ export class AgentProfileService implements IAgentProfileService {
}
private setActiveTools(names: readonly string[]): void {
- this.wireRecord.append({ type: 'tools.set_active_tools', names: [...names] });
+ this.record.append({ type: 'tools.set_active_tools', names: [...names] });
this.applyActiveToolNames(names);
this.initializeBuiltinTools();
}
@@ -319,7 +318,7 @@ export class AgentProfileService implements IAgentProfileService {
return;
}
if (!this.hasModel()) return;
- this.events.emit({
+ this.record.signal({
type: 'agent.status.updated',
model: this.modelAlias,
maxContextTokens: this.getModelCapabilities().max_context_tokens,
@@ -381,6 +380,11 @@ export class AgentProfileService implements IAgentProfileService {
}
}
+function stripConfigMeta(record: AgentRecord<'config.update'>): ProfileUpdateData {
+ const { type: _type, time: _time, ...changed } = record;
+ return changed;
+}
+
registerScopedService(
LifecycleScope.Agent,
IAgentProfileService,
diff --git a/packages/agent-core-v2/src/agent/prompt/promptService.ts b/packages/agent-core-v2/src/agent/prompt/promptService.ts
index e68ad3be5..779f630ae 100644
--- a/packages/agent-core-v2/src/agent/prompt/promptService.ts
+++ b/packages/agent-core-v2/src/agent/prompt/promptService.ts
@@ -8,9 +8,8 @@ import {
USER_PROMPT_ORIGIN,
type ContextMessage,
} from '#/agent/contextMemory';
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import { IAgentTurnService, type Turn } from '#/agent/turn';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
import { IAgentPromptService } from './prompt';
export class AgentPromptService implements IAgentPromptService {
@@ -21,8 +20,7 @@ export class AgentPromptService implements IAgentPromptService {
constructor(
@IAgentContextMemoryService private readonly context: IAgentContextMemoryService,
@IAgentTurnService private readonly turnService: IAgentTurnService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
) {
turnService.hooks.beforeStep.register('prompt-service-steer-before-step', async (_ctx, next) => {
this.flushSteerQueue();
@@ -89,7 +87,7 @@ export class AgentPromptService implements IAgentPromptService {
}
}
- if (!this.wireRecord.restoring && (removedUserCount < count || stoppedAtCompaction)) {
+ if (!this.record.restoring && (removedUserCount < count || stoppedAtCompaction)) {
throw new KimiError(
ErrorCodes.REQUEST_INVALID,
formatUndoUnavailableMessage(count, removedUserCount, stoppedAtCompaction),
@@ -142,7 +140,7 @@ export class AgentPromptService implements IAgentPromptService {
private emitBusyIfActive(): boolean {
const activeTurn = this.turnService.getActiveTurn();
if (activeTurn === undefined) return false;
- this.events.emit({
+ this.record.signal({
type: 'error',
...makeErrorPayload(
ErrorCodes.TURN_AGENT_BUSY,
diff --git a/packages/agent-core-v2/src/agent/record/index.ts b/packages/agent-core-v2/src/agent/record/index.ts
new file mode 100644
index 000000000..d25c6f461
--- /dev/null
+++ b/packages/agent-core-v2/src/agent/record/index.ts
@@ -0,0 +1,8 @@
+/**
+ * `record` domain barrel — re-exports the record service contract and
+ * implementation. Importing this barrel registers the `IAgentRecordService`
+ * binding into the scope registry.
+ */
+
+export * from './record';
+export * from './recordService';
diff --git a/packages/agent-core-v2/src/agent/record/record.ts b/packages/agent-core-v2/src/agent/record/record.ts
new file mode 100644
index 000000000..f39a839f3
--- /dev/null
+++ b/packages/agent-core-v2/src/agent/record/record.ts
@@ -0,0 +1,106 @@
+/**
+ * `record` domain (L3) — `IAgentRecordService` contract.
+ *
+ * Single entry point for recording facts that happen inside an agent. One
+ * `append(record)` call fans out to every facet declared for the record type:
+ * durable persistence (for resume), live broadcast (`AgentEvent` to the edge),
+ * and replay capture. `signal(event)` emits a live-only event that is never
+ * recorded (deltas / progress). Bound at Agent scope.
+ */
+
+import type { AgentEvent } from '@moonshot-ai/protocol';
+
+import type { IDisposable } from '#/_base/di';
+import { createDecorator } from '#/_base/di';
+import type {
+ IAgentWireRecordService,
+ WireRecord,
+ WireRecordBlobSelector,
+ WireRecordMap,
+ WireRecordRestoreOptions,
+ WireRecordRestoreResult,
+ PersistedWireRecord,
+ WireRecordRestoringContext,
+} from '#/agent/wireRecord';
+import type { AgentReplayRecord, AgentReplayRecordPayload } from '#/agent/replayBuilder/types';
+
+/**
+ * Per-agent fact registry. Each domain augments `WireRecordMap` (in
+ * `#/agent/wireRecord`) to declare the payload shape of the records it owns;
+ * `AgentRecordMap` extends it so the same record types are accepted by the
+ * `record` facade and by the persistence layer (`PersistedWireRecord`) without
+ * duplication. New domains may augment either interface.
+ */
+export interface AgentRecordMap extends WireRecordMap {}
+
+export type AgentRecord = {
+ [T in K]: { readonly type: T; readonly time?: number } & Readonly;
+}[K];
+
+/**
+ * Facets declared for a record type. Every facet is optional: a record with no
+ * facets is still persisted (durable) but broadcasts nothing and captures no
+ * replay. Restore-time behavior (resume + replay capture) is driven entirely
+ * by these facets, so business code never re-implements the live/restore split.
+ */
+export interface RecordFacets {
+ /**
+ * Live projection: `record → AgentEvent` broadcast on the wire. When omitted
+ * the record is not broadcast live. Automatically suppressed while restoring.
+ */
+ readonly toLive?: (record: AgentRecord) => AgentEvent;
+ /**
+ * Replay projection: `record →` zero/one/many replay records. Runs on both
+ * the live `append` path and the restore path; the replay layer gates by
+ * phase so live appends are only captured when `captureLiveRecords` is set.
+ */
+ readonly toReplay?: (
+ record: AgentRecord,
+ ) => AgentReplayRecordPayload | readonly AgentReplayRecordPayload[] | undefined;
+ /** Resumer: rebuild in-memory state from a restored record. */
+ readonly resume?: (record: AgentRecord) => void | Promise;
+ /**
+ * Blob offload/rehydrate selector for large content parts (e.g. images in
+ * context messages). Forwarded to the durable store so oversized parts are
+ * offloaded to the blob store on append and rehydrated on restore.
+ */
+ readonly blobs?: WireRecordBlobSelector>;
+}
+
+export interface IAgentRecordService {
+ readonly _serviceBrand: undefined;
+
+ /** Record a fact: persists it and fans out to its declared facets. */
+ append(record: AgentRecord): void;
+ /**
+ * Subscribe to the live `AgentEvent` stream. Suppressed while restoring so
+ * edge consumers never receive historical events. This is the contract the
+ * `server-v2` edge depends on (`on(handler)` → `AgentEvent`).
+ */
+ on(handler: (event: AgentEvent) => void): IDisposable;
+ /** Emit a live-only event that is never persisted (deltas / progress). */
+ signal(event: AgentEvent): void;
+ /**
+ * Declare the facets for a record type. Call from a domain's constructor.
+ * Returns a disposable that unregisters the facets (and the resumer).
+ */
+ define(type: K, facets: RecordFacets): IDisposable;
+
+ restore(
+ records?: readonly PersistedWireRecord[],
+ options?: WireRecordRestoreOptions,
+ ): Promise;
+ flush(): Promise;
+ close(): Promise;
+
+ /** Replay result built from restored (and optionally live) records. */
+ buildReplay(): readonly AgentReplayRecord[];
+ /** When true, live `append` calls also feed the replay buffer. */
+ captureLiveRecords: boolean;
+
+ readonly restoring: WireRecordRestoringContext | null;
+ readonly postRestoring: boolean;
+ readonly hooks: IAgentWireRecordService['hooks'];
+}
+
+export const IAgentRecordService = createDecorator('agentRecordService');
diff --git a/packages/agent-core-v2/src/agent/record/recordService.ts b/packages/agent-core-v2/src/agent/record/recordService.ts
new file mode 100644
index 000000000..f059ed6d5
--- /dev/null
+++ b/packages/agent-core-v2/src/agent/record/recordService.ts
@@ -0,0 +1,176 @@
+/**
+ * `record` domain (L3) — `IAgentRecordService` implementation.
+ *
+ * Owns the unified `append` / `signal` / `define` API: one `append(record)`
+ * fans out to durable persistence (delegated to `wireRecord`), live broadcast
+ * (an owned `Emitter`), and replay capture (delegated to
+ * `replayBuilder`). `signal(event)` emits a live-only event that is never
+ * recorded. Live emission is suppressed while restoring, so edge consumers
+ * never receive historical events. The former `eventSink` service is folded
+ * into this class; `wireRecord` / `replayBuilder` remain registered backends
+ * that this service coordinates.
+ */
+
+import { Disposable, toDisposable } from '#/_base/di';
+import type { IDisposable } from '#/_base/di';
+import { Emitter } from '#/_base/event';
+import { InstantiationType } from '#/_base/di/extensions';
+import { LifecycleScope, registerScopedService } from '#/_base/di/scope';
+import type { AgentEvent } from '@moonshot-ai/protocol';
+
+import {
+ IAgentWireRecordService,
+ type WireRecord,
+ type WireRecordBlobSelector,
+ type WireRecordMap,
+ type WireRecordRestoreOptions,
+ type WireRecordRestoreResult,
+ type PersistedWireRecord,
+ type WireRecordRestoringContext,
+} from '#/agent/wireRecord';
+import { IAgentReplayBuilderService } from '#/agent/replayBuilder';
+import type { AgentReplayRecord } from '#/agent/replayBuilder/types';
+
+import {
+ IAgentRecordService,
+ type AgentRecord,
+ type AgentRecordMap,
+ type RecordFacets,
+} from './record';
+
+export class AgentRecordService extends Disposable implements IAgentRecordService {
+ declare readonly _serviceBrand: undefined;
+ private readonly facets = new Map>();
+ private readonly liveEmitter = this._register(new Emitter());
+
+ constructor(
+ @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
+ @IAgentReplayBuilderService private readonly replayBuilder: IAgentReplayBuilderService,
+ ) {
+ super();
+ // Restore-time replay capture: every restored record runs its `toReplay`
+ // facet. `replayBuilder.push` gates by phase, so this only records while
+ // restoring / post-restoring.
+ this._register(
+ wireRecord.hooks.onRestoredRecord.register('record-replay', async (ctx, next) => {
+ await next();
+ this.runReplayFacet(ctx.record as unknown as AgentRecord);
+ }),
+ );
+ }
+
+ append(record: AgentRecord): void {
+ this.wireRecord.append(record as unknown as WireRecord);
+ const facet = this.facets.get(record.type);
+ if (facet?.toLive !== undefined) {
+ this.emitLive(facet.toLive(record));
+ }
+ if (facet?.toReplay !== undefined) {
+ this.runReplayFacet(record);
+ }
+ }
+
+ on(handler: (event: AgentEvent) => void): IDisposable {
+ return this.liveEmitter.event(handler);
+ }
+
+ signal(event: AgentEvent): void {
+ this.emitLive(event);
+ }
+
+ define(type: K, facets: RecordFacets): IDisposable {
+ // Merge live/replay facets rather than overwriting: the record type's owner
+ // supplies `toLive`/`toReplay`, while secondary resumers (other domains that
+ // listen to the same type, e.g. microCompaction listening to
+ // `full_compaction.complete`) only contribute a `resume`. First writer wins
+ // for live/replay; every `resume` is kept (the durable store supports
+ // multiple resumers per type).
+ const previous = this.facets.get(type);
+ this.facets.set(type, {
+ toLive: facets.toLive ?? previous?.toLive,
+ toReplay: facets.toReplay ?? previous?.toReplay,
+ } as RecordFacets);
+ const resumeReg =
+ facets.resume === undefined
+ ? undefined
+ : this.wireRecord.register(
+ type as unknown as keyof WireRecordMap,
+ (record) => facets.resume!(record as unknown as AgentRecord),
+ facets.blobs === undefined
+ ? undefined
+ : {
+ blobs: facets.blobs as unknown as WireRecordBlobSelector<
+ WireRecord
+ >,
+ },
+ );
+ return toDisposable(() => {
+ resumeReg?.dispose();
+ });
+ }
+
+ restore(
+ records?: readonly PersistedWireRecord[],
+ options?: WireRecordRestoreOptions,
+ ): Promise {
+ return this.wireRecord.restore(records, options);
+ }
+
+ flush(): Promise {
+ return this.wireRecord.flush();
+ }
+
+ close(): Promise {
+ return this.wireRecord.close();
+ }
+
+ buildReplay(): readonly AgentReplayRecord[] {
+ return this.replayBuilder.buildResult();
+ }
+
+ get captureLiveRecords(): boolean {
+ return this.replayBuilder.captureLiveRecords;
+ }
+
+ set captureLiveRecords(value: boolean) {
+ this.replayBuilder.captureLiveRecords = value;
+ }
+
+ get restoring(): WireRecordRestoringContext | null {
+ return this.wireRecord.restoring;
+ }
+
+ get postRestoring(): boolean {
+ return this.replayBuilder.postRestoring;
+ }
+
+ get hooks(): IAgentWireRecordService['hooks'] {
+ return this.wireRecord.hooks;
+ }
+
+ private emitLive(event: AgentEvent): void {
+ // Suppress live emission while restoring so edge consumers never receive
+ // historical events (matches the former `eventSink.emit` guard).
+ if (this.wireRecord.restoring !== null) return;
+ this.liveEmitter.fire(event);
+ }
+
+ private runReplayFacet(record: AgentRecord): void {
+ const facet = this.facets.get(record.type);
+ if (facet?.toReplay === undefined) return;
+ const out = facet.toReplay(record);
+ if (out === undefined) return;
+ const list = Array.isArray(out) ? out : [out];
+ for (const replayRecord of list) {
+ this.replayBuilder.push(replayRecord);
+ }
+ }
+}
+
+registerScopedService(
+ LifecycleScope.Agent,
+ IAgentRecordService,
+ AgentRecordService,
+ InstantiationType.Delayed,
+ 'record',
+);
diff --git a/packages/agent-core-v2/src/agent/rpc/rpcService.ts b/packages/agent-core-v2/src/agent/rpc/rpcService.ts
index 5caae2ae2..9fb03ef32 100644
--- a/packages/agent-core-v2/src/agent/rpc/rpcService.ts
+++ b/packages/agent-core-v2/src/agent/rpc/rpcService.ts
@@ -8,7 +8,7 @@ import { IAgentContextSizeService } from '#/agent/contextSize';
import { IAgentFileToolsService } from '#/agent/fileTools';
import { IAgentFullCompactionService } from '#/agent/fullCompaction';
import { IAgentGoalService } from '#/agent/goal';
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import { ErrorCodes, KimiError } from '#/errors';
import { userCancellationReason } from '#/_base/utils/abort';
import { IAgentPermissionGate } from '#/agent/permissionGate';
@@ -96,7 +96,7 @@ export class AgentRPCService implements IAgentRPCService {
@IAgentUsageService private readonly usage: IAgentUsageService,
@ITelemetryService private readonly telemetry: ITelemetryService,
@IAgentGoalService private readonly goal: IAgentGoalService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentQuestionToolsService private readonly questionTools: IAgentQuestionToolsService,
@IAgentWebService private readonly web: IAgentWebService,
@IPluginService private readonly plugins: IPluginService,
@@ -149,12 +149,12 @@ export class AgentRPCService implements IAgentRPCService {
else if (update.kind === 'stderr') stderr += update.text ?? '';
else return;
if (payload.commandId !== undefined) {
- this.events.emit({ type: 'shell.output', commandId: payload.commandId, update });
+ this.record.signal({ type: 'shell.output', commandId: payload.commandId, update });
}
},
onForegroundTaskStart: (taskId: string) => {
if (payload.commandId !== undefined) {
- this.events.emit({ type: 'shell.started', commandId: payload.commandId, taskId });
+ this.record.signal({ type: 'shell.started', commandId: payload.commandId, taskId });
}
},
});
@@ -316,7 +316,7 @@ export class AgentRPCService implements IAgentRPCService {
commandArgs: payload.args,
trigger: 'user-slash' as const,
};
- this.events.emit({
+ this.record.signal({
type: 'plugin_command.activated',
activationId: origin.activationId,
pluginId: origin.pluginId,
diff --git a/packages/agent-core-v2/src/agent/skill/skillService.ts b/packages/agent-core-v2/src/agent/skill/skillService.ts
index 92bb83d82..fff92e064 100644
--- a/packages/agent-core-v2/src/agent/skill/skillService.ts
+++ b/packages/agent-core-v2/src/agent/skill/skillService.ts
@@ -19,12 +19,11 @@ import {
isUserActivatableSkillType,
type SkillDefinition,
} from '#/app/globalSkillCatalog/types';
-import { IAgentEventSinkService } from '#/agent/eventSink';
import { IAgentPromptService } from '#/agent/prompt';
import { ITelemetryService } from '#/app/telemetry';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
import type { Turn } from '#/agent/turn';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentRecordService } from '#/agent/record';
import {
IAgentSkillService,
type ModelSkillActivationInput,
@@ -47,15 +46,16 @@ export class AgentSkillService extends Disposable implements IAgentSkillService
constructor(
@ISessionSkillCatalog private readonly skillCatalog: ISessionSkillCatalog,
@IAgentPromptService private readonly prompt: IAgentPromptService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
+ @IAgentRecordService private readonly records: IAgentRecordService,
@ITelemetryService private readonly telemetry: ITelemetryService,
@IAgentToolRegistryService toolRegistry: IAgentToolRegistryService,
) {
super();
this._register(
- this.wireRecord.register('skill.activate', (record) => {
- this.publishActivation(record.origin);
+ records.define('skill.activate', {
+ resume: (r) => {
+ this.publishActivation(r.origin);
+ },
}),
);
this._register(toolRegistry.register(new SkillTool(this)));
@@ -162,7 +162,7 @@ export class AgentSkillService extends Disposable implements IAgentSkillService
input?: readonly ContentPart[],
delivery: 'prompt' | 'steer' = 'prompt',
): Turn | undefined {
- this.wireRecord.append({ type: 'skill.activate', origin });
+ this.records.append({ type: 'skill.activate', origin });
this.publishActivation(origin);
if (input === undefined) return undefined;
@@ -180,7 +180,7 @@ export class AgentSkillService extends Disposable implements IAgentSkillService
}
private publishActivation(origin: SkillActivationOrigin): void {
- this.events.emit({
+ this.records.signal({
type: 'skill.activated',
activationId: origin.activationId,
skillName: origin.skillName,
@@ -189,7 +189,7 @@ export class AgentSkillService extends Disposable implements IAgentSkillService
skillPath: origin.skillPath,
skillSource: origin.skillSource,
});
- if (this.wireRecord.restoring !== null) return;
+ if (this.records.restoring !== null) return;
this.telemetry.track('skill_invoked', {
skill_name: origin.skillName,
trigger: origin.trigger,
diff --git a/packages/agent-core-v2/src/agent/swarm/subagentBatch.ts b/packages/agent-core-v2/src/agent/swarm/subagentBatch.ts
index 21e28bfde..fbc439208 100644
--- a/packages/agent-core-v2/src/agent/swarm/subagentBatch.ts
+++ b/packages/agent-core-v2/src/agent/swarm/subagentBatch.ts
@@ -22,7 +22,7 @@ import {
} from '#/agent/agentTool';
import type { IAgentLifecycleService } from '#/session/agent-lifecycle';
import type { ISessionMetadata } from '#/session/session-metadata';
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import { isUserCancellation } from '#/_base/utils/abort';
/*
@@ -719,7 +719,7 @@ export function runChildAgentQueued({
retryChildAgent({ lifecycle, parentAgentId, metadata, agentId, ...options }),
suspended: (event) => {
const parent = lifecycle.getHandle(parentAgentId);
- parent?.accessor.get(IAgentEventSinkService)?.emit({
+ parent?.accessor.get(IAgentRecordService)?.signal({
type: 'subagent.suspended',
subagentId: event.agentId,
reason: event.reason,
diff --git a/packages/agent-core-v2/src/agent/swarm/swarmService.ts b/packages/agent-core-v2/src/agent/swarm/swarmService.ts
index 50698bb27..58c48c3ad 100644
--- a/packages/agent-core-v2/src/agent/swarm/swarmService.ts
+++ b/packages/agent-core-v2/src/agent/swarm/swarmService.ts
@@ -11,12 +11,11 @@
import { Disposable } from '#/_base/di';
import { InstantiationType } from '#/_base/di/extensions';
import { LifecycleScope, registerScopedService } from '#/_base/di/scope';
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import { IAgentScopeContext } from '#/agent/scopeContext';
import { IAgentSystemReminderService } from '#/agent/systemReminder';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
import { IAgentTurnService } from '#/agent/turn';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
import { IAgentLifecycleService } from '#/session/agent-lifecycle';
import SWARM_MODE_ENTER_REMINDER from './enter-reminder.md?raw';
import SWARM_MODE_EXIT_REMINDER from './exit-reminder.md?raw';
@@ -42,8 +41,7 @@ export class AgentSwarmService extends Disposable implements IAgentSwarmService
constructor(
runQueued: AgentSwarmToolHost['runQueued'] | undefined,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentSystemReminderService private readonly reminders: IAgentSystemReminderService,
@IAgentTurnService turnService: IAgentTurnService,
@IAgentToolRegistryService toolRegistry: IAgentToolRegistryService,
@@ -52,13 +50,17 @@ export class AgentSwarmService extends Disposable implements IAgentSwarmService
) {
super();
this._register(
- wireRecord.register('swarm_mode.enter', (record) => {
- this.restoreEnter(record.trigger);
+ record.define('swarm_mode.enter', {
+ resume: (r) => {
+ this.restoreEnter(r.trigger);
+ },
}),
);
this._register(
- wireRecord.register('swarm_mode.exit', () => {
- this.applyExit(false);
+ record.define('swarm_mode.exit', {
+ resume: () => {
+ this.applyExit(false);
+ },
}),
);
this._register(
@@ -79,13 +81,13 @@ export class AgentSwarmService extends Disposable implements IAgentSwarmService
enter(trigger: SwarmModeTrigger): void {
if (this._active !== null) return;
- this.wireRecord.append({ type: 'swarm_mode.enter', trigger });
+ this.record.append({ type: 'swarm_mode.enter', trigger });
this.applyEnter(trigger, true);
}
exit(): void {
if (this._active === null) return;
- this.wireRecord.append({ type: 'swarm_mode.exit' });
+ this.record.append({ type: 'swarm_mode.exit' });
this.applyExit(true);
}
@@ -124,7 +126,7 @@ export class AgentSwarmService extends Disposable implements IAgentSwarmService
}
private emitChanged(): void {
- this.events.emit({ type: 'agent.status.updated', swarmMode: this.isActive });
+ this.record.signal({ type: 'agent.status.updated', swarmMode: this.isActive });
}
}
diff --git a/packages/agent-core-v2/src/agent/toolStore/toolStoreService.ts b/packages/agent-core-v2/src/agent/toolStore/toolStoreService.ts
index a3c3a3a1f..e537d84a8 100644
--- a/packages/agent-core-v2/src/agent/toolStore/toolStoreService.ts
+++ b/packages/agent-core-v2/src/agent/toolStore/toolStoreService.ts
@@ -5,7 +5,7 @@ import { InstantiationType } from '#/_base/di/extensions';
import { LifecycleScope, registerScopedService } from '#/_base/di/scope';
import { OrderedHookSlot } from '#/hooks';
import { IAgentToolStoreService, type ToolStoreData, type ToolStoreKey } from './toolStore';
-import { IAgentWireRecordService, type WireRecord } from '#/agent/wireRecord';
+import { IAgentRecordService, type AgentRecord } from '#/agent/record';
declare module '#/agent/wireRecord' {
interface WireRecordMap {
@@ -27,13 +27,13 @@ export class AgentToolStoreService extends Disposable implements IAgentToolStore
}>(),
};
- constructor(
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
- ) {
+ constructor(@IAgentRecordService private readonly records: IAgentRecordService) {
super();
this._register(
- wireRecord.register('tools.update_store', (record) => {
- this.apply(record.key, record.value);
+ records.define('tools.update_store', {
+ resume: (r) => {
+ this.apply(r.key, r.value);
+ },
}),
);
}
@@ -43,12 +43,12 @@ export class AgentToolStoreService extends Disposable implements IAgentToolStore
}
set(key: K, value: ToolStoreData[K]): void {
- const record: WireRecord<'tools.update_store'> = {
+ const record: AgentRecord<'tools.update_store'> = {
type: 'tools.update_store',
key,
value,
};
- this.wireRecord.append(record);
+ this.records.append(record);
this.apply(key, value);
}
@@ -58,7 +58,7 @@ export class AgentToolStoreService extends Disposable implements IAgentToolStore
private apply(key: K, value: ToolStoreData[K]): void {
this.store[key] = value;
- if (!this.wireRecord.restoring) {
+ if (!this.records.restoring) {
void this.hooks.onUpdated.run({ key, value });
}
}
diff --git a/packages/agent-core-v2/src/agent/turn/turnService.ts b/packages/agent-core-v2/src/agent/turn/turnService.ts
index 526055c91..d8d7f322a 100644
--- a/packages/agent-core-v2/src/agent/turn/turnService.ts
+++ b/packages/agent-core-v2/src/agent/turn/turnService.ts
@@ -4,12 +4,11 @@ import { toKimiErrorPayload, type KimiErrorPayload } from "#/errors";
import { isUserCancellation } from "#/_base/utils/abort";
import type { ContextMessage, PromptOrigin } from '#/agent/contextMemory';
import { IAgentContextMemoryService, USER_PROMPT_ORIGIN } from '#/agent/contextMemory';
-import { IAgentEventSinkService } from '#/agent/eventSink';
import { IAgentExternalHooksService } from '#/agent/externalHooks';
import { OrderedHookSlot } from '#/hooks';
import { IAgentLoopService } from '#/agent/loop';
import { ITelemetryService } from '#/app/telemetry';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentRecordService } from '#/agent/record';
import type {
Turn,
TurnContextOverflowContext,
@@ -50,20 +49,21 @@ export class AgentTurnService implements IAgentTurnService {
constructor(
@IAgentLoopService private readonly loop: IAgentLoopService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
+ @IAgentRecordService private readonly record: IAgentRecordService,
@IAgentContextMemoryService private readonly context: IAgentContextMemoryService,
@IAgentExternalHooksService private readonly externalHooks: IAgentExternalHooksService,
@ITelemetryService private readonly telemetry: ITelemetryService,
) {
- wireRecord.register('turn.launch', (record) => {
- this.restoreLaunch(record.turnId);
+ record.define('turn.launch', {
+ resume: (r) => {
+ this.restoreLaunch(r.turnId);
+ },
});
this.hooks.beforeStep.register('turn-before-step-event', async (ctx, next) => {
await next();
this.resolveReady(ctx.turn);
});
- this.events.on((event) => {
+ this.record.on((event) => {
if (event.type === 'agent.status.updated' && event.planMode !== undefined) {
this.planModeActive = event.planMode;
return;
@@ -86,7 +86,7 @@ export class AgentTurnService implements IAgentTurnService {
this.lastEndedReasonValue = undefined;
const turnId = this.nextTurnId;
- this.wireRecord.append({ type: 'turn.launch', turnId, origin, promptMessageId });
+ this.record.append({ type: 'turn.launch', turnId, origin, promptMessageId });
this.restoreLaunch(turnId);
const abortController = new AbortController();
const ready = createControlledPromise();
@@ -120,7 +120,7 @@ export class AgentTurnService implements IAgentTurnService {
let result: TurnResult | undefined;
try {
this.telemetry.track('turn_started', { mode: telemetryMode });
- this.events.emit({
+ this.record.signal({
type: 'turn.started',
turnId: turn.id,
origin,
@@ -163,9 +163,9 @@ export class AgentTurnService implements IAgentTurnService {
) {
this.externalHooks.triggerInterrupt({ turnId: turn.id, reason: 'cancelled' });
}
- this.events.emit(ended);
+ this.record.signal(ended);
if (ended.error !== undefined) {
- this.events.emit({ type: 'error', ...ended.error });
+ this.record.signal({ type: 'error', ...ended.error });
}
if (ended.reason !== 'completed') {
this.trackTurnInterrupted(turn.id, 0);
@@ -210,7 +210,7 @@ export class AgentTurnService implements IAgentTurnService {
toolCalls: [],
origin: { kind: 'hook_result', event: hookResult.event, blocked: true },
});
- this.events.emit({
+ this.record.signal({
type: 'hook.result',
turnId: turn.id,
hookEvent: hookResult.event,
@@ -227,7 +227,7 @@ export class AgentTurnService implements IAgentTurnService {
toolCalls: [],
origin: { kind: 'hook_result', event: hookResult.event },
});
- this.events.emit({
+ this.record.signal({
type: 'hook.result',
turnId: turn.id,
hookEvent: hookResult.event,
diff --git a/packages/agent-core-v2/src/agent/usage/usageService.ts b/packages/agent-core-v2/src/agent/usage/usageService.ts
index 2a2c294eb..37ac7f96c 100644
--- a/packages/agent-core-v2/src/agent/usage/usageService.ts
+++ b/packages/agent-core-v2/src/agent/usage/usageService.ts
@@ -5,10 +5,9 @@ import { InstantiationType } from '#/_base/di/extensions';
import { LifecycleScope, registerScopedService } from '#/_base/di/scope';
import { Disposable } from '#/_base/di/lifecycle';
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import type { UsageRecordContext, UsageStatus } from './usage';
import { IAgentUsageService } from './usage';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
declare module '#/agent/wireRecord' {
interface WireRecordMap {
@@ -26,20 +25,19 @@ export class AgentUsageService extends Disposable implements IAgentUsageService
private currentTurnId: number | undefined;
private currentTurn: TokenUsage | undefined;
- constructor(
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
- @IAgentEventSinkService private readonly events: IAgentEventSinkService,
- ) {
+ constructor(@IAgentRecordService private readonly records: IAgentRecordService) {
super();
this._register(
- wireRecord.register('usage.record', (record) => {
- this.apply(record.model, record.usage, record.context);
+ records.define('usage.record', {
+ resume: (r) => {
+ this.apply(r.model, r.usage, r.context);
+ },
}),
);
}
record(model: string, usage: TokenUsage, context?: UsageRecordContext): void {
- this.wireRecord.append({
+ this.records.append({
type: 'usage.record',
model,
usage,
@@ -78,7 +76,7 @@ export class AgentUsageService extends Disposable implements IAgentUsageService
private publishChanged(): void {
const status = this.status();
if (status === undefined) return;
- this.events.emit({ type: 'agent.status.updated', usage: status });
+ this.records.signal({ type: 'agent.status.updated', usage: status });
}
private byModelSnapshot(): Record {
diff --git a/packages/agent-core-v2/src/agent/userTool/userToolService.ts b/packages/agent-core-v2/src/agent/userTool/userToolService.ts
index 39488c332..543ec029a 100644
--- a/packages/agent-core-v2/src/agent/userTool/userToolService.ts
+++ b/packages/agent-core-v2/src/agent/userTool/userToolService.ts
@@ -12,7 +12,7 @@ import type {
import { ISessionInteractionService } from '#/session/interaction';
import { IAgentProfileService } from '#/agent/profile';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
-import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentRecordService } from '#/agent/record';
import {
IAgentUserToolService,
type UserToolRegistration,
@@ -44,29 +44,33 @@ export class AgentUserToolService extends Disposable implements IAgentUserToolSe
constructor(
@IAgentToolRegistryService private readonly registry: IAgentToolRegistryService,
@IAgentProfileService private readonly profile: IAgentProfileService,
- @IAgentWireRecordService private readonly wireRecord: IAgentWireRecordService,
+ @IAgentRecordService private readonly records: IAgentRecordService,
@ISessionInteractionService private readonly interaction: ISessionInteractionService,
) {
super();
this._register(
- wireRecord.register('tools.register_user_tool', (record) => {
- this.applyRegister(record);
+ records.define('tools.register_user_tool', {
+ resume: (r) => {
+ this.applyRegister(r);
+ },
}),
);
this._register(
- wireRecord.register('tools.unregister_user_tool', (record) => {
- this.applyUnregister(record.name);
+ records.define('tools.unregister_user_tool', {
+ resume: (r) => {
+ this.applyUnregister(r.name);
+ },
}),
);
}
register(input: UserToolRegistration): void {
- this.wireRecord.append({ type: 'tools.register_user_tool', ...input });
+ this.records.append({ type: 'tools.register_user_tool', ...input });
this.applyRegister(input);
}
unregister(name: string): void {
- this.wireRecord.append({ type: 'tools.unregister_user_tool', name });
+ this.records.append({ type: 'tools.unregister_user_tool', name });
this.applyUnregister(name);
}
diff --git a/packages/agent-core-v2/src/app/gateway/gatewayService.ts b/packages/agent-core-v2/src/app/gateway/gatewayService.ts
index 412bc7b44..12773d3e0 100644
--- a/packages/agent-core-v2/src/app/gateway/gatewayService.ts
+++ b/packages/agent-core-v2/src/app/gateway/gatewayService.ts
@@ -7,13 +7,13 @@
*
* WS event fan-out (sequencing, journaling, replay, per-connection dispatch)
* is a transport concern and lives in the edge package (`packages/server-v2`)
- * on top of `IEventService` + `IAgentEventSinkService` — not here.
+ * on top of `IEventService` + `IAgentRecordService` — not here.
*/
import { InstantiationType } from '#/_base/di/extensions';
-import { type IScopeHandle, LifecycleScope, registerScopedService } from '#/_base/di/scope';
+import { type IAgentScopeHandle, LifecycleScope, registerScopedService } from '#/_base/di/scope';
import { IAgentLifecycleService } from '#/session/agent-lifecycle/agentLifecycle';
-import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentRecordService } from '#/agent/record';
import { ILogService, ISessionLogService } from '#/app/log';
import { IAgentPromptService } from '#/agent/prompt';
import { ISessionLifecycleService } from '#/app/session-lifecycle';
@@ -29,7 +29,7 @@ export class RestGateway implements IRestGateway {
@ILogService private readonly log: ILogService,
) {}
- private agent(sessionId: string, agentId: string): IScopeHandle {
+ private agent(sessionId: string, agentId: string): IAgentScopeHandle {
const session = this.sessions.get(sessionId);
if (session === undefined) throw new Error(`unknown session '${sessionId}'`);
const agents = session.accessor.get(IAgentLifecycleService);
@@ -92,7 +92,7 @@ export class WSGateway implements IWSGateway {
constructor(
@ISessionLifecycleService _sessions: ISessionLifecycleService,
- @IAgentEventSinkService _event: IAgentEventSinkService,
+ @IAgentRecordService _record: IAgentRecordService,
) {}
connect(connectionId: string): void {
diff --git a/packages/agent-core-v2/src/index.ts b/packages/agent-core-v2/src/index.ts
index b27edeb4f..3815b5790 100644
--- a/packages/agent-core-v2/src/index.ts
+++ b/packages/agent-core-v2/src/index.ts
@@ -49,7 +49,6 @@ export * from '#/session/session-context';
export * from '#/session/session-activity';
export * from '#/session/session';
-export * from '#/agent/eventSink';
import '#/session/approval';
export { ISessionApprovalService } from '#/session/approval';
export * from '#/session/question';
@@ -76,7 +75,6 @@ export * from '#/agent/systemReminder';
export * from '#/agent/contextProjector';
export * from '#/agent/contextSize';
export * from '#/agent/contextInjector';
-export * from '#/agent/eventSink';
export * from '#/agent/externalHooks';
export * from '#/agent/fullCompaction';
export * from '#/agent/llmRequestLog';
@@ -92,9 +90,11 @@ export * from '#/agent/prompt';
export * from '#/agent/promptLegacy';
export * from '#/app/messageLegacy';
export * from '#/agent/replayBuilder';
+export * from '#/agent/record';
export * from '#/agent/rpc';
export * from '#/agent/scopeContext';
export * from '#/agent/agentTool';
+export * from '#/session/btw';
export * from '#/agent/todoList';
export * from '#/agent/tool';
export * from '#/agent/toolExecutor';
diff --git a/packages/agent-core-v2/test/harness/agent.ts b/packages/agent-core-v2/test/harness/agent.ts
index de8edb800..14d473da1 100644
--- a/packages/agent-core-v2/test/harness/agent.ts
+++ b/packages/agent-core-v2/test/harness/agent.ts
@@ -147,6 +147,7 @@ import type {
WireRecordRestoreOptions,
WireRecordRestoreResult,
} from '#/agent/wireRecord';
+import { IAgentRecordService } from '#/agent/record';
import type { PathAccessOperation } from '#/session/workspaceContext';
import { createFakeAgentFs, createFakeHostEnvironment, createFakeProcessRunner } from '../tools/fixtures/fake-exec';
@@ -213,12 +214,11 @@ class TestAgentSkillService extends AgentSkillService {
constructor(
@ISessionSkillCatalog skillCatalog: ISessionSkillCatalog,
@IAgentPromptService prompt: IAgentPromptService,
- @IAgentEventSinkService events: IAgentEventSinkService,
- @IAgentWireRecordService wireRecord: IAgentWireRecordService,
+ @IAgentRecordService records: IAgentRecordService,
@ITelemetryService telemetry: ITelemetryService,
@IAgentToolRegistryService toolRegistry: IAgentToolRegistryService,
) {
- super(skillCatalog, prompt, events, wireRecord, telemetry, toolRegistry);
+ super(skillCatalog, prompt, records, telemetry, toolRegistry);
if (skillCatalog.catalog.listInvocableSkills().length > 0) {
this._register(toolRegistry.register(new ModelSkillTool(this)));
}
diff --git a/packages/agent-core-v2/test/record/recordService.test.ts b/packages/agent-core-v2/test/record/recordService.test.ts
new file mode 100644
index 000000000..c53eda864
--- /dev/null
+++ b/packages/agent-core-v2/test/record/recordService.test.ts
@@ -0,0 +1,207 @@
+import { describe, expect, it, vi } from 'vitest';
+
+import { DisposableStore, toDisposable } from '#/_base/di';
+import { createServices } from '#/_base/di/test';
+import { OrderedHookSlot } from '#/hooks';
+import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentReplayBuilderService } from '#/agent/replayBuilder';
+import { IAgentWireRecordService } from '#/agent/wireRecord';
+import type { WireRecord, WireRecordRestoredContext } from '#/agent/wireRecord';
+import {
+ AgentRecordService,
+ IAgentRecordService,
+ type AgentRecord,
+} from '#/agent/record';
+import type { AgentEvent } from '@moonshot-ai/protocol';
+import type { AgentReplayRecord, AgentReplayRecordPayload } from '#/agent/replayBuilder/types';
+
+declare module '#/agent/record' {
+ interface AgentRecordMap {
+ 'test.fact': { value: number };
+ }
+}
+
+interface StubHost {
+ readonly record: IAgentRecordService;
+ readonly wire: ReturnType;
+ readonly eventSink: ReturnType;
+ readonly replay: ReturnType;
+ readonly dispose: () => void;
+}
+
+function createWireStub() {
+ const appended: unknown[] = [];
+ const resumers = new Map void | Promise>();
+ const hooks = {
+ onRestoredRecord: new OrderedHookSlot(),
+ onResumeEnded: new OrderedHookSlot<{}>(),
+ };
+ let restoring: { time?: number } | null = null;
+ return {
+ appended,
+ resumers,
+ hooks,
+ append: vi.fn((record: unknown) => appended.push(record)),
+ register: vi.fn((type: string, resumer: (record: unknown) => void | Promise) => {
+ resumers.set(type, resumer);
+ return toDisposable(() => resumers.delete(type));
+ }),
+ restore: vi.fn(async () => ({}) as { warning?: string }),
+ flush: vi.fn(async () => {}),
+ close: vi.fn(async () => {}),
+ get restoring() {
+ return restoring;
+ },
+ setRestoring(value: { time?: number } | null) {
+ restoring = value;
+ },
+ postRestoring: false,
+ };
+}
+
+function createEventSinkStub() {
+ const emitted: AgentEvent[] = [];
+ return {
+ emitted,
+ emit: vi.fn((event: AgentEvent) => emitted.push(event)),
+ on: vi.fn(() => toDisposable(() => {})),
+ };
+}
+
+function createReplayStub() {
+ const records: AgentReplayRecord[] = [];
+ return {
+ records,
+ push: vi.fn((record: AgentReplayRecordPayload) =>
+ records.push(record as unknown as AgentReplayRecord),
+ ),
+ buildResult: vi.fn(() => records),
+ captureLiveRecords: false,
+ postRestoring: false,
+ };
+}
+
+function createHost(): StubHost {
+ const wire = createWireStub();
+ const eventSink = createEventSinkStub();
+ const replay = createReplayStub();
+ const disposables = new DisposableStore();
+ const services = createServices(disposables, {
+ additionalServices: (reg) => {
+ reg.definePartialInstance(IAgentWireRecordService, wire);
+ reg.definePartialInstance(IAgentEventSinkService, eventSink);
+ reg.definePartialInstance(IAgentReplayBuilderService, replay);
+ reg.define(IAgentRecordService, AgentRecordService);
+ },
+ });
+ return {
+ record: services.get(IAgentRecordService),
+ wire,
+ eventSink,
+ replay,
+ dispose: () => disposables.dispose(),
+ };
+}
+
+describe('AgentRecordService facade', () => {
+ it('append fans out to durable + live + replay facets', () => {
+ const host = createHost();
+ host.record.define('test.fact', {
+ toLive: (r) => ({ type: 'test.live', value: r.value }) as unknown as AgentEvent,
+ toReplay: (r) =>
+ ({ type: 'message', value: r.value }) as unknown as AgentReplayRecordPayload,
+ });
+
+ host.record.append({ type: 'test.fact', value: 42 });
+
+ expect(host.wire.append).toHaveBeenCalledWith(
+ expect.objectContaining({ type: 'test.fact', value: 42 }),
+ );
+ expect(host.eventSink.emitted).toContainEqual(
+ expect.objectContaining({ type: 'test.live', value: 42 }),
+ );
+ expect(host.replay.records).toContainEqual(
+ expect.objectContaining({ type: 'message', value: 42 }),
+ );
+ host.dispose();
+ });
+
+ it('append omits facets that are not declared', () => {
+ const host = createHost();
+ host.record.define('test.fact', {});
+
+ host.record.append({ type: 'test.fact', value: 1 });
+
+ expect(host.wire.append).toHaveBeenCalledTimes(1);
+ expect(host.eventSink.emit).not.toHaveBeenCalled();
+ expect(host.replay.push).not.toHaveBeenCalled();
+ host.dispose();
+ });
+
+ it('signal emits live only and never persists or captures replay', () => {
+ const host = createHost();
+ host.record.signal({ type: 'test.delta', delta: 'x' } as unknown as AgentEvent);
+
+ expect(host.eventSink.emitted).toContainEqual(
+ expect.objectContaining({ type: 'test.delta', delta: 'x' }),
+ );
+ expect(host.wire.append).not.toHaveBeenCalled();
+ expect(host.replay.push).not.toHaveBeenCalled();
+ host.dispose();
+ });
+
+ it('define registers the resumer with wireRecord and forwards the record', async () => {
+ const host = createHost();
+ const resume = vi.fn();
+ host.record.define('test.fact', { resume });
+
+ expect(host.wire.register).toHaveBeenCalledWith('test.fact', expect.any(Function));
+ const registered = host.wire.resumers.get('test.fact');
+ expect(registered).toBeDefined();
+ await registered?.({ type: 'test.fact', value: 7 });
+
+ expect(resume).toHaveBeenCalledWith(expect.objectContaining({ value: 7 }));
+ host.dispose();
+ });
+
+ it('captures replay during restore through the onRestoredRecord hook', async () => {
+ const host = createHost();
+ host.record.define('test.fact', {
+ toReplay: (r) =>
+ ({ type: 'message', value: r.value }) as unknown as AgentReplayRecordPayload,
+ });
+ host.wire.setRestoring({ time: 123 });
+
+ await host.wire.hooks.onRestoredRecord.run({
+ record: { type: 'test.fact', value: 5 } as unknown as WireRecord,
+ stop: false,
+ });
+
+ expect(host.replay.records).toContainEqual(
+ expect.objectContaining({ type: 'message', value: 5 }),
+ );
+ host.dispose();
+ });
+
+ it('on() delegates to the live event sink', () => {
+ const host = createHost();
+ const handler = vi.fn();
+ host.record.on(handler);
+ expect(host.eventSink.on).toHaveBeenCalledWith(handler);
+ host.dispose();
+ });
+
+ it('dispose returned by define unregisters the resumer and facets', () => {
+ const host = createHost();
+ const subscription = host.record.define('test.fact', { resume: vi.fn() });
+ expect(host.wire.resumers.has('test.fact')).toBe(true);
+
+ subscription.dispose();
+
+ expect(host.wire.resumers.has('test.fact')).toBe(false);
+ // Facet removed: append should no longer fan out (still persists).
+ host.record.append({ type: 'test.fact', value: 1 });
+ expect(host.wire.append).toHaveBeenCalledTimes(1);
+ host.dispose();
+ });
+});
diff --git a/packages/agent-core-v2/test/skill/skill.test.ts b/packages/agent-core-v2/test/skill/skill.test.ts
index e37cd548b..902ed67a3 100644
--- a/packages/agent-core-v2/test/skill/skill.test.ts
+++ b/packages/agent-core-v2/test/skill/skill.test.ts
@@ -20,6 +20,8 @@ import { ITelemetryService } from '#/app/telemetry';
import { IAgentToolRegistryService } from '#/agent/toolRegistry';
import type { Turn } from '#/agent/turn';
import { IAgentWireRecordService } from '#/agent/wireRecord';
+import { IAgentReplayBuilderService } from '#/agent/replayBuilder';
+import { AgentRecordService, IAgentRecordService } from '#/agent/record';
import { stubWireRecord } from '../contextMemory/stubs';
import { executeTool } from '../tools/fixtures/execute-tool';
import { stubSkill } from './stubs';
@@ -71,6 +73,13 @@ describe('AgentSkillService', () => {
on: () => ({ dispose: () => {} }),
});
reg.defineInstance(IAgentWireRecordService, stubWireRecord());
+ reg.definePartialInstance(IAgentReplayBuilderService, {
+ push: () => {},
+ buildResult: () => [],
+ captureLiveRecords: false,
+ postRestoring: false,
+ });
+ reg.define(IAgentRecordService, AgentRecordService);
reg.definePartialInstance(ITelemetryService, { track: () => {} });
reg.definePartialInstance(IAgentToolRegistryService, {
register: () => ({ dispose: () => {} }),
@@ -171,6 +180,13 @@ describe('SkillTool', () => {
on: () => ({ dispose: () => {} }),
});
reg.defineInstance(IAgentWireRecordService, stubWireRecord());
+ reg.definePartialInstance(IAgentReplayBuilderService, {
+ push: () => {},
+ buildResult: () => [],
+ captureLiveRecords: false,
+ postRestoring: false,
+ });
+ reg.define(IAgentRecordService, AgentRecordService);
reg.definePartialInstance(ITelemetryService, { track: () => {} });
reg.definePartialInstance(IAgentToolRegistryService, {
register: () => ({ dispose: () => {} }),
diff --git a/packages/agent-core-v2/test/usage/usage.test.ts b/packages/agent-core-v2/test/usage/usage.test.ts
index 4653c45f3..bd61f5078 100644
--- a/packages/agent-core-v2/test/usage/usage.test.ts
+++ b/packages/agent-core-v2/test/usage/usage.test.ts
@@ -3,10 +3,14 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest';
import { toDisposable } from '#/_base/di';
import { DisposableStore } from '#/_base/di/lifecycle';
import { createServices, type TestInstantiationService } from '#/_base/di/test';
+import { OrderedHookSlot } from '#/hooks';
import { IAgentEventSinkService } from '#/agent/eventSink';
+import { IAgentReplayBuilderService } from '#/agent/replayBuilder';
+import { AgentRecordService, IAgentRecordService } from '#/agent/record';
import { IAgentUsageService, type UsageStatus } from '#/agent/usage';
import { AgentUsageService } from '#/agent/usage/usageService';
import { IAgentWireRecordService, type WireRecord } from '#/agent/wireRecord';
+import type { WireRecordRestoredContext } from '#/agent/wireRecord';
let disposables: DisposableStore;
@@ -221,6 +225,10 @@ function createUsageHarness(): {
reg.definePartialInstance(IAgentWireRecordService, {
restoring: null,
postRestoring: false,
+ hooks: {
+ onRestoredRecord: new OrderedHookSlot(),
+ onResumeEnded: new OrderedHookSlot<{}>(),
+ },
append: (record) => {
records.push(record);
},
@@ -235,6 +243,13 @@ function createUsageHarness(): {
},
on: () => toDisposable(() => {}),
});
+ reg.definePartialInstance(IAgentReplayBuilderService, {
+ push: () => {},
+ buildResult: () => [],
+ captureLiveRecords: false,
+ postRestoring: false,
+ });
+ reg.define(IAgentRecordService, AgentRecordService);
reg.define(IAgentUsageService, AgentUsageService);
},
});
diff --git a/packages/server-v2/src/transport/ws/eventMap.ts b/packages/server-v2/src/transport/ws/eventMap.ts
index 0f72be2a8..4920eee00 100644
--- a/packages/server-v2/src/transport/ws/eventMap.ts
+++ b/packages/server-v2/src/transport/ws/eventMap.ts
@@ -7,12 +7,12 @@
* Core `events` — process-wide `DomainEvent` bus (`IEventService`)
* Session `interactions` — pending human-in-the-loop requests (`ISessionInteractionService.onDidChange`)
* Session `interactions:resolved` — request resolutions (`ISessionInteractionService.onDidResolve`)
- * Agent `events` — per-agent `AgentEvent` stream (`IAgentEventSinkService`)
+ * Agent `events` — per-agent `AgentEvent` stream (`IAgentRecordService`)
*/
import {
IEventService,
- IAgentEventSinkService,
+ IAgentRecordService,
ISessionInteractionService,
type DomainEvent,
type IDisposable,
@@ -58,7 +58,7 @@ export const eventMap: Record> = {
},
agent: {
events: {
- subscribe: (scope, listener) => scope.accessor.get(IAgentEventSinkService).on(listener),
+ subscribe: (scope, listener) => scope.accessor.get(IAgentRecordService).on(listener),
},
},
};
diff --git a/packages/server-v2/src/transport/ws/v1/sessionEventBroadcaster.ts b/packages/server-v2/src/transport/ws/v1/sessionEventBroadcaster.ts
index b1d420168..2da3760f6 100644
--- a/packages/server-v2/src/transport/ws/v1/sessionEventBroadcaster.ts
+++ b/packages/server-v2/src/transport/ws/v1/sessionEventBroadcaster.ts
@@ -1,13 +1,13 @@
/**
* `SessionEventBroadcaster` — per-session single fan-out point that turns
- * agent `IAgentEventSinkService` emissions into a sequenced, journaled, replayable
+ * agent `IAgentRecordService` emissions into a sequenced, journaled, replayable
* `/api/v1/ws` event stream (the `{seq, epoch}` watermark).
*
* Port of v1's `WSBroadcastService` (`packages/server/.../wsBroadcastService.ts`),
- * adapted to v2 where agent events live on per-agent `IAgentEventSinkService`s (not a Core
+ * adapted to v2 where agent events live on per-agent `IAgentRecordService`s (not a Core
* firehose). For each session it:
*
- * 1. Subscribes to every agent's `IAgentEventSinkService` via `IAgentLifecycleService`
+ * 1. Subscribes to every agent's `IAgentRecordService` via `IAgentLifecycleService`
* reach-down-via-handle (and `onDidCreate`/`onDidDispose` for late agents).
* 2. Attaches `agentId`/`sessionId` to build the wire `Event`.
* 3. Classifies durable vs volatile (`VOLATILE_EVENT_TYPES`).
@@ -24,10 +24,16 @@
* the journal is continuous from first activation onward.
*/
-import type { DomainEvent, IDisposable, IScopeHandle, Scope } from '@moonshot-ai/agent-core-v2';
+import type {
+ DomainEvent,
+ IAgentScopeHandle,
+ IDisposable,
+ ISessionScopeHandle,
+ Scope,
+} from '@moonshot-ai/agent-core-v2';
import {
IAgentLifecycleService,
- IAgentEventSinkService,
+ IAgentRecordService,
IEventService,
ISessionLifecycleService,
} from '@moonshot-ai/agent-core-v2';
@@ -247,11 +253,11 @@ export class SessionEventBroadcaster {
state.queue = state.queue.then(() => this.dispatch(state, event)).catch(() => {});
}
- private attachAgents(sessionId: string, session: IScopeHandle, state: SessionState): void {
+ private attachAgents(sessionId: string, session: ISessionScopeHandle, state: SessionState): void {
const agents = session.accessor.get(IAgentLifecycleService);
- const subscribeAgent = (handle: IScopeHandle): void => {
+ const subscribeAgent = (handle: IAgentScopeHandle): void => {
if (state.agentDisposables.has(handle.id)) return;
- const sink = handle.accessor.get(IAgentEventSinkService);
+ const sink = handle.accessor.get(IAgentRecordService);
const d = sink.on((agentEvent) => this.onAgentEvent(sessionId, handle.id, agentEvent));
state.agentDisposables.set(handle.id, d);
};