session.ts 35 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836
  1. // Sessions remain resident after creation so their open Remote sources keep running off-screen.
  2. import type { Context } from '@deepseek-ai/cordis'
  3. import type { InboxState } from '@deepseek-ai/dsh-agent/types'
  4. import { randomUUID } from '@deepseek-ai/dsh-util-crypto'
  5. import type { AttachmentIdType, FileAttachmentRef, ImageAttachmentRef } from '@deepseek-ai/dsh-attachment'
  6. import type { SubagentAddress } from '@deepseek-ai/dsh-subagent/client'
  7. import type { MessageId } from '@deepseek-ai/dsh-llm/brand'
  8. import { SessionLogOffset, SessionSeq, type SessionId } from '@deepseek-ai/dsh-session/types'
  9. import { SessionEventStream } from '../transport.ts'
  10. import type { SessionJournalChange } from '../transport.ts'
  11. import type {
  12. PromptContentPart,
  13. QueueAction,
  14. SessionAddress,
  15. SessionAssistantStreamBaseline,
  16. SessionProjectionBaseline,
  17. SessionRequestId,
  18. } from '../../types.ts'
  19. import type {
  20. BeginSubmissionInput, PendingSubmissionRetirement, SessionFace, SubmissionHandle,
  21. } from '../contract/session.ts'
  22. import type {
  23. OpenState, PendingSubmission, PromptError, SessionSnapshot,
  24. } from '../contract/snapshot.ts'
  25. import { MutableSessionEventSource } from '../contract/events.ts'
  26. import type {
  27. SessionEventLikeEntry, SessionLiveEventEntry,
  28. } from '../contract/events.ts'
  29. import { Notifier } from './notifier.ts'
  30. import { isRemoteFailure } from '@deepseek-ai/dsh-api-gateway/client'
  31. import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
  32. import type { RemoteFailure, RemoteResult } from '@deepseek-ai/dsh-typert-protocol'
  33. import type { SessionRemotes } from './remotes.ts'
  34. import { ProjectionValueStore } from './projection-store.ts'
  35. import type { ProjectionsBaseline } from './projection-store.ts'
  36. import { resolvedClientTimeZone } from '../time-zone.ts'
  37. import {
  38. ClientAssistantStream,
  39. type ClientAssistantStreamResult,
  40. } from './assistant-stream.ts'
  41. function projectionsBaseline(value: SessionProjectionBaseline): ProjectionsBaseline {
  42. return {
  43. ...value,
  44. asOfSeq: value.asOfSeq === -1 ? -1 : SessionSeq(value.asOfSeq),
  45. }
  46. }
  47. /** Messages requested per history page. */
  48. export const PAGE_MESSAGES = 50
  49. /** Messages requested per page while a turn jump loops backwards (fewer, larger round trips). */
  50. export const JUMP_PAGE_MESSAGES = 200
  51. /** Manager-owned observers of a Session object's local state edges. */
  52. export interface SessionOptions {
  53. /** Catalog-discovered address selecting non-activating subagent transport. */
  54. address?: SubagentAddress
  55. /** Whether the exact direct parent Agent was live at the latest catalog read; absent before that read. */
  56. parentAvailable?: boolean
  57. /**
  58. * First ACCEPTED prompt on a blank session (fires at most once, on the
  59. * prompt RPC's success response): the manager mirrors the blank→false flip
  60. * into its list row so the session surfaces without waiting for a host
  61. * frame. Acceptance is the flip point because it proves the user message
  62. * is in the host log; a rejected first prompt keeps the session blank
  63. * (hidden, still reusable by connectWorkspace).
  64. */
  65. onEngaged?(session: Session): void
  66. /**
  67. * Manager-owned projection value store to adopt (frames route through the
  68. * manager and values outlive instantiation); omitted, the Session owns a
  69. * private store (bare object-layer construction).
  70. */
  71. projections?: ProjectionValueStore
  72. }
  73. /**
  74. * Owns a session's event window, lifecycle state, and observable
  75. * snapshot. React bindings remain outside this data layer. Features see only
  76. * the {@link SessionFace} slice (ISession verbs + the snapshot source); the
  77. * remaining public members are Session Controller internals.
  78. */
  79. export class Session implements SessionFace {
  80. // ---- Window and derived state (all private; the snapshot is the only read API) ----
  81. private baseSeq = SessionLogOffset(0)
  82. private hasMore = false
  83. private openState: OpenState = 'cold'
  84. private openError: RemoteFailure | null = null
  85. private openPromise: Promise<void> | null = null
  86. /** Bumped by stream replacement to invalidate an in-flight doOpen. Stale
  87. * passes drop all writes once the generation moves on. */
  88. private openGeneration = 0
  89. private loadingOlder = false
  90. /** Shared low-water target of the running jump loop; null when no jump is paging. */
  91. private jumpTargetSeq: SessionSeq | null = null
  92. /** The running jump loop's completion, shared by retargeting callers. */
  93. private jumpPromise: Promise<void> | null = null
  94. private readonly stopObservingInbox: () => void
  95. private readonly assistantStream = new ClientAssistantStream()
  96. private running = false
  97. private address: SubagentAddress | undefined
  98. private parentAvailable: boolean | undefined
  99. /**
  100. * Sticky send marker, private input of the composerPhase derivation: set
  101. * synchronously before prompt()'s first await, never reset — the blank →
  102. * engaging edge of the phase machine (see ComposerPhase).
  103. */
  104. private promptAttempted = false
  105. /** A first accepted prompt stays in the engaging phase until its turn is observable. */
  106. private firstPromptPendingTurn = false
  107. /** Empty-log mirror (see ConversationSnapshot.blank); unknown bare sessions begin conservatively blank. */
  108. private blankBit = true
  109. private removed = false
  110. private promptError: PromptError | null = null
  111. private lastAgentError: string | null = null
  112. /** Local submission echoes, insertion-ordered (see SessionSnapshot.pendingSubmissions). */
  113. private pendingSubmissions: readonly PendingSubmission[] = []
  114. /** Per-echo settlement state; `retiring` latches the first observation so a
  115. * Inbox projection and its durable event cannot both retire one echo. */
  116. private readonly submissionSettlements = new Map<SessionRequestId, {
  117. readonly onRetire?: ((retirement: PendingSubmissionRetirement) => void) | undefined
  118. retiring: boolean
  119. }>()
  120. /** Owns the addressed page/follow lifecycle while this Session is open. */
  121. private events: SessionEventStream | undefined
  122. /**
  123. * Per-session projection value store (push model; see the session-projection
  124. * subsystem page, docs/subsystems/session-projection.md): finished whole
  125. * values computed on the Host, seeded by the tail page's
  126. * projections block and updated by Session Controller control frames under the
  127. * one higher-seq-wins rule. Keys are read via `projections.faceOf(key)`
  128. * (the useProjection resolution face); the conversation snapshot never
  129. * carries projection values, and no client-side domain folding exists.
  130. * Manager-owned when constructed through SessionManager (frames route and
  131. * the store outlives instantiation, the title-snapshot precedent); a bare
  132. * construction gets a private store.
  133. */
  134. readonly projections: ProjectionValueStore
  135. /** Contiguous history and live tail consumed by Conversation assembly. */
  136. readonly eventSource = new MutableSessionEventSource()
  137. private snapshotCache: SessionSnapshot
  138. private readonly notifier: Notifier
  139. /**
  140. * Agent-scoped cordis context, bound once by ClientSessions when it
  141. * mints the scope (the client mirror of the host Agent's loopCtx). The
  142. * Session dispatches its own scoped events through it; undefined means
  143. * unbound (bare object-layer construction) or already pruned — both skip
  144. * dispatch-dependent behavior rather than fail.
  145. */
  146. private actx: Context | undefined
  147. /**
  148. * @param sessionId - Host session identity (client sessions are always Host-born).
  149. * @param remote - generated Remote namespaces this session calls.
  150. * @param options - optional manager-owned state observers.
  151. */
  152. constructor(
  153. readonly sessionId: SessionId,
  154. private readonly remote: SessionRemotes,
  155. private readonly options: SessionOptions = {},
  156. ) {
  157. this.projections = options.projections ?? new ProjectionValueStore()
  158. this.address = options.address
  159. this.parentAvailable = options.parentAvailable
  160. this.notifier = new Notifier(() => {
  161. this.snapshotCache = this.buildSnapshot()
  162. })
  163. this.snapshotCache = this.buildSnapshot()
  164. this.stopObservingInbox = this.projections.faceOf('inbox').subscribe(() => {
  165. this.observeSubmissionInbox()
  166. })
  167. }
  168. /**
  169. * Bind the Agent-scoped context minted by ClientSessions (single write;
  170. * a second bind is a wiring error and throws). Direction stays one-way at
  171. * this binding boundary: consumers still reach the Session via `sessions.sessionOf`,
  172. * while the Session holds its own dispatch point (host Agent.loopCtx
  173. * mirror).
  174. * @param actx - the agent's scoped context.
  175. */
  176. bindScope(actx: Context): void {
  177. if (this.actx !== undefined) throw new Error(`session ${this.sessionId} already has a bound scope`)
  178. this.actx = actx
  179. }
  180. /** Release the bound scope at prune time (a later rebind accompanies a freshly minted scope). */
  181. unbindScope(): void {
  182. this.actx = undefined
  183. }
  184. // ---- Operations ----
  185. /**
  186. * Register one local submission echo (see the ISession declaration).
  187. * Synchronous through markDirty: the echo is in the very next snapshot, so
  188. * the conversation can paint it before the caller starts serializing.
  189. * @param input - echo content and the optional settlement callback.
  190. * @returns the minted identity for {@link prompt} plus the pre-prompt abandon path.
  191. */
  192. beginSubmission(input: BeginSubmissionInput): SubmissionHandle {
  193. const requestId = randomUUID() as SessionRequestId
  194. this.pendingSubmissions = [...this.pendingSubmissions, {
  195. requestId,
  196. placement: this.running
  197. ? input.mode === 'steer' ? 'steering' : 'queued'
  198. : 'transcript',
  199. time: Date.now(),
  200. text: input.text,
  201. attachments: input.attachments,
  202. }]
  203. this.submissionSettlements.set(requestId, { onRetire: input.onRetire, retiring: false })
  204. // The blank → engaging edge flips here, ahead of prompt(): the composer
  205. // docks and the echo renders on the click's own frame.
  206. this.promptAttempted = true
  207. this.notifier.markDirty()
  208. return { requestId, abandon: () => { this.retireFailedSubmission(requestId) } }
  209. }
  210. /**
  211. * Send (queue/steer passed through 1:1); failures land in the snapshot's promptError.
  212. * @param content - text, browser-owned temporary image uploads, and staged-file receipts.
  213. * @param mode - queue appends after the current turn; steer interrupts it.
  214. * @param signal - optional caller cancellation for the complete admission round-trip.
  215. * @param requestId - identity from {@link beginSubmission}; a failed identified prompt retires its echo.
  216. * @returns the prompt result (also mirrored into promptError on failure).
  217. */
  218. async prompt(
  219. content: PromptContentPart[],
  220. mode: 'queue' | 'steer',
  221. signal?: AbortSignal,
  222. requestId?: SessionRequestId,
  223. ): Promise<RemoteResult<{ accepted: true }>> {
  224. this.promptError = null
  225. this.lastAgentError = null
  226. // Synchronous, before the first await: the blank → engaging edge must be
  227. // visible on the session area's very first frame when a caller sends
  228. // ahead of navigation (first-send flow).
  229. this.promptAttempted = true
  230. if (this.blankBit) this.firstPromptPendingTurn = true
  231. this.notifier.markDirty()
  232. let result: RemoteResult<{ accepted: true }>
  233. if (this.address === undefined) {
  234. const clientTimeZone = resolvedClientTimeZone()
  235. result = await this.remote.session.prompt({
  236. requestId: requestId ?? randomUUID() as SessionRequestId,
  237. sessionId: this.sessionId,
  238. mode,
  239. content,
  240. clientTimeZone,
  241. }, signal)
  242. } else if (content.some(part => part.type === 'file')) {
  243. result = {
  244. ok: false,
  245. error: new RemoteError(
  246. 'subagent/attachment-invalid',
  247. 'subagent continuation does not accept files',
  248. { reason: 'SUBAGENT_FILE_UNSUPPORTED' },
  249. ),
  250. }
  251. } else {
  252. // The preceding branch rejects file parts before the narrower subagent
  253. // wire type is used; this array is not filtered or reordered.
  254. const routedContent = content as Exclude<PromptContentPart, { readonly type: 'file' }>[]
  255. const routed = await this.remote.subagents.prompt({
  256. requestId: randomUUID() as SessionRequestId,
  257. parentSessionId: this.address.parentSessionId,
  258. childSessionId: this.address.childSessionId,
  259. mode: 'continuable',
  260. delivery: mode,
  261. content: routedContent,
  262. clientTimeZone: resolvedClientTimeZone(),
  263. }, signal)
  264. result = routed.ok ? { ok: true, value: { accepted: true } } : routed
  265. }
  266. if (!result.ok) {
  267. if (requestId !== undefined) this.retireFailedSubmission(requestId)
  268. this.promptError = { op: 'send', error: result.error }
  269. this.notifier.markDirty()
  270. return result
  271. }
  272. // Blank flips on ACCEPTANCE, not attempt: an accepted prompt starts the
  273. // conversation's first turn on the host (the host criterion — a logged
  274. // turn/start — is fact, not optimism; standalone command and projection
  275. // events never flip it), while a rejected first prompt must keep the
  276. // session blank — the client-side blank mirror only ever lowers, so
  277. // flipping early on a failure would surface the session forever and
  278. // strip its connectWorkspace reuse eligibility against the host's
  279. // authority.
  280. if (this.blankBit) {
  281. this.blankBit = false
  282. this.options.onEngaged?.(this)
  283. this.notifier.markDirty()
  284. }
  285. return result
  286. }
  287. /**
  288. * Resolve one image referenced by this session into browser-consumable bytes.
  289. * @param attachmentId - opaque id found in the folded session log.
  290. * @returns the authenticated reference and decoded bytes.
  291. */
  292. async readAttachment(
  293. attachmentId: AttachmentIdType,
  294. ): Promise<RemoteResult<{ attachment: ImageAttachmentRef; data: Uint8Array }>> {
  295. const result = await this.remote.session.attachment({
  296. sessionId: this.sessionId,
  297. attachmentId,
  298. })
  299. if (!result.ok) return result
  300. const binary = atob(result.value.data)
  301. const data = Uint8Array.from(binary, char => char.charCodeAt(0))
  302. return { ok: true, value: { attachment: result.value.attachment, data } }
  303. }
  304. /** Apply one operation to a still-pending queue occurrence. */
  305. async updateQueue(itemId: MessageId, action: QueueAction): Promise<RemoteResult<{ accepted: true }>> {
  306. return this.remote.session.updateQueue({ sessionId: this.sessionId, itemId, action })
  307. }
  308. /**
  309. * Stop the active turn while the Host preserves pending inbox work; failures
  310. * land in promptError (same error-strip display slot). A subagent address
  311. * routes through `subagents.interruptByParent`, whose durable parent-address
  312. * authority works without a live parent Agent.
  313. * @returns the cancel result.
  314. */
  315. async cancel(): Promise<RemoteResult<{ accepted: true }>> {
  316. const address = this.address
  317. const result = address !== undefined
  318. ? await this.remote.subagents.interruptByParent(
  319. address.childSessionId,
  320. address.parentSessionId,
  321. 'continuable',
  322. )
  323. : await this.remote.session.cancel({ sessionId: this.sessionId })
  324. if (!result.ok) {
  325. this.promptError = { op: 'stop', error: result.error }
  326. this.notifier.markDirty()
  327. }
  328. return result
  329. }
  330. /**
  331. * Rename: contract session.rename 1:1. On success settle the 'title'
  332. * projection cell from the response's `{title, seq}` under the store's
  333. * higher-seq-wins rule (the push frame arriving later is a no-op replay),
  334. * so the list row and any useProjection('title') reader update without
  335. * waiting for the control-stream projection update.
  336. * @param title - raw title text (the host normalizes acceptance).
  337. * @returns the rename result (normalized accepted title + title event seq).
  338. */
  339. async rename(title: string): Promise<RemoteResult<{ title: string; seq: SessionSeq }>> {
  340. const result = await this.remote.session.rename({ sessionId: this.sessionId, title })
  341. if (!result.ok) return result
  342. const seq = SessionSeq(result.value.seq)
  343. this.projections.apply('title', result.value.title, seq)
  344. return { ok: true, value: { title: result.value.title, seq } }
  345. }
  346. /**
  347. * Execute one slash-command line against this session's agent — pure
  348. * admission semantics (the host executor durably logs the lifecycle;
  349. * outcomes render as flow nodes, never as a response echo).
  350. * @param line - the full command line, leading slash included.
  351. * @returns the admission result.
  352. */
  353. async command(line: string): Promise<RemoteResult<{ matched: boolean }>> {
  354. const result = await this.remote.commands.execute(this.sessionId, line, [])
  355. if (!result.ok) return result
  356. return { ok: true, value: { matched: result.value !== undefined } }
  357. }
  358. /** First open: pull the tail page (idempotent — in-flight/already-open returns the existing promise). */
  359. open(): Promise<void> {
  360. if (this.openState === 'open') return Promise.resolve()
  361. if (this.openPromise !== null) return this.openPromise
  362. const promise = this.doOpen(this.openGeneration).finally(() => {
  363. // Identity-guarded: a superseded open must not null out the promise resync just started.
  364. if (this.openPromise === promise) this.openPromise = null
  365. })
  366. this.openPromise = promise
  367. return promise
  368. }
  369. /** Page up: pull one earlier page with the window's first seq as beforeSeq and prepend. */
  370. async loadOlder(): Promise<void> {
  371. if (this.openState !== 'open' || !this.hasMore || this.loadingOlder) return
  372. const events = this.events
  373. if (events === undefined) return
  374. this.loadingOlder = true
  375. this.notifier.markDirty()
  376. try {
  377. await events.prepend({ beforeSeq: this.baseSeq, maxMessages: PAGE_MESSAGES })
  378. } catch (error) {
  379. if (!isRemoteFailure(error)) {
  380. console.error('[session-controller] loadOlder failed:', error)
  381. }
  382. } finally {
  383. this.loadingOlder = false
  384. this.notifier.markDirty()
  385. }
  386. }
  387. /** Jump loader: page backwards until the window covers seq (see ISession.loadThrough). */
  388. loadThrough(seq: SessionSeq): Promise<void> {
  389. if (this.openState !== 'open' || !this.hasMore || this.baseSeq <= seq) return Promise.resolve()
  390. if (this.jumpPromise !== null) {
  391. // Retarget the running loop to the lowest requested seq.
  392. this.jumpTargetSeq = SessionSeq(Math.min(this.jumpTargetSeq ?? seq, seq))
  393. return this.jumpPromise
  394. }
  395. // A plain single-page pull owns the busy flag; the jump does not queue
  396. // behind it (the caller retries once it settles) and must leave no
  397. // target behind — only the loop's finally clears that field, and no
  398. // loop starts here.
  399. if (this.loadingOlder) return Promise.resolve()
  400. this.jumpTargetSeq = seq
  401. this.loadingOlder = true
  402. this.notifier.markDirty()
  403. // Stale-pass guard (the doOpen pattern): a resync mid-loop replaces the
  404. // stream generation; this pass then stops instead of paging the new
  405. // generation toward its old target.
  406. const generation = this.openGeneration
  407. this.jumpPromise = (async () => {
  408. try {
  409. while (this.hasMore && this.jumpTargetSeq !== null && this.baseSeq > this.jumpTargetSeq) {
  410. if (generation !== this.openGeneration) return
  411. const events = this.events
  412. if (events === undefined) return
  413. const before = this.baseSeq
  414. await events.prepend({ beforeSeq: this.baseSeq, maxMessages: JUMP_PAGE_MESSAGES })
  415. // No-progress guard: an empty or dropped page that still claims more
  416. // history must end the loop, not spin it.
  417. if (this.baseSeq >= before) return
  418. }
  419. } catch (error) {
  420. if (!isRemoteFailure(error)) {
  421. console.error('[session-controller] loadThrough failed:', error)
  422. }
  423. } finally {
  424. this.jumpTargetSeq = null
  425. this.jumpPromise = null
  426. this.loadingOlder = false
  427. this.notifier.markDirty()
  428. }
  429. })()
  430. return this.jumpPromise
  431. }
  432. /** Rebuild an opened history source after address replacement.
  433. * Invalidates any in-flight open first; projection state belongs to the independently
  434. * reconnecting control stream and remains untouched. */
  435. async resync(): Promise<void> {
  436. if (this.openState === 'cold') return // never opened: no window to rebuild (doOpen flips to 'loading' synchronously, so cold implies no in-flight open)
  437. this.openGeneration++
  438. const events = this.events
  439. this.events = undefined
  440. await events?.dispose()
  441. this.openPromise = null
  442. this.openState = 'cold'
  443. this.openError = null
  444. this.baseSeq = SessionLogOffset(0)
  445. this.notifier.markDirty()
  446. await this.open()
  447. }
  448. // ---- Subscription API (useSyncExternalStore direct wiring) ----
  449. /**
  450. * uSES subscription entry.
  451. * @param listener - change callback.
  452. * @returns the unsubscribe function.
  453. */
  454. subscribe(listener: () => void): () => void {
  455. return this.notifier.subscribe(listener)
  456. }
  457. /**
  458. * Cached Session snapshot (rebuilt lazily when dirty with no listeners).
  459. * @returns the cached reference (stable until the next flush).
  460. */
  461. getSnapshot(): SessionSnapshot {
  462. this.notifier.ensureFresh()
  463. return this.snapshotCache
  464. }
  465. // ---- Manager-only entry points (@internal; never called by the UI) ----
  466. /**
  467. * Running-bit relay from the host stream (list entry and snapshot stay consistent).
  468. * @param running - the new running state.
  469. */
  470. handleRunning(running: boolean): void {
  471. // Turn-start conversion: a blank session never runs, so the first
  472. // running:true proves another side's first message landed.
  473. if (running && this.blankBit) {
  474. this.blankBit = false
  475. this.notifier.markDirty()
  476. }
  477. if (running) this.firstPromptPendingTurn = false
  478. if (this.running === running) return
  479. this.running = running
  480. this.notifier.markDirty()
  481. }
  482. /**
  483. * Install or clear the catalog-discovered transport address. A changed
  484. * address rebuilds an already-open window through its new history route.
  485. * @param address - direct parent/child address, or undefined for ordinary transport.
  486. * @param parentAvailable - latest exact-parent availability hint, or undefined before a catalog read.
  487. */
  488. configureSubagent(address: SubagentAddress | undefined, parentAvailable?: boolean): void {
  489. const same = this.address?.parentSessionId === address?.parentSessionId
  490. && this.address?.childSessionId === address?.childSessionId
  491. && this.address?.mode === address?.mode
  492. this.address = address
  493. this.parentAvailable = parentAvailable
  494. if (!same && this.openState !== 'cold') void this.resync()
  495. else this.notifier.markDirty()
  496. }
  497. /**
  498. * Update only the parent availability hint from a catalog refresh.
  499. * @param available - whether the exact direct parent is live.
  500. */
  501. handleSubagentParentAvailable(available: boolean): void {
  502. if (this.parentAvailable === available) return
  503. this.parentAvailable = available
  504. this.notifier.markDirty()
  505. }
  506. /**
  507. * Blank-bit relay from the authoritative summary source (`session.list` and
  508. * `api-session/added`). Monotone: once any signal (local first send,
  509. * running flip, an earlier summary) cleared it, a stale true never
  510. * re-blanks.
  511. * @param blank - the summary's derived empty-log bit.
  512. */
  513. handleBlank(blank: boolean): void {
  514. if (blank === this.blankBit) return
  515. if (blank && (this.promptAttempted || this.running)) return
  516. this.blankBit = blank
  517. this.notifier.markDirty()
  518. }
  519. /** `api-session/removed` relay: flag the snapshot while retaining the resident instance. */
  520. handleRemoved(): void {
  521. this.removed = true
  522. this.notifier.markDirty()
  523. }
  524. /**
  525. * `api-session/error` relay: the outlet for live failures with no turn position.
  526. * @param message - the stringified error.
  527. */
  528. handleAgentError(message: string): void {
  529. this.lastAgentError = message
  530. this.notifier.markDirty()
  531. }
  532. /**
  533. * Stop the Session's live Remote source.
  534. * @returns when the Remote iterator has completed teardown.
  535. */
  536. async dispose(): Promise<void> {
  537. this.stopObservingInbox()
  538. // Unsettled echoes retire as failed so their owners can restore or
  539. // release browser resources; echoes already scheduled as observed keep
  540. // that settlement.
  541. for (const requestId of [...this.submissionSettlements.keys()]) {
  542. this.retireFailedSubmission(requestId)
  543. }
  544. this.openGeneration++
  545. const events = this.events
  546. this.events = undefined
  547. await events?.dispose()
  548. }
  549. // ---- Private ----
  550. /** @param generation - openGeneration at launch; stale passes cannot publish after replacement. */
  551. private async doOpen(generation: number): Promise<void> {
  552. this.openState = 'loading'
  553. this.openError = null
  554. this.notifier.markDirty()
  555. const events = new SessionEventStream(this.remote, this.sessionAddress(), {
  556. publish: (change) => {
  557. if (generation !== this.openGeneration || this.events !== events) return
  558. this.acceptEventChange(change)
  559. },
  560. failed: (error) => {
  561. this.failEventStream(events, generation, error)
  562. },
  563. })
  564. this.events = events
  565. try {
  566. await events.open({ maxMessages: PAGE_MESSAGES })
  567. if (generation !== this.openGeneration || this.events !== events) return
  568. this.openState = 'open'
  569. } catch (error) {
  570. if (generation !== this.openGeneration || this.events !== events) return
  571. if (!isRemoteFailure(error)) throw error
  572. this.events = undefined
  573. this.openState = 'error'
  574. this.openError = error
  575. } finally {
  576. if (generation === this.openGeneration) this.notifier.markDirty()
  577. }
  578. }
  579. /** Apply one contiguous journal update already reconciled by the Remote stream. */
  580. private acceptEventChange(change: SessionJournalChange): void {
  581. switch (change.type) {
  582. case 'replace':
  583. this.installWindow(
  584. change.entries,
  585. change.hasMore,
  586. change.page.projections === undefined ? undefined : projectionsBaseline(change.page.projections),
  587. change.page.assistantStream,
  588. )
  589. return
  590. case 'prepend':
  591. this.prependWindow(change.entries, change.hasMore)
  592. return
  593. case 'append':
  594. this.publishAssistantEntry(this.assistantStream.acceptDurable(change.entry))
  595. return
  596. case 'assistant-stream':
  597. this.publishAssistantEntry(this.assistantStream.acceptFrame(change.frame))
  598. }
  599. }
  600. /** Replace the complete contiguous window and apply page-owned projection metadata. */
  601. private installWindow(
  602. entries: readonly SessionEventLikeEntry[],
  603. hasMore: boolean,
  604. projections?: ProjectionsBaseline,
  605. assistantStream?: SessionAssistantStreamBaseline,
  606. ): void {
  607. // A durable gap-repair page has no assistant baseline. Clearing transient
  608. // attempts makes a held notification reopen follow once for an atomic
  609. // page/baseline pair instead of applying it to an unrelated repair cut.
  610. const visible = this.assistantStream.replace(entries, assistantStream)
  611. this.baseSeq = SessionLogOffset(entries[0]?.event.seq ?? 0)
  612. this.hasMore = hasMore
  613. if (visible.some(entry => entry.event.type === 'turn/start')) this.firstPromptPendingTurn = false
  614. if (projections !== undefined) this.projections.seed(projections)
  615. this.eventSource.replace(visible, hasMore)
  616. for (const entry of visible) this.observeSubmissionEvent(entry.event)
  617. this.notifier.markDirty()
  618. }
  619. private publishAssistantEntry(result: ClientAssistantStreamResult): void {
  620. if (result?.type === 'rebaseline') {
  621. const events = this.events
  622. queueMicrotask(() => {
  623. if (events !== undefined && this.events === events) events.restart()
  624. })
  625. return
  626. }
  627. if (result?.type === 'settlement') {
  628. this.eventSource.settleAssistant(result.attemptId, result.entry)
  629. this.observeSubmissionEvent(result.entry.event)
  630. this.notifier.markDirty()
  631. return
  632. }
  633. if (result?.type === 'abandonment') {
  634. this.eventSource.settleAssistant(result.attemptId)
  635. this.notifier.markDirty()
  636. return
  637. }
  638. if (result?.type === 'publish' && this.appendLive(result.entry)) {
  639. this.notifier.markDirty()
  640. } else if (result?.type === 'transient') {
  641. this.eventSource.append(result.entry)
  642. this.notifier.markDirty()
  643. }
  644. }
  645. /** Prepend one stream-validated history page. */
  646. private prependWindow(entries: readonly SessionEventLikeEntry[], hasMore: boolean): void {
  647. this.baseSeq = entries[0] === undefined ? this.baseSeq : SessionLogOffset(entries[0].event.seq)
  648. this.hasMore = hasMore
  649. this.eventSource.prepend(entries, hasMore)
  650. }
  651. /** Append one stream-validated live event. */
  652. private appendLive(entry: SessionLiveEventEntry): boolean {
  653. const event = entry.event
  654. const awaitingFirstTurn = this.firstPromptPendingTurn
  655. if (event.type === 'turn/start') this.firstPromptPendingTurn = false
  656. this.eventSource.append(entry)
  657. // After the feed append: the conversation assembly's animation frame is
  658. // registered by the feed subscribers above, so the echo-retirement frame
  659. // scheduled here always runs after the durable node became renderable.
  660. this.observeSubmissionEvent(event)
  661. return awaitingFirstTurn !== this.firstPromptPendingTurn
  662. }
  663. /** Observe durable acceptance even when insertion and claim share one projection notification. */
  664. private observeSubmissionEvent(event: { readonly type: string; readonly data?: unknown }): void {
  665. if (this.submissionSettlements.size === 0) return
  666. if (event.type === 'agent/inbox/spliced') {
  667. const splice = event.data as { readonly inserted?: unknown } | undefined
  668. if (Array.isArray(splice?.inserted)) {
  669. for (const message of splice.inserted) this.observeSubmissionEvent({ type: 'user/message', data: message })
  670. }
  671. return
  672. }
  673. if (event.type !== 'user/message') return
  674. // Structural read: window entries may be compact history records, so the
  675. // fields are narrowed rather than trusted (same posture as Conversation
  676. // assembly matchers).
  677. const data = event.data as { readonly source?: unknown; readonly content?: unknown } | undefined
  678. const source = data?.source as { readonly kind?: unknown; readonly rpcId?: unknown } | undefined
  679. if (source?.kind !== 'user' || typeof source.rpcId !== 'string') return
  680. this.scheduleObservedRetirement(source.rpcId as SessionRequestId, attachmentRefsIn(data?.content))
  681. }
  682. /** Retire local echoes when their accepted messages appear in the durable Inbox projection. */
  683. private observeSubmissionInbox(): void {
  684. if (this.submissionSettlements.size === 0) return
  685. const inbox = this.projections.get('inbox') as InboxState | undefined
  686. if (inbox === undefined) return
  687. for (const message of [...inbox['next-turn'], ...inbox['next-step']]) {
  688. const source = message.source
  689. if (source.kind === 'user' && 'rpcId' in source) {
  690. this.scheduleObservedRetirement(source.rpcId, attachmentRefsIn(message.content))
  691. }
  692. }
  693. }
  694. /**
  695. * Latch one observed settlement and remove the echo an animation frame
  696. * later. The delay keeps the echo in the snapshot until the frame in which
  697. * the durable node (whose assembly frame was registered first) is
  698. * renderable; the render-time rpcId dedupe hides the one-frame overlap.
  699. */
  700. private scheduleObservedRetirement(
  701. requestId: SessionRequestId,
  702. attachments: readonly (ImageAttachmentRef | FileAttachmentRef)[],
  703. ): void {
  704. const settlement = this.submissionSettlements.get(requestId)
  705. if (settlement === undefined || settlement.retiring) return
  706. settlement.retiring = true
  707. scheduleFrame(() => { this.finishSubmission(requestId, { reason: 'observed', attachments }) })
  708. }
  709. /** Remove one unsettled echo immediately (prompt rejection, abort, or disposal). */
  710. private retireFailedSubmission(requestId: SessionRequestId): void {
  711. const settlement = this.submissionSettlements.get(requestId)
  712. if (settlement === undefined || settlement.retiring) return
  713. settlement.retiring = true
  714. this.finishSubmission(requestId, { reason: 'failed' })
  715. }
  716. /** Single removal point: drop the echo, publish, then notify the owner. */
  717. private finishSubmission(requestId: SessionRequestId, retirement: PendingSubmissionRetirement): void {
  718. const settlement = this.submissionSettlements.get(requestId)
  719. /* v8 ignore next -- retiring latches before every schedule, so one settlement never finishes twice. */
  720. if (settlement === undefined) return
  721. this.submissionSettlements.delete(requestId)
  722. this.pendingSubmissions = this.pendingSubmissions.filter(echo => echo.requestId !== requestId)
  723. this.notifier.markDirty()
  724. settlement.onRetire?.(retirement)
  725. }
  726. /** Publish a terminal background failure only while this stream still owns the Session. */
  727. private failEventStream(events: SessionEventStream, generation: number, error: unknown): void {
  728. if (generation !== this.openGeneration || this.events !== events) return
  729. if (!isRemoteFailure(error)) throw error
  730. this.openGeneration++
  731. this.events = undefined
  732. this.openPromise = null
  733. this.openState = 'error'
  734. this.openError = error
  735. void events.dispose()
  736. this.notifier.markDirty()
  737. }
  738. private buildSnapshot(): SessionSnapshot {
  739. return {
  740. sessionId: this.sessionId,
  741. pendingSubmissions: this.pendingSubmissions,
  742. running: this.running,
  743. subagent: this.address === undefined
  744. ? null
  745. : {
  746. address: this.address,
  747. ...(this.parentAvailable === undefined ? {} : { parentAvailable: this.parentAvailable }),
  748. },
  749. removed: this.removed,
  750. openState: this.openState,
  751. openError: this.openError,
  752. hasMore: this.hasMore,
  753. loadingOlder: this.loadingOlder,
  754. promptError: this.promptError,
  755. blank: this.blankBit,
  756. lastAgentError: this.lastAgentError,
  757. promptAttempted: this.promptAttempted,
  758. awaitingFirstTurn: this.firstPromptPendingTurn,
  759. }
  760. }
  761. private sessionAddress(): SessionAddress {
  762. return this.address === undefined
  763. ? { kind: 'session', sessionId: this.sessionId }
  764. : { kind: 'subagent', ...this.address }
  765. }
  766. }
  767. /** Run one callback on the next animation frame, or a macrotask where no frame clock exists. */
  768. function scheduleFrame(fn: () => void): void {
  769. if (typeof requestAnimationFrame === 'function') requestAnimationFrame(() => { fn() })
  770. else setTimeout(fn, 0)
  771. }
  772. /** Attachment references in one structurally-read content block list, in block order. */
  773. function attachmentRefsIn(content: unknown): readonly (ImageAttachmentRef | FileAttachmentRef)[] {
  774. if (!Array.isArray(content)) return []
  775. const refs: Array<ImageAttachmentRef | FileAttachmentRef> = []
  776. for (const block of content) {
  777. if (typeof block !== 'object' || block === null) continue
  778. const candidate = block as { readonly type?: unknown; readonly attachment?: unknown }
  779. if ((candidate.type === 'image' || candidate.type === 'file')
  780. && typeof candidate.attachment === 'object' && candidate.attachment !== null) {
  781. refs.push(candidate.attachment as ImageAttachmentRef | FileAttachmentRef)
  782. }
  783. }
  784. return refs
  785. }