operations.ts 9.9 KB

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