1
/**2
* Process-local provider for the background-job capability seam3
* (`ctx.jobs`). It keeps every job — lifecycle state, the bounded output4
* ring, and the model cursor — in memory and hands out fresh projections and5
* chunk copies, never live state.6
*7
* Registrations outlive producer and controller fibers. Agent or service8
* disposal cancels live work and awaits compliant producers; a throwing9
* teardown cancel force-fails only the record and reports a possible orphan.10
* @module @deepseek-ai/dsh-jobs-local11
*/13
import { Context } from '@deepseek-ai/cordis'14
import z from '@deepseek-ai/schemastery'15
import type { Agent } from '@deepseek-ai/dsh-agent'16
import { ScopedLayers, scopeOf } from '@deepseek-ai/dsh-scope'17
import type { SessionId } from '@deepseek-ai/dsh-session'18
import { deadline, timeoutOf } from '@deepseek-ai/dsh-timeout'19
import { JobRegistry, JobId } from '@deepseek-ai/dsh-jobs'20
import type {21
JobAppendOptions, JobEvent, JobEvents, JobHandle, JobKind, JobOutcome, JobOutputRead, JobOutputSource,22
JobRead, JobSettleCause, JobSpec, JobStatus, JobView,23
} from '@deepseek-ai/dsh-jobs'24
import { JobEventHub, JobLayer } from './events.ts'25
import { startPump } from './pump.ts'26
import type { PumpHandle } from './pump.ts'27
import { OutputRing } from './ring.ts'29
/** Timeout code that distinguishes a bounded wait from caller cancellation. */30
export const TASK_WAIT_TIMEOUT = 'TASK_WAIT_TIMEOUT'32
/** Default maximum number of active jobs in one exact-owner bucket. */33
const DEFAULT_MAX_CONCURRENT_JOBS_PER_OWNER = 1035
/** Default live ring retention per job, in UTF-8 bytes. */36
const DEFAULT_RETAIN_BYTES = 256 * 102438
/** Default ring retention kept after settlement, in UTF-8 bytes. */39
const DEFAULT_SETTLED_RETAIN_BYTES = 16 * 102441
/** Default poll interval for pull sources, in milliseconds. */42
const DEFAULT_PUMP_POLL_MS = 15044
/** Configuration for the process-local job registry. */45
export 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?: number51
/** Live ring retention per job in UTF-8 bytes; omission defaults to 262144. */52
retainBytes?: number53
/**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 of56
* this cap; the first terminal model read then trims to it.57
*/58
settledRetainBytes?: number59
/** Poll interval for a job's pull sources, in milliseconds; omission defaults to 150. */60
pumpPollMs?: number61
}63
/**64
* Producer-written state shared between the {@link JobHandle} and the65
* registered record: the starter call writes through it before the commit,66
* the same object serves the job for its whole life afterwards.67
*/68
interface ProducerState {69
/** Live progress line until settlement clears it. */70
progress: string | undefined71
/** The committed registry record; undefined exactly during the starter call. */72
job: TrackedJob | undefined73
}75
/** The registry's mutable per-job record (never handed out — see {@link LocalJobRegistry.view}). */76
interface TrackedJob {77
id: JobId78
kind: JobKind79
label: string80
outputLimitBytes: number | undefined81
/** Exact lifecycle owner; session-id authorization is derived from it. */82
owner: Agent | undefined83
cancel: (reason?: string) => void84
status: JobStatus85
ring: OutputRing86
/** The model's consuming cursor; {@link JobRegistry.readAt} never moves it. */87
modelCursor: number88
/** Whether the first post-settlement read already handed out `result`. */89
resultDelivered: boolean90
/** Producer-shared progress line and commit binding. */91
state: ProducerState92
/** Terminal reason; a recorded kill reason is merged in at settlement. */93
detail: string | undefined94
result: string | undefined95
startedAt: number96
finishedAt: number | undefined97
/** Reason recorded by {@link JobRegistry.kill}, merged into a `killed` settlement's detail. */98
killReason: string | undefined99
/** Set once a kill or teardown cancel ran; settlement reports it as the cause. */100
settleCause: JobSettleCause | undefined101
/** 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: () => void105
/** 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 | undefined109
/**110
* The spill file each pull source reported on its latest read, by source111
* index; an entry is undefined while that source keeps none. Source112
* metadata rather than per-chunk metadata, so it survives ring eviction and113
* follows a source that withdraws its file.114
*/115
spillPaths: (string | undefined)[]116
}118
/** True for the three terminal {@link JobStatus} values. */119
function isTerminal(status: JobStatus): boolean {120
return status === 'completed' || status === 'killed' || status === 'failed'121
}123
/**124
* The in-memory `jobs` registry. See the Service Definition contract in125
* `@deepseek-ai/dsh-jobs` for the ownership, isolation, and lifecycle126
* semantics this implementation honors.127
*/128
export 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
})152
/** Schemastery-defaulted active-job limit. */153
private readonly maxConcurrentJobsPerOwner: number154
/** Schemastery-defaulted live ring retention cap. */155
private readonly retainBytes: number156
/** Schemastery-defaulted settled ring retention cap. */157
private readonly settledRetainBytes: number158
/** Schemastery-defaulted pull-source poll interval. */159
private readonly pumpPollMs: number160
private store = new Map<JobId, TrackedJob>()161
private counters = new Map<string, number>()162
/**163
* Controllers and scoped subscriptions layered by the scope that registered164
* them, in the tools-registry shape: a contribution files into its165
* registering context's scope, and a read unions the global layer with the166
* owner's scope chain.167
*168
* The registry is one process-wide instance serving every composition, so a169
* flat table would answer a per-owner question process-wide: one preset's170
* job controls would hold `start()` open for an agent whose own composition171
* loads none, and one settlement would reach every preset's notice listener.172
* Layers make both reads owner-relative. Nothing derives a cache from a173
* layer, so change notification is a no-op.174
*/175
private readonly layers = new ScopedLayers<JobLayer>(() => new JobLayer(), () => {})176
private readonly hub: JobEventHub177
/** 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: Context182
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.maxConcurrentJobsPerOwner187
this.retainBytes = resolved.retainBytes188
this.settledRetainBytes = resolved.settledRetainBytes189
this.pumpPollMs = resolved.pumpPollMs190
this.selfCtx = ctx191
this.hub = new JobEventHub(this.layers, (message) => { ctx.logger.warn(message) })192
ctx.effect(() => () => this.disposeAll(), 'jobs teardown')193
}195
/**196
* The event stream bound to the accessing context: a subscription is an197
* effect of that context, and `{ owners: 'scope' }` names its scope.198
*/199
get events(): JobEvents {200
const registrar = this.ctx201
return {202
subscribe: (filter, listener) => this.hub.subscribe(registrar, filter, listener),203
}204
}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 !== undefined214
&& (!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)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
}226
// The id is issued before the starter runs so the producer face can carry227
// it; a throwing starter still leaves nothing registered — its ordinal is228
// simply skipped.229
const count = (this.counters.get(spec.kind) ?? 0) + 1230
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)241
let markSettled!: () => void242
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 inside268
// the starter are already in the ring and `state`, and every later handle269
// call reaches the registered record for its terminal checks and signals.270
state.job = job271
this.store.set(id, job)272
// Registration is complete and cannot fail from here, so the visible set273
// has genuinely changed. The announcement precedes the pump because the274
// pump drains its sources once synchronously, and that drain may append275
// and announce output: a job's first event is always `registered`.276
this.emit({ type: 'registered', job: this.view(job) }, owner)278
// The producer's settlement or a registry-forced one ends the pump; the279
// 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.done301
this.settle(job, outcome, job.settleCause ?? 'producer')302
})303
return id304
}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
}312
get(id: JobId, caller?: SessionId): JobView {313
return this.view(this.expect(id, caller))314
}316
read(id: JobId, caller?: SessionId): JobRead {317
return this.readJob(this.expect(id, caller))318
}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
}328
kill(id: JobId, caller?: SessionId, reason?: string): 'requested' | 'already-finished' {329
return this.killJob(this.expect(id, caller), reason)330
}332
async wait(id: JobId, timeoutMs: number, caller?: SessionId, signal?: AbortSignal): Promise<JobView> {333
return this.waitJob(this.expect(id, caller), timeoutMs, signal)334
}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
}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
}352
/**353
* Resolve a spec's owner session to its live Agent. An owned registration354
* needs the agent registry, and the session must currently have a live355
* 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 undefined359
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 owner368
}370
/**371
* Whether an attached job controller can collect and stop work owned by372
* `owner`. The global layer holds every controller attached from an unscoped373
* context — a host composition's own controls — and therefore serves every374
* 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 true380
return this.layers.chainLayers(owner === undefined ? undefined : scopeOf(owner.ctx))381
.some(layer => !layer.controllers.isEmpty())382
}384
/** Count authoritative active records for one exact owner or the shared unowned bucket. */385
private activeJobCount(owner: Agent | undefined): number {386
let count = 0387
for (const job of this.store.values()) {388
if (job.owner === owner && (job.status === 'running' || job.status === 'stopping')) count += 1389
}390
return count391
}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 job399
}401
/**402
* The isolation fence: a job with an owner is reachable only by callers403
* whose session id matches (`!== undefined` semantics — an unowned job is404
* 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
}412
/** Project a fresh read-only view from the mutable record. */413
private view(job: TrackedJob): JobView {414
const owner = job.owner?.id415
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
}435
private emit(event: JobEvent, owner: Agent | undefined): void {436
this.hub.emit(event, owner)437
}439
/**440
* Consume the ring from the model cursor; the result rides the first read441
* after settlement. A terminal read is the point the settled stream drops442
* 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.total447
const result = isTerminal(job.status) && !job.resultDelivered ? job.result : undefined448
if (result !== undefined) job.resultDelivered = true449
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
}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 = reason465
job.settleCause = 'kill'466
this.emit({ type: 'stopping', job: this.view(job) }, job.owner)467
return 'requested'468
}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 from477
// 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 waiter488
// before it announces, and each released waiter detaches this489
// listener in the same synchronous span, so nothing that reacts to a490
// 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
}504
/**505
* Append one chunk to the ring. A producer chunk against a settled job is506
* logged and dropped; the registry's own pump drains silently after507
* settlement (a forced settlement may precede the producer's). A chunk508
* staged inside the starter call is retained and signals no observer — the509
* 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.job519
if (job !== undefined && isTerminal(job.status)) {520
if (writer === 'producer') this.selfCtx.logger.warn(`jobs: append to settled job ${job.id} dropped`)521
return522
}523
if (!ring.append(text, options, this.retainBytes)) return524
if (job !== undefined) this.emitOutput(job)525
}527
/**528
* Contain a failing pull source: the first throw is logged, and the source529
* reads as exhausted from then on, so the job runs to its own settlement530
* with whatever the ring holds instead of freezing on a pump failure.531
*/532
private guardSource(job: TrackedJob, source: JobOutputSource): JobOutputSource {533
let failed = false534
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 = true542
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
}549
/** Announce that one job's ring advanced (append or settlement). */550
private emitOutput(job: TrackedJob): void {551
const owner = job.owner?.id552
this.emit({ type: 'output', id: job.id, ...owner !== undefined ? { owner } : {}, total: job.ring.total }, job.owner)553
}555
/**556
* Replace the live progress line through a producer face; a write against a557
* settled job is logged and dropped. A write staged inside the starter call558
* seeds the registered projection and signals no observer.559
*/560
private updateProgress(state: ProducerState, line: string): void {561
const job = state.job562
if (job !== undefined && isTerminal(job.status)) {563
this.selfCtx.logger.warn(`jobs: progress update on settled job ${job.id} dropped`)564
return565
}566
state.progress = line567
if (job !== undefined) this.emit({ type: 'progress', job: this.view(job) }, job.owner)568
}570
/**571
* Record the first terminal outcome, release waiters, then announce the572
* settlement. First-wins preserves a teardown force-failure against late573
* producer settlement. The settled event follows every released waiter and574
* reports whether it released one: a timed-out or aborted wait has already575
* 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)) return579
job.status = outcome.status580
// A killed settlement carries the recorded kill reason in its detail:581
// producer facts first (`signal: SIGTERM; cancelled by the user`). A job582
// that outran its kill request (settled `completed`/`failed`) keeps the583
// producer detail alone — the reason describes a kill that never landed.584
if (outcome.status === 'killed' && job.killReason !== undefined) {585
job.detail = outcome.detail !== undefined586
? `${outcome.detail}; ${job.killReason}`587
: job.killReason588
} else if (outcome.detail !== undefined) {589
job.detail = outcome.detail590
}591
job.state.progress = undefined592
job.result = outcome.result593
job.finishedAt = Date.now()594
// Settlement ends the stream: trim to the settled cap before any observer595
// reads the terminal projection, but never below the bytes the model596
// cursor has not consumed. A job that finishes before its first model597
// 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 committed605
// settlement so an observer that wakes on it reads the terminal state.606
this.emitOutput(job)607
}609
/**610
* Attach one awaited cleanup through the exact owner's scope. This survives611
* producer reloads and joins agent quiescence; the retained disposer lets612
* service teardown detach the cross-fiber effect.613
*/614
private ensureOwnerCleanup(owner: Agent): void {615
if (this.ownerCleanups.has(owner)) return616
// 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
}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
}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
}640
/**641
* Cancel live jobs, await settlement, drop every record, and detach owner642
* 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 rows649
// stream registers from its own context — is still reachable here.650
// Without the removals it keeps the rows it last received after a651
// 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
}659
/**660
* Cancel jobs during teardown with per-job containment. A throwing cancel661
* force-fails the record and reports a possible orphan; a cancel that returns662
* 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)) continue667
// Whatever settles this job from here on, its owner or the service is668
// being destroyed: the settlement announces `teardown` so a completion669
// 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 a675
// slow stop can defer; announcing the transition here is what keeps an676
// 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
}687
export default LocalJobRegistry