operations.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286
  1. /**
  2. * Tool operation orchestration over session-query service capabilities.
  3. *
  4. * @module @deepseek-ai/dsh-tool-session-query/operations
  5. */
  6. import type { Context } from '@deepseek-ai/cordis'
  7. import { HarnessError } from '@deepseek-ai/dsh-llm'
  8. import { SessionSeq } from '@deepseek-ai/dsh-session'
  9. import type { SessionId } from '@deepseek-ai/dsh-session'
  10. import {
  11. SessionQueryError,
  12. type SessionEventSearchPage,
  13. type SessionEventSurface,
  14. type SessionRecord,
  15. type SessionSearchCursor,
  16. } from '@deepseek-ai/dsh-session-query'
  17. import type { ToolRunContext } from '@deepseek-ai/dsh-tools'
  18. import { toolInput } from './input.ts'
  19. import { presentation } from './presentation.ts'
  20. import { serviceBoundary } from './service-boundary.ts'
  21. import { workspaceAccess } from './workspace-access.ts'
  22. type SessionSearchArgs = Parameters<typeof toolInput.buildSessionFilters>[0]
  23. interface EventSearchArgs {
  24. session_id?: string
  25. query: string
  26. seq_from?: number
  27. seq_to?: number
  28. time_from?: string
  29. time_to?: string
  30. event_types?: string[]
  31. surfaces?: SessionEventSurface[]
  32. }
  33. interface SessionTargetArgs {
  34. session_id?: string
  35. }
  36. interface EventTargetArgs extends SessionTargetArgs {
  37. seq: number
  38. }
  39. interface EventReadArgs extends EventTargetArgs {
  40. before?: number
  41. after?: number
  42. }
  43. interface SearchCollection<T> {
  44. readonly items: T[]
  45. readonly capped: boolean
  46. }
  47. async function executeSessionSearch(
  48. ctx: Context,
  49. args: SessionSearchArgs,
  50. exec: ToolRunContext,
  51. maxResults: number,
  52. ): Promise<string> {
  53. const caller = workspaceAccess.callerOf(exec, ctx)
  54. const cwd = caller.header.cwd
  55. if (cwd === undefined) {
  56. throw new HarnessError(
  57. 'cross-session search is unavailable because the caller session has no workspace',
  58. 'SESSION_QUERY_TOOL_UNAUTHORIZED',
  59. )
  60. }
  61. const query = toolInput.normalizeQuery(args.query)
  62. const sessionFilters = toolInput.buildSessionFilters(args)
  63. const eventFilters = toolInput.buildEventFilters({
  64. seqFrom: args.event_seq_from,
  65. seqTo: args.event_seq_to,
  66. timeFrom: args.event_time_from,
  67. timeTo: args.event_time_to,
  68. eventTypes: args.event_types,
  69. surfaces: args.event_surfaces,
  70. })
  71. const requestedParentIds = toolInput.materializeParentSessionIds(args.parent_session_ids)
  72. if (requestedParentIds !== undefined || args.include_root_sessions === true) {
  73. const authorizedParentIds = requestedParentIds === undefined
  74. ? new Set<SessionId>()
  75. : await workspaceAccess.authorizeSessionIds(ctx, caller, requestedParentIds, exec.signal)
  76. const parentValues: Array<SessionId | null> = requestedParentIds
  77. ?.filter(id => authorizedParentIds.has(id)) ?? []
  78. if (args.include_root_sessions === true) parentValues.push(null)
  79. if (parentValues.length === 0) return presentation.formatEmptySessionSearch()
  80. sessionFilters.push({ kind: 'parent', values: parentValues })
  81. }
  82. sessionFilters.push({ kind: 'cwd', values: [cwd] })
  83. const collected = await collectPages(
  84. maxResults,
  85. exec.signal,
  86. cursor => serviceBoundary.call(ctx, exec.signal, 'session search', () =>
  87. ctx.sessionQuery.searchSessions({
  88. query,
  89. sessionFilters,
  90. eventFilters,
  91. ...cursor === undefined ? {} : { cursor },
  92. }, { signal: exec.signal })),
  93. hit => hit.header.id !== caller.id && workspaceAccess.recordAuthorized(hit, caller),
  94. )
  95. const parentIds = collected.items
  96. .map(hit => hit.header.parentSession)
  97. .filter((id): id is SessionId => id !== undefined)
  98. const authorizedParents = await workspaceAccess.authorizeSessionIds(ctx, caller, parentIds, exec.signal)
  99. const titles = await workspaceAccess.readTitles(
  100. ctx,
  101. caller,
  102. collected.items.map(hit => hit.header.id),
  103. exec.signal,
  104. )
  105. return presentation.formatSessionSearch(collected, titles, authorizedParents)
  106. }
  107. async function executeEventSearch(
  108. ctx: Context,
  109. args: EventSearchArgs,
  110. exec: ToolRunContext,
  111. maxResults: number,
  112. ): Promise<string> {
  113. const caller = workspaceAccess.callerOf(exec, ctx)
  114. const sessionId = workspaceAccess.targetId(args, caller)
  115. await workspaceAccess.authorizeTarget(ctx, caller, sessionId, exec.signal)
  116. const query = toolInput.normalizeQuery(args.query)
  117. const range = toolInput.sequenceRange(args.seq_from, args.seq_to)
  118. if (sessionId === caller.id) {
  119. const stepStartSeq = caller.boundary?.lastStepStartSeq
  120. if (stepStartSeq === undefined) {
  121. throw new HarnessError(
  122. 'current-session search requires an active step boundary',
  123. 'SESSION_QUERY_TOOL_NO_CURRENT_STEP',
  124. )
  125. }
  126. // `null` (no step started yet) caps the range to a degenerate `to` the
  127. // filter validation rejects — the loop never runs tools outside a step.
  128. range.to = Math.min(range.to ?? Number.MAX_SAFE_INTEGER, (stepStartSeq ?? 0) - 1)
  129. }
  130. const title = await workspaceAccess.readTitle(ctx, caller, sessionId, exec.signal)
  131. if (range.from !== undefined && range.to !== undefined && range.from > range.to) {
  132. return presentation.formatEventSearch(sessionId, title, { items: [], capped: false })
  133. }
  134. const filters = toolInput.buildEventFilters({
  135. seqFrom: range.from,
  136. seqTo: range.to,
  137. timeFrom: args.time_from,
  138. timeTo: args.time_to,
  139. eventTypes: args.event_types,
  140. surfaces: args.surfaces,
  141. })
  142. const collected = await collectPages(
  143. maxResults,
  144. exec.signal,
  145. async (cursor): Promise<SessionEventSearchPage> => {
  146. const page = await serviceBoundary.call(ctx, exec.signal, 'event search', () =>
  147. ctx.sessionQuery.searchEvents({
  148. sessionId,
  149. query,
  150. filters,
  151. ...cursor === undefined ? {} : { cursor },
  152. }, { signal: exec.signal }))
  153. workspaceAccess.assertObservedTargetAuthorized(caller, sessionId, page.session)
  154. return page
  155. },
  156. () => true,
  157. )
  158. return presentation.formatEventSearch(sessionId, title, collected)
  159. }
  160. async function executeSessionTrace(
  161. ctx: Context,
  162. args: SessionTargetArgs,
  163. exec: ToolRunContext,
  164. ): Promise<string> {
  165. const caller = workspaceAccess.callerOf(exec, ctx)
  166. const sessionId = workspaceAccess.targetId(args, caller)
  167. await workspaceAccess.authorizeTarget(ctx, caller, sessionId, exec.signal)
  168. const trace = await serviceBoundary.call(ctx, exec.signal, 'session lineage trace', () =>
  169. ctx.sessionQuery.traceSession(sessionId, exec.signal))
  170. workspaceAccess.assertObservedTargetAuthorized(caller, sessionId, trace.target.header)
  171. const ancestors: SessionRecord[] = []
  172. let ancestorBoundary = false
  173. for (const ancestor of trace.ancestors) {
  174. if (!workspaceAccess.recordAuthorized(ancestor, caller)) {
  175. ancestorBoundary = true
  176. break
  177. }
  178. ancestors.push(ancestor)
  179. }
  180. if (ancestors.length === trace.ancestors.length && !trace.complete) ancestorBoundary = true
  181. const descendants = workspaceAccess.authorizeDescendants(trace.descendants, caller)
  182. const visibleIds = [
  183. trace.target.header.id,
  184. ...ancestors.map(record => record.header.id),
  185. ...workspaceAccess.descendantIds(descendants),
  186. ]
  187. const titles = await workspaceAccess.readTitles(ctx, caller, visibleIds, exec.signal)
  188. return presentation.formatSessionTrace(trace, ancestors, ancestorBoundary, descendants, titles)
  189. }
  190. async function executeEventTrace(
  191. ctx: Context,
  192. args: EventTargetArgs,
  193. exec: ToolRunContext,
  194. ): Promise<string> {
  195. toolInput.assertNonNegativeSafeInteger('seq', args.seq)
  196. const seq = SessionSeq(args.seq)
  197. const caller = workspaceAccess.callerOf(exec, ctx)
  198. const sessionId = workspaceAccess.targetId(args, caller)
  199. await workspaceAccess.authorizeTarget(ctx, caller, sessionId, exec.signal)
  200. const trace = await serviceBoundary.call(ctx, exec.signal, 'event trace', () =>
  201. ctx.sessionQuery.traceEvent({ sessionId, seq }, exec.signal))
  202. workspaceAccess.assertObservedTargetAuthorized(caller, sessionId, trace.session)
  203. const title = await workspaceAccess.readTitle(ctx, caller, sessionId, exec.signal)
  204. return presentation.formatEventTrace(sessionId, title, trace)
  205. }
  206. async function executeEventRead(
  207. ctx: Context,
  208. args: EventReadArgs,
  209. exec: ToolRunContext,
  210. ): Promise<string> {
  211. toolInput.assertNonNegativeSafeInteger('seq', args.seq)
  212. const seq = SessionSeq(args.seq)
  213. if (args.before !== undefined) toolInput.assertNonNegativeSafeInteger('before', args.before)
  214. if (args.after !== undefined) toolInput.assertNonNegativeSafeInteger('after', args.after)
  215. const caller = workspaceAccess.callerOf(exec, ctx)
  216. const sessionId = workspaceAccess.targetId(args, caller)
  217. await workspaceAccess.authorizeTarget(ctx, caller, sessionId, exec.signal)
  218. const window = await serviceBoundary.call(ctx, exec.signal, 'event read', () =>
  219. ctx.sessionQuery.readEvent({
  220. sessionId,
  221. seq,
  222. ...args.before === undefined ? {} : { before: args.before },
  223. ...args.after === undefined ? {} : { after: args.after },
  224. }, exec.signal))
  225. workspaceAccess.assertObservedTargetAuthorized(caller, sessionId, window.session)
  226. const title = await workspaceAccess.readTitle(ctx, caller, sessionId, exec.signal)
  227. return presentation.formatEventRead(sessionId, title, window)
  228. }
  229. async function collectPages<T>(
  230. maxResults: number,
  231. signal: AbortSignal,
  232. request: (cursor?: SessionSearchCursor) => Promise<{
  233. readonly items: readonly T[]
  234. readonly nextCursor?: SessionSearchCursor
  235. }>,
  236. accept: (item: T) => boolean,
  237. ): Promise<SearchCollection<T>> {
  238. const items: T[] = []
  239. const seen = new Set<SessionSearchCursor>()
  240. let cursor: SessionSearchCursor | undefined
  241. while (true) {
  242. signal.throwIfAborted()
  243. const page = await request(cursor)
  244. signal.throwIfAborted()
  245. for (const item of page.items) {
  246. if (!accept(item)) continue
  247. if (items.length === maxResults) {
  248. return { items, capped: true }
  249. }
  250. items.push(item)
  251. }
  252. if (page.nextCursor === undefined) return { items, capped: false }
  253. if (seen.has(page.nextCursor)) {
  254. throw new SessionQueryError(
  255. 'session-search provider repeated a continuation cursor',
  256. 'SESSION_QUERY_INVALID_CURSOR',
  257. )
  258. }
  259. seen.add(page.nextCursor)
  260. cursor = page.nextCursor
  261. }
  262. }
  263. /** Five model-facing session-query operation implementations. */
  264. export const operations = {
  265. executeSessionSearch,
  266. executeEventSearch,
  267. executeSessionTrace,
  268. executeEventTrace,
  269. executeEventRead,
  270. }