From 2e13432b0a378bd90975bd623ee84a6d614b16ae Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 10 Sep 2026 15:29:04 +0800 Subject: [PATCH] Show cron state/labels in sidebar and push archive history. Surface lifecycle badges and editable monitor fields in the UI, and POST terminal runs to persist_history archive endpoints on settle/tick. Co-authored-by: Cursor --- lib/archive.js | 113 +++++++++++++++++++++ lib/client.js | 231 +++++++++++++++++++++++++++++++++++++++++-- lib/host.js | 79 +++++++++++++-- test/archive.test.js | 90 +++++++++++++++++ 4 files changed, 500 insertions(+), 13 deletions(-) create mode 100644 lib/archive.js create mode 100644 test/archive.test.js diff --git a/lib/archive.js b/lib/archive.js new file mode 100644 index 0000000..51010a1 --- /dev/null +++ b/lib/archive.js @@ -0,0 +1,113 @@ +/** + * Push terminal cron runs to an external archive endpoint. + * Endpoint receives JSON: { plugin, archivedAt, job, runs }. + */ + +export function runsNeedingArchive(runs, jobId) { + return (runs || []).filter((run) => ( + run + && run.jobId === jobId + && run.status !== 'queued' + && run.status !== 'running' + && !run.archivedAt + )) +} + +/** + * @param {string} endpoint + * @param {{ job: object, runs: object[] }} payload + * @param {{ fetchImpl?: typeof fetch, now?: () => number }} [opts] + */ +export async function pushArchiveBatch(endpoint, payload, opts = {}) { + const url = String(endpoint || '').trim() + if (!url) { + const error = new Error('archive endpoint is empty') + error.code = 'INVALID_ARCHIVE' + throw error + } + if (!/^https?:\/\//i.test(url)) { + const error = new Error('archive endpoint must be http(s) URL') + error.code = 'INVALID_ARCHIVE' + throw error + } + const fetchImpl = opts.fetchImpl || globalThis.fetch + if (typeof fetchImpl !== 'function') { + const error = new Error('fetch is not available for archive push') + error.code = 'ARCHIVE_UNAVAILABLE' + throw error + } + const body = { + plugin: 'dsh-ops-cron', + archivedAt: new Date((opts.now || Date.now)()).toISOString(), + job: payload.job, + runs: payload.runs, + } + const res = await fetchImpl(url, { + method: 'POST', + headers: { + 'content-type': 'application/json; charset=utf-8', + accept: 'application/json', + }, + body: JSON.stringify(body), + }) + if (!res.ok) { + const text = typeof res.text === 'function' ? await res.text().catch(() => '') : '' + const error = new Error(`archive endpoint returned ${res.status}${text ? `: ${text.slice(0, 200)}` : ''}`) + error.code = 'ARCHIVE_FAILED' + error.status = res.status + throw error + } + let parsed = null + try { + parsed = typeof res.json === 'function' ? await res.json() : null + } catch { + parsed = null + } + return { ok: true, status: res.status, body: parsed } +} + +/** + * Mark runs archived in state (pure). + */ +export function markRunsArchived(state, runIds, now) { + const ids = new Set(runIds || []) + if (!ids.size) return state + return { + ...state, + runs: (state.runs || []).map((run) => ( + ids.has(run.id) ? { ...run, archivedAt: now } : run + )), + } +} + +/** + * After archive, drop old archived terminals for archive-policy jobs, + * keeping the newest `keep` archived + all non-archived + all active. + */ +export function pruneArchivedRuns(state, keep = 50) { + const jobs = state.jobs || [] + const byJob = new Map(jobs.map((job) => [job.id, job])) + const active = [] + const pending = [] + const archivedByJob = new Map() + for (const run of state.runs || []) { + if (!run) continue + if (run.status === 'queued' || run.status === 'running') { + active.push(run) + continue + } + const job = byJob.get(run.jobId) + if (job?.persistHistory?.kind === 'archive' && run.archivedAt) { + if (!archivedByJob.has(run.jobId)) archivedByJob.set(run.jobId, []) + archivedByJob.get(run.jobId).push(run) + continue + } + pending.push(run) + } + const keptArchived = [] + for (const rows of archivedByJob.values()) { + rows.sort((a, b) => (b.archivedAt || b.actualAt || 0) - (a.archivedAt || a.actualAt || 0)) + keptArchived.push(...rows.slice(0, Math.max(0, keep))) + } + return { ...state, runs: [...active, ...pending, ...keptArchived] } +} diff --git a/lib/client.js b/lib/client.js index 82ed51d..2d0a9c1 100644 --- a/lib/client.js +++ b/lib/client.js @@ -232,7 +232,7 @@ window.__ModuleLoader__.load({ } function jobStamp(rows) { - return (rows || []).map((row) => `${row.id}:${row.updatedAt}:${row.nextRunAt}:${row.enabled}:${row.lastStatus}:${row.model}:${row.provider}:${row.delivery?.kind}:${row.delivery?.targetId}`).join('|') + return (rows || []).map((row) => `${row.id}:${row.updatedAt}:${row.nextRunAt}:${row.enabled}:${row.lastStatus}:${row.state}:${row.stuck}:${row.model}:${row.provider}:${row.delivery?.kind}:${row.delivery?.targetId}:${JSON.stringify(row.labels || {})}`).join('|') } function runStamp(rows) { @@ -352,6 +352,12 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba .dsh-ct-cwdHint{color:var(--dsw-alias-label-tertiary);font-size:12px;line-height:1.4} .dsh-ct-check{flex-direction:row!important;align-items:center;gap:8px!important} .dsh-ct-check input{width:auto;margin:0} +.dsh-ct-state{display:inline-flex;align-items:center;gap:4px;font-size:11px;font-weight:600;letter-spacing:.02em;padding:1px 6px;border-radius:999px;border:1px solid var(--dsw-alias-border-l2);color:var(--dsw-alias-label-secondary);vertical-align:middle;margin-left:6px} +.dsh-ct-state[data-state=running],.dsh-ct-state[data-state=pending]{color:var(--dsw-alias-state-business-primary);border-color:currentColor} +.dsh-ct-state[data-state=succeeded]{color:var(--dsw-alias-state-success-primary, #1a7f37);border-color:currentColor} +.dsh-ct-state[data-state=failed],.dsh-ct-state[data-stuck=true]{color:var(--dsw-alias-state-error-primary);border-color:currentColor} +.dsh-ct-state[data-state=paused]{opacity:.75} +.dsh-ct-labels{color:var(--dsw-alias-label-tertiary);font-size:11px;font-weight:400;margin-left:6px} .dsh-ct-editorActs{display:flex;gap:8px;margin-top:8px;flex-wrap:wrap} .dsh-ct-primary{appearance:none;font:inherit;border:none;border-radius:8px;padding:8px 14px;background:var(--dsw-alias-label-primary);color:var(--dsw-alias-bg-layer-3);cursor:pointer;font-size:13px} .dsh-ct-secondary{appearance:none;font:inherit;border:1px solid var(--dsw-alias-border-l2);border-radius:8px;padding:8px 14px;background:transparent;color:var(--dsw-alias-label-secondary);cursor:pointer;font-size:13px} @@ -377,6 +383,27 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba expandRuns: '展开记录', collapseRuns: '收起记录', search: '搜索', searchPlaceholder: '搜索任务', searchClear: '清除搜索', searchEmpty: '没有匹配的任务。', paused: '已暂停', + state: '状态', + stateIdle: '空闲', + statePending: '排队', + stateRunning: '运行中', + stateSucceeded: '成功', + stateFailed: '失败', + statePaused: '已暂停', + stateStuck: '疑似卡死', + labels: '标签', + labelsHint: '每行 key=value,或 JSON 对象。用于批量监控与发现。', + labelsPlaceholder: 'role=worker\ntask=theory', + persistHistory: '历史保留', + persistRetain: '保留最近 N 条(跟随全局)', + persistForever: '永久保留', + persistArchive: '归档到远程 endpoint', + archiveEndpoint: '归档 URL', + archiveEndpointHint: 'persist_history=archive 时 POST 终态运行到该 http(s) 地址。', + watch: '监听绑定 (JSON)', + watchHint: '可选。例 {"taskId":"..."} 或 {"labels":{"task":"theory"},"match":"all"}', + progressFile: '进度文件', + progressFileHint: '相对工作目录或绝对路径;任务运行时写入 JSON 快照。', scheduleTz: '时区', cwd: '工作目录', timeout: '超时(分钟)', cwdRecent: '归属工作区(默认)', cwdCustom: '自定义路径…', cwdPlaceholder: '/absolute/path', @@ -423,6 +450,27 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba expandRuns: 'Show runs', collapseRuns: 'Hide runs', search: 'Search', searchPlaceholder: 'Search jobs', searchClear: 'Clear search', searchEmpty: 'No matching jobs.', paused: 'Paused', + state: 'State', + stateIdle: 'Idle', + statePending: 'Pending', + stateRunning: 'Running', + stateSucceeded: 'Succeeded', + stateFailed: 'Failed', + statePaused: 'Paused', + stateStuck: 'Stuck', + labels: 'Labels', + labelsHint: 'One key=value per line, or a JSON object. Used for discovery and batch monitoring.', + labelsPlaceholder: 'role=worker\ntask=theory', + persistHistory: 'History retention', + persistRetain: 'Retain last N (global setting)', + persistForever: 'Keep forever', + persistArchive: 'Archive to remote endpoint', + archiveEndpoint: 'Archive URL', + archiveEndpointHint: 'When archive is selected, terminal runs are POSTed to this http(s) URL.', + watch: 'Watch binding (JSON)', + watchHint: 'Optional. e.g. {"taskId":"..."} or {"labels":{"task":"theory"},"match":"all"}', + progressFile: 'Progress file', + progressFileHint: 'Path relative to cwd or absolute; the job writes a JSON snapshot while running.', scheduleTz: 'Time zone', cwd: 'Working directory', timeout: 'Timeout (minutes)', cwdRecent: 'Owner workspace (default)', cwdCustom: 'Custom path…', cwdPlaceholder: '/absolute/path', @@ -705,6 +753,11 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba imBotId: '', imTargetId: '', mirrorToSession: false, + labelsText: '', + persistKind: 'retain', + archiveEndpoint: '', + watchText: '', + progressFile: '', } } @@ -726,6 +779,11 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba imBotId: job.delivery?.botId || '', imTargetId: job.delivery?.targetId || '', mirrorToSession: job.mirrorToSession === true, + labelsText: formatLabelsText(job.labels), + persistKind: persistHistoryKind(job), + archiveEndpoint: job.persistHistory?.kind === 'archive' ? (job.persistHistory.endpoint || '') : '', + watchText: job.watch ? JSON.stringify(job.watch, null, 2) : '', + progressFile: job.progress?.channel?.file || '', } } @@ -736,6 +794,72 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba return t('deliveryDsh') } + function jobStateKey(job) { + if (job?.enabled === false) return 'paused' + return String(job?.state || 'idle') + } + + function jobStateLabel(job, t) { + const key = jobStateKey(job) + const map = { + idle: t('stateIdle'), + pending: t('statePending'), + running: t('stateRunning'), + succeeded: t('stateSucceeded'), + failed: t('stateFailed'), + paused: t('statePaused'), + } + const base = map[key] || key + return job?.stuck ? `${base} · ${t('stateStuck')}` : base + } + + function formatLabelsText(labels) { + if (!labels || typeof labels !== 'object') return '' + return Object.entries(labels) + .filter(([k]) => k) + .map(([k, v]) => `${k}=${v}`) + .join('\n') + } + + function parseLabelsText(raw) { + const text = String(raw || '').trim() + if (!text) return {} + if (text.startsWith('{')) { + try { + const parsed = JSON.parse(text) + if (parsed && typeof parsed === 'object' && !Array.isArray(parsed)) { + const out = {} + for (const [k, v] of Object.entries(parsed)) { + if (!k) continue + out[String(k)] = String(v ?? '') + } + return out + } + } catch { /* fall through */ } + } + const out = {} + for (const line of text.split(/[\n,]+/)) { + const part = line.trim() + if (!part) continue + const idx = part.indexOf('=') + if (idx <= 0) continue + out[part.slice(0, idx).trim()] = part.slice(idx + 1).trim() + } + return out + } + + function persistHistoryKind(job) { + const kind = job?.persistHistory?.kind + if (kind === 'forever' || kind === 'archive') return kind + return 'retain' + } + + function labelsInline(job) { + const labels = job?.labels + if (!labels || !Object.keys(labels).length) return '' + return Object.entries(labels).map(([k, v]) => `${k}=${v}`).join(' ') + } + function modelKey(provider, model) { if (!provider || !model) return '' return `${provider}::${model}` @@ -771,6 +895,8 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba job.model, job.ownerEmpNo, job.ownerDisplayName, + job.state, + labelsInline(job), ...(runs.filter((run) => run.jobId === job.id).map((run) => run.summary || run.status)), ].join(' ').toLowerCase() return hay.includes(needle) @@ -799,7 +925,11 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba 'data-on': selectedRun ? 'true' : 'false', onClick: () => onSelectRun(job.id, run), }, - h('span', { className: skin.title }, (run.summary || run.status || '').split('\n')[0] || run.status), + h('span', { className: skin.title }, + (run.state || run.status || '').toString(), + ' · ', + (run.summary || '').split('\n')[0] || '', + ), h('span', { className: skin.time }, formatTime(run.actualAt || run.scheduledAt, job.schedule?.timezone)), ) }) @@ -834,9 +964,18 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba h('span', { className: joinClass(skin.arrow, open ? skin.arrowOpen : '') }, chevronIcon() || '▸')), h('span', { className: skin.projectText }, - h('span', { className: skin.title }, paused - ? `${job.name} · ${t('paused')}` - : `${job.name} · ${jobDeliveryLabel(job, t)}${jobModelLabel(job) ? ` · ${jobModelLabel(job)}` : ''}`)), + h('span', { className: skin.title }, + job.name, + h('span', { + className: 'dsh-ct-state', + 'data-state': jobStateKey(job), + 'data-stuck': job.stuck ? 'true' : 'false', + }, jobStateLabel(job, t)), + labelsInline(job) + ? h('span', { className: 'dsh-ct-labels' }, labelsInline(job)) + : null, + ), + ), h('span', { className: skin.rowActs, onClick: (e) => e.stopPropagation() }, h('button', { type: 'button', className: skin.rowIcon, title: t('runNow'), onClick: () => onRun(job.id) }, playIcon() || '▶'), @@ -1158,11 +1297,22 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba ) } - function JobEditor({ t, form, setForm, isNew, error, lastOutput, nextRunAt, lastRunAt, workspaces, catalog, presets, imCatalog, viewer, onSave, onRun, onPause, onRemove }) { + function JobEditor({ t, form, setForm, isNew, error, lastOutput, nextRunAt, lastRunAt, jobMeta, workspaces, catalog, presets, imCatalog, viewer, onSave, onRun, onPause, onRemove }) { return h('div', { className: 'dsh-ct-editor' }, h('h1', null, isNew ? t('newJob') : (form.name || t('title'))), h('p', { className: 'dsh-ct-editorLead' }, t('editorLead')), error ? h('p', { className: 'dsh-ct-error' }, error) : null, + !isNew && jobMeta + ? h('p', { className: 'dsh-ct-editorLead' }, + h('span', { + className: 'dsh-ct-state', + 'data-state': jobStateKey(jobMeta), + 'data-stuck': jobMeta.stuck ? 'true' : 'false', + }, `${t('state')}: ${jobStateLabel(jobMeta, t)}`), + jobMeta.watch ? ` · watch` : '', + jobMeta.progress?.channel?.file ? ` · progress` : '', + ) + : null, h('label', null, t('name'), h('input', { value: form.name, onChange: (e) => setForm({ ...form, name: e.target.value }) })), h('label', null, t('prompt'), h('textarea', { value: form.prompt, onChange: (e) => setForm({ ...form, prompt: e.target.value }) })), h('div', { className: 'dsh-ct-editorRow' }, @@ -1191,6 +1341,52 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba h(CwdField, { t, form, setForm, workspaces, viewer }), h(ModelField, { t, form, setForm, catalog }), h(PresetField, { t, form, setForm, presets }), + h('label', null, t('labels'), + h('textarea', { + value: form.labelsText || '', + placeholder: t('labelsPlaceholder'), + onChange: (e) => setForm({ ...form, labelsText: e.target.value }), + rows: 3, + }), + ), + h('span', { className: 'dsh-ct-cwdHint' }, t('labelsHint')), + h('label', null, t('persistHistory'), + h('select', { + value: form.persistKind || 'retain', + onChange: (e) => setForm({ ...form, persistKind: e.target.value }), + }, + h('option', { value: 'retain' }, t('persistRetain')), + h('option', { value: 'forever' }, t('persistForever')), + h('option', { value: 'archive' }, t('persistArchive')), + ), + ), + form.persistKind === 'archive' + ? h(React.Fragment, null, + h('label', null, t('archiveEndpoint'), h('input', { + value: form.archiveEndpoint || '', + placeholder: 'https://example.com/cron-archive', + onChange: (e) => setForm({ ...form, archiveEndpoint: e.target.value }), + })), + h('span', { className: 'dsh-ct-cwdHint' }, t('archiveEndpointHint')), + ) + : null, + h('label', null, t('watch'), + h('textarea', { + value: form.watchText || '', + placeholder: '{"labels":{"task":"theory"},"match":"all"}', + onChange: (e) => setForm({ ...form, watchText: e.target.value }), + rows: 3, + }), + ), + h('span', { className: 'dsh-ct-cwdHint' }, t('watchHint')), + h('label', null, t('progressFile'), + h('input', { + value: form.progressFile || '', + placeholder: 'progress.json', + onChange: (e) => setForm({ ...form, progressFile: e.target.value }), + }), + ), + h('span', { className: 'dsh-ct-cwdHint' }, t('progressFileHint')), h('label', null, t('delivery'), h('select', { value: form.deliveryKind || 'dsh', @@ -1474,6 +1670,28 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba ? { kind: 'im', botId: form.imBotId, targetId: form.imTargetId } : { kind: 'dsh' }, mirrorToSession: form.mirrorToSession === true, + labels: parseLabelsText(form.labelsText), + persistHistory: form.persistKind === 'forever' + ? { kind: 'forever' } + : form.persistKind === 'archive' + ? { kind: 'archive', endpoint: String(form.archiveEndpoint || '').trim() } + : { kind: 'retain' }, + } + const watchRaw = String(form.watchText || '').trim() + if (watchRaw) { + try { + payload.watch = JSON.parse(watchRaw) + } catch { + throw new Error('watch must be valid JSON') + } + } else { + payload.watch = null + } + const progressFile = String(form.progressFile || '').trim() + if (progressFile) { + payload.progress = { channel: { file: progressFile }, metrics: [] } + } else { + payload.progress = null } if (selection.type === 'job' && selection.jobId) { await api(`/jobs/${selection.jobId}`, { method: 'PATCH', body: JSON.stringify(payload) }) @@ -1519,6 +1737,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba t, form, setForm, isNew: selection.type === 'new', error, lastOutput, nextRunAt: selectedJob?.nextRunAt, lastRunAt: selectedJob?.lastRunAt, + jobMeta: selectedJob || null, workspaces, catalog, presets, diff --git a/lib/host.js b/lib/host.js index 3204280..2381e2e 100644 --- a/lib/host.js +++ b/lib/host.js @@ -18,6 +18,12 @@ import { readProgressSnapshot, selectJobsByWatch, } from './monitor.js' +import { + markRunsArchived, + pruneArchivedRuns, + pushArchiveBatch, + runsNeedingArchive, +} from './archive.js' import { decideDispatch, nextFire, validateSchedule } from './scheduler.js' import { createJobRecord, @@ -369,6 +375,7 @@ function runView(run) { summary: run.summary, reason: run.reason, output_ref: run.outputRef || (run.sessionId ? `session:${run.sessionId}` : null), + ...run.archivedAt ? { archivedAt: run.archivedAt } : {}, } } @@ -481,7 +488,57 @@ export function createHostService(options = {}) { await maybeDeliverIm(job, run) await maybeMirrorResult(job, run) await maybeNotifyWatchers(job, run) - return { job: jobView(job, state.runs), run: runView(run), decision: claimedDecision } + if (job?.persistHistory?.kind === 'archive' && job.persistHistory.endpoint) { + try { + await archiveJobRuns(jobId, { identity: null }) + } catch (error) { + logger.warn?.(`[dsh-ops-cron] archive push failed for ${jobId}: ${error instanceof Error ? error.message : error}`) + } + } + return { job: jobView(job, (await snapshot()).runs), run: runView(run), decision: claimedDecision } + } + + /** + * Push unarchived terminal runs for archive-policy jobs to their endpoint. + * @param {string|null} [jobId] limit to one job; null = all archive jobs + * @param {{ identity?: object|null, fetchImpl?: typeof fetch, limit?: number }} [opts] + */ + async function archiveJobRuns(jobId = null, opts = {}) { + const identity = opts.identity || null + const limit = Number.isInteger(opts.limit) && opts.limit > 0 ? Math.min(200, opts.limit) : 50 + const state = await snapshot() + let jobs = listJobs(state).filter((job) => ( + job?.persistHistory?.kind === 'archive' && job.persistHistory.endpoint + )) + if (jobId) jobs = jobs.filter((job) => job.id === jobId) + if (identity) { + for (const job of jobs) assertCanAccessJob(job, identity) + jobs = filterJobsForIdentity(jobs, identity) + } + const results = [] + for (const job of jobs) { + const live = await snapshot() + const batch = runsNeedingArchive(live.runs, job.id) + .sort((a, b) => (a.actualAt || a.scheduledAt || 0) - (b.actualAt || b.scheduledAt || 0)) + .slice(0, limit) + if (!batch.length) { + results.push({ jobId: job.id, archived: 0, skipped: true }) + continue + } + await pushArchiveBatch(job.persistHistory.endpoint, { + job: jobView(job, live.runs), + runs: batch.map(runView), + }, { fetchImpl: opts.fetchImpl || options.fetchImpl, now }) + const mark = now() + const ids = batch.map((run) => run.id) + await withState((current) => pruneArchivedRuns(markRunsArchived(current, ids, mark), 50)) + results.push({ jobId: job.id, archived: batch.length, endpoint: job.persistHistory.endpoint }) + } + return { + ok: true, + archived: results.reduce((sum, row) => sum + (row.archived || 0), 0), + results, + } } async function createJob(input, identity = null, opts = {}) { @@ -721,6 +778,11 @@ export function createHostService(options = {}) { if (result.run) fired.push(result) } await tickWatcherReports() + try { + await archiveJobRuns(null, { limit: 20 }) + } catch (error) { + logger.warn?.(`[dsh-ops-cron] archive tick failed: ${error instanceof Error ? error.message : error}`) + } return fired } finally { ticking = false @@ -1080,12 +1142,14 @@ export function createHostService(options = {}) { if (path === `${API_PREFIX}/archive` && method === 'POST') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return - // Reserved: persist_history=archive:endpoint will push here later. - write(200, { - ok: true, - archived: false, - message: 'archive endpoint reserved; set persist_history=archive: on jobs for future use', + const body = await readJsonBody(req).catch(() => ({})) + const jobId = typeof body.jobId === 'string' ? body.jobId.trim() + : (typeof body.task_id === 'string' ? body.task_id.trim() : '') + const result = await archiveJobRuns(jobId || null, { + identity, + limit: Number(body.limit) > 0 ? Number(body.limit) : 50, }) + write(200, { ok: true, ...result }) return } @@ -1118,7 +1182,7 @@ export function createHostService(options = {}) { const locale = resolveLocale(req) const code = error && error.code if (code === 'NOT_FOUND') return write(404, { ...apiError('not_found', locale), error: error.message || 'not_found' }) - if (code === 'INVALID_CRON' || code === 'INVALID_AT' || code === 'INVALID_SCHEDULE' || code === 'INVALID_JOB' || code === 'INVALID_TIMEZONE' || code === 'INVALID_CWD' || code === 'INVALID_DELIVERY' || code === 'IM_DELIVERY_FORBIDDEN' || code === 'IM_TARGET_FORBIDDEN' || code === 'IM_JOB_MISSING_CWD' || code === 'INVALID_WATCH' || code === 'INVALID_PROGRESS' || code === 'INVALID_REPORT' || code === 'INVALID_PERSIST_HISTORY' || code === 'INVALID_PROGRESS_PATH') { + if (code === 'INVALID_CRON' || code === 'INVALID_AT' || code === 'INVALID_SCHEDULE' || code === 'INVALID_JOB' || code === 'INVALID_TIMEZONE' || code === 'INVALID_CWD' || code === 'INVALID_DELIVERY' || code === 'IM_DELIVERY_FORBIDDEN' || code === 'IM_TARGET_FORBIDDEN' || code === 'IM_JOB_MISSING_CWD' || code === 'INVALID_WATCH' || code === 'INVALID_PROGRESS' || code === 'INVALID_REPORT' || code === 'INVALID_PERSIST_HISTORY' || code === 'INVALID_PROGRESS_PATH' || code === 'INVALID_ARCHIVE' || code === 'ARCHIVE_FAILED' || code === 'ARCHIVE_UNAVAILABLE') { return write(400, { ok: false, error: error.message, code, message: error.message }) } if (code === 'PAYLOAD_TOO_LARGE') return write(413, apiError('payload_too_large', locale)) @@ -1304,6 +1368,7 @@ export function createHostService(options = {}) { queryJobs, listRuns, getProgress, + archiveJobRuns, async listJobs(identity = null, query = null) { const state = await snapshot() let jobs = listJobs(state) diff --git a/test/archive.test.js b/test/archive.test.js new file mode 100644 index 0000000..e8196a6 --- /dev/null +++ b/test/archive.test.js @@ -0,0 +1,90 @@ +import assert from 'node:assert/strict' +import { test } from 'node:test' +import { + markRunsArchived, + pruneArchivedRuns, + pushArchiveBatch, + runsNeedingArchive, +} from '../lib/archive.js' +import { createHostService } from '../lib/host.js' +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' + +test('runsNeedingArchive skips active and already archived', () => { + const runs = [ + { id: '1', jobId: 'j', status: 'running' }, + { id: '2', jobId: 'j', status: 'succeeded' }, + { id: '3', jobId: 'j', status: 'failed', archivedAt: 1 }, + { id: '4', jobId: 'other', status: 'succeeded' }, + ] + assert.deepEqual(runsNeedingArchive(runs, 'j').map((r) => r.id), ['2']) +}) + +test('pushArchiveBatch posts JSON and mark/prune archive local copy', async () => { + const posts = [] + const result = await pushArchiveBatch('https://archive.example/cron', { + job: { id: 'j1', name: 'n' }, + runs: [{ id: 'r1', status: 'succeeded' }], + }, { + now: () => Date.parse('2026-09-10T08:00:00.000Z'), + fetchImpl: async (url, init) => { + posts.push({ url, init }) + return { ok: true, status: 200, async json() { return { received: 1 } } } + }, + }) + assert.equal(result.ok, true) + assert.equal(posts[0].url, 'https://archive.example/cron') + assert.match(posts[0].init.body, /"plugin":"dsh-ops-cron"/) + + let state = { + jobs: [{ id: 'j1', persistHistory: { kind: 'archive', endpoint: 'https://archive.example/cron' } }], + runs: Array.from({ length: 60 }, (_, i) => ({ + id: `r${i}`, + jobId: 'j1', + status: 'succeeded', + actualAt: i, + archivedAt: i < 55 ? 1000 + i : undefined, + })), + } + state = markRunsArchived(state, ['r59'], 9999) + state = pruneArchivedRuns(state, 10) + const archived = state.runs.filter((r) => r.archivedAt) + assert.ok(archived.length <= 10) + assert.ok(state.runs.some((r) => r.id === 'r59')) +}) + +test('host archiveJobRuns pushes and marks runs', async (t) => { + const dir = await mkdtemp(join(tmpdir(), 'dsh-cron-arch-')) + t.after(() => rm(dir, { recursive: true, force: true })) + const posts = [] + const service = createHostService({ + filePath: join(dir, 'store.json'), + now: () => Date.parse('2026-09-10T08:00:00.000Z'), + fetchImpl: async (url, init) => { + posts.push({ url, body: JSON.parse(init.body) }) + return { ok: true, status: 204, async json() { return null } } + }, + sessionPort: { + async createAndPrompt() { + return { sessionId: 's-arch', status: 'succeeded', summary: 'done' } + }, + }, + }) + const job = await service.createJob({ + name: 'arch-job', + prompt: 'work', + schedule: { kind: 'at', at: new Date(Date.parse('2026-09-10T08:01:00.000Z')).toISOString(), timezone: 'Asia/Shanghai' }, + cwd: dir, + persistHistory: { kind: 'archive', endpoint: 'https://archive.example/intake' }, + }) + await service.dispatchRun(job.id, 'run-now') + assert.ok(posts.length >= 1, 'settle should auto-archive') + assert.equal(posts[0].url, 'https://archive.example/intake') + assert.equal(posts[0].body.job.id, job.id) + const snap = await service.snapshot() + const run = (snap.runs || []).find((row) => row.jobId === job.id && row.status === 'succeeded') + assert.ok(run?.archivedAt) + const again = await service.archiveJobRuns(job.id, { limit: 10 }) + assert.equal(again.archived, 0) +})