返回源码地图

packages/jobs/jobs-local/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 13 / 13;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Process-local provider for the background-job capability seam
3 * (`ctx.jobs`). It keeps every job — lifecycle state, the bounded output
4 * ring, and the model cursor — in memory and hands out fresh projections and
5 * chunk copies, never live state.
6 *
7 * Registrations outlive producer and controller fibers. Agent or service
8 * disposal cancels live work and awaits compliant producers; a throwing
9 * teardown cancel force-fails only the record and reports a possible orphan.
10 * @module @deepseek-ai/dsh-jobs-local
11 */
12
13import { Context } from '@deepseek-ai/cordis'
14import z from '@deepseek-ai/schemastery'
15import type { Agent } from '@deepseek-ai/dsh-agent'
16import { ScopedLayers, scopeOf } from '@deepseek-ai/dsh-scope'
17import type { SessionId } from '@deepseek-ai/dsh-session'
18import { deadline, timeoutOf } from '@deepseek-ai/dsh-timeout'
19import { JobRegistry, JobId } from '@deepseek-ai/dsh-jobs'
20import type {
21 JobAppendOptions, JobEvent, JobEvents, JobHandle, JobKind, JobOutcome, JobOutputRead, JobOutputSource,
22 JobRead, JobSettleCause, JobSpec, JobStatus, JobView,
23} from '@deepseek-ai/dsh-jobs'
24import { JobEventHub, JobLayer } from './events.ts'
25import { startPump } from './pump.ts'
26import type { PumpHandle } from './pump.ts'
27import { OutputRing } from './ring.ts'
28
29/** Timeout code that distinguishes a bounded wait from caller cancellation. */
30export const TASK_WAIT_TIMEOUT = 'TASK_WAIT_TIMEOUT'
31
32/** Default maximum number of active jobs in one exact-owner bucket. */
33const DEFAULT_MAX_CONCURRENT_JOBS_PER_OWNER = 10
34
35/** Default live ring retention per job, in UTF-8 bytes. */
36const DEFAULT_RETAIN_BYTES = 256 * 1024
37
38/** Default ring retention kept after settlement, in UTF-8 bytes. */
39const DEFAULT_SETTLED_RETAIN_BYTES = 16 * 1024
40
41/** Default poll interval for pull sources, in milliseconds. */
42const DEFAULT_PUMP_POLL_MS = 150
43
44/** Configuration for the process-local job registry. */
45export interface Config {
46 /**
47 * Maximum `running` plus `stopping` jobs per exact owner or in the shared unowned bucket;
48 * omission defaults to 10.
49 */
50 maxConcurrentJobsPerOwner?: number
51 /** Live ring retention per job in UTF-8 bytes; omission defaults to 262144. */
52 retainBytes?: number
53 /**
54 * Ring retention kept after a job settles, in UTF-8 bytes; omission defaults to 16384.
55 * Settlement keeps every byte the model cursor has not consumed on top of
56 * this cap; the first terminal model read then trims to it.
57 */
58 settledRetainBytes?: number
59 /** Poll interval for a job's pull sources, in milliseconds; omission defaults to 150. */
60 pumpPollMs?: number
61}
62
63/**
64 * Producer-written state shared between the {@link JobHandle} and the
65 * registered record: the starter call writes through it before the commit,
66 * the same object serves the job for its whole life afterwards.
67 */
68interface ProducerState {
69 /** Live progress line until settlement clears it. */
70 progress: string | undefined
71 /** The committed registry record; undefined exactly during the starter call. */
72 job: TrackedJob | undefined
73}
74
75/** The registry's mutable per-job record (never handed out — see {@link LocalJobRegistry.view}). */
76interface TrackedJob {
77 id: JobId
78 kind: JobKind
79 label: string
80 outputLimitBytes: number | undefined
81 /** Exact lifecycle owner; session-id authorization is derived from it. */
82 owner: Agent | undefined
83 cancel: (reason?: string) => void
84 status: JobStatus
85 ring: OutputRing
86 /** The model's consuming cursor; {@link JobRegistry.readAt} never moves it. */
87 modelCursor: number
88 /** Whether the first post-settlement read already handed out `result`. */
89 resultDelivered: boolean
90 /** Producer-shared progress line and commit binding. */
91 state: ProducerState
92 /** Terminal reason; a recorded kill reason is merged in at settlement. */
93 detail: string | undefined
94 result: string | undefined
95 startedAt: number
96 finishedAt: number | undefined
97 /** Reason recorded by {@link JobRegistry.kill}, merged into a `killed` settlement's detail. */
98 killReason: string | undefined
99 /** Set once a kill or teardown cancel ran; settlement reports it as the cause. */
100 settleCause: JobSettleCause | undefined
101 /** Resolves once the terminal record is committed and announced. */
102 settled: Promise<void>
103 /** Resolver for {@link settled}, called by the first effective settlement. */
104 markSettled: () => void
105 /** Removable resolvers for live waits; timeout/abort unregister before the job settles. */
106 waitResolvers: Set<() => void>
107 /** The registry-owned pump over the spec's pull sources, when it named any. */
108 pump: PumpHandle | undefined
109 /**
110 * The spill file each pull source reported on its latest read, by source
111 * index; an entry is undefined while that source keeps none. Source
112 * metadata rather than per-chunk metadata, so it survives ring eviction and
113 * follows a source that withdraws its file.
114 */
115 spillPaths: (string | undefined)[]
116}
117
118/** True for the three terminal {@link JobStatus} values. */
119function isTerminal(status: JobStatus): boolean {
120 return status === 'completed' || status === 'killed' || status === 'failed'
121}
122
123/**
124 * The in-memory `jobs` registry. See the Service Definition contract in
125 * `@deepseek-ai/dsh-jobs` for the ownership, isolation, and lifecycle
126 * semantics this implementation honors.
127 */
128export class LocalJobRegistry extends JobRegistry {
129 static Config: z<Config> = z.object({
130 maxConcurrentJobsPerOwner: z.number()
131 .step(1)
132 .min(1)
133 .max(Number.MAX_SAFE_INTEGER)
134 .default(DEFAULT_MAX_CONCURRENT_JOBS_PER_OWNER),
135 retainBytes: z.number()
136 .step(1)
137 .min(1)
138 .max(Number.MAX_SAFE_INTEGER)
139 .default(DEFAULT_RETAIN_BYTES),
140 settledRetainBytes: z.number()
141 .step(1)
142 .min(1)
143 .max(Number.MAX_SAFE_INTEGER)
144 .default(DEFAULT_SETTLED_RETAIN_BYTES),
145 pumpPollMs: z.number()
146 .step(1)
147 .min(1)
148 .max(Number.MAX_SAFE_INTEGER)
149 .default(DEFAULT_PUMP_POLL_MS),
150 })
151
152 /** Schemastery-defaulted active-job limit. */
153 private readonly maxConcurrentJobsPerOwner: number
154 /** Schemastery-defaulted live ring retention cap. */
155 private readonly retainBytes: number
156 /** Schemastery-defaulted settled ring retention cap. */
157 private readonly settledRetainBytes: number
158 /** Schemastery-defaulted pull-source poll interval. */
159 private readonly pumpPollMs: number
160 private store = new Map<JobId, TrackedJob>()
161 private counters = new Map<string, number>()
162 /**
163 * Controllers and scoped subscriptions layered by the scope that registered
164 * them, in the tools-registry shape: a contribution files into its
165 * registering context's scope, and a read unions the global layer with the
166 * owner's scope chain.
167 *
168 * The registry is one process-wide instance serving every composition, so a
169 * flat table would answer a per-owner question process-wide: one preset's
170 * job controls would hold `start()` open for an agent whose own composition
171 * loads none, and one settlement would reach every preset's notice listener.
172 * Layers make both reads owner-relative. Nothing derives a cache from a
173 * layer, so change notification is a no-op.
174 */
175 private readonly layers = new ScopedLayers<JobLayer>(() => new JobLayer(), () => {})
176 private readonly hub: JobEventHub
177 /** Owner agents with attached scope cleanup, mapped to the exact disposer. */
178 private ownerCleanups = new Map<Agent, () => Promise<void> | void>()
179 /** Service context used by detached settlement continuations and teardown. */
180 private readonly selfCtx: Context
181
182 constructor(ctx: Context, config: Config) {
183 super(ctx)
184 // Schemastery validates and fills the defaults before constructing the service.
185 const resolved = config as Required<Config>
186 this.maxConcurrentJobsPerOwner = resolved.maxConcurrentJobsPerOwner
187 this.retainBytes = resolved.retainBytes
188 this.settledRetainBytes = resolved.settledRetainBytes
189 this.pumpPollMs = resolved.pumpPollMs
190 this.selfCtx = ctx
191 this.hub = new JobEventHub(this.layers, (message) => { ctx.logger.warn(message) })
192 ctx.effect(() => () => this.disposeAll(), 'jobs teardown')
193 }
194
195 /**
196 * The event stream bound to the accessing context: a subscription is an
197 * effect of that context, and `{ owners: 'scope' }` names its scope.
198 */
199 get events(): JobEvents {
200 const registrar = this.ctx
201 return {
202 subscribe: (filter, listener) => this.hub.subscribe(registrar, filter, listener),
203 }
204 }
205
206 start(spec: JobSpec): JobId {
207 const owner = this.resolveOwner(spec.owner)
208 if (!this.servesOwner(owner)) {
209 throw new Error('background jobs unavailable: no job controller serves this agent (load @deepseek-ai/dsh-tool-jobs in its composition)')
210 }
211 if (spec.kind.length === 0) throw new Error('invalid job kind: expected a non-empty string')
212 if (spec.label.length === 0) throw new Error('invalid job label: expected a non-empty string')
213 if (spec.outputLimitBytes !== undefined
214 && (!Number.isSafeInteger(spec.outputLimitBytes) || spec.outputLimitBytes <= 0)) {
215 throw new Error(`invalid outputLimitBytes: expected a positive safe integer, got ${JSON.stringify(spec.outputLimitBytes)}`)
216 }
217 if (owner !== undefined) this.ensureOwnerCleanup(owner)
218
219 const active = this.activeJobCount(owner)
220 if (active >= this.maxConcurrentJobsPerOwner) {
221 throw new Error(
222 `background job limit reached for this owner (limit: ${this.maxConcurrentJobsPerOwner}); use job_kill to stop an unneeded job, wait for it to finish, then retry`,
223 )
224 }
225
226 // The id is issued before the starter runs so the producer face can carry
227 // it; a throwing starter still leaves nothing registered — its ordinal is
228 // simply skipped.
229 const count = (this.counters.get(spec.kind) ?? 0) + 1
230 this.counters.set(spec.kind, count)
231 const id = JobId(`${spec.kind}-${count}`)
232 const ring = new OutputRing()
233 const state: ProducerState = { progress: undefined, job: undefined }
234 const handle: JobHandle = {
235 id,
236 append: (text, options) => { this.appendRing(state, ring, text, options, 'producer') },
237 updateProgress: (line) => { this.updateProgress(state, line) },
238 }
239 const hooks = spec.run(handle)
240
241 let markSettled!: () => void
242 const settled = new Promise<void>((resolve) => { markSettled = resolve })
243 const job: TrackedJob = {
244 id,
245 kind: spec.kind,
246 label: spec.label,
247 outputLimitBytes: spec.outputLimitBytes,
248 owner,
249 cancel: hooks.cancel.bind(hooks),
250 status: 'running',
251 ring,
252 modelCursor: 0,
253 resultDelivered: false,
254 state,
255 detail: undefined,
256 result: undefined,
257 startedAt: Date.now(),
258 finishedAt: undefined,
259 killReason: undefined,
260 settleCause: undefined,
261 settled,
262 markSettled,
263 waitResolvers: new Set(),
264 pump: undefined,
265 spillPaths: [],
266 }
267 // Binding the shared producer state is the commit: writes staged inside
268 // the starter are already in the ring and `state`, and every later handle
269 // call reaches the registered record for its terminal checks and signals.
270 state.job = job
271 this.store.set(id, job)
272 // Registration is complete and cannot fail from here, so the visible set
273 // has genuinely changed. The announcement precedes the pump because the
274 // pump drains its sources once synchronously, and that drain may append
275 // and announce output: a job's first event is always `registered`.
276 this.emit({ type: 'registered', job: this.view(job) }, owner)
277
278 // The producer's settlement or a registry-forced one ends the pump; the
279 // pump's final drain then lands before this registry trims the ring.
280 const producerDone = hooks.done.then(
281 outcome => outcome,
282 (error: unknown): JobOutcome => {
283 // Contain a producer contract violation (`done` rejected) so cleanup and waiters cannot hang.
284 this.selfCtx.logger.warn(`jobs: job ${job.id} producer done promise rejected (producer contract violation): ${String(error)}`)
285 return { status: 'failed', detail: String(error) }
286 },
287 )
288 if (spec.output !== undefined && spec.output.length > 0) {
289 job.pump = startPump(
290 spec.output.map(source => this.guardSource(job, source)),
291 {
292 append: (text, options) => { this.appendRing(state, ring, text, options, 'pump') },
293 spill: (index, path) => { job.spillPaths[index] = path },
294 },
295 this.pumpPollMs,
296 Promise.race([producerDone, settled]),
297 )
298 }
299 void producerDone.then(async (outcome) => {
300 if (job.pump !== undefined) await job.pump.done
301 this.settle(job, outcome, job.settleCause ?? 'producer')
302 })
303 return id
304 }
305
306 list(caller?: SessionId): JobView[] {
307 return [...this.store.values()]
308 .filter(job => job.owner === undefined || job.owner.id === caller)
309 .map(job => this.view(job))
310 }
311
312 get(id: JobId, caller?: SessionId): JobView {
313 return this.view(this.expect(id, caller))
314 }
315
316 read(id: JobId, caller?: SessionId): JobRead {
317 return this.readJob(this.expect(id, caller))
318 }
319
320 readAt(id: JobId, from: number, caller?: SessionId): JobOutputRead {
321 const job = this.expect(id, caller)
322 if (!Number.isSafeInteger(from) || from < 0) {
323 throw new Error(`invalid output read offset: expected a non-negative safe integer, got ${JSON.stringify(from)}`)
324 }
325 return job.ring.readFrom(from)
326 }
327
328 kill(id: JobId, caller?: SessionId, reason?: string): 'requested' | 'already-finished' {
329 return this.killJob(this.expect(id, caller), reason)
330 }
331
332 async wait(id: JobId, timeoutMs: number, caller?: SessionId, signal?: AbortSignal): Promise<JobView> {
333 return this.waitJob(this.expect(id, caller), timeoutMs, signal)
334 }
335
336 remove(id: JobId, caller?: SessionId): void {
337 const job = this.expect(id, caller)
338 if (!isTerminal(job.status)) throw new Error(`job ${id} is still ${job.status}; kill it and wait for settlement before removing it`)
339 this.drop([job])
340 }
341
342 attachController(name: string): () => void {
343 // One token per call keeps duplicate labels independently disposable.
344 const token = Symbol(name)
345 return this.layers.effect(
346 this.ctx,
347 layer => layer.controllers.append(token),
348 { label: 'jobs.attachController()' },
349 )
350 }
351
352 /**
353 * Resolve a spec's owner session to its live Agent. An owned registration
354 * needs the agent registry, and the session must currently have a live
355 * instance: that instance's disposal is what cancels and drops the job.
356 */
357 private resolveOwner(session: SessionId | undefined): Agent | undefined {
358 if (session === undefined) return undefined
359 const agents = this.selfCtx.get('agents')
360 if (agents === undefined) {
361 throw new Error('background job ownership requires the agent registry (load @deepseek-ai/dsh-agent)')
362 }
363 const owner = agents.get(session)
364 if (owner === undefined) {
365 throw new Error(`session "${session}" has no live agent (background job owner must be live)`)
366 }
367 return owner
368 }
369
370 /**
371 * Whether an attached job controller can collect and stop work owned by
372 * `owner`. The global layer holds every controller attached from an unscoped
373 * context — a host composition's own controls — and therefore serves every
374 * owner; a scoped controller serves exactly the agents composed under it.
375 * @param owner - the job's owner, or undefined for unowned work.
376 * @returns whether some reachable controller serves the owner.
377 */
378 private servesOwner(owner?: Agent): boolean {
379 if (!this.layers.global.controllers.isEmpty()) return true
380 return this.layers.chainLayers(owner === undefined ? undefined : scopeOf(owner.ctx))
381 .some(layer => !layer.controllers.isEmpty())
382 }
383
384 /** Count authoritative active records for one exact owner or the shared unowned bucket. */
385 private activeJobCount(owner: Agent | undefined): number {
386 let count = 0
387 for (const job of this.store.values()) {
388 if (job.owner === owner && (job.status === 'running' || job.status === 'stopping')) count += 1
389 }
390 return count
391 }
392
393 /** Look up a job and enforce caller access. */
394 private expect(id: JobId, caller?: SessionId): TrackedJob {
395 const job = this.store.get(id)
396 if (job === undefined) throw new Error(`unknown job ${id}`)
397 this.assertAccess(job, caller)
398 return job
399 }
400
401 /**
402 * The isolation fence: a job with an owner is reachable only by callers
403 * whose session id matches (`!== undefined` semantics — an unowned job is
404 * open, and a caller-less view can never match an owned one).
405 */
406 private assertAccess(job: TrackedJob, caller: SessionId | undefined): void {
407 if (job.owner !== undefined && job.owner.id !== caller) {
408 throw new Error(`job ${job.id} belongs to another session`)
409 }
410 }
411
412 /** Project a fresh read-only view from the mutable record. */
413 private view(job: TrackedJob): JobView {
414 const owner = job.owner?.id
415 const spillPaths = [...new Set(job.spillPaths.filter((path): path is string => path !== undefined))]
416 return {
417 id: job.id,
418 kind: job.kind,
419 label: job.label,
420 ...owner !== undefined ? { owner } : {},
421 ...job.outputLimitBytes !== undefined ? { outputLimitBytes: job.outputLimitBytes } : {},
422 status: job.status,
423 ...job.state.progress !== undefined ? { progress: job.state.progress } : {},
424 ...job.detail !== undefined ? { detail: job.detail } : {},
425 startedAt: job.startedAt,
426 ...job.finishedAt !== undefined ? { finishedAt: job.finishedAt } : {},
427 output: {
428 total: job.ring.total,
429 earliest: job.ring.earliest,
430 ...spillPaths.length > 0 ? { spillPaths } : {},
431 },
432 }
433 }
434
435 private emit(event: JobEvent, owner: Agent | undefined): void {
436 this.hub.emit(event, owner)
437 }
438
439 /**
440 * Consume the ring from the model cursor; the result rides the first read
441 * after settlement. A terminal read is the point the settled stream drops
442 * to the settled cap: settlement kept every unconsumed byte for it.
443 */
444 private readJob(job: TrackedJob): JobRead {
445 const read = job.ring.readFrom(job.modelCursor)
446 job.modelCursor = job.ring.total
447 const result = isTerminal(job.status) && !job.resultDelivered ? job.result : undefined
448 if (result !== undefined) job.resultDelivered = true
449 if (isTerminal(job.status)) job.ring.trim(this.settledRetainBytes)
450 return {
451 chunks: read.chunks,
452 lossy: read.lossy,
453 ...result !== undefined ? { result } : {},
454 job: this.view(job),
455 }
456 }
457
458 private killJob(job: TrackedJob, reason?: string): 'requested' | 'already-finished' {
459 if (isTerminal(job.status)) return 'already-finished'
460 // Cancel first so a throw leaves lifecycle state unchanged.
461 job.cancel(reason)
462 job.status = 'stopping'
463 // Last writer wins on purpose: the detail reports the latest kill intent.
464 if (reason !== undefined) job.killReason = reason
465 job.settleCause = 'kill'
466 this.emit({ type: 'stopping', job: this.view(job) }, job.owner)
467 return 'requested'
468 }
469
470 private async waitJob(job: TrackedJob, timeoutMs: number, signal?: AbortSignal): Promise<JobView> {
471 if (!Number.isFinite(timeoutMs) || timeoutMs <= 0) {
472 throw new Error(`invalid wait timeout: expected a positive number of milliseconds, got ${JSON.stringify(timeoutMs)}`)
473 }
474 if (!isTerminal(job.status)) {
475 if (signal?.aborted) throw new Error('wait aborted')
476 // The scoped deadline distinguishes a successful wait timeout from
477 // caller cancellation and clears its timer on every exit.
478 using d = deadline(signal, timeoutMs, TASK_WAIT_TIMEOUT)
479 await new Promise<void>((resolve, reject) => {
480 const onSettled = (): void => {
481 job.waitResolvers.delete(onSettled)
482 d.signal.removeEventListener('abort', onAbort)
483 resolve()
484 }
485 const onAbort = (): void => {
486 job.waitResolvers.delete(onSettled)
487 // A settled job cannot reach here: settlement releases every waiter
488 // before it announces, and each released waiter detaches this
489 // listener in the same synchronous span, so nothing that reacts to a
490 // settlement can abort a wait the settlement already owed.
491 if (timeoutOf(d.signal, TASK_WAIT_TIMEOUT) !== undefined) {
492 resolve()
493 } else {
494 reject(new Error('wait aborted'))
495 }
496 }
497 job.waitResolvers.add(onSettled)
498 d.signal.addEventListener('abort', onAbort, { once: true })
499 })
500 }
501 return this.view(job)
502 }
503
504 /**
505 * Append one chunk to the ring. A producer chunk against a settled job is
506 * logged and dropped; the registry's own pump drains silently after
507 * settlement (a forced settlement may precede the producer's). A chunk
508 * staged inside the starter call is retained and signals no observer — the
509 * registration commit publishes it.
510 */
511 private appendRing(
512 state: ProducerState,
513 ring: OutputRing,
514 text: string,
515 options: JobAppendOptions | undefined,
516 writer: 'producer' | 'pump',
517 ): void {
518 const job = state.job
519 if (job !== undefined && isTerminal(job.status)) {
520 if (writer === 'producer') this.selfCtx.logger.warn(`jobs: append to settled job ${job.id} dropped`)
521 return
522 }
523 if (!ring.append(text, options, this.retainBytes)) return
524 if (job !== undefined) this.emitOutput(job)
525 }
526
527 /**
528 * Contain a failing pull source: the first throw is logged, and the source
529 * reads as exhausted from then on, so the job runs to its own settlement
530 * with whatever the ring holds instead of freezing on a pump failure.
531 */
532 private guardSource(job: TrackedJob, source: JobOutputSource): JobOutputSource {
533 let failed = false
534 return {
535 ...source.channel !== undefined ? { channel: source.channel } : {},
536 read: (fromByte) => {
537 if (failed) return { text: '', nextOffset: fromByte, lossy: false }
538 try {
539 return source.read(fromByte)
540 } catch (error: unknown) {
541 failed = true
542 this.selfCtx.logger.warn(`jobs: output source for ${job.id} failed; its stream stops here: ${String(error)}`)
543 return { text: '', nextOffset: fromByte, lossy: false }
544 }
545 },
546 }
547 }
548
549 /** Announce that one job's ring advanced (append or settlement). */
550 private emitOutput(job: TrackedJob): void {
551 const owner = job.owner?.id
552 this.emit({ type: 'output', id: job.id, ...owner !== undefined ? { owner } : {}, total: job.ring.total }, job.owner)
553 }
554
555 /**
556 * Replace the live progress line through a producer face; a write against a
557 * settled job is logged and dropped. A write staged inside the starter call
558 * seeds the registered projection and signals no observer.
559 */
560 private updateProgress(state: ProducerState, line: string): void {
561 const job = state.job
562 if (job !== undefined && isTerminal(job.status)) {
563 this.selfCtx.logger.warn(`jobs: progress update on settled job ${job.id} dropped`)
564 return
565 }
566 state.progress = line
567 if (job !== undefined) this.emit({ type: 'progress', job: this.view(job) }, job.owner)
568 }
569
570 /**
571 * Record the first terminal outcome, release waiters, then announce the
572 * settlement. First-wins preserves a teardown force-failure against late
573 * producer settlement. The settled event follows every released waiter and
574 * reports whether it released one: a timed-out or aborted wait has already
575 * left the set, so only a wait still owed the projection counts.
576 */
577 private settle(job: TrackedJob, outcome: JobOutcome, cause: JobSettleCause): void {
578 if (isTerminal(job.status)) return
579 job.status = outcome.status
580 // A killed settlement carries the recorded kill reason in its detail:
581 // producer facts first (`signal: SIGTERM; cancelled by the user`). A job
582 // that outran its kill request (settled `completed`/`failed`) keeps the
583 // producer detail alone — the reason describes a kill that never landed.
584 if (outcome.status === 'killed' && job.killReason !== undefined) {
585 job.detail = outcome.detail !== undefined
586 ? `${outcome.detail}; ${job.killReason}`
587 : job.killReason
588 } else if (outcome.detail !== undefined) {
589 job.detail = outcome.detail
590 }
591 job.state.progress = undefined
592 job.result = outcome.result
593 job.finishedAt = Date.now()
594 // Settlement ends the stream: trim to the settled cap before any observer
595 // reads the terminal projection, but never below the bytes the model
596 // cursor has not consumed. A job that finishes before its first model
597 // read keeps everything the live cap retained until that read.
598 job.ring.trim(Math.max(this.settledRetainBytes, job.ring.total - job.modelCursor))
599 const waitResolvers = [...job.waitResolvers]
600 job.waitResolvers.clear()
601 for (const resolveWait of waitResolvers) resolveWait()
602 job.markSettled()
603 this.emit({ type: 'settled', job: this.view(job), cause, awaited: waitResolvers.length > 0 }, job.owner)
604 // The ring's stream ends with settlement; the signal follows the committed
605 // settlement so an observer that wakes on it reads the terminal state.
606 this.emitOutput(job)
607 }
608
609 /**
610 * Attach one awaited cleanup through the exact owner's scope. This survives
611 * producer reloads and joins agent quiescence; the retained disposer lets
612 * service teardown detach the cross-fiber effect.
613 */
614 private ensureOwnerCleanup(owner: Agent): void {
615 if (this.ownerCleanups.has(owner)) return
616 // Record only after attach succeeds; a disposing scope rejects new effects.
617 const detach = owner.ctx.effect(() => async () => {
618 this.ownerCleanups.delete(owner)
619 await this.disposeOwned(owner)
620 }, 'jobs.ownerCleanup()')
621 this.ownerCleanups.set(owner, detach)
622 }
623
624 /** Cancel, await terminal records, and drop every job owned by one exact agent lifecycle. */
625 private async disposeOwned(owner: Agent): Promise<void> {
626 const owned = [...this.store.values()].filter(job => job.owner === owner)
627 this.cancelForTeardown(owned, 'owner disposed')
628 await Promise.all(owned.map(job => job.settled))
629 this.drop(owned)
630 }
631
632 /** Drop settled records and announce each removal, the one visible-set change no per-job record carries. */
633 private drop(jobs: readonly TrackedJob[]): void {
634 for (const job of jobs) {
635 this.store.delete(job.id)
636 this.emit({ type: 'removed', job: this.view(job) }, job.owner)
637 }
638 }
639
640 /**
641 * Cancel live jobs, await settlement, drop every record, and detach owner
642 * effects. Throwing cancels are force-failed to avoid teardown deadlock.
643 */
644 private async disposeAll(): Promise<void> {
645 const all = [...this.store.values()]
646 this.cancelForTeardown(all, 'jobs service disposed')
647 await Promise.all(all.map(job => job.settled))
648 // A subscriber mounted outside this service — the job controller's rows
649 // stream registers from its own context — is still reachable here.
650 // Without the removals it keeps the rows it last received after a
651 // registry reload.
652 this.drop(all)
653 // Detach cross-fiber owner effects after the shared store is quiescent.
654 const ownerCleanups = [...this.ownerCleanups.values()]
655 this.ownerCleanups.clear()
656 await Promise.all(ownerCleanups.map(cleanup => Promise.resolve(cleanup())))
657 }
658
659 /**
660 * Cancel jobs during teardown with per-job containment. A throwing cancel
661 * force-fails the record and reports a possible orphan; a cancel that returns
662 * without settling remains indistinguishable from a slow stop and may stall.
663 */
664 private cancelForTeardown(jobs: TrackedJob[], reason: string): void {
665 for (const job of jobs) {
666 if (isTerminal(job.status)) continue
667 // Whatever settles this job from here on, its owner or the service is
668 // being destroyed: the settlement announces `teardown` so a completion
669 // reporter does not address a reader that no longer exists.
670 job.settleCause = 'teardown'
671 try {
672 job.cancel(reason)
673 job.status = 'stopping'
674 // Teardown reaches settlement only after the producer releases, which a
675 // slow stop can defer; announcing the transition here is what keeps an
676 // observer from showing `running` for that whole window.
677 this.emit({ type: 'stopping', job: this.view(job) }, job.owner)
678 } catch (error: unknown) {
679 const detail = `cancel threw during teardown; work may be orphaned: ${String(error)}`
680 this.selfCtx.logger.warn(`jobs: cancel of ${job.id} threw during teardown; job record forced failed and work may be orphaned: ${String(error)}`)
681 this.settle(job, { status: 'failed', detail }, 'teardown')
682 }
683 }
684 }
685}
686
687export default LocalJobRegistry