index.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349
  1. /**
  2. * Model-facing `task_output`, `task_list`, and `task_kill` tools over
  3. * `ctx.tasks`. Loading the plugin attaches the control surface required by
  4. * producers. It also injects unreported completions as durable context for the
  5. * owner's next request; notices do not wake idle agents.
  6. * @module @deepseek-ai/dsh-tool-tasks
  7. */
  8. import type { Context } from 'cordis'
  9. import z from 'schemastery'
  10. import { boundContextSummary, createUserMessage, type ContentBlock } from '@deepseek-ai/dsh-llm'
  11. import { TextRetainer } from '@deepseek-ai/dsh-retention'
  12. import { defineTool } from '@deepseek-ai/dsh-tools'
  13. import type { GenericCallView, ToolDefinition, ToolExecution } from '@deepseek-ai/dsh-tools'
  14. import { TaskId } from '@deepseek-ai/dsh-tasks'
  15. import type { TaskSnapshot } from '@deepseek-ai/dsh-tasks'
  16. import type {} from '@deepseek-ai/dsh-system-prompt'
  17. export const name = 'tool-tasks'
  18. export const inject = ['tools', 'tasks', 'systemPrompt']
  19. /** Configures bounded `task_output` waits. */
  20. export interface Config {
  21. /** Wait duration applied when `task_output` sets `wait` without `timeout_ms` (default 30s). */
  22. waitTimeoutMs?: number
  23. /** Hard cap on any single wait; a larger model-supplied `timeout_ms` is clamped down to it (default 10min). */
  24. maxWaitTimeoutMs?: number
  25. }
  26. export const Config: z<Config> = z.object({
  27. waitTimeoutMs: z.number().min(1).default(30_000),
  28. maxWaitTimeoutMs: z.number().min(1).default(600_000),
  29. })
  30. /** Task state safe for model-authored programs; ownership/bookkeeping fields are omitted. */
  31. export interface PublicTaskSnapshot {
  32. id: string
  33. kind: string
  34. label: string
  35. status: TaskSnapshot['status']
  36. detail?: string
  37. startedAt: number
  38. finishedAt?: number
  39. }
  40. /** Shared schema for task-control outputs. */
  41. const PUBLIC_TASK_SCHEMA = {
  42. type: 'object',
  43. additionalProperties: false,
  44. properties: {
  45. id: { type: 'string', required: true },
  46. kind: { type: 'string', required: true },
  47. label: { type: 'string', required: true },
  48. status: {
  49. type: 'string',
  50. required: true,
  51. enum: ['running', 'stopping', 'completed', 'killed', 'failed'],
  52. },
  53. detail: { type: 'string' },
  54. startedAt: { type: 'integer', required: true },
  55. finishedAt: { type: 'integer' },
  56. },
  57. } as const
  58. /** Remove task ownership and notification bookkeeping from a registry snapshot. */
  59. function publicTask(snapshot: TaskSnapshot): PublicTaskSnapshot {
  60. return {
  61. id: snapshot.id,
  62. kind: snapshot.kind,
  63. label: snapshot.label,
  64. status: snapshot.status,
  65. ...snapshot.detail !== undefined ? { detail: snapshot.detail } : {},
  66. startedAt: snapshot.startedAt,
  67. ...snapshot.finishedAt !== undefined ? { finishedAt: snapshot.finishedAt } : {},
  68. }
  69. }
  70. /**
  71. * Render generic status with optional producer detail.
  72. * @param snapshot - task state to render.
  73. * @returns a bracketed status line.
  74. */
  75. export function statusLine(snapshot: Pick<TaskSnapshot, 'status' | 'detail'>): string {
  76. return snapshot.detail !== undefined
  77. ? `[status: ${snapshot.status}, ${snapshot.detail}]`
  78. : `[status: ${snapshot.status}]`
  79. }
  80. const encoder = new TextEncoder()
  81. function retainTail(text: string, maxBytes: number): string {
  82. const retainer = new TextRetainer({ kind: 'tail', maxBytes })
  83. retainer.push(text)
  84. return retainer.finish().text
  85. }
  86. function retainHead(text: string, maxBytes: number): string {
  87. const retainer = new TextRetainer({ kind: 'head', maxBytes })
  88. retainer.push(text)
  89. return retainer.finish().text
  90. }
  91. function fitWithSuffix(
  92. content: string,
  93. suffix: string,
  94. maxBytes: number | undefined,
  95. omitted: string,
  96. ): string {
  97. const complete = `${content}${suffix}`
  98. if (maxBytes === undefined || encoder.encode(complete).byteLength <= maxBytes) return complete
  99. const fixed = `${content.endsWith(omitted.trimStart()) ? '' : omitted}${suffix}`
  100. const fixedBytes = encoder.encode(fixed).byteLength
  101. if (fixedBytes >= maxBytes) return retainTail(fixed, maxBytes)
  102. return `${retainTail(content, maxBytes - fixedBytes)}${fixed}`
  103. }
  104. /**
  105. * One-line account of a settled task for the `notice` form's collapsed row.
  106. * @param snapshot - the settled task.
  107. * @returns its kind, label, and status, bounded like every notice summary.
  108. */
  109. function completionSummary(snapshot: TaskSnapshot): string {
  110. return boundContextSummary(`${snapshot.kind} ${snapshot.label} ${statusLine(snapshot)}`)
  111. }
  112. function fitCompletionNotice(snapshot: TaskSnapshot): string {
  113. const prefix = `background task ${snapshot.id}`
  114. const detail = ` (${snapshot.kind}: ${snapshot.label}) finished ${statusLine(snapshot)}`
  115. const action = '\nDone; task_output.'
  116. const complete = `${prefix}${detail}. Read its output with task_output.`
  117. const maxBytes = snapshot.outputLimitBytes
  118. if (maxBytes === undefined || encoder.encode(complete).byteLength <= maxBytes) return complete
  119. const omitted = '\n[notice truncated]'
  120. const fixed = `${prefix}${omitted}${action}`
  121. const fixedBytes = encoder.encode(fixed).byteLength
  122. if (fixedBytes <= maxBytes) {
  123. return fixedBytes === maxBytes
  124. ? fixed
  125. : `${prefix}${retainHead(detail, maxBytes - fixedBytes)}${omitted}${action}`
  126. }
  127. const compact = `${prefix}${action}`
  128. const compactBytes = encoder.encode(compact).byteLength
  129. if (compactBytes <= maxBytes) return compact
  130. const actionBytes = encoder.encode(action).byteLength
  131. if (actionBytes >= maxBytes) return retainTail(action, maxBytes)
  132. return `${retainHead(prefix, maxBytes - actionBytes)}${action}`
  133. }
  134. function rawSingleText(content: readonly ContentBlock[]): string | undefined {
  135. if (content.length !== 1) return undefined
  136. const block = content[0]
  137. if (block?.type !== 'text') return undefined
  138. return block.text
  139. }
  140. function boundSingleText(content: readonly ContentBlock[], maxBytes: number): ContentBlock[] | undefined {
  141. const text = rawSingleText(content)
  142. if (text === undefined) return undefined
  143. return [{
  144. type: 'text',
  145. text: fitWithSuffix(text, '', maxBytes, '\n[result truncated]'),
  146. }]
  147. }
  148. function visibleOutputLimit(ctx: Context, exec: ToolExecution): number | undefined {
  149. if (exec.name !== 'task_output' && exec.name !== 'task_kill') return undefined
  150. const taskId = (exec.arguments as { task_id?: unknown } | null | undefined)?.task_id
  151. if (typeof taskId !== 'string' || taskId.length === 0) return undefined
  152. return ctx.tasks.list(exec.agent).find(snapshot => snapshot.id === taskId)?.outputLimitBytes
  153. }
  154. /** Validate the non-empty constraint that ParameterSchemaSpec cannot express. */
  155. function validateTaskId(value: string): TaskId {
  156. if (value.length === 0) {
  157. throw new Error(`invalid task_id: expected a non-empty string, got ${JSON.stringify(value)}`)
  158. }
  159. return TaskId(value)
  160. }
  161. /** Pending presentation shared by the three generic task controls. */
  162. function presentTaskCall(title: string, kind: 'read' | 'execute', rawInput?: string): GenericCallView {
  163. return { card: 'generic', title, kind, ...rawInput !== undefined ? { rawInput } : {} }
  164. }
  165. export function apply(ctx: Context, config: Config): void {
  166. const waitDefault = config.waitTimeoutMs ?? 30_000
  167. const waitCap = config.maxWaitTimeoutMs ?? 600_000
  168. if (waitDefault > waitCap) {
  169. throw new Error(`tool-tasks: waitTimeoutMs (${waitDefault}) exceeds maxWaitTimeoutMs (${waitCap})`)
  170. }
  171. const outputLimits = new WeakMap<ToolExecution, number>()
  172. ctx.on('tools/pre-execute', (exec, next) => {
  173. const maxBytes = visibleOutputLimit(ctx, exec)
  174. if (maxBytes !== undefined) outputLimits.set(exec, maxBytes)
  175. return next()
  176. }, { prepend: true })
  177. const finalizeTaskContent: NonNullable<ToolDefinition['finalizeContent']> = (exec, result) => {
  178. const maxBytes = outputLimits.get(exec) ?? visibleOutputLimit(ctx, exec)
  179. outputLimits.delete(exec)
  180. if (maxBytes === undefined) return undefined
  181. if (exec.name === 'task_output' && !result.isError) {
  182. // This definition owns and schema-validates the canonical value. Preserve
  183. // its output/status split only while policy left the default rendering intact.
  184. const value = result.value as unknown as { text: string; task: PublicTaskSnapshot }
  185. const body = value.text.length > 0 ? value.text : '(no new output)'
  186. const content = body.endsWith('\n') ? body.slice(0, -1) : body
  187. const suffix = `\n${statusLine(value.task)}`
  188. if (rawSingleText(result.content) === `${content}${suffix}`) {
  189. return [{
  190. type: 'text',
  191. text: fitWithSuffix(content, suffix, maxBytes, '\n[output truncated]'),
  192. }]
  193. }
  194. }
  195. return boundSingleText(result.content, maxBytes)
  196. }
  197. // Producers may start work only while a control surface is attached.
  198. ctx.tasks.attachSurface('tool-tasks')
  199. // Cross-call guidance follows the bash section and precedes product sections.
  200. ctx.systemPrompt.section({
  201. name: 'tool:tasks',
  202. order: 106,
  203. text: 'Track every background task id you start. You are notified in-session when a task finishes — do not busy-poll or sleep on one; keep working on independent steps and do not duplicate a running task\'s work. Before giving a final answer, collect every still-relevant task with task_output (set wait: true only when you are genuinely blocked on it), and task_kill tasks that stopped mattering.',
  204. })
  205. // Use the exact lifecycle owner; reusable ids could resolve to a replacement.
  206. // Delivery targets the exact lifecycle owner. The notice waits in its
  207. // next-step inbox until another step claims it; disposal before that
  208. // boundary discards it with the owner.
  209. ctx.tasks.onTaskDone((snapshot, owner) => {
  210. if (snapshot.reported || owner === undefined) return
  211. owner.inject(createUserMessage({
  212. content: [{
  213. type: 'text',
  214. text: fitCompletionNotice(snapshot),
  215. }],
  216. source: {
  217. kind: 'plugin',
  218. plugin: 'tool-tasks',
  219. form: 'notice',
  220. summary: completionSummary(snapshot),
  221. },
  222. }))
  223. })
  224. ctx.tools.register(defineTool({
  225. name: 'task_output',
  226. description: 'Read a background task. Stream tasks return only output since the previous read; '
  227. + 'final-output tasks return their result after settlement. Every response ends with '
  228. + '`[status: ...]`. Reads are non-blocking unless `wait: true`, which waits up to the configured cap.',
  229. // A timed-out wait returns task state rather than a TOOL_TIMEOUT error, so
  230. // this tool owns its deadline instead of using ToolDefinition.timeoutMs.
  231. parameters: {
  232. task_id: { type: 'string', required: true, description: 'Task id returned by the tool that started the background work.' },
  233. wait: { type: 'boolean', description: 'Block until the task reaches a terminal status or the timeout expires. A timed-out wait returns [status: running] and leaves the task alive.' },
  234. timeout_ms: { type: 'number', description: 'Max wait in milliseconds (only meaningful with wait: true). Defaults to the configured wait timeout; capped by the configured maximum.' },
  235. },
  236. finalizeContent: finalizeTaskContent,
  237. output: {
  238. schema: {
  239. type: 'object',
  240. additionalProperties: false,
  241. properties: {
  242. text: { type: 'string', required: true },
  243. task: { ...PUBLIC_TASK_SCHEMA, required: true },
  244. },
  245. },
  246. render: (_args, value) => {
  247. const body = value.text.length > 0 ? value.text : '(no new output)'
  248. const separator = body.endsWith('\n') ? '' : '\n'
  249. return [{ type: 'text', text: `${body}${separator}${statusLine(value.task)}` }]
  250. },
  251. },
  252. async execute(args, exec) {
  253. const id = validateTaskId(args.task_id)
  254. if (args.wait === true) {
  255. const timeout = Math.min(args.timeout_ms ?? waitDefault, waitCap)
  256. await ctx.tasks.wait(id, timeout, exec.agent, exec.signal)
  257. }
  258. const read = ctx.tasks.read(id, exec.agent)
  259. return { text: read.text, task: publicTask(read.snapshot) }
  260. },
  261. presentCall: args => presentTaskCall(`Read output from background task ${args.task_id}`, 'read', args.task_id),
  262. }))
  263. ctx.tools.register(defineTool({
  264. name: 'task_list',
  265. description: 'List your background tasks (running and finished) with their ids, kinds, and statuses.',
  266. parameters: {},
  267. output: {
  268. schema: { type: 'array', items: PUBLIC_TASK_SCHEMA },
  269. render: (_args, tasks) => [{
  270. type: 'text',
  271. text: tasks.length === 0
  272. ? '(no background tasks)'
  273. : tasks.map(t => `${t.id} [${t.kind}] ${t.status} — ${t.label}`).join('\n'),
  274. }],
  275. },
  276. execute(_args, exec) {
  277. const tasks = ctx.tasks.list(exec.agent)
  278. return Promise.resolve(tasks.map(publicTask))
  279. },
  280. presentCall: () => presentTaskCall('List background tasks', 'read'),
  281. }))
  282. ctx.tools.register(defineTool({
  283. name: 'task_kill',
  284. description: 'Request cancellation of a running background task by task id. Returns immediately; the task settles as killed once its work actually stops.',
  285. parameters: {
  286. task_id: { type: 'string', required: true, description: 'Task id returned by the tool that started the background work.' },
  287. reason: { type: 'string', description: 'Optional short reason, recorded in the log and forwarded to the task.' },
  288. },
  289. finalizeContent: finalizeTaskContent,
  290. output: {
  291. schema: {
  292. type: 'object',
  293. additionalProperties: false,
  294. properties: {
  295. outcome: {
  296. type: 'string',
  297. required: true,
  298. enum: ['cancellation-requested', 'already-finished'],
  299. },
  300. task: { ...PUBLIC_TASK_SCHEMA, required: true },
  301. },
  302. },
  303. render: (_args, value) => [{
  304. type: 'text',
  305. text: value.outcome === 'already-finished'
  306. ? `task ${value.task.id} had already finished ${statusLine(value.task)}`
  307. : `requested cancellation of task ${value.task.id}`,
  308. }],
  309. },
  310. execute(args, exec) {
  311. const id = validateTaskId(args.task_id)
  312. const result = ctx.tasks.kill(id, exec.agent, args.reason)
  313. // A snapshot describes current state without consuming pending output.
  314. const snapshot = publicTask(ctx.tasks.get(id, exec.agent))
  315. return Promise.resolve({
  316. outcome: result === 'already-finished' ? 'already-finished' as const : 'cancellation-requested' as const,
  317. task: snapshot,
  318. })
  319. },
  320. presentCall: args => presentTaskCall(`Kill background task ${args.task_id}`, 'execute', args.task_id),
  321. }))
  322. }