返回源码地图

packages/experimental/agent-team/src/journal.ts

main snapshot · da00f7f5358f · 正文引用章节 13;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/** Serialized Team transactions over the exact live Lead Session log. */
2
3import type { Agent } from '@deepseek-ai/dsh-agent'
4import type { Context } from '@deepseek-ai/cordis'
5import type { SessionEventMap, SessionId } from '@deepseek-ai/dsh-session'
6import type { TeamEventType, TeamState } from './projection.ts'
7
8type AppendTeamEvent = <T extends TeamEventType>(type: T, data: SessionEventMap[T]) => void
9type MutableTeamEventType = 'team/member' | 'team/task' | 'team/message/queued' | 'team/message/delivered'
10
11/** Owns per-Lead transaction order and committed Team event publication. */
12export class TeamJournal {
13 private readonly tails = new Map<SessionId, Promise<void>>()
14
15 /**
16 * @param ctx - Team service context with the injected Session service.
17 * @param onCommit - synchronous notification after the Team event flush succeeds.
18 */
19 constructor(
20 private readonly ctx: Context,
21 private readonly onCommit: (root: Agent) => void,
22 ) {}
23
24 /**
25 * Read authoritative Team state for one exact live Lead.
26 * @param root - exact live Team Lead.
27 * @returns current projected state selected by the Lead Team id.
28 */
29 state(root: Agent): TeamState {
30 const projection = this.ctx.sessionProjections.stateOf(root.session, 'agentTeam')
31 if (projection === undefined) throw new Error('Agent Teams projection is not registered')
32 if (projection.failure !== undefined) throw new Error(projection.failure)
33 return projection
34 }
35
36 /**
37 * Serialize one Lead's asynchronous mutation operation.
38 * @param rootId - Lead Session identity selecting the transaction queue.
39 * @param operation - complete read-check-append operation.
40 * @returns the operation result.
41 */
42 async transact<T>(rootId: SessionId, operation: () => Promise<T>): Promise<T> {
43 const prior = this.tails.get(rootId) ?? Promise.resolve()
44 const run = prior.then(operation, operation)
45 const tail = run.then(() => undefined, () => undefined)
46 this.tails.set(rootId, tail)
47 try {
48 return await run
49 } finally {
50 if (this.tails.get(rootId) === tail) this.tails.delete(rootId)
51 }
52 }
53
54 /**
55 * Append and checkpoint one root-owned Team event before publication.
56 * @param root - exact live Lead whose Session owns the event.
57 * @param type - Team event discriminant.
58 * @param data - payload correlated with the event type.
59 */
60 async appendAndFlush<T extends MutableTeamEventType>(
61 root: Agent,
62 type: T,
63 data: SessionEventMap[T],
64 ): Promise<void> {
65 // Team events never enter the conversation surface. This narrower local
66 // capability removes Session.append's conditional surface argument while
67 // preserving the event-key/payload correlation.
68 const append = root.session.append.bind(root.session) as unknown as AppendTeamEvent
69 append(type, data)
70 await this.ctx.sessions.flush(root.session)
71 this.onCommit(root)
72 }
73}