index.ts 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359
  1. /**
  2. * Service Definition for combined session-history reads, traces, filters, and full-text search.
  3. *
  4. * @module @deepseek-ai/dsh-session-query
  5. */
  6. import { Context, Service } from '@deepseek-ai/cordis'
  7. import { Session, snapshotSessionEvent, type SessionId } from '@deepseek-ai/dsh-session'
  8. import { foldSessionTitle } from '@deepseek-ai/dsh-session-title'
  9. import type { SessionTitleSnapshot } from '@deepseek-ai/dsh-session-title'
  10. import type {
  11. SessionEventResultFilter,
  12. SessionEventSearchPage,
  13. SessionEventReadRequest,
  14. SessionEventRecord,
  15. SessionEventSearchDocument,
  16. SessionEventSearchRequest,
  17. SessionEventTraceObservation,
  18. SessionEventTraceRequest,
  19. SessionEventWindow,
  20. SessionLineageTrace,
  21. SessionLogSnapshot,
  22. SessionRecord,
  23. SessionResultFilter,
  24. SessionSearchExecContext,
  25. SessionSearchHit,
  26. SessionSearchPage,
  27. SessionSearchRequest,
  28. SessionSurfaceSnapshot,
  29. SessionTitleObservation,
  30. SessionTitleObservationResult,
  31. } from './types.ts'
  32. import {
  33. SESSION_QUERY_DEFAULT_PERSISTED_INSPECT_CONCURRENCY,
  34. SESSION_QUERY_READ_WINDOW_MAX,
  35. SessionQueryError,
  36. type Config,
  37. } from './config.ts'
  38. import { SessionCorpus } from './corpus.ts'
  39. import { buildSessionEventSearchDocuments } from './documents.ts'
  40. import {
  41. filterSessionEventDocuments,
  42. filterSessionResults,
  43. materializeSessionEventResultFilters,
  44. materializeSessionResultFilters,
  45. } from './filters.ts'
  46. import * as tracing from './tracing.ts'
  47. export type * from './types.ts'
  48. export { SessionSearchCursor } from './cursor.ts'
  49. export type { Config, SessionQueryErrorCode } from './config.ts'
  50. export {
  51. SESSION_QUERY_DEFAULT_PERSISTED_INSPECT_CONCURRENCY,
  52. SESSION_QUERY_READ_WINDOW_MAX,
  53. SessionQueryError,
  54. } from './config.ts'
  55. export { extractSessionEventText } from './extraction.ts'
  56. export { buildSessionEventRecords, buildSessionEventSearchDocuments } from './documents.ts'
  57. export {
  58. compileSessionTextFilter,
  59. filterSessionEventDocuments,
  60. filterSessionResults,
  61. materializeSessionEventResultFilters,
  62. materializeSessionResultFilters,
  63. } from './filters.ts'
  64. export { assertSessionHeadersCompatible } from './sources.ts'
  65. declare module '@deepseek-ai/cordis' {
  66. interface Context {
  67. sessionQuery: SessionQueryEngine
  68. }
  69. }
  70. /**
  71. * Unified live-preferred session query service.
  72. *
  73. * Exact reads, filters, and traces are backend-independent concrete behavior.
  74. * A backend implements full-text observation, reconciliation, ranking, cursor
  75. * generations, and query execution on the same `ctx.sessionQuery` service.
  76. */
  77. export abstract class SessionQueryEngine extends Service {
  78. static inject = ['sessions']
  79. private readonly _readWindowMax: number
  80. private readonly _corpus: SessionCorpus
  81. constructor(ctx: Context, config: Config = {}) {
  82. super(ctx, 'sessionQuery')
  83. this._readWindowMax = config.readWindowMax ?? SESSION_QUERY_READ_WINDOW_MAX
  84. if (!Number.isInteger(this._readWindowMax) || this._readWindowMax < 0) {
  85. throw new SessionQueryError(
  86. 'session-query: readWindowMax must be a non-negative integer',
  87. 'SESSION_QUERY_INVALID_CONFIG',
  88. )
  89. }
  90. const persistedInspectConcurrency = config.persistedInspectConcurrency
  91. ?? SESSION_QUERY_DEFAULT_PERSISTED_INSPECT_CONCURRENCY
  92. if (!Number.isSafeInteger(persistedInspectConcurrency) || persistedInspectConcurrency < 1) {
  93. throw new SessionQueryError(
  94. 'session-query: persistedInspectConcurrency must be a positive safe integer',
  95. 'SESSION_QUERY_INVALID_CONFIG',
  96. )
  97. }
  98. this._corpus = new SessionCorpus(ctx, persistedInspectConcurrency)
  99. }
  100. /**
  101. * Search the live-preferred logical corpus and group by session.
  102. * @param request - query text, metadata filters, page size, and cursor.
  103. * @param exec - optional cancellation control.
  104. * @returns session hits ranked by their strongest matching event.
  105. */
  106. abstract searchSessions(
  107. request: SessionSearchRequest,
  108. exec?: SessionSearchExecContext,
  109. ): Promise<SessionSearchPage<SessionSearchHit>>
  110. /**
  111. * Search events within one live-preferred logical session.
  112. * @param request - target session, query text, filters, page size, and cursor.
  113. * @param exec - optional cancellation control.
  114. * @returns matching event hits and their target header from one indexed generation.
  115. */
  116. abstract searchEvents(
  117. request: SessionEventSearchRequest,
  118. exec?: SessionSearchExecContext,
  119. ): Promise<SessionEventSearchPage>
  120. /**
  121. * List the complete logical corpus using live-preferred records.
  122. * @param signal - optional cancellation for persistence listing.
  123. * @returns deterministic newest-first cloned session records.
  124. */
  125. listSessions(signal?: AbortSignal): Promise<SessionRecord[]> {
  126. return this._corpus.listSessions(signal)
  127. }
  128. /**
  129. * Read and replay-validate one complete logical session log without making it live.
  130. * @param sessionId - live or persisted session id to read.
  131. * @returns cloned header and complete raw event log from one observation.
  132. * @throws when persistence, header compatibility, or replay validation fails.
  133. */
  134. async readSession(sessionId: SessionId): Promise<SessionLogSnapshot> {
  135. const loaded = await this._corpus.load(sessionId)
  136. Session.create(sessionId, loaded.events, loaded.header)
  137. return {
  138. session: structuredClone(loaded.header),
  139. events: loaded.events.map(snapshotSessionEvent),
  140. }
  141. }
  142. /**
  143. * Filter the complete logical corpus with provider-independent predicates.
  144. * @param filters - ANDed session metadata and availability clauses.
  145. * @param signal - optional cancellation for persistence listing.
  146. * @returns matching cloned records in deterministic newest-first order.
  147. */
  148. async filterSessions(
  149. filters: readonly SessionResultFilter[],
  150. signal?: AbortSignal,
  151. ): Promise<SessionRecord[]> {
  152. const ownedFilters = materializeSessionResultFilters(filters)
  153. return this._filterSessions(ownedFilters, signal)
  154. }
  155. /**
  156. * Fold the latest log-backed title from one live-preferred logical session.
  157. * @param sessionId - live or persisted session id to read.
  158. * @param signal - optional cancellation for source resolution and title folding.
  159. * @returns latest title snapshot, or `undefined` when the log has no title event.
  160. */
  161. async readTitle(
  162. sessionId: SessionId,
  163. signal?: AbortSignal,
  164. ): Promise<SessionTitleSnapshot | undefined> {
  165. return (await this.readTitleSnapshot(sessionId, signal)).title
  166. }
  167. /**
  168. * Fold the latest title and return its source header from one corpus observation.
  169. * @param sessionId - live or persisted session id to read.
  170. * @param signal - optional cancellation for source resolution and title folding.
  171. * @returns cloned source header and optional latest title snapshot.
  172. */
  173. async readTitleSnapshot(
  174. sessionId: SessionId,
  175. signal?: AbortSignal,
  176. ): Promise<SessionTitleObservation> {
  177. const result = (await this.readTitleSnapshots([sessionId], signal))[0] as SessionTitleObservationResult
  178. if (result.status === 'rejected') throw result.reason
  179. return result.value
  180. }
  181. /**
  182. * Fold titles for unique sessions from one cancellable corpus observation.
  183. *
  184. * Results preserve first-occurrence input order. Operational failures stay
  185. * isolated per session, while cancellation rejects the complete operation.
  186. * @param sessionIds - live or persisted session ids to observe.
  187. * @param signal - optional cancellation shared by all source reads.
  188. * @returns one fulfilled or rejected result per unique requested id.
  189. */
  190. async readTitleSnapshots(
  191. sessionIds: readonly SessionId[],
  192. signal?: AbortSignal,
  193. ): Promise<SessionTitleObservationResult[]> {
  194. return this._corpus.projectMany(sessionIds, (source): SessionTitleObservation => {
  195. const title = foldSessionTitle(source.events)
  196. return {
  197. session: structuredClone(source.header),
  198. ...title === undefined ? {} : { title },
  199. }
  200. }, signal)
  201. }
  202. /**
  203. * List lightweight raw-log event records for one logical session.
  204. * @param sessionId - live-preferred session id to read.
  205. * @returns event records in ascending seq order.
  206. */
  207. async listEvents(sessionId: SessionId): Promise<SessionEventRecord[]> {
  208. const loaded = await this._corpus.load(sessionId)
  209. return tracing.eventRecords(sessionId, loaded.events)
  210. }
  211. /**
  212. * Scan first-party semantic event documents with provider-independent filters.
  213. * @param sessionId - live-preferred session id to scan.
  214. * @param filters - ANDed metadata and literal-text predicates.
  215. * @returns matching semantic documents in ascending seq order.
  216. */
  217. async filterEvents(
  218. sessionId: SessionId,
  219. filters: readonly SessionEventResultFilter[],
  220. ): Promise<SessionEventSearchDocument[]> {
  221. const ownedFilters = materializeSessionEventResultFilters(filters)
  222. return this._filterEvents(sessionId, ownedFilters)
  223. }
  224. private async _filterSessions(
  225. filters: readonly SessionResultFilter[],
  226. signal?: AbortSignal,
  227. ): Promise<SessionRecord[]> {
  228. return filterSessionResults(await this._corpus.listSessions(signal), filters)
  229. }
  230. private async _filterEvents(
  231. sessionId: SessionId,
  232. filters: readonly SessionEventResultFilter[],
  233. ): Promise<SessionEventSearchDocument[]> {
  234. const loaded = await this._corpus.load(sessionId)
  235. const documents = buildSessionEventSearchDocuments(sessionId, loaded.events)
  236. return filterSessionEventDocuments(documents, filters)
  237. }
  238. /**
  239. * Read one session's complete current model surface from one corpus observation.
  240. * @param sessionId - live-preferred session id to read.
  241. * @returns cloned header, current surface, and the last sequence number included in the raw-log capture.
  242. * @throws when source resolution fails or the session surface is invalid.
  243. */
  244. async readSurface(sessionId: SessionId): Promise<SessionSurfaceSnapshot> {
  245. const loaded = await this._corpus.load(sessionId)
  246. return {
  247. session: structuredClone(loaded.header),
  248. capturedThroughSeq: loaded.events.at(-1)?.seq ?? null,
  249. events: tracing.currentSurfaceEvents(sessionId, loaded.events),
  250. }
  251. }
  252. /**
  253. * Trace known ancestry and descendants from one corpus observation.
  254. * @param sessionId - logical session id to trace.
  255. * @param signal - optional cancellation for persistence listing.
  256. * @returns a complete lineage or the first parent that could not be resolved.
  257. * @throws when corpus resolution fails, the target is absent, or its known ancestry cycles.
  258. */
  259. async traceSession(sessionId: SessionId, signal?: AbortSignal): Promise<SessionLineageTrace> {
  260. const records = await this._corpus.listSessions(signal)
  261. signal?.throwIfAborted()
  262. return tracing.traceSession(records, sessionId)
  263. }
  264. /**
  265. * Trace one event's direct positional replacements and cited source events.
  266. * @param request - target session id and event seq.
  267. * @param signal - optional cancellation for persisted source resolution.
  268. * @returns source header, direct links, and the target's positional replacement chain.
  269. * @throws when source resolution fails, the target is absent, or surface/source-event validation fails.
  270. */
  271. async traceEvent(request: SessionEventTraceRequest, signal?: AbortSignal): Promise<SessionEventTraceObservation> {
  272. const loaded = await this._corpus.load(request.sessionId, signal)
  273. signal?.throwIfAborted()
  274. return {
  275. session: loaded.header,
  276. ...tracing.traceEvent(request.sessionId, loaded.events, request.seq),
  277. }
  278. }
  279. /**
  280. * Read one full event plus a bounded raw-log context window.
  281. * @param request - target session/seq and context sizes.
  282. * @param signal - optional cancellation for persisted source resolution.
  283. * @returns cloned target and neighboring events.
  284. */
  285. async readEvent(request: SessionEventReadRequest, signal?: AbortSignal): Promise<SessionEventWindow> {
  286. const before = this._readWindow('before', request.before)
  287. const after = this._readWindow('after', request.after)
  288. const sessionId = request.sessionId
  289. const seq = request.seq
  290. return this._readEvent(sessionId, seq, before, after, signal)
  291. }
  292. private async _readEvent(
  293. sessionId: SessionId,
  294. seq: number,
  295. before: number,
  296. after: number,
  297. signal?: AbortSignal,
  298. ): Promise<SessionEventWindow> {
  299. const loaded = await this._corpus.load(sessionId, signal)
  300. signal?.throwIfAborted()
  301. const target = loaded.events[seq]
  302. if (target === undefined || target.seq !== seq) {
  303. throw new SessionQueryError(
  304. `session "${sessionId}" has no event at seq ${seq}`,
  305. 'SESSION_QUERY_EVENT_NOT_FOUND',
  306. )
  307. }
  308. const startSeq = Math.max(0, seq - before)
  309. const endSeq = Math.min(loaded.events.length - 1, seq + after)
  310. const targetSnapshot = snapshotSessionEvent(target)
  311. const events = loaded.events.slice(startSeq, endSeq + 1)
  312. .map(event => event === target
  313. ? targetSnapshot
  314. : snapshotSessionEvent(event))
  315. return {
  316. session: structuredClone(loaded.header),
  317. target: targetSnapshot,
  318. events,
  319. startSeq,
  320. endSeq,
  321. }
  322. }
  323. private _readWindow(name: 'before' | 'after', value: number | undefined): number {
  324. if (value === undefined) return 0
  325. if (!Number.isInteger(value) || value < 0 || value > this._readWindowMax) {
  326. throw new SessionQueryError(
  327. `${name} must be an integer between 0 and ${this._readWindowMax}`,
  328. 'SESSION_QUERY_INVALID_WINDOW',
  329. )
  330. }
  331. return value
  332. }
  333. }
  334. export default SessionQueryEngine