queue-mirror.ts 2.5 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071
  1. import type { ContentBlock } from '@deepseek-ai/dsh-llm/types'
  2. import type { SessionQueuedItem } from '../../types.ts'
  3. import type { SessionEvent } from '@deepseek-ai/dsh-session/types'
  4. import type { QueuedMessage } from '../contract/snapshot.ts'
  5. const QUEUE_PREVIEW_CHARS = 200
  6. // Attachment blocks are excluded: queue presentation renders them from
  7. // `content`, so the text preview covers only what has no visual form.
  8. function previewOf(content: readonly ContentBlock[]): string {
  9. const flat = content
  10. .filter(block => block.type !== 'image' && block.type !== 'file')
  11. .map(block => (block.type === 'text' ? block.text : `[${block.type}]`))
  12. .join(' ').replace(/\s+/g, ' ').trim()
  13. const chars = Array.from(flat)
  14. return chars.length > QUEUE_PREVIEW_CHARS ? `${chars.slice(0, QUEUE_PREVIEW_CHARS).join('')}…` : flat
  15. }
  16. function textOf(content: readonly ContentBlock[]): string | null {
  17. if (!content.every(block => block.type === 'text')) return null
  18. return content.map(block => block.text).join('')
  19. }
  20. type QueueItems = readonly SessionQueuedItem[]
  21. /** Authoritative transient queue projection and durable steering handoff. */
  22. export class SessionQueueMirror {
  23. private current: readonly QueuedMessage[] = []
  24. /**
  25. * Return the current immutable queue projection.
  26. * @returns current queue rows.
  27. */
  28. snapshot(): readonly QueuedMessage[] {
  29. return this.current
  30. }
  31. /**
  32. * Replace from one authoritative stream queue frame.
  33. * @param items - complete host queue snapshot.
  34. */
  35. replace(items: QueueItems): void {
  36. this.current = items.map((item) => {
  37. const content = item.message.content as unknown as readonly ContentBlock[]
  38. return {
  39. id: item.id,
  40. messageId: item.message.id,
  41. placement: item.placement,
  42. ...(item.rpcId === undefined ? {} : { rpcId: item.rpcId }),
  43. content,
  44. preview: previewOf(content),
  45. text: textOf(content),
  46. }
  47. })
  48. }
  49. /**
  50. * Retire a transient steering row once its durable message enters the log.
  51. * @param event - newly contiguous durable Session event.
  52. * @returns whether the projection changed.
  53. */
  54. acceptDurable(event: SessionEvent): boolean {
  55. if (event.type !== 'user/message') return false
  56. const messageId = event.data.id
  57. const index = this.current.findIndex(item =>
  58. item.placement === 'steering' && item.messageId === messageId)
  59. if (index < 0) return false
  60. this.current = this.current.filter((_item, candidate) => candidate !== index)
  61. return true
  62. }
  63. }