codec.ts 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343
  1. /**
  2. * Schema-18 physical chunk-row codec. This package owns the durable tags,
  3. * validation, and row-size limits independently from other persistence formats.
  4. * @module @deepseek-ai/dsh-session-persistence-sqlite/codec
  5. */
  6. import type { StreamChunk } from '@deepseek-ai/dsh-llm'
  7. import type { SessionEvent } from '@deepseek-ai/dsh-session'
  8. /* jscpd:ignore-start -- schema 18 deliberately owns a frozen physical codec;
  9. * importing or sharing the JSONL codec would let that format mutate this database interpreter. */
  10. type DeltaKind = 'text-delta' | 'reasoning-delta' | 'tool-call-delta'
  11. type DeltaEvent = SessionEvent<'assistant/chunk'>
  12. interface RunDataBase {
  13. readonly turn: number
  14. readonly step: number
  15. readonly index: number
  16. readonly dt: number[]
  17. }
  18. interface TextRunData extends RunDataBase {
  19. readonly texts: string[]
  20. }
  21. interface ToolCallRunData extends RunDataBase {
  22. readonly id: Extract<StreamChunk, { type: 'tool-call-delta' }>['id']
  23. readonly name?: string
  24. readonly args: string[]
  25. }
  26. /** One schema-18 packed physical record. */
  27. export type ChunkRow =
  28. | { readonly type: 'text-chunks'; readonly seq0: number; readonly time0: number; readonly data: TextRunData }
  29. | { readonly type: 'reasoning-chunks'; readonly seq0: number; readonly time0: number; readonly data: TextRunData }
  30. | { readonly type: 'tool-call-chunks'; readonly seq0: number; readonly time0: number; readonly data: ToolCallRunData }
  31. /** One scalar event or schema-18 packed physical record. */
  32. export type StorageRecord = SessionEvent | ChunkRow
  33. /** Minimum eligible members in a packed physical record. */
  34. export const MIN_PACKED_ROW_MEMBERS = 3
  35. /** Maximum logical members represented by one packed physical record. */
  36. export const MAX_PACKED_ROW_MEMBERS = 1_024
  37. /** Maximum UTF-8 bytes in one packed physical record's data column. */
  38. export const MAX_PACKED_DATA_BYTES = 1_048_576
  39. function isRecord(value: unknown): value is Record<string, unknown> {
  40. return typeof value === 'object' && value !== null
  41. }
  42. function hasExactKeys(value: object, keys: readonly string[]): boolean {
  43. return Object.keys(value).length === keys.length && keys.every(key => Object.hasOwn(value, key))
  44. }
  45. function classify(event: SessionEvent): DeltaKind | undefined {
  46. if (event.type !== 'assistant/chunk') return undefined
  47. if (!hasExactKeys(event, ['type', 'seq', 'time', 'data'])) return undefined
  48. if (!Number.isSafeInteger(event.seq) || event.seq < 0 || !Number.isSafeInteger(event.time)) return undefined
  49. const data: unknown = event.data
  50. if (!isRecord(data) || !hasExactKeys(data, ['turn', 'step', 'chunk'])) return undefined
  51. if (typeof data.turn !== 'number' || typeof data.step !== 'number') return undefined
  52. const chunk = data.chunk
  53. if (!isRecord(chunk) || typeof chunk.index !== 'number') return undefined
  54. switch (chunk.type) {
  55. case 'text-delta':
  56. case 'reasoning-delta':
  57. return hasExactKeys(chunk, ['type', 'index', 'text']) && typeof chunk.text === 'string'
  58. ? chunk.type
  59. : undefined
  60. case 'tool-call-delta': {
  61. const validKeys = hasExactKeys(chunk, ['type', 'index', 'id', 'argumentsDelta'])
  62. || (hasExactKeys(chunk, ['type', 'index', 'id', 'name', 'argumentsDelta'])
  63. && typeof chunk.name === 'string')
  64. return validKeys && typeof chunk.id === 'string' && typeof chunk.argumentsDelta === 'string'
  65. ? chunk.type
  66. : undefined
  67. }
  68. default:
  69. return undefined
  70. }
  71. }
  72. function toolCallOf(event: DeltaEvent): { readonly id: string; readonly name?: string } {
  73. return event.data.chunk as { readonly id: string; readonly name?: string }
  74. }
  75. function indexOf(event: DeltaEvent): number {
  76. return (event.data.chunk as { readonly index: number }).index
  77. }
  78. function continues(previous: DeltaEvent, next: DeltaEvent, kind: DeltaKind): boolean {
  79. if (next.seq !== previous.seq + 1 || !Number.isSafeInteger(next.time - previous.time)) return false
  80. if (next.data.turn !== previous.data.turn || next.data.step !== previous.data.step) return false
  81. if (indexOf(next) !== indexOf(previous)) return false
  82. if (kind !== 'tool-call-delta') return true
  83. const left = toolCallOf(previous)
  84. const right = toolCallOf(next)
  85. return left.id === right.id
  86. && Object.hasOwn(left, 'name') === Object.hasOwn(right, 'name')
  87. && left.name === right.name
  88. }
  89. function buildRow(kind: DeltaKind, run: readonly DeltaEvent[]): ChunkRow {
  90. const first = run[0] as DeltaEvent
  91. const base = {
  92. turn: first.data.turn,
  93. step: first.data.step,
  94. index: indexOf(first),
  95. dt: run.slice(1).map((event, index) => event.time - (run[index] as DeltaEvent).time),
  96. }
  97. const envelope = { seq0: first.seq, time0: first.time }
  98. if (kind === 'tool-call-delta') {
  99. const call = toolCallOf(first)
  100. return {
  101. type: 'tool-call-chunks',
  102. ...envelope,
  103. data: {
  104. ...base,
  105. id: call.id as Extract<StreamChunk, { type: 'tool-call-delta' }>['id'],
  106. ...Object.hasOwn(call, 'name') ? { name: call.name as string } : {},
  107. args: run.map(event => (event.data.chunk as { readonly argumentsDelta: string }).argumentsDelta),
  108. },
  109. }
  110. }
  111. const data = {
  112. ...base,
  113. texts: run.map(event => (event.data.chunk as { readonly text: string }).text),
  114. }
  115. return kind === 'text-delta'
  116. ? { type: 'text-chunks', ...envelope, data }
  117. : { type: 'reasoning-chunks', ...envelope, data }
  118. }
  119. function packedDataBytes(row: ChunkRow): number {
  120. return Buffer.byteLength(JSON.stringify(row.data))
  121. }
  122. function emitBoundedRun(out: StorageRecord[], kind: DeltaKind, completeRun: readonly DeltaEvent[]): void {
  123. let offset = 0
  124. while (completeRun.length - offset >= MIN_PACKED_ROW_MEMBERS) {
  125. let low = MIN_PACKED_ROW_MEMBERS
  126. let high = Math.min(completeRun.length - offset, MAX_PACKED_ROW_MEMBERS)
  127. const largest = buildRow(kind, completeRun.slice(offset, offset + high))
  128. if (packedDataBytes(largest) <= MAX_PACKED_DATA_BYTES) {
  129. out.push(largest)
  130. offset += high
  131. continue
  132. }
  133. high -= 1
  134. let accepted = 0
  135. let acceptedRow: ChunkRow | undefined
  136. while (low <= high) {
  137. const middle = Math.floor((low + high) / 2)
  138. const candidate = buildRow(kind, completeRun.slice(offset, offset + middle))
  139. if (packedDataBytes(candidate) <= MAX_PACKED_DATA_BYTES) {
  140. accepted = middle
  141. acceptedRow = candidate
  142. low = middle + 1
  143. } else {
  144. high = middle - 1
  145. }
  146. }
  147. if (accepted === 0) {
  148. out.push(completeRun[offset] as DeltaEvent)
  149. offset += 1
  150. continue
  151. }
  152. /* v8 ignore next -- accepted is set only with its same-branch candidate. */
  153. out.push(acceptedRow ?? malformed(kind, 'bounded encoder lost its accepted row'))
  154. offset += accepted
  155. }
  156. out.push(...completeRun.slice(offset))
  157. }
  158. /**
  159. * Pack eligible logical chunk runs into bounded schema-18 records.
  160. * @param events - logical events in sequence order.
  161. * @returns scalar and packed physical records in equivalent order.
  162. */
  163. export function packChunkRuns(events: readonly SessionEvent[]): StorageRecord[] {
  164. const out: StorageRecord[] = []
  165. let kind: DeltaKind | undefined
  166. let run: DeltaEvent[] = []
  167. const flush = (): void => {
  168. if (kind === undefined) out.push(...run)
  169. else emitBoundedRun(out, kind, run)
  170. kind = undefined
  171. run = []
  172. }
  173. for (const event of events) {
  174. const nextKind = classify(event)
  175. if (nextKind === undefined) {
  176. flush()
  177. out.push(event)
  178. continue
  179. }
  180. const delta = event as DeltaEvent
  181. const previous = run.at(-1)
  182. if (nextKind === kind && previous !== undefined && continues(previous, delta, nextKind)) {
  183. run.push(delta)
  184. continue
  185. }
  186. flush()
  187. kind = nextKind
  188. run = [delta]
  189. }
  190. flush()
  191. return out
  192. }
  193. function malformed(tag: string, reason: string): never {
  194. throw new Error(`malformed ${tag} storage row: ${reason}`)
  195. }
  196. function validateRunData(
  197. tag: string,
  198. data: Record<string, unknown>,
  199. payloadKey: 'texts' | 'args',
  200. serializedBytes?: number,
  201. ): string[] {
  202. if (typeof data.turn !== 'number' || typeof data.step !== 'number' || typeof data.index !== 'number') {
  203. malformed(tag, 'turn/step/index must be numbers')
  204. }
  205. const payload = data[payloadKey]
  206. if (!Array.isArray(payload)
  207. || payload.length < MIN_PACKED_ROW_MEMBERS
  208. || payload.length > MAX_PACKED_ROW_MEMBERS
  209. || payload.some(member => typeof member !== 'string')) {
  210. malformed(tag, `${payloadKey} must contain ${MIN_PACKED_ROW_MEMBERS}..${MAX_PACKED_ROW_MEMBERS} strings`)
  211. }
  212. const gaps = data.dt
  213. if (!Array.isArray(gaps) || gaps.some(gap => !Number.isSafeInteger(gap))) {
  214. malformed(tag, 'dt must be an array of safe integers')
  215. }
  216. if (gaps.length !== payload.length - 1) malformed(tag, 'dt length must match the member count')
  217. if ((serializedBytes ?? Buffer.byteLength(JSON.stringify(data))) > MAX_PACKED_DATA_BYTES) {
  218. malformed(tag, `data exceeds ${MAX_PACKED_DATA_BYTES} UTF-8 bytes`)
  219. }
  220. return payload as string[]
  221. }
  222. function validateRow(
  223. value: Record<string, unknown>,
  224. tag: ChunkRow['type'],
  225. serializedBytes?: number,
  226. ): ChunkRow {
  227. if (!hasExactKeys(value, ['type', 'seq0', 'time0', 'data'])) malformed(tag, 'invalid envelope fields')
  228. if (!Number.isSafeInteger(value.seq0) || (value.seq0 as number) < 0) malformed(tag, 'seq0 must be non-negative')
  229. if (!Number.isSafeInteger(value.time0)) malformed(tag, 'time0 must be a safe integer')
  230. const data = value.data
  231. if (!isRecord(data)) malformed(tag, 'data must be an object')
  232. let payload: string[]
  233. if (tag === 'tool-call-chunks') {
  234. const withName = hasExactKeys(data, ['turn', 'step', 'index', 'id', 'name', 'dt', 'args'])
  235. if (!withName && !hasExactKeys(data, ['turn', 'step', 'index', 'id', 'dt', 'args'])) {
  236. malformed(tag, 'invalid tool-call data fields')
  237. }
  238. if (typeof data.id !== 'string' || (withName && typeof data.name !== 'string')) {
  239. malformed(tag, 'id and optional name must be strings')
  240. }
  241. payload = validateRunData(tag, data, 'args', serializedBytes)
  242. } else {
  243. if (!hasExactKeys(data, ['turn', 'step', 'index', 'dt', 'texts'])) malformed(tag, 'invalid text data fields')
  244. payload = validateRunData(tag, data, 'texts', serializedBytes)
  245. }
  246. if (!Number.isSafeInteger((value.seq0 as number) + payload.length - 1)) malformed(tag, 'member seqs exceed safe integers')
  247. let time = value.time0 as number
  248. for (const gap of data.dt as number[]) {
  249. time += gap
  250. if (!Number.isSafeInteger(time)) malformed(tag, 'member times exceed safe integers')
  251. }
  252. return value as unknown as ChunkRow
  253. }
  254. function expandRow(row: ChunkRow): SessionEvent[] {
  255. const members = row.type === 'tool-call-chunks' ? row.data.args : row.data.texts
  256. const events: SessionEvent[] = []
  257. let time = row.time0
  258. for (let index = 0; index < members.length; index += 1) {
  259. if (index > 0) time += row.data.dt[index - 1] as number
  260. let chunk: StreamChunk
  261. switch (row.type) {
  262. case 'text-chunks':
  263. chunk = { type: 'text-delta', index: row.data.index, text: members[index] as string }
  264. break
  265. case 'reasoning-chunks':
  266. chunk = { type: 'reasoning-delta', index: row.data.index, text: members[index] as string }
  267. break
  268. case 'tool-call-chunks':
  269. chunk = {
  270. type: 'tool-call-delta',
  271. index: row.data.index,
  272. id: row.data.id,
  273. ...Object.hasOwn(row.data, 'name') ? { name: row.data.name as string } : {},
  274. argumentsDelta: members[index] as string,
  275. }
  276. break
  277. }
  278. events.push({
  279. type: 'assistant/chunk',
  280. seq: row.seq0 + index,
  281. time,
  282. data: { turn: row.data.turn, step: row.data.step, chunk },
  283. })
  284. }
  285. return events
  286. }
  287. /**
  288. * Decode one scalar or packed schema-18 record.
  289. * @param value - parsed physical-record value.
  290. * @returns the represented logical events.
  291. */
  292. export function decodeStorageRecord(value: unknown): SessionEvent[] {
  293. if (!isRecord(value)) return [value as SessionEvent]
  294. const tag = value.type
  295. if (tag !== 'text-chunks' && tag !== 'reasoning-chunks' && tag !== 'tool-call-chunks') {
  296. return [value as SessionEvent]
  297. }
  298. return expandRow(validateRow(value, tag))
  299. }
  300. /**
  301. * Decode one packed row from its exact uncompressed data value. The byte bound
  302. * rejects oversized input before JSON parsing and avoids serializing it again.
  303. * @param tag - validated packed physical type.
  304. * @param seq0 - first represented logical sequence number.
  305. * @param time0 - first represented logical timestamp.
  306. * @param serializedData - decoded SQLite data-column text.
  307. * @returns the represented logical events.
  308. */
  309. export function decodeSerializedChunkRow(
  310. tag: ChunkRow['type'],
  311. seq0: number,
  312. time0: number,
  313. serializedData: string,
  314. ): SessionEvent[] {
  315. const bytes = Buffer.byteLength(serializedData)
  316. if (bytes > MAX_PACKED_DATA_BYTES) malformed(tag, `data exceeds ${MAX_PACKED_DATA_BYTES} UTF-8 bytes`)
  317. return expandRow(validateRow({ type: tag, seq0, time0, data: JSON.parse(serializedData) as unknown }, tag, bytes))
  318. }
  319. /* jscpd:ignore-end */