1
/** Serialized Team transactions over the exact live Lead Session log. */3
import type { Agent } from '@deepseek-ai/dsh-agent'4
import type { Context } from '@deepseek-ai/cordis'5
import type { SessionEventMap, SessionId } from '@deepseek-ai/dsh-session'6
import type { TeamEventType, TeamState } from './projection.ts'8
type AppendTeamEvent = <T extends TeamEventType>(type: T, data: SessionEventMap[T]) => void9
type MutableTeamEventType = 'team/member' | 'team/task' | 'team/message/queued' | 'team/message/delivered'11
/** Owns per-Lead transaction order and committed Team event publication. */12
export class TeamJournal {13
private readonly tails = new Map<SessionId, Promise<void>>()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
) {}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 projection34
}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 run49
} finally {50
if (this.tails.get(rootId) === tail) this.tails.delete(rootId)51
}52
}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 local66
// capability removes Session.append's conditional surface argument while67
// preserving the event-key/payload correlation.68
const append = root.session.append.bind(root.session) as unknown as AppendTeamEvent69
append(type, data)70
await this.ctx.sessions.flush(root.session)71
this.onCommit(root)72
}73
}