File size: 6,757 Bytes
d04f74a | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 | /** Live Session queue, jobs, and projection state with reconnect baselines. */
import type { Context } from '@deepseek-ai/cordis'
import type { Agent, InboxState } from '@deepseek-ai/dsh-agent'
import { Deque } from '@deepseek-ai/dsh-deque'
import type { JobSnapshot } from '@deepseek-ai/dsh-jobs'
import type {
Session, SessionId, UserMessage,
} from '@deepseek-ai/dsh-session'
import type { JsonValue } from '@deepseek-ai/dsh-util-values'
import type {
SessionControlBaseline,
SessionControlFrame,
SessionJob,
SessionProjectionBaseline,
SessionProjectionValues,
SessionQueuedItem,
} from './types.ts'
/** Owns the Host-wide Session control stream. */
export class SessionControlController {
private readonly streams = new Set<ControlQueue>()
/** @param ctx - Host context carrying live Agent, projection, and jobs services. */
constructor(private readonly ctx: Context) {
ctx.sessionProjections.onChanged((session, key, value, seq) => {
this.broadcast({
type: 'projection',
sessionId: session.id,
key,
value: value as JsonValue,
seq,
})
if (key !== 'inbox') return
const agent = this.ctx.agents.get(session.id)
if (agent?.session !== session) return
this.broadcast({
type: 'queue',
sessionId: session.id,
items: queueItemsFromInbox(value as InboxState),
})
})
ctx.inject(['jobs'], (jobsCtx) => {
jobsCtx.jobs.onJobsChanged((owner) => { this.onJobsChanged(owner) })
})
ctx.on('session/created', (session) => {
const jobs = this.jobsFor(this.ctx.agents.get(session.id))
if (jobs.length > 0) this.broadcast({ type: 'jobs', sessionId: session.id, jobs })
})
ctx.effect(() => () => {
for (const stream of this.streams) stream.end()
this.streams.clear()
}, 'session-controller.control')
}
/**
* Open one generation of Host-wide live control state.
* @param signal - Remote stream cancellation.
* @returns one complete baseline followed by live replacement frames.
*/
async *control(signal: AbortSignal): AsyncIterable<SessionControlFrame> {
signal.throwIfAborted()
const queue = new ControlQueue()
this.streams.add(queue)
try {
yield { type: 'baseline', value: this.baseline() }
yield* queue.iterate(signal)
} finally {
this.streams.delete(queue)
queue.end()
}
}
private baseline(): SessionControlBaseline {
const sessions = this.ctx.sessions.list()
const queues = Object.create(null) as Record<SessionId, readonly SessionQueuedItem[]>
const jobs = Object.create(null) as Record<SessionId, readonly SessionJob[]>
for (const session of sessions) {
const agent = this.ctx.agents.get(session.id)
queues[session.id] = agent?.session === session ? queueItems(agent) : []
jobs[session.id] = this.jobsFor(agent)
}
return {
queues,
jobs,
projections: this.projectionBaseline(sessions),
}
}
private projectionBaseline(
sessions: readonly Session[],
): Readonly<Record<SessionId, SessionProjectionBaseline>> {
const blocks = Object.create(null) as Record<SessionId, SessionProjectionBaseline>
for (const session of sessions) {
const snapshot = this.ctx.sessionProjections.snapshot(session)
blocks[session.id] = {
asOfSeq: snapshot.asOfSeq,
// Every projection definition validates its value before snapshot publication.
values: snapshot.values as SessionProjectionValues,
}
}
return blocks
}
private onJobsChanged(owner: Agent | undefined): void {
if (owner !== undefined) {
this.broadcast({ type: 'jobs', sessionId: owner.id, jobs: this.jobsFor(owner) })
return
}
for (const session of this.ctx.sessions.list()) {
this.broadcast({
type: 'jobs',
sessionId: session.id,
jobs: this.jobsFor(this.ctx.agents.get(session.id)),
})
}
}
private jobsFor(agent: Agent | undefined): SessionJob[] {
const jobs = this.ctx.get('jobs')
return jobs === undefined ? [] : jobs.list(agent).map(jobView)
}
private broadcast(frame: SessionControlFrame): void {
for (const stream of this.streams) stream.push(frame)
}
}
class ControlQueue {
private readonly buffer = new Deque<SessionControlFrame>()
private wake: (() => void) | undefined
private done = false
push(frame: SessionControlFrame): void {
if (this.done) return
this.buffer.pushBack(frame)
const wake = this.wake
this.wake = undefined
wake?.()
}
end(): void {
if (this.done) return
this.done = true
const wake = this.wake
this.wake = undefined
wake?.()
}
async *iterate(signal: AbortSignal): AsyncIterable<SessionControlFrame> {
const onAbort = (): void => { this.end() }
signal.addEventListener('abort', onAbort, { once: true })
try {
while (!this.done && !signal.aborted) {
const frame = this.buffer.popFront()
if (frame !== undefined) {
yield frame
continue
}
await new Promise<void>((resolve) => { this.wake = resolve })
}
while (this.buffer.size > 0 && !signal.aborted) yield this.buffer.popFront() as SessionControlFrame
} finally {
signal.removeEventListener('abort', onAbort)
this.end()
}
}
}
function queueItems(agent: Agent): SessionQueuedItem[] {
return queueItemsFromInbox({
'next-turn': agent.inbox.nextTurn,
'next-step': agent.inbox.nextStep,
})
}
function queueItemsFromInbox(inbox: InboxState): SessionQueuedItem[] {
return [
...inbox['next-turn'].map(message => ({
id: message.id,
placement: 'queued' as const,
...promptRpcId(message),
message: { id: message.id, content: message.content as unknown as JsonValue[] },
})),
...inbox['next-step'].map(message => ({
id: message.id,
placement: message.source.kind === 'user' ? 'steering' as const : 'context' as const,
...promptRpcId(message),
message: { id: message.id, content: message.content as unknown as JsonValue[] },
})),
]
}
/** Prompt-RPC identity carried by a browser-submitted message's user source. */
function promptRpcId(message: UserMessage): Pick<SessionQueuedItem, 'rpcId'> {
const source = message.source
return source.kind === 'user' && 'rpcId' in source ? { rpcId: source.rpcId } : {}
}
function jobView(job: JobSnapshot): SessionJob {
return {
id: job.id,
kind: job.kind,
label: job.label,
status: job.status,
...(job.detail === undefined ? {} : { detail: job.detail }),
startedAt: job.startedAt,
...(job.finishedAt === undefined ? {} : { finishedAt: job.finishedAt }),
}
}
|