list.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389
  1. /** Cold-safe Session list and search projection. */
  2. import { stat } from 'node:fs/promises'
  3. import type { Context } from '@deepseek-ai/cordis'
  4. import type {} from '@deepseek-ai/dsh-agent-presets'
  5. import type { ImageAttachmentLimits } from '@deepseek-ai/dsh-attachment'
  6. import type { Session, SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
  7. import type {} from '@deepseek-ai/dsh-session-projection'
  8. import type {} from '@deepseek-ai/dsh-session-projection-cache'
  9. import { SessionQueryError, type SessionSearchCursor } from '@deepseek-ai/dsh-session-query'
  10. import { TypertRemoteFailure } from '@deepseek-ai/dsh-typert-protocol'
  11. import { z } from 'zod'
  12. import {
  13. SESSION_SEARCH_RESULT_LIMIT,
  14. SESSION_SEARCH_SNIPPET_MAX_CODE_POINTS,
  15. } from './types.ts'
  16. import type {
  17. SessionListMetadata, SessionProjectionHints, SessionProjectionValues, SessionSearchItem,
  18. SessionSearchValue, SessionSummary,
  19. } from './types.ts'
  20. /** Default maximum artifact size eligible for one cold projection observation. */
  21. export const DEFAULT_COLD_BLANK_PROBE_MAX_BYTES = 1024
  22. const COLD_SUMMARY_BATCH_SIZE = 16
  23. const SEARCH_PROVIDER_CALL_LIMIT = 100
  24. const SESSION_SEARCH_QUERY_MAX_CHARS = 500
  25. const MESSAGE_TYPES = new Set(['user/message', 'assistant/message'])
  26. const sessionListMetadataSchema: z.ZodType<SessionListMetadata> = z.object({
  27. blank: z.boolean(),
  28. lastPromptAt: z.number().nullable(),
  29. })
  30. const imageLimitsSchema = z.object({
  31. maxImageBytes: z.number().int().positive(),
  32. maxImagesPerMessage: z.number().int().positive(),
  33. maxMessageImageBytes: z.number().int().positive(),
  34. maxImagePixels: z.number().int().positive(),
  35. maxImageDimension: z.number().int().positive(),
  36. mediaTypes: z.array(z.string()),
  37. }) as unknown as z.ZodType<ImageAttachmentLimits>
  38. /**
  39. * Advance the Session-list metadata projection by one committed event.
  40. * @param state - metadata before the event.
  41. * @param event - next committed Session event.
  42. * @returns the original or advanced metadata value.
  43. */
  44. export function applySessionListMetadata(
  45. state: SessionListMetadata,
  46. event: SessionEvent,
  47. ): SessionListMetadata {
  48. const blank = state.blank && event.type !== 'turn/start'
  49. const lastPromptAt = event.type === 'user/message' && event.data.source.kind === 'user'
  50. ? event.time
  51. : state.lastPromptAt
  52. return blank === state.blank && lastPromptAt === state.lastPromptAt
  53. ? state
  54. : { blank, lastPromptAt }
  55. }
  56. /**
  57. * Return the longest prefix containing at most `maximum` Unicode code points.
  58. * @param value - source text.
  59. * @param maximum - maximum number of Unicode code points.
  60. * @returns the source text or its longest allowed prefix.
  61. */
  62. export function truncateUnicodeCodePoints(value: string, maximum: number): string {
  63. let count = 0
  64. let end = 0
  65. for (const codePoint of value) {
  66. if (count === maximum) return value.slice(0, end)
  67. count++
  68. end += codePoint.length
  69. }
  70. return value
  71. }
  72. /** Owns list projection registration, bounded cold summaries, and authorized search. */
  73. export class ApiSessionList {
  74. /**
  75. * @param ctx - Host context carrying Session, query, persistence, and projection services.
  76. * @param coldBlankProbeMaxBytes - maximum physical artifact size eligible for a full observation.
  77. */
  78. constructor(
  79. private readonly ctx: Context,
  80. private readonly coldBlankProbeMaxBytes: number,
  81. ) {
  82. ctx.inject(['sessionProjections'], (projectionCtx) => {
  83. projectionCtx.sessionProjections.register<'sessionListMetadata', SessionListMetadata>({
  84. key: 'sessionListMetadata',
  85. stateSchema: sessionListMetadataSchema,
  86. init: () => ({ blank: true, lastPromptAt: null }),
  87. apply: applySessionListMetadata,
  88. wire: { viewSchema: sessionListMetadataSchema, view: state => state },
  89. stateVersion: 1,
  90. })
  91. })
  92. ctx.inject(['sessionProjections', 'attachments'], (projectionCtx) => {
  93. projectionCtx.sessionProjections.register<'imageLimits', null>({
  94. key: 'imageLimits',
  95. stateSchema: z.null(),
  96. init: () => null,
  97. apply: state => state,
  98. wire: {
  99. viewSchema: imageLimitsSchema,
  100. view: () => projectionCtx.attachments.imageLimits,
  101. },
  102. stateVersion: 1,
  103. })
  104. })
  105. }
  106. /**
  107. * Build one current attached-Session summary.
  108. * @param session - attached Session to summarize.
  109. * @returns current list metadata and available projections.
  110. */
  111. summaryFor(session: Session): SessionSummary {
  112. const projections = this.projectionsFor(session.header, session)
  113. const metadata = projections?.values.sessionListMetadata
  114. return {
  115. sessionId: session.id,
  116. updatedAt: updatedAt(session.header, metadata),
  117. running: this.ctx.agents.get(session.id)?.status === 'running',
  118. blank: metadata?.blank ?? session.seq === 0,
  119. ...listFields(session.header),
  120. ...(projections === undefined ? {} : { projections }),
  121. }
  122. }
  123. /**
  124. * Read every visible attached and persisted Session without activating an Agent.
  125. * @param signal - optional cancellation for persistence reads.
  126. * @returns visible Session summaries ordered by activity.
  127. */
  128. async list(signal?: AbortSignal): Promise<SessionSummary[]> {
  129. signal?.throwIfAborted()
  130. const records = await this.ctx.sessionQuery.listSessions(signal)
  131. signal?.throwIfAborted()
  132. const items: SessionSummary[] = []
  133. const cold: SessionHeader[] = []
  134. for (const record of records) {
  135. const live = this.ctx.sessions.get(record.header.id)
  136. if (live !== undefined) {
  137. items.push(this.summaryFor(live))
  138. continue
  139. }
  140. if (record.header.cwd === undefined) continue
  141. cold.push(record.header)
  142. }
  143. for (let offset = 0; offset < cold.length; offset += COLD_SUMMARY_BATCH_SIZE) {
  144. const settled = await Promise.allSettled(cold.slice(offset, offset + COLD_SUMMARY_BATCH_SIZE)
  145. .map(header => this.summarizeCold(header, signal)))
  146. for (const result of settled) {
  147. if (result.status === 'rejected') throw result.reason
  148. items.push(result.value)
  149. }
  150. }
  151. items.sort((left, right) => right.updatedAt - left.updatedAt)
  152. return items
  153. }
  154. private async summarizeCold(
  155. header: SessionHeader,
  156. signal: AbortSignal | undefined,
  157. ): Promise<SessionSummary> {
  158. const cached = this.projectionsFor(header, undefined)
  159. const projections = cached?.values.sessionListMetadata?.blank === false
  160. ? cached
  161. : await this.probeSmallCold(header, signal) ?? cached
  162. const raced = this.ctx.sessions.get(header.id)
  163. if (raced !== undefined) return this.summaryFor(raced)
  164. const metadata = projections?.values.sessionListMetadata
  165. return {
  166. sessionId: header.id,
  167. updatedAt: updatedAt(header, metadata),
  168. running: false,
  169. // A large or inaccessible cache miss remains unknown and visible.
  170. blank: metadata?.blank ?? false,
  171. ...listFields(header),
  172. ...(projections === undefined ? {} : { projections }),
  173. }
  174. }
  175. private async probeSmallCold(
  176. header: SessionHeader,
  177. signal: AbortSignal | undefined,
  178. ): Promise<SessionProjectionHints | undefined> {
  179. if (this.coldBlankProbeMaxBytes === 0) return undefined
  180. const persistence = this.ctx.get('sessionPersistence')
  181. const location = persistence?.locate(header)
  182. if (location === undefined) return undefined
  183. signal?.throwIfAborted()
  184. try {
  185. if ((await stat(location.path)).size > this.coldBlankProbeMaxBytes) return undefined
  186. } catch {
  187. signal?.throwIfAborted()
  188. return undefined
  189. }
  190. try {
  191. using observation = await this.ctx.sessionQuery.observeSession(header.id, {
  192. ...(signal === undefined ? {} : { signal }),
  193. projectionMode: 'all',
  194. })
  195. const block = observation.projections
  196. return block === undefined
  197. ? undefined
  198. : { asOfSeq: block.asOfSeq, values: block.values as SessionProjectionValues }
  199. } catch (error: unknown) {
  200. signal?.throwIfAborted()
  201. this.ctx.logger.warn(
  202. `api-session.list: small cold observation for "${header.id}" failed; serving it as visible: ${String(error)}`,
  203. )
  204. return undefined
  205. }
  206. }
  207. /**
  208. * Search current visible message content without activating any matching Session.
  209. * @param query - literal message-content query.
  210. * @param signal - cancellation for list and search reads.
  211. * @returns authorized bounded Session search results.
  212. */
  213. async search(query: string, signal: AbortSignal): Promise<SessionSearchValue> {
  214. const normalizedQuery = normalizeSearchQuery(query)
  215. signal.throwIfAborted()
  216. const provider = this.ctx.get('sessionQuery')
  217. if (provider === undefined) {
  218. reject(
  219. 'internal',
  220. 'session search is unavailable: this deployment does not mount @deepseek-ai/dsh-session-query',
  221. {},
  222. )
  223. }
  224. try {
  225. const visible = await provider.listSessions(signal)
  226. signal.throwIfAborted()
  227. const visibleIds = new Set(visible
  228. .filter(record => record.header.cwd !== undefined)
  229. .map(record => record.header.id))
  230. if (visibleIds.size === 0) return { items: [], hasMore: false }
  231. const authorized: SessionSearchItem[] = []
  232. const acceptedIds = new Set<SessionId>()
  233. const seenCursors = new Set<SessionSearchCursor>()
  234. let cursor: SessionSearchCursor | undefined
  235. let providerCalls = 0
  236. let pageLimit = SESSION_SEARCH_RESULT_LIMIT
  237. while (authorized.length <= SESSION_SEARCH_RESULT_LIMIT) {
  238. signal.throwIfAborted()
  239. if (providerCalls >= SEARCH_PROVIDER_CALL_LIMIT) {
  240. throw new Error(`session search provider exceeded the ${SEARCH_PROVIDER_CALL_LIMIT}-call work budget`)
  241. }
  242. providerCalls++
  243. const requestedCursor = cursor
  244. const requestedLimit = pageLimit
  245. let page
  246. try {
  247. page = await provider.searchSessions({
  248. query: normalizedQuery,
  249. eventFilters: [
  250. { kind: 'type', values: ['user/message', 'assistant/message'] },
  251. { kind: 'surface', values: ['current'] },
  252. ],
  253. limit: requestedLimit,
  254. ...(requestedCursor === undefined ? {} : { cursor: requestedCursor }),
  255. }, { signal })
  256. signal.throwIfAborted()
  257. } catch (error: unknown) {
  258. signal.throwIfAborted()
  259. if (requestedCursor === undefined
  260. && error instanceof SessionQueryError
  261. && error.code === 'SESSION_QUERY_INVALID_LIMIT'
  262. && requestedLimit > 1) {
  263. pageLimit = Math.max(1, Math.floor(requestedLimit / 2))
  264. continue
  265. }
  266. if (requestedCursor !== undefined
  267. && error instanceof SessionQueryError
  268. && error.code === 'SESSION_QUERY_STALE_CURSOR') {
  269. authorized.length = 0
  270. acceptedIds.clear()
  271. seenCursors.clear()
  272. cursor = undefined
  273. continue
  274. }
  275. throw error
  276. }
  277. if (page.items.length > requestedLimit) {
  278. throw new Error(`session search provider returned ${String(page.items.length)} items; maximum is ${String(requestedLimit)}`)
  279. }
  280. for (const hit of page.items) {
  281. if (authorized.length > SESSION_SEARCH_RESULT_LIMIT) continue
  282. if (!visibleIds.has(hit.header.id)
  283. || hit.bestMatch.sessionId !== hit.header.id
  284. || hit.bestMatch.surface !== 'current'
  285. || !MESSAGE_TYPES.has(hit.bestMatch.type)
  286. || acceptedIds.has(hit.header.id)) continue
  287. acceptedIds.add(hit.header.id)
  288. authorized.push({
  289. sessionId: hit.header.id,
  290. snippet: truncateUnicodeCodePoints(hit.bestMatch.snippet, SESSION_SEARCH_SNIPPET_MAX_CODE_POINTS),
  291. })
  292. }
  293. if (page.nextCursor !== undefined) {
  294. if (seenCursors.has(page.nextCursor)) {
  295. throw new Error('session search provider repeated a continuation cursor')
  296. }
  297. seenCursors.add(page.nextCursor)
  298. }
  299. if (authorized.length > SESSION_SEARCH_RESULT_LIMIT || page.nextCursor === undefined) break
  300. cursor = page.nextCursor
  301. }
  302. return {
  303. items: authorized.slice(0, SESSION_SEARCH_RESULT_LIMIT),
  304. hasMore: authorized.length > SESSION_SEARCH_RESULT_LIMIT,
  305. }
  306. } catch (error: unknown) {
  307. signal.throwIfAborted()
  308. if (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED') {
  309. reject('cancelled', 'session search was aborted', {})
  310. }
  311. reject('internal', `session search failed: ${String(error)}`, {})
  312. }
  313. }
  314. private projectionsFor(
  315. header: SessionHeader,
  316. session: Session | undefined,
  317. ): SessionProjectionHints | undefined {
  318. try {
  319. const block = session === undefined
  320. ? this.ctx.get('sessionProjectionCache')?.cachedSnapshot(header)
  321. : this.ctx.get('sessionProjections')?.cachedSnapshot(session)
  322. return block !== undefined && Object.keys(block.values).length > 0
  323. ? {
  324. asOfSeq: block.asOfSeq,
  325. // Listing hints contain every currently cached wire value but remain
  326. // partial: missing cells and cache rows are never materialized here.
  327. values: block.values as SessionProjectionValues,
  328. }
  329. : undefined
  330. } catch (error) {
  331. this.ctx.logger.warn(
  332. `api-session.list: projection column for "${header.id}" failed; serving the row without it: ${String(error)}`,
  333. )
  334. return undefined
  335. }
  336. }
  337. }
  338. function normalizeSearchQuery(query: string): string {
  339. const normalized = query.trim()
  340. if (normalized.length === 0) {
  341. reject('bad-request', 'session search query must not be empty', {})
  342. }
  343. if (normalized.length > SESSION_SEARCH_QUERY_MAX_CHARS) {
  344. reject(
  345. 'bad-request',
  346. `session search query must contain at most ${SESSION_SEARCH_QUERY_MAX_CHARS} UTF-16 code units`,
  347. {},
  348. )
  349. }
  350. if (normalized.includes('\0')) {
  351. reject('bad-request', 'session search query must not contain NUL', {})
  352. }
  353. return normalized
  354. }
  355. function reject(code: string, message: string, details: object): never {
  356. throw new TypertRemoteFailure({ code, message, details })
  357. }
  358. function updatedAt(header: SessionHeader, metadata: SessionListMetadata | undefined): number {
  359. return Math.max(header.createdAt, metadata?.lastPromptAt ?? 0)
  360. }
  361. function listFields(header: SessionHeader): {
  362. readonly parentSessionId?: SessionId
  363. readonly origin?: 'subagent'
  364. readonly cwd?: string
  365. } {
  366. return {
  367. ...(header.parentSession === undefined ? {} : { parentSessionId: header.parentSession }),
  368. ...(header.origin === undefined ? {} : { origin: header.origin }),
  369. ...(header.cwd === undefined ? {} : { cwd: header.cwd }),
  370. }
  371. }