1
/** Process-local assistant state retained for reconnecting Web followers. */3
import type { AssistantStreamFrame } from '@deepseek-ai/dsh-agent'4
import { AssistantStreamAccumulator } from '@deepseek-ai/dsh-llm'5
import type { SessionSeqCursor } from '@deepseek-ai/dsh-session'6
import type { JsonValue } from '@deepseek-ai/dsh-util-values'7
import type {8
SessionAssistantStreamAttempt,9
SessionAssistantStreamBaseline,10
} from './types.ts'12
interface MutableAttempt {13
readonly attemptId: SessionAssistantStreamAttempt['attemptId']14
readonly startedAfterSeq: SessionSeqCursor15
readonly turn: number16
readonly step: number17
readonly stream: AssistantStreamAccumulator18
nextIndex: number19
}21
const EMPTY_BASELINE: SessionAssistantStreamBaseline = { revision: 0 }23
/**24
* Folds dense Agent frames and materializes one shared immutable reconnect25
* baseline per accepted revision.26
*/27
export class SessionAssistantStreamAccumulator {28
private activeAttempt: MutableAttempt | undefined29
private revision = 030
private snapshotValue: SessionAssistantStreamBaseline = EMPTY_BASELINE31
private dirty = false33
/**34
* Fold one trusted frame from the current attached Agent lifecycle.35
* @param frame - next dense process-local Assistant frame.36
* @param durableCursor - last committed Session seq when this frame was observed.37
*/38
accept(frame: AssistantStreamFrame, durableCursor: SessionSeqCursor): void {39
if (frame.type === 'start' && frame.revision === 1 && this.revision !== 0) {40
this.activeAttempt = undefined41
this.revision = 042
}43
if (frame.revision !== this.revision + 1) {44
this.activeAttempt = undefined45
this.revision = frame.revision46
this.dirty = true47
return48
}49
this.revision = frame.revision50
switch (frame.type) {51
case 'start':52
this.activeAttempt = {53
attemptId: frame.attemptId,54
startedAfterSeq: durableCursor,55
turn: frame.turn,56
step: frame.step,57
stream: new AssistantStreamAccumulator(),58
nextIndex: 0,59
}60
break61
case 'chunk': {62
const attempt = this.activeAttempt63
if (attempt === undefined64
|| attempt.attemptId !== frame.attemptId65
|| frame.index !== attempt.nextIndex) {66
this.activeAttempt = undefined67
break68
}69
attempt.stream.push({ time: frame.time, chunk: frame.chunk })70
attempt.nextIndex += 171
break72
}73
case 'end':74
this.activeAttempt = undefined75
break76
}77
this.dirty = true78
}80
/**81
* Read the cached reconnect baseline, materializing it after a state change.82
* @returns the identity-stable baseline for the latest accepted revision.83
*/84
snapshot(): SessionAssistantStreamBaseline {85
if (!this.dirty) return this.snapshotValue86
this.snapshotValue = {87
revision: this.revision,88
...this.activeAttempt === undefined ? {} : {89
activeAttempt: {90
attemptId: this.activeAttempt.attemptId,91
startedAfterSeq: this.activeAttempt.startedAfterSeq,92
turn: this.activeAttempt.turn,93
step: this.activeAttempt.step,94
nextIndex: this.activeAttempt.nextIndex,95
stream: this.activeAttempt.stream.snapshot() as unknown as readonly JsonValue[],96
},97
},98
}99
this.dirty = false100
return this.snapshotValue101
}102
}