1
/**2
* Same-session goal domain: event-sourced state, compare-and-set mutations,3
* and process-local continuation activation.4
* @module @deepseek-ai/dsh-goal5
*/7
import { randomUUID } from 'node:crypto'8
import { Context } from '@deepseek-ai/cordis'9
import z from '@deepseek-ai/schemastery'10
import { z as zod } from 'zod'11
import type { ZodType } from 'zod'12
import { agentEvents } from '@deepseek-ai/dsh-agent'13
import type { Agent } from '@deepseek-ai/dsh-agent'14
import { SessionSeq } from '@deepseek-ai/dsh-session'15
import type { Session, SessionEvent, SessionLogOffset } from '@deepseek-ai/dsh-session'16
import { TypertRemoteService, Remote } from '@deepseek-ai/dsh-typert-protocol'17
import type {} from '@deepseek-ai/dsh-session-projection'18
import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'19
import {20
applyGoalEvent,21
goalChangeRef,22
} from './fold.ts'23
import type { GoalFoldState } from './fold.ts'24
import {25
GOAL_CHANGE_VERSION,26
GoalError,27
GoalId,28
} from './runtime.ts'29
import type {30
CreateGoalRequest,31
CreateGoalResult,32
EditGoalRequest,33
GoalActivation,34
GoalBlockReason,35
GoalPhase,36
GoalProjection,37
GoalProjectionState,38
GoalRef,39
GoalSnapshot,40
GoalView,41
} from './types.ts'42
import type {43
GoalChangeMeta,44
GoalChanged,45
GoalClearChangeMeta,46
GoalOperation,47
GoalSnapshotChangeMeta,48
} from './domain.ts'50
// The pure payload outlet (./types.ts, ONE home of the `goal` projection-key51
// declaration) re-exported onto the package root keeps the module edge in52
// the emitted index.d.ts, so aggregate programs consuming the declarations53
// still receive the SessionProjectionStateMap merge.54
export type * from './types.ts'55
export type * from './domain.ts'56
export { GOAL_CHANGE_VERSION, GoalError, GoalId } from './runtime.ts'57
export { decodeGoalChange, foldGoal, goalChangeRef } from './fold.ts'59
declare module '@deepseek-ai/cordis' {60
interface Context {61
goals: GoalService62
}63
}65
/** Wire payload schema of the `goal` projection (current goal or pre-create/cleared null). */66
const goalProjectionSchema: ZodType<GoalProjection | null> = zod.union([67
zod.object({68
goal: zod.object({69
id: zod.string().min(1),70
revision: zod.number().int().positive(),71
objective: zod.string().min(1),72
phase: zod.union([zod.literal('active'), zod.literal('paused'), zod.literal('blocked'), zod.literal('complete')]),73
blockedReason: zod.object({ code: zod.string(), message: zod.string() }).optional(),74
maxGoalRounds: zod.number().int().positive(),75
}),76
roundsStarted: zod.number().int().nonnegative(),77
createdAt: zod.number(),78
updatedAt: zod.number(),79
}),80
zod.null(),81
]) as ZodType<GoalProjection | null>83
const goalProjectionStateSchema: ZodType<GoalProjectionState> = zod.object({84
current: goalProjectionSchema,85
seenGoalIds: zod.array(zod.string().min(1)).refine(86
ids => new Set(ids).size === ids.length,87
{ message: 'seen goal ids must be unique' },88
),89
failure: zod.string().min(1).nullable(),90
}).strict().superRefine((state, context) => {91
if (state.current === null) return92
if (!state.seenGoalIds.includes(state.current.goal.id)) {93
context.addIssue({ code: 'custom', message: 'current goal id must be retained among seen goal ids' })94
}95
if (state.current.updatedAt < state.current.createdAt) {96
context.addIssue({ code: 'custom', message: 'current goal update cannot precede its creation' })97
}98
if (state.current.roundsStarted > state.current.goal.maxGoalRounds) {99
context.addIssue({ code: 'custom', message: 'current goal rounds cannot exceed its configured limit' })100
}101
}) as unknown as ZodType<GoalProjectionState>103
/** Build strict fold state from one checkpoint-safe projection state. */104
function goalFoldState(state: GoalProjectionState): GoalFoldState {105
return {106
goal: state.current?.goal,107
roundsStarted: state.current?.roundsStarted ?? 0,108
createdAt: state.current?.createdAt,109
updatedAt: state.current?.updatedAt,110
lastRef: undefined,111
seenGoalIds: new Set(state.seenGoalIds),112
}113
}115
/** Convert strict fold state into checkpoint-safe projection state. */116
function goalProjectionState(state: GoalFoldState): GoalProjectionState {117
let current: GoalProjection | null = null118
if (state.goal !== undefined) {119
const { createdAt, updatedAt } = state120
if (createdAt === undefined || updatedAt === undefined) {121
throw new Error('current goal fold lacks timestamps')122
}123
current = {124
goal: state.goal,125
roundsStarted: state.roundsStarted,126
createdAt,127
updatedAt,128
}129
}130
return {131
current,132
seenGoalIds: [...state.seenGoalIds],133
failure: null,134
}135
}137
/**138
* Fold durable goal events through the strict replay rules without throwing139
* from the projection registry's event drive. The first invalid owned event140
* is retained in `failure`; host goal access rejects that state while the141
* client view remains at the last valid goal.142
* @param state - the projection covering all prior events.143
* @param event - the next committed session event.144
* @returns the next projection (same reference when the event is unrelated).145
*/146
export function applyGoalProjection(state: GoalProjectionState, event: SessionEvent): GoalProjectionState {147
if (state.failure !== null) return state148
if (event.type !== 'goal/change'149
&& (event.type !== 'user/message' || event.data.source.kind !== 'goal')) return state150
const folded = goalFoldState(state)151
try {152
applyGoalEvent(folded, event)153
return goalProjectionState(folded)154
} catch (error: unknown) {155
/* v8 ignore next -- the strict goal fold throws Error instances. */156
const message = error instanceof Error ? error.message : String(error)157
return { ...state, failure: `goal replay failed at session event ${event.seq}: ${message}` }158
}159
}161
/** Strict host goal state with the existing cropped client value. */162
export const goalProjectionDefinition = {163
key: 'goal',164
stateSchema: goalProjectionStateSchema,165
init: (): GoalProjectionState => ({ current: null, seenGoalIds: [], failure: null }),166
apply: applyGoalProjection,167
wire: { viewSchema: goalProjectionSchema, view: state => state.current },168
stateVersion: 6,169
} satisfies ProjectionDefinition<'goal', GoalProjectionState>171
/** Deployment defaults for goal creation. */172
export interface Config {173
/** Total rounds used when a create request omits its own cap. */174
defaultMaxGoalRounds?: number175
}177
/** Resolved defaults. */178
export interface ResolvedConfig {179
/** Validated positive safe-integer default round cap. */180
defaultMaxGoalRounds: number181
}183
/** Process-local activation state crossing the synchronous append boundary. */184
interface GoalRuntimeState {185
activation: GoalActivation186
pendingActivation: {187
readonly offset: SessionLogOffset188
readonly activation: GoalActivation189
} | undefined190
}192
/** Validated create input with every deployment default materialized. */193
interface ResolvedCreateGoal {194
readonly objective: string195
readonly maxGoalRounds: number196
}198
/** Validate a caller-visible positive safe-integer round cap. */199
function resolveMaxGoalRounds(value: number): number {200
if (!Number.isSafeInteger(value) || value < 1) {201
throw new GoalError('maxGoalRounds must be a positive safe integer', 'GOAL_INVALID_MAX_ROUNDS')202
}203
return value204
}206
/** Validate and normalize an objective at the domain boundary. */207
function resolveObjective(value: string): string {208
if (typeof value !== 'string' || value.trim().length === 0) {209
throw new GoalError('goal objective must be a non-empty string', 'GOAL_INVALID_OBJECTIVE')210
}211
return value.trim()212
}214
/** Materialize deployment defaults and validate one create request. */215
function resolveCreateGoal(request: CreateGoalRequest, defaultMaxGoalRounds: number): ResolvedCreateGoal {216
return {217
objective: resolveObjective(request.objective),218
maxGoalRounds: resolveMaxGoalRounds(request.maxGoalRounds ?? defaultMaxGoalRounds),219
}220
}222
/** Validate and detach one policy-owned blocker explanation. */223
function resolveBlockReason(reason: unknown): GoalBlockReason {224
const record = typeof reason === 'object' && reason !== null && !Array.isArray(reason)225
? reason as Record<string, unknown>226
: undefined227
const code = record?.['code']228
const message = record?.['message']229
if (typeof code !== 'string' || !/^[a-z][a-z0-9]*(?:-[a-z0-9]+)*$/.test(code)230
|| typeof message !== 'string' || message.trim().length === 0) {231
throw new GoalError(232
'goal block reason requires a lower-kebab-case code and a non-empty message',233
'GOAL_INVALID_BLOCK_REASON',234
)235
}236
return { code, message: message.trim() }237
}239
/** Goal service (`ctx.goals`) backed exclusively by the owning session log. */240
export class GoalService extends TypertRemoteService {241
static inject = ['agents', 'sessionProjections']243
static Config: z<Config> = z.object({244
defaultMaxGoalRounds: z.number().default(256),245
})247
private readonly resolved: ResolvedConfig248
private readonly runtimeStates = new WeakMap<Session, GoalRuntimeState>()250
constructor(ctx: Context, config: Config = {}) {251
super(ctx, 'goals')252
this.resolved = {253
defaultMaxGoalRounds: resolveMaxGoalRounds(config.defaultMaxGoalRounds ?? 256),254
}255
ctx.on('agent/created', ({ agent }) => {256
this.setActivation(agent.session, 'disarmed')257
})258
ctx.sessionProjections.register(goalProjectionDefinition)259
ctx.on('session/event', (session, event) => {260
if (event.type !== 'goal/change') return261
const runtime = this.runtimeState(session)262
const activation = runtime.pendingActivation !== undefined263
&& SessionSeq(runtime.pendingActivation.offset) === event.seq264
? runtime.pendingActivation.activation265
: 'disarmed'266
this.setActivation(session, activation)267
})268
}270
/**271
* Read the current goal for one exact live agent.272
* @param agent - owning live agent.273
* @returns a fresh view or `undefined` when no goal is current.274
* @throws {@link GoalError} when the agent is not the registry's live instance.275
*/276
@Remote('get')277
get(agent: Agent): GoalView | undefined {278
this.assertLive(agent)279
return this.view(this.state(agent.session), this.runtimeState(agent.session))280
}282
/**283
* Remove process-local continuation authority without changing durable goal284
* phase or revision. Lifecycle owners use this before unloading a driver;285
* a later human-authorized {@link resume} records the new activation edge.286
* @param agent - owning live agent.287
* @returns a fresh disarmed view, or `undefined` when no goal is current.288
*/289
disarm(agent: Agent): GoalView | undefined {290
this.assertLive(agent)291
this.setActivation(agent.session, 'disarmed')292
const runtime = this.runtimeState(agent.session)293
return this.view(this.state(agent.session), runtime)294
}296
/**297
* Create and arm a goal. A completed goal may be replaced; every other298
* current phase must be cleared or resumed instead.299
* @param agent - owning live agent.300
* @param request - objective and optional round cap.301
* @returns the created live view.302
*/303
create(agent: Agent, request: CreateGoalRequest): GoalView {304
const spec = resolveCreateGoal(request, this.resolved.defaultMaxGoalRounds)305
const [state, runtime] = this.prepareMutation(agent)306
const current = state?.goal307
if (current !== undefined && current.phase !== 'complete') {308
throw new GoalError(`goal "${current.id}" already exists with phase "${current.phase}"`, 'GOAL_ALREADY_EXISTS')309
}310
const now = Date.now()311
const goal: GoalSnapshot = {312
id: GoalId(`goal-${randomUUID()}`),313
revision: 1,314
objective: spec.objective,315
phase: 'active',316
maxGoalRounds: spec.maxGoalRounds,317
}318
return this.commitSnapshot(agent, runtime, 'create', goal, 0, now, now, 'armed')319
}321
/**322
* Edit objective and/or round cap without changing phase.323
* @param agent - owning live agent.324
* @param ref - expected current revision.325
* @param request - at least one replacement field.326
* @returns the edited view.327
*/328
@Remote('edit')329
edit(agent: Agent, ref: GoalRef, request: EditGoalRequest): GoalView {330
const [state, runtime] = this.prepareMutation(agent)331
const currentState = this.expectCurrent(state, ref)332
const current = currentState.goal333
if (request.objective === undefined && request.maxGoalRounds === undefined) {334
throw new GoalError('goal edit requires objective and/or maxGoalRounds', 'GOAL_INVALID_EDIT')335
}336
const goal: GoalSnapshot = {337
...current,338
revision: current.revision + 1,339
...request.objective === undefined ? {} : { objective: resolveObjective(request.objective) },340
...request.maxGoalRounds === undefined ? {} : { maxGoalRounds: resolveMaxGoalRounds(request.maxGoalRounds) },341
}342
return this.commitCurrent(agent, currentState, runtime, 'edit', goal, runtime.activation)343
}345
/**346
* Pause an active goal and disarm automatic continuation.347
* @param agent - owning live agent.348
* @param ref - expected current revision.349
* @returns the paused view.350
*/351
@Remote('pause')352
pause(agent: Agent, ref: GoalRef): GoalView {353
return this.transition(agent, ref, 'pause', ['active'], 'paused', 'disarmed')354
}356
/**357
* Resume and arm a stopped goal, or rearm an active goal after a358
* session-start edge, while its round budget still has capacity.359
* @param agent - owning live agent.360
* @param ref - expected current revision.361
* @returns the active view.362
*/363
@Remote('resume')364
resume(agent: Agent, ref: GoalRef): GoalView {365
const [state, runtime] = this.prepareMutation(agent)366
const currentState = this.expectCurrent(state, ref)367
const current = currentState.goal368
const resumable: readonly GoalPhase[] = ['active', 'paused', 'blocked']369
if (!resumable.includes(current.phase)) {370
throw this.transitionError(current, 'resume', resumable)371
}372
if (current.phase === 'active' && runtime.activation === 'armed') {373
throw new GoalError(`goal "${current.id}" is already active and armed`, 'GOAL_INVALID_TRANSITION')374
}375
if (currentState.roundsStarted >= current.maxGoalRounds) {376
throw new GoalError(377
`goal "${current.id}" exhausted ${current.maxGoalRounds} goal rounds; increase maxGoalRounds before resuming`,378
'GOAL_INVALID_TRANSITION',379
)380
}381
return this.commitCurrent(agent, currentState, runtime, 'resume', this.withPhase(current, 'active'), 'armed')382
}384
/**385
* Mark a current non-complete goal complete and disarm it.386
* @param agent - owning live agent.387
* @param ref - expected current revision.388
* @returns the completed view.389
*/390
@Remote('complete')391
complete(agent: Agent, ref: GoalRef): GoalView {392
return this.transition(393
agent,394
ref,395
'complete',396
['active', 'paused', 'blocked'],397
'complete',398
'disarmed',399
)400
}402
/**403
* Mark an active goal blocked and disarm it.404
* @param agent - owning live agent.405
* @param ref - expected current revision.406
* @param reason - policy-owned stable code and human-readable explanation.407
* @returns the blocked view with its durable reason.408
*/409
block(agent: Agent, ref: GoalRef, reason: GoalBlockReason): GoalView {410
const [state, runtime] = this.prepareMutation(agent)411
const currentState = this.expectCurrent(state, ref)412
const current = currentState.goal413
if (current.phase !== 'active') {414
throw this.transitionError(current, 'block', ['active'])415
}416
return this.commitCurrent(417
agent,418
currentState,419
runtime,420
'block',421
{ ...this.withPhase(current, 'blocked'), blockedReason: resolveBlockReason(reason) },422
'disarmed',423
)424
}426
/**427
* Clear the current goal while retaining a durable tombstone and history.428
* @param agent - owning live agent.429
* @param ref - expected current revision.430
* @returns the tombstone ref whose revision is one past the cleared snapshot.431
*/432
@Remote('clear')433
clear(agent: Agent, ref: GoalRef): GoalRef {434
const [state, runtime] = this.prepareMutation(agent)435
const currentState = this.expectCurrent(state, ref)436
const current = currentState.goal437
const tombstone: GoalRef = { id: current.id, revision: current.revision + 1 }438
const change: GoalClearChangeMeta = {439
kind: 'goal/change',440
version: GOAL_CHANGE_VERSION,441
operation: 'clear',442
cleared: tombstone,443
clearedAt: this.nextMutationTime(currentState),444
}445
this.commit(agent, runtime, change, 'disarmed')446
return { ...tombstone }447
}449
/** Resolve the durable and process-local state used by a mutation. */450
private prepareMutation(agent: Agent): readonly [GoalProjection | null, GoalRuntimeState] {451
this.assertLive(agent)452
return [this.state(agent.session), this.runtimeState(agent.session)]453
}455
/** Reject stale or missing current-state refs. */456
private expectCurrent(state: GoalProjection | null, ref: GoalRef): GoalProjection {457
if (state === null) throw new GoalError('no current goal', 'GOAL_NOT_FOUND')458
const current = state.goal459
if (ref.id !== current.id || ref.revision !== current.revision) {460
throw new GoalError(461
`stale goal ref "${ref.id}" revision ${ref.revision}; current is "${current.id}" revision ${current.revision}`,462
'GOAL_STALE_REVISION',463
)464
}465
return state466
}468
/** Enforce exact live-agent identity rather than trusting a matching id. */469
private assertLive(agent: Agent): void {470
if (this.ctx.agents.get(agent.id) !== agent) {471
throw new GoalError(`agent "${agent.id}" is not live in this registry`, 'GOAL_AGENT_NOT_LIVE')472
}473
}475
/** Read the current durable projection maintained by the registry. */476
private state(session: Session): GoalProjection | null {477
const state = this.ctx.sessionProjections.stateOf(session, 'goal')478
if (state === undefined) throw new Error('goal projection is not registered')479
if (state.failure !== null) throw new Error(state.failure)480
return state.current481
}483
/** Return the process-local activation state, initially disarmed. */484
private runtimeState(session: Session): GoalRuntimeState {485
let runtime = this.runtimeStates.get(session)486
if (runtime !== undefined) return runtime487
runtime = {488
activation: 'disarmed',489
pendingActivation: undefined,490
}491
this.runtimeStates.set(session, runtime)492
return runtime493
}495
/** Publish one process-local activation edge when it actually changes. */496
private setActivation(session: Session, activation: GoalActivation): void {497
const runtime = this.runtimeState(session)498
if (runtime.activation === activation) return499
runtime.activation = activation500
const state = this.ctx.sessionProjections.stateOf(session, 'goal')501
/* v8 ignore next -- static inject requires the projection registry before this service activates. */502
if (state === undefined) return503
if (state.failure !== null) return504
const goal = this.view(state.current, runtime)505
this.ctx.emit('goal/activation-changed', {506
sessionId: session.id,507
...goal === undefined ? {} : {508
goal: {509
id: goal.id,510
revision: goal.revision,511
activation: goal.activation,512
},513
},514
})515
}517
/** Build a new revision with one replacement phase. */518
private withPhase(current: GoalSnapshot, phase: GoalPhase): GoalSnapshot {519
return {520
id: current.id,521
revision: current.revision + 1,522
objective: current.objective,523
phase,524
maxGoalRounds: current.maxGoalRounds,525
}526
}528
/** Shared validated phase transition. */529
private transition(530
agent: Agent,531
ref: GoalRef,532
operation: Exclude<GoalOperation, 'create' | 'edit' | 'clear'>,533
allowed: readonly GoalPhase[],534
phase: GoalPhase,535
activation: GoalActivation,536
): GoalView {537
const [state, runtime] = this.prepareMutation(agent)538
const currentState = this.expectCurrent(state, ref)539
const current = currentState.goal540
if (!allowed.includes(current.phase)) throw this.transitionError(current, operation, allowed)541
return this.commitCurrent(agent, currentState, runtime, operation, this.withPhase(current, phase), activation)542
}544
/** Render a stable invalid-transition error. */545
private transitionError(current: GoalSnapshot, operation: GoalOperation, allowed: readonly GoalPhase[]): GoalError {546
return new GoalError(547
`cannot ${operation} goal "${current.id}" from phase "${current.phase}"; expected ${allowed.join(' or ')}`,548
'GOAL_INVALID_TRANSITION',549
)550
}552
/** Commit a mutation that retains the current goal's derived counters/times. */553
private commitCurrent(554
agent: Agent,555
state: GoalProjection,556
runtime: GoalRuntimeState,557
operation: Exclude<GoalOperation, 'create' | 'clear'>,558
goal: GoalSnapshot,559
activation: GoalActivation,560
): GoalView {561
return this.commitSnapshot(562
agent,563
runtime,564
operation,565
goal,566
state.roundsStarted,567
state.createdAt,568
this.nextMutationTime(state),569
activation,570
)571
}573
/** Clamp a current goal's next timestamp across backward wall-clock movement. */574
private nextMutationTime(state: GoalProjection): number {575
return Math.max(Date.now(), state.updatedAt)576
}578
/** Build and commit one full-snapshot mutation. */579
private commitSnapshot(580
agent: Agent,581
runtime: GoalRuntimeState,582
operation: Exclude<GoalOperation, 'clear'>,583
goal: GoalSnapshot,584
roundsStarted: number,585
createdAt: number,586
updatedAt: number,587
activation: GoalActivation,588
): GoalView {589
const change: GoalSnapshotChangeMeta = {590
kind: 'goal/change',591
version: GOAL_CHANGE_VERSION,592
operation,593
goal,594
roundsStarted,595
createdAt,596
updatedAt,597
}598
this.commit(agent, runtime, change, activation)599
return {600
...goal,601
roundsStarted,602
createdAt,603
updatedAt,604
activation: runtime.activation,605
}606
}608
/** Commit one mutation into the goal log and live event stream. */609
private commit(agent: Agent, runtime: GoalRuntimeState, change: GoalChangeMeta, activation: GoalActivation): void {610
const ref = goalChangeRef(change)611
runtime.pendingActivation = { offset: agent.session.seq, activation }612
try {613
const event = agent.session.append('goal/change', change)614
/* v8 ignore next -- Session.append returns the event committed at the pre-append seq. */615
if (SessionSeq(runtime.pendingActivation.offset) === event.seq) runtime.activation = activation616
} finally {617
runtime.pendingActivation = undefined618
}619
const goal = this.view(this.state(agent.session), runtime)620
const notification: GoalChanged = {621
operation: change.operation,622
ref: { ...ref },623
...goal === undefined ? {} : { goal },624
}625
agentEvents(this.ctx, agent).emit('goal/changed', { change: notification })626
}628
/** Build a detached current view. */629
private view(state: GoalProjection | null, runtime: GoalRuntimeState): GoalView | undefined {630
if (state === null) return undefined631
return {632
...state.goal,633
roundsStarted: state.roundsStarted,634
createdAt: state.createdAt,635
updatedAt: state.updatedAt,636
activation: runtime.activation,637
}638
}640
/**641
* Create one Goal through the remote boundary.642
* @param agent - exact live Agent resolved from the wire identity.643
* @param request - objective and optional round cap.644
* @returns the created Goal identity.645
*/646
@Remote('create')647
remoteExportCreate(agent: Agent, request: CreateGoalRequest): CreateGoalResult {648
const view = this.create(agent, request)649
return { ref: { id: view.id, revision: view.revision } }650
}651
}653
export default GoalService