From 1cc5c951ceb1e57054cc9f5355464e0606540218 Mon Sep 17 00:00:00 2001 From: barlind Date: Tue, 9 Jun 2026 14:33:15 +0200 Subject: [PATCH] Enhance job handling with provider-private payload support and recovery tests --- static/capabilities/jobs.js | 113 ++++++++++++++++++++++++++++-- tests/js/jobs_diagnostics.test.js | 24 +++++++ tests/js/jobs_domain.test.js | 63 +++++++++++++++++ tests/js/jobs_test_harness.js | 5 ++ 4 files changed, 198 insertions(+), 7 deletions(-) diff --git a/static/capabilities/jobs.js b/static/capabilities/jobs.js index d4e1864..d58f166 100644 --- a/static/capabilities/jobs.js +++ b/static/capabilities/jobs.js @@ -54,6 +54,7 @@ const outcomes = []; const bridgeHits = []; const pendingRecoverableRefs = new Map(); + const providerPrivatePayloads = new Map(); const memoryStorage = new Map(); function _now() { return new Date().toISOString(); } @@ -355,6 +356,7 @@ state: STATES.QUEUED, progress: { mode: 'indeterminate', percent: null, step: '', message: '', updatedAt: now }, attempts: [], + externallyManaged: !!args.externallyManaged, history: [], retryable: false, safeReason: null, @@ -384,6 +386,7 @@ progress: _clone(job.progress), actionsAvailable: _actionsAvailable(job), retryable: !!job.retryable, + externallyManaged: !!job.externallyManaged, attempts: job.attempts.map(attempt => ({ attemptId: attempt.attemptId, attemptNumber: attempt.attemptNumber, @@ -463,12 +466,14 @@ function _handled(payload = {}) { return _result('handled', payload); } - function _callProvider(provider, operation, payload) { + function _callProvider(provider, operation, payload, options = {}) { const handlers = provider && provider.operationHandlers ? provider.operationHandlers : {}; const handler = handlers[operation] || handlers[operation.replace(/^job\./, '')]; if (typeof handler !== 'function') return null; try { - return handler(_safeValue(payload)); + const request = _safeValue(payload) || {}; + if (Object.prototype.hasOwnProperty.call(options, 'providerPayload')) request.providerPayload = options.providerPayload; + return handler(request); } catch (err) { return _providerException(operation, err); } @@ -604,6 +609,7 @@ priority: job.priority, safeLabel: job.safeLabel, state: recoveryState, + externallyManaged: !!job.externallyManaged, persistedAt: _now(), }; } @@ -645,6 +651,7 @@ state, progress: { mode: 'indeterminate', percent: null, step: 'recovered', message: 'Recovered after reload', updatedAt: now }, attempts: [], + externallyManaged: !!ref.externallyManaged, history: [], retryable: false, safeReason: state === STATES.PROVIDER_UNAVAILABLE || state === STATES.ORPHANED ? 'Provider recovery unavailable after reload' : 'Recovered after reload', @@ -681,6 +688,17 @@ pendingRecoverableRefs.delete(jobId); continue; } + const recoveryOutcome = recoveryResult && recoveryResult.outcome; + const recoveryState = recoveryResult && (recoveryResult.terminalStatus || recoveryResult.status || recoveryResult.state); + const terminalStatus = ['completed', 'cancelled', 'failed', 'timeout', 'provider-unavailable', 'orphaned'].includes(recoveryState) ? recoveryState : recoveryOutcome; + if (['completed', 'cancelled', 'failed', 'timeout', 'provider-unavailable', 'orphaned'].includes(terminalStatus)) { + const job = _jobFromRecoveryRef(ref, STATES.QUEUED); + jobs.set(job.jobId, job); + _terminal(job, terminalStatus, recoveryResult.category, recoveryResult.safeReason || recoveryResult.reason || 'Recovered terminal job', !!recoveryResult.retryable, recoveryResult.resultSummary || recoveryResult.summary || 'Recovered terminal job'); + _emitJobs(terminalStatus === 'completed' ? 'completed' : (terminalStatus === 'cancelled' ? 'cancelled' : (terminalStatus === 'provider-unavailable' ? 'provider-unavailable' : (terminalStatus === 'orphaned' ? 'orphaned' : 'failed'))), { job: _jobSummary(job, { includeHistory: false }) }); + pendingRecoverableRefs.delete(jobId); + continue; + } if (recoveryResult && recoveryResult.state && [STATES.QUEUED, STATES.RUNNING, STATES.PAUSED].includes(recoveryResult.state)) nextState = recoveryResult.state; const job = _jobFromRecoveryRef(ref, nextState); jobs.set(job.jobId, job); @@ -752,6 +770,8 @@ } const job = _newJob(args, provider.providerId); jobs.set(job.jobId, job); + if (Object.prototype.hasOwnProperty.call(args, 'providerPayload')) providerPrivatePayloads.set(job.jobId, args.providerPayload); + else if (Object.prototype.hasOwnProperty.call(args, 'privatePayload')) providerPrivatePayloads.set(job.jobId, args.privatePayload); _history(job, 'event', 'Job queued'); _rememberOutcome('enqueue', 'queued', { jobId: job.jobId, providerId: provider.providerId, requesterId: job.requesterId }); _emitJobs('queued', { job: _jobSummary(job, { includeHistory: false }) }); @@ -762,6 +782,66 @@ return _result(job.state, { job: _jobSummary(job) }); } + function _adoptCommand(ctx) { + const args = _plainObject(ctx.payload); + const providerId = _safeId(args.providerId, ''); + const jobType = _safeId(args.jobType, ''); + if (!providerId || !jobType) return _result('validation-failed', {}, 'providerId and jobType are required'); + const provider = providers.get(providerId); + if (!provider) return _result('no-owner', {}, `Provider ${providerId} is not registered`); + if (!provider.jobTypes.includes(jobType)) return _result('no-handler', {}, `Provider ${providerId} does not handle ${jobType}`); + if (!args.requester && !args.requesterId && !ctx.requester) args.requester = provider.pluginId || providerId; + const jobId = _safeId(args.jobId, '') || _id('job'); + let job = jobs.get(jobId); + const desiredState = _safeId(args.state || args.status || 'running', 'running'); + const activeState = [STATES.QUEUED, STATES.RUNNING, STATES.PAUSED, STATES.CANCELLATION_REQUESTED].includes(desiredState) ? desiredState : STATES.RUNNING; + const isTerminal = ['completed', 'cancelled', 'failed', 'timeout', 'provider-unavailable', 'orphaned'].includes(desiredState); + if (job && TERMINAL_STATES.has(job.state)) return _result('stale', { job: _jobSummary(job) }, 'job is terminal'); + if (!job) { + job = _newJob({ ...args, jobId, externallyManaged: true }, providerId); + jobs.set(job.jobId, job); + _history(job, 'event', 'Job adopted from provider-owned work'); + } else { + job.externallyManaged = true; + job.safeLabel = _safeText(args.safeLabel || args.label || job.safeLabel, job.safeLabel, MAX_LABEL); + job.updatedAt = _now(); + } + if (args.progress) { + const source = _plainObject(args.progress); + const mode = source.mode === 'determinate' || source.mode === 'indeterminate' || source.mode === 'step-only' ? source.mode : (source.percent == null ? 'indeterminate' : 'determinate'); + job.progress = { + mode, + percent: mode === 'determinate' ? Math.max(0, Math.min(100, _number(source.percent, 0) || 0)) : null, + step: _safeText(source.step || source.currentStep || '', '', 80), + message: _safeText(source.message || source.safeMessage || '', '', MAX_REASON), + updatedAt: _now(), + }; + } + if (isTerminal) { + _terminal(job, desiredState, args.category, args.safeReason || args.reason || 'Provider reported terminal state', !!args.retryable, args.resultSummary || args.summary || 'Provider job settled'); + _rememberOutcome('adopt', desiredState === 'timeout' ? 'timeout' : desiredState, { jobId: job.jobId, providerId, requesterId: job.requesterId, category: args.category, safeReason: job.safeReason }); + _emitJobs(desiredState === 'completed' ? 'completed' : (desiredState === 'cancelled' ? 'cancelled' : (desiredState === 'provider-unavailable' ? 'provider-unavailable' : (desiredState === 'orphaned' ? 'orphaned' : 'failed'))), { job: _jobSummary(job, { includeHistory: false }) }); + _scheduleProvider(providerId); + return _result(desiredState === 'timeout' ? 'timeout' : desiredState, { job: _jobSummary(job) }); + } + job.state = activeState; + job.safeReason = args.safeReason || args.reason ? _safeText(args.safeReason || args.reason) : null; + if (activeState === STATES.RUNNING) job.startedAt = job.startedAt || _now(); + if (activeState === STATES.QUEUED) job.queuedAt = job.queuedAt || _now(); + job.updatedAt = _now(); + const attempt = job.attempts[job.attempts.length - 1]; + if (attempt) { + attempt.state = activeState; + attempt.updatedAt = job.updatedAt; + if (activeState === STATES.RUNNING) attempt.startedAt = attempt.startedAt || job.updatedAt; + } + _history(job, 'event', `Job adopted as ${activeState}`); + _rememberOutcome('adopt', 'handled', { jobId: job.jobId, providerId, requesterId: job.requesterId, safeReason: job.safeReason }); + _emitJobs(activeState === STATES.PAUSED ? 'paused' : (activeState === STATES.RUNNING || activeState === STATES.CANCELLATION_REQUESTED ? 'started' : 'queued'), { job: _jobSummary(job, { includeHistory: false }) }); + _persistRecoverableRefs(); + return _handled({ job: _jobSummary(job) }); + } + function _listCommand(ctx) { const filters = _plainObject(ctx.payload); const includeTerminal = filters.includeTerminal !== false; @@ -770,6 +850,7 @@ if (filters.jobType && job.jobType !== filters.jobType) return false; if (filters.state && job.state !== filters.state) return false; if (filters.requesterId && job.requesterId !== filters.requesterId) return false; + if (filters.externallyManaged != null && !!job.externallyManaged !== !!filters.externallyManaged) return false; if (!includeTerminal && TERMINAL_STATES.has(job.state)) return false; return true; }).map(job => _jobSummary(job, { includeHistory: false })); @@ -934,7 +1015,7 @@ if (!provider || !_providerCanRun(provider)) return; let load = _providerLoad(providerId); const queued = Array.from(jobs.values()) - .filter(job => job.providerId === providerId && job.state === STATES.QUEUED) + .filter(job => job.providerId === providerId && job.state === STATES.QUEUED && !job.externallyManaged) .sort((a, b) => _priorityRank(a) - _priorityRank(b) || String(a.queuedAt || '').localeCompare(String(b.queuedAt || '')) || a.jobId.localeCompare(b.jobId)); for (const job of queued) { if (load.running >= provider.capacity.maxRunning) { @@ -960,7 +1041,9 @@ attempt.updatedAt = _now(); } _history(job, 'event', 'Job started'); - const result = _callProvider(provider, 'job.enqueue', { job: _jobSummary(job, { includeHistory: false }) }); + const hasProviderPayload = providerPrivatePayloads.has(job.jobId); + const result = _callProvider(provider, 'job.enqueue', { job: _jobSummary(job, { includeHistory: false }) }, hasProviderPayload ? { providerPayload: providerPrivatePayloads.get(job.jobId) } : {}); + providerPrivatePayloads.delete(job.jobId); if (_providerCallFailed(result)) { _terminalProviderFailure(job, 'enqueue', result); return; @@ -993,6 +1076,7 @@ attempt.terminalAt = job.terminalAt; attempt.terminalOutcome = _clone(job.terminalOutcome); } + providerPrivatePayloads.delete(job.jobId); if (!terminalJobIds.includes(job.jobId)) terminalJobIds.push(job.jobId); while (terminalJobIds.length > MAX_TERMINAL_JOBS) terminalJobIds.shift(); _history(job, 'event', `Job ${outcomeStatus}`); @@ -1048,6 +1132,17 @@ return _result('completed', { job: _jobSummary(job) }); } + function cancelled(providerId, jobId, result = {}) { + const job = jobs.get(_safeId(jobId, '')); + if (!job || job.providerId !== providerId) return _result('no-target', {}, 'job not found'); + if (TERMINAL_STATES.has(job.state)) return _result('stale', { job: _jobSummary(job) }, 'job is terminal'); + _terminal(job, 'cancelled', 'cancellation', result.safeReason || result.reason || 'Provider reported cancellation', !!result.retryable, result.resultSummary || result.summary || 'Cancelled'); + _rememberOutcome('cancelled', 'cancelled', { jobId: job.jobId, providerId: job.providerId, requesterId: job.requesterId, category: 'cancellation', safeReason: job.safeReason }); + _emitJobs('cancelled', { job: _jobSummary(job, { includeHistory: false }) }); + _scheduleProvider(providerId); + return _result('cancelled', { job: _jobSummary(job) }); + } + function fail(providerId, jobId, result = {}) { const job = jobs.get(_safeId(jobId, '')); if (!job || job.providerId !== providerId) return _result('no-target', {}, 'job not found'); @@ -1171,6 +1266,7 @@ outcomes.length = 0; bridgeHits.length = 0; pendingRecoverableRefs.clear(); + providerPrivatePayloads.clear(); if (options.clearStorage !== false) { _storageRemove(SELECTED_PROVIDER_STORAGE_KEY); _storageRemove(RECOVERY_STORAGE_KEY); @@ -1189,8 +1285,8 @@ compatibility: 'shim-allowed', ownership: 'multi-provider', safety: 'privileged', - commands: ['register-provider', 'unregister-provider', 'list-providers', 'enqueue', 'list', 'inspect', 'cancel', 'pause', 'resume', 'retry', 'record-bridge-hit'], - operations: ['job.enqueue', 'job.status', 'job.cancel', 'job.pause', 'job.resume', 'job.retry', 'job.recover'], + commands: ['register-provider', 'unregister-provider', 'list-providers', 'enqueue', 'adopt', 'list', 'inspect', 'cancel', 'pause', 'resume', 'retry', 'record-bridge-hit'], + operations: ['job.enqueue', 'job.status', 'job.cancel', 'job.pause', 'job.resume', 'job.retry', 'job.recover', 'job.adopt'], events: ['provider-registered', 'provider-unregistered', 'provider-unavailable', 'queued', 'started', 'progress', 'log', 'paused', 'resumed', 'cancellation-requested', 'cancelled', 'completed', 'failed', 'retried', 'orphaned', 'bridge-hit'], description: 'Owns the privileged jobs provider registry, scheduling, lifecycle state, recovery, bridge hits, and redaction-safe diagnostics.', provider_policy: { providerId: 'jobs', kind: 'core', safety: 'privileged' }, @@ -1199,6 +1295,7 @@ 'unregister-provider': _unregisterProviderCommand, 'list-providers': _listProvidersCommand, enqueue: _enqueueCommand, + adopt: _adoptCommand, list: _listCommand, inspect: _inspectCommand, cancel: _cancelCommand, @@ -1221,11 +1318,13 @@ reportProgress: updateProgress, log, complete, + cancelled, fail, markProviderUnavailable, simulateReload, resetForTests, - _test: { reset: resetForTests, providers, jobs, selectedProviders, pendingRecoverableRefs, storage: memoryStorage }, + adopt(args) { return _adoptCommand({ payload: args || {}, requester: 'api' }); }, + _test: { reset: resetForTests, providers, jobs, selectedProviders, pendingRecoverableRefs, providerPrivatePayloads, storage: memoryStorage }, }; window.slopsmith.jobs = api; diff --git a/tests/js/jobs_diagnostics.test.js b/tests/js/jobs_diagnostics.test.js index 2210f44..df27a67 100644 --- a/tests/js/jobs_diagnostics.test.js +++ b/tests/js/jobs_diagnostics.test.js @@ -115,6 +115,30 @@ test('async recovery rejections become provider-unavailable terminal jobs', asyn assert.doesNotMatch(JSON.stringify(snapshot), /Users\/example|recovery\.db/); }); +test('recovery handlers can restore terminal provider-owned jobs', async () => { + const window = loadJobs(); + const initial = makeProvider({ providerId: 'provider.recover-terminal', recoverySupport: { queued: true, running: true, paused: true } }); + await dispatch(window, 'register-provider', { provider: initial.provider }); + await dispatch(window, 'enqueue', enqueuePayload({ logicalJobKey: 'recover-terminal', safeLabel: 'Terminal Recover' })); + + window.slopsmith.jobs.resetForTests({ clearStorage: false }); + const recovered = makeProvider({ + providerId: 'provider.recover-terminal', + recoverySupport: { queued: true, running: true, paused: true }, + operationHandlers: { + 'job.recover': async () => ({ outcome: 'handled', state: 'completed', resultSummary: 'Backend finished while reloading' }), + }, + }); + await dispatch(window, 'register-provider', { provider: recovered.provider }); + const snapshot = diagnosticsSnapshot(window); + + assert.equal(snapshot.jobs.active.length, 0); + assert.equal(snapshot.jobs.queued.length, 0); + assert.equal(snapshot.jobs.recentTerminal.length, 1); + assert.equal(snapshot.jobs.recentTerminal[0].state, 'completed'); + assert.equal(snapshot.jobs.recentTerminal[0].terminalOutcome.resultSummary, 'Backend finished while reloading'); +}); + test('reload marks non-recoverable jobs orphaned or provider-unavailable without restoring raw payloads', async () => { const window = loadJobs(); const { provider } = makeProvider({ providerId: 'provider.no-recover', recoverySupport: { queued: false, running: false, paused: false } }); diff --git a/tests/js/jobs_domain.test.js b/tests/js/jobs_domain.test.js index 459ea61..c905216 100644 --- a/tests/js/jobs_domain.test.js +++ b/tests/js/jobs_domain.test.js @@ -77,6 +77,69 @@ test('approved enqueue queues and starts with redaction-safe public job fields', assert.doesNotMatch(JSON.stringify(diagnosticsSnapshot(window)), /Secret\.sloppak|token|secret/); }); +test('provider-private enqueue payload reaches only the provider callback', async () => { + const window = loadJobs(); + let privatePayload = null; + const { provider } = makeProvider({ + operationHandlers: { + 'job.enqueue': request => { + privatePayload = request.providerPayload; + return { outcome: 'handled' }; + }, + }, + }); + await dispatch(window, 'register-provider', { provider }); + + const result = await dispatch(window, 'enqueue', enqueuePayload({ + providerPayload: { filename: '/Users/example/DLC/Secret Song.psarc', token: 'abc123' }, + target: { safeRef: 'target-secret-song' }, + inputs: { safeFingerprint: 'input-secret-song' }, + })); + + assert.equal(result.status, 'applied'); + assert.deepEqual(privatePayload, { filename: '/Users/example/DLC/Secret Song.psarc', token: 'abc123' }); + assert.doesNotMatch(JSON.stringify(result.payload.job), /Secret Song|abc123|filename/); + assert.doesNotMatch(JSON.stringify(diagnosticsSnapshot(window)), /Secret Song|abc123|filename/); +}); + +test('adopted provider-owned jobs do not invoke enqueue handlers', async () => { + const window = loadJobs(); + const { calls, provider } = makeProvider({ providerId: 'provider.backend', capacity: { maxRunning: 1, maxQueued: 10 } }); + await dispatch(window, 'register-provider', { provider }); + + const adopted = await dispatch(window, 'adopt', enqueuePayload({ + providerId: 'provider.backend', + jobId: 'backend-job-1', + state: 'queued', + logicalJobKey: 'legacy-backend-job-1', + safeLabel: 'Convert existing backend row', + })); + const native = await dispatch(window, 'enqueue', enqueuePayload({ providerId: 'provider.backend', logicalJobKey: 'native-after-adopt' })); + + assert.equal(adopted.status, 'applied'); + assert.equal(adopted.payload.job.externallyManaged, true); + assert.equal(adopted.payload.job.state, 'queued'); + assert.equal(native.status, 'applied'); + assert.equal(native.payload.job.state, 'running'); + assert.equal(calls.length, 1); + assert.equal(calls[0][1].job.jobId, native.payload.job.jobId); +}); + +test('provider can report a running job as cancelled', async () => { + const window = loadJobs(); + const { provider } = makeProvider({ providerId: 'provider.cancelled' }); + await dispatch(window, 'register-provider', { provider }); + const enqueued = await dispatch(window, 'enqueue', enqueuePayload({ providerId: provider.providerId })); + + const result = window.slopsmith.jobs.cancelled(provider.providerId, enqueued.payload.job.jobId, { safeReason: 'User stopped backend conversion' }); + const snapshot = diagnosticsSnapshot(window); + + assert.equal(result.outcome, 'cancelled'); + assert.equal(snapshot.jobs.active.length, 0); + assert.equal(snapshot.jobs.recentTerminal[0].state, 'cancelled'); + assert.equal(snapshot.jobs.recentTerminal[0].terminalOutcome.category, 'cancellation'); +}); + test('list and inspect are prompt-free and do not invoke provider callbacks', async () => { const window = loadJobs(); const { calls, provider } = makeProvider(); diff --git a/tests/js/jobs_test_harness.js b/tests/js/jobs_test_harness.js index d05228e..fad79e6 100644 --- a/tests/js/jobs_test_harness.js +++ b/tests/js/jobs_test_harness.js @@ -100,8 +100,13 @@ function enqueuePayload(overrides = {}) { target: overrides.target || { targetRef: 'song-1' }, inputs: overrides.inputs || { safeFingerprint: 'input-1' }, safeLabel: overrides.safeLabel || 'Build playable cache', + jobId: overrides.jobId, logicalJobKey: overrides.logicalJobKey, providerId: overrides.providerId, + state: overrides.state, + status: overrides.status, + providerPayload: overrides.providerPayload, + privatePayload: overrides.privatePayload, privileged: overrides.privileged, approvalScopeKey: overrides.approvalScopeKey, };