1
/** Process-local assistant attempt framing and durable stream accumulation. */3
import {4
AssistantStreamAccumulator,5
BlockAssembler,6
LlmAttemptId,7
type AssistantStreamRecord,8
type ContentBlock,9
type FinishReason,10
type ReplayEnvelope,11
type StreamChunk,12
type TokenUsage,13
} from '@deepseek-ai/dsh-llm'14
import type { AssistantStreamFrame } from '@deepseek-ai/dsh-agent'15
import type { SessionEventMap, SessionId, SessionSeq } from '@deepseek-ai/dsh-session'17
/** Folds one model attempt into one compact stream plus ordered transient frames. */18
export class AssistantStreamAttempt {19
private readonly accumulator = new AssistantStreamAccumulator()20
private readonly assembler = new BlockAssembler()21
private index = 022
private terminal = false23
/** Attempt identity unique within this Agent lifecycle. */24
readonly attemptId: LlmAttemptId26
/** Whether this started attempt has emitted its terminal frame. */27
get ended(): boolean { return this.terminal }29
/**30
* @param sessionId - identity embedded only in the Agent-lifecycle-local attempt id.31
* @param attempt - attached-Session-local attempt counter.32
* @param nextRevision - allocates the next emitted frame revision.33
* @param turn - durable turn owning the request.34
* @param step - durable step owning the request.35
* @param emit - agent-scoped notification publisher.36
*/37
constructor(38
sessionId: SessionId,39
attempt: number,40
private readonly nextRevision: () => number,41
readonly turn: number,42
readonly step: number,43
private readonly emit: (frame: AssistantStreamFrame) => void,44
) {45
this.attemptId = LlmAttemptId(`${sessionId}:${attempt}`)46
}48
/** Publish the opening marker before the first delivered chunk. */49
start(): void {50
this.emit({51
type: 'start',52
attemptId: this.attemptId,53
revision: this.nextRevision(),54
turn: this.turn,55
step: this.step,56
})57
}59
/** Snapshot one chunk once, then feed durable compaction, assembly, and live publication. */60
push(chunk: StreamChunk): void {61
const timed = this.accumulator.push({ time: Date.now(), chunk })62
this.assembler.push(timed.chunk)63
this.emit({64
type: 'chunk',65
attemptId: this.attemptId,66
revision: this.nextRevision(),67
index: this.index++,68
time: timed.time,69
chunk: timed.chunk,70
})71
}73
/**74
* Publish terminal settlement after the matching durable event commits.75
* @param eventType - durable settlement type.76
* @param append - synchronous durable append returning its committed seq.77
*/78
settle(79
eventType: 'assistant/message' | 'assistant/attempt',80
append: () => SessionSeq,81
): void {82
let seq: SessionSeq83
try {84
seq = append()85
} catch (error: unknown) {86
this.abandon()87
throw error88
}89
this.terminal = true90
this.emit({91
type: 'end',92
attemptId: this.attemptId,93
revision: this.nextRevision(),94
index: this.index,95
outcome: { kind: 'committed', eventType, seq },96
})97
}99
/** Publish abandonment when no durable attempt event can be committed. */100
abandon(): void {101
this.terminal = true102
this.emit({103
type: 'end',104
attemptId: this.attemptId,105
revision: this.nextRevision(),106
index: this.index,107
outcome: { kind: 'abandoned' },108
})109
}111
/** Exact compact stream for the final durable event. */112
get stream(): SessionEventMap['assistant/attempt']['stream'] {113
return [...this.accumulator.snapshot()] as AssistantStreamRecord[]114
}116
/** Canonical completed-message blocks from the same chunks. */117
blocks(): ContentBlock[] {118
return this.assembler.blocks()119
}121
/** Safe visible prefix when cancellation interrupts the attempt. */122
interruptedBlocks(): ContentBlock[] {123
return this.assembler.interruptedBlocks()124
}126
/** Latest adapter-reported usage in the stream. */127
get usage(): TokenUsage | undefined {128
return this.assembler.usage129
}131
/** Terminal reason, defaulting to stop when the stream omitted one. */132
get finish(): FinishReason {133
return this.assembler.finish134
}136
/** Replay state carried by the terminal finish record. */137
get replayState(): ReplayEnvelope | undefined {138
return this.assembler.replayState139
}140
}