/** * Host service used by apply() and by tests. * Owns the durable store, the single timer, CRUD, and the fire path. */ import { randomUUID } from 'node:crypto' import { mkdir } from 'node:fs/promises' import { homedir } from 'node:os' import { basename, join } from 'node:path' import { claimOccurrence, executeClaimedRun, extractAssistantText, interruptActiveRuns, publicJob, settleRun, TITLE_PREFIX } from './fire.js' import { deliverRunToIm, mergeDeliveryMention, mirrorRunToSession, normalizeOrigin } from './delivery.js' import { workspaceVisibleIds } from './isolation.js' import { decideDispatch, nextFire, validateSchedule } from './scheduler.js' import { createJobRecord, createStore, DEFAULT_SETTINGS, getJob, listHistory, listJobs, normalizeSettings, removeJob, storePath, upsertJob, } from './store.js' import { assertCanAccessJob, canViewAllJobs, claimUnassignedForViewer, filterJobsForIdentity, filterRunsForJobs, migrateJobOwners, normalizeOwnerEmpNo, UNASSIGNED_OWNER, viewerPayload, } from './ownership.js' export const PLUGIN_NAME = 'dsh-ops-cron' export const API_PREFIX = '/dsh-ops-cron' export function dshHome() { return process.env.DSH_HOME || join(homedir(), '.dsh') } export function defaultCwd() { return join(dshHome(), 'ops-cron', 'workspace') } function json(res, status, body) { if (res.headersSent || res.writableEnded) return res.writeHead(status, { 'content-type': 'application/json; charset=utf-8', 'cache-control': 'no-store', }) res.end(JSON.stringify(body)) } function parseUrl(req) { try { return new URL(req.url || '/', 'http://dsh.local') } catch { return new URL('http://dsh.local/') } } function parseCookieHeader(header, name) { if (!header || typeof header !== 'string') return null for (const part of header.split(';')) { const idx = part.indexOf('=') if (idx < 0) continue if (part.slice(0, idx).trim() !== name) continue try { return decodeURIComponent(part.slice(idx + 1).trim()) } catch { return part.slice(idx + 1).trim() } } return null } function empNoFromUserWorkspacePath(cwd) { const norm = String(cwd || '').replace(/\\/g, '/') const match = norm.match(/\/user-workspaces\/([^/]+)(?:\/|$)/) return match ? decodeURIComponent(match[1]) : null } function inferIdentityFromJobInput(input, getUdsAuth) { const uds = typeof getUdsAuth === 'function' ? getUdsAuth() : null if (!uds) return null const origin = input?.origin const sessionId = origin?.kind === 'web' && typeof origin.sessionId === 'string' ? origin.sessionId.trim() : '' if (sessionId && typeof uds.getSessionOwner === 'function') { const empNo = uds.getSessionOwner(sessionId) if (empNo) { return { empNo: String(empNo), displayName: String(empNo), permissions: {}, } } } const fromCwd = empNoFromUserWorkspacePath(input?.cwd) if (fromCwd) { return { empNo: fromCwd, displayName: fromCwd, permissions: {}, } } return null } function browserEmpNo(request) { const cookie = request?.headers?.cookie || '' return parseCookieHeader(cookie, 'PORTALSSOUser') || parseCookieHeader(cookie, 'ZTEDPGSSOUser') || parseCookieHeader(cookie, 'UDS_FALLBACK_USER') || parseCookieHeader(cookie, 'UDS_FALLBACK_UI') } /** * Resolve browser identity via uds-auth. Returns null when missing/unavailable. * @param {object} request * @param {() => object|undefined} getUdsAuth */ async function resolveBrowserIdentity(request, getUdsAuth) { const uds = typeof getUdsAuth === 'function' ? getUdsAuth() : null if (!uds || typeof uds.resolveRequestIdentity !== 'function') return null try { const identity = await uds.resolveRequestIdentity(request) if (!identity?.empNo) return null const workspacePath = typeof uds.getProvisionedWorkspacePath === 'function' ? uds.getProvisionedWorkspacePath(identity.empNo) : null return { ...identity, workspacePath: workspacePath || identity.workspacePath || null } } catch { return null } } /** * @returns {Promise} identity or null after writing error response */ async function requireIdentity(request, write, getUdsAuth) { if (!isTrustedApiRequest(request)) { write(403, { ok: false, error: 'forbidden' }) return null } const uds = typeof getUdsAuth === 'function' ? getUdsAuth() : null if (!uds || typeof uds.resolveRequestIdentity !== 'function') { write(503, { ok: false, error: 'auth_unavailable', message: 'uds-auth 未就绪,无法使用定时任务' }) return null } const identity = await resolveBrowserIdentity(request, getUdsAuth) if (!identity?.empNo) { write(401, { ok: false, error: 'login_required', message: '登录后才能使用定时任务' }) return null } return identity } function sanitizeJobCwd(cwd, identity, getUdsAuth, { forOwnerEmpNo } = {}) { const raw = typeof cwd === 'string' ? cwd.trim() : '' const uds = typeof getUdsAuth === 'function' ? getUdsAuth() : null const owner = forOwnerEmpNo || identity?.empNo const ownerPath = owner && uds?.getProvisionedWorkspacePath ? uds.getProvisionedWorkspacePath(owner) : null if (!raw) return ownerPath || '' if (canViewAllJobs(identity)) return raw // Without uds-auth, keep the caller cwd (tools may stamp owner from path only). if (!uds) return raw if (owner && uds?.isUserPath?.(owner, raw)) return raw if (ownerPath) return ownerPath return '' } function isTrustedApiRequest(request) { const host = request.headers.host ?? '' if (!host) return false const hostname = host.split(':')[0].replace(/^\[|\]$/g, '') if ((request.headers['sec-fetch-site'] ?? '') === 'cross-site') return false // Allow same-origin LAN / non-loopback Host (login still required separately). const site = request.headers['sec-fetch-site'] ?? '' if (site === 'same-origin' || site === 'same-site') return true const origin = request.headers.origin if (origin !== undefined && origin !== 'null') { try { if (new URL(origin).host !== host) return false } catch { return false } } return hostname === '127.0.0.1' || hostname === 'localhost' || hostname === '::1' || hostname === '0.0.0.0' } async function readJsonBody(req) { const chunks = [] let size = 0 const max = 256 * 1024 for await (const chunk of req) { size += chunk.length if (size > max) { const error = new Error('payload too large') error.code = 'PAYLOAD_TOO_LARGE' throw error } chunks.push(chunk) } const text = Buffer.concat(chunks).toString('utf8') if (!text) return {} return JSON.parse(text) } function jobView(job) { return publicJob(job) } function runView(run) { if (!run) return null return { id: run.id, jobId: run.jobId, scheduledAt: run.scheduledAt, actualAt: run.actualAt, status: run.status, trigger: run.trigger, sessionId: run.sessionId, error: run.error, summary: run.summary, reason: run.reason, } } /** * @param {object} options * @param {() => number} [options.now] * @param {string} [options.filePath] * @param {object} [options.sessionPort] { createAndPrompt, archiveSession, waitForTurn } * @param {() => object|undefined} [options.getDshIm] soft-injected proactive IM API * @param {() => object|undefined} [options.getAgents] soft-injected agents service for session mirror * @param {{ warn?: Function, info?: Function }} [options.logger] * @param {number} [options.tickIntervalMs] */ export function createHostService(options = {}) { const now = options.now || (() => Date.now()) const filePath = options.filePath || storePath(dshHome()) const store = options.store || createStore({ filePath }) const sessionPort = options.sessionPort || null const getDshIm = typeof options.getDshIm === 'function' ? options.getDshIm : () => undefined const getAgents = typeof options.getAgents === 'function' ? options.getAgents : () => undefined const getUdsAuth = typeof options.getUdsAuth === 'function' ? options.getUdsAuth : () => undefined const logger = options.logger || console const tickIntervalMs = Number(options.tickIntervalMs) > 0 ? Number(options.tickIntervalMs) : 15_000 let timer = null let ticking = false async function withState(fn) { return store.mutate(fn) } async function snapshot() { return store.read() } async function maybeDeliverIm(job, run) { if (!job || job.delivery?.kind !== 'im') return try { await deliverRunToIm(job, run?.summary, { dshIm: getDshIm() }) logger.info?.(`[dsh-ops-cron] im delivery sent for job ${job.id} → ${job.delivery.targetId}`) } catch (error) { logger.warn?.(`[dsh-ops-cron] im delivery failed for job ${job.id}: ${error instanceof Error ? error.message : error}`) } } async function maybeMirrorResult(job, run) { if (!job || !run || run.status === 'skipped') return try { const result = await mirrorRunToSession(job, run.summary || run.error || '', { dshIm: getDshIm(), getAgents, pluginName: PLUGIN_NAME, newId: () => randomUUID(), }) if (result?.mirrored) { logger.info?.(`[dsh-ops-cron] mirrored run ${run.id} → session ${result.sessionId} via ${result.via} (${result.method}${result.resumed ? ', resumed' : ''})`) } else if (result?.reason && result.reason !== 'no_target') { logger.warn?.(`[dsh-ops-cron] mirror skipped for run ${run.id}: ${result.reason}${result.sessionId ? ` session=${result.sessionId}` : ''}${result.error ? ` (${result.error})` : ''}`) } else if (!result?.skipped) { logger.warn?.(`[dsh-ops-cron] mirror incomplete for run ${run.id}: ${result?.reason || 'unknown'}`) } } catch (error) { logger.warn?.(`[dsh-ops-cron] mirror failed for job ${job.id}: ${error instanceof Error ? error.message : error}`) } } async function settleAndNotify(jobId, runId, claimedDecision) { const state = await snapshot() const run = (state.runs || []).find((row) => row.id === runId) const job = getJob(state, jobId) await maybeDeliverIm(job, run) await maybeMirrorResult(job, run) return { job: jobView(job), run: runView(run), decision: claimedDecision } } async function createJob(input, identity = null) { const t = now() // Ignore client-supplied ids on create — otherwise POST/tools can overwrite. const { id: _ignoredId, ownerEmpNo: _ignoreOwner, ...safeInput } = input && typeof input === 'object' ? input : {} let ownerIdentity = identity?.empNo ? identity : inferIdentityFromJobInput(safeInput, getUdsAuth) if (ownerIdentity?.empNo) { safeInput.ownerEmpNo = ownerIdentity.empNo safeInput.ownerDisplayName = ownerIdentity.displayName || ownerIdentity.empNo safeInput.cwd = sanitizeJobCwd(safeInput.cwd, ownerIdentity, getUdsAuth, { forOwnerEmpNo: ownerIdentity.empNo }) } else if (safeInput.ownerEmpNo) { safeInput.ownerEmpNo = normalizeOwnerEmpNo(safeInput.ownerEmpNo) } else { safeInput.ownerEmpNo = UNASSIGNED_OWNER } let created await withState((current) => { created = createJobRecord(safeInput, current, t) return upsertJob(current, created) }) return jobView(created) } async function updateJob(jobId, patch, identity = null) { const t = now() const state = await withState((current) => { const job = getJob(current, jobId) if (!job) { const error = new Error('job not found') error.code = 'NOT_FOUND' throw error } if (identity) patch = { ...patch, _identity: identity } const nextInput = { id: job.id, name: patch.name !== undefined ? patch.name : job.name, prompt: patch.prompt !== undefined ? patch.prompt : job.prompt, enabled: patch.enabled !== undefined ? patch.enabled : job.enabled, schedule: patch.schedule !== undefined ? patch.schedule : job.schedule, cwd: patch.cwd !== undefined ? patch.cwd : job.cwd, timeoutMinutes: patch.timeoutMinutes !== undefined ? patch.timeoutMinutes : job.timeoutMinutes, provider: patch.provider !== undefined ? patch.provider : job.provider, model: patch.model !== undefined ? patch.model : job.model, reasoningEffort: patch.reasoningEffort !== undefined ? patch.reasoningEffort : job.reasoningEffort, agentPreset: patch.agentPreset !== undefined ? patch.agentPreset : job.agentPreset, delivery: patch.delivery !== undefined ? mergeDeliveryMention(job.delivery, patch.delivery) : job.delivery, origin: patch.origin !== undefined ? (normalizeOrigin(patch.origin) || undefined) : job.origin, ownerEmpNo: job.ownerEmpNo || UNASSIGNED_OWNER, ownerDisplayName: job.ownerDisplayName || '', } if (patch._identity) { nextInput.cwd = sanitizeJobCwd( nextInput.cwd, patch._identity, getUdsAuth, { forOwnerEmpNo: nextInput.ownerEmpNo }, ) } if (patch.ownerEmpNo !== undefined && canViewAllJobs(patch._identity)) { nextInput.ownerEmpNo = normalizeOwnerEmpNo(patch.ownerEmpNo) if (patch.ownerDisplayName !== undefined) { nextInput.ownerDisplayName = String(patch.ownerDisplayName || '').trim().slice(0, 80) } else if (nextInput.ownerEmpNo !== job.ownerEmpNo) { nextInput.ownerDisplayName = nextInput.ownerEmpNo === UNASSIGNED_OWNER ? '' : (patch.ownerDisplayName || nextInput.ownerEmpNo) } } const record = createJobRecord(nextInput, current, t) record.createdAt = job.createdAt record.lastRunAt = job.lastRunAt record.lastStatus = job.lastStatus if (patch.origin === undefined && job.origin) record.origin = job.origin if (patch.schedule === undefined && patch.enabled === undefined) { record.nextRunAt = job.nextRunAt } return upsertJob(current, record) }) return jobView(getJob(state, jobId)) } async function pauseJob(jobId, enabled, identity = null) { return updateJob(jobId, { enabled }, identity) } async function deleteJob(jobId) { await withState((current) => { if (!getJob(current, jobId)) { const error = new Error('job not found') error.code = 'NOT_FOUND' throw error } return removeJob(current, jobId) }) return { ok: true, id: jobId } } async function updateSettings(patch) { const state = await withState((current) => ({ ...current, settings: normalizeSettings({ ...current.settings, ...patch }, current.settings), })) return state.settings } async function dispatchRun(jobId, trigger) { let claimed const t = now() await withState((current) => { claimed = claimOccurrence(current, jobId, t, trigger, current.settings) if (claimed.decision.action === 'wait') return current return claimed.state }) if (!claimed.run) return { job: claimed.job, run: null, decision: claimed.decision } if (claimed.run.status === 'skipped') { return { job: claimed.job, run: runView(claimed.run), decision: claimed.decision } } const executed = await store.mutate(async (current) => { const result = await executeClaimedRun(current, claimed.run.id, { now, createAndPrompt: sessionPort?.createAndPrompt, archiveSession: sessionPort?.archiveSession, }) return result.state }) let run = (executed.runs || []).find((row) => row.id === claimed.run.id) if (run?.status === 'running' && typeof sessionPort?.waitForTurn === 'function') { let terminal try { terminal = await sessionPort.waitForTurn(run.sessionId) } catch (error) { const message = error instanceof Error ? error.message : String(error) terminal = { status: 'failed', error: message, summary: message } } const settled = await store.mutate((current) => settleRun(current, run.id, terminal, now())) run = (settled.runs || []).find((row) => row.id === claimed.run.id) return settleAndNotify(jobId, run.id, claimed.decision) } return settleAndNotify(jobId, claimed.run.id, claimed.decision) } async function tick() { if (ticking) return [] ticking = true const fired = [] try { const state = await snapshot() if (state.settings?.enabled === false) return fired const t = now() for (const job of listJobs(state)) { if (job.enabled === false) continue const runs = (state.runs || []).filter((run) => run.jobId === job.id) const decision = decideDispatch({ job, runs, now: t, overlapPolicy: state.settings.overlapPolicy, misfirePolicy: state.settings.misfirePolicy, }) if (decision.action === 'wait') continue const result = await dispatchRun(job.id, 'schedule') if (result.run) fired.push(result) } return fired } finally { ticking = false } } function startTimer() { if (timer) return timer = setInterval(() => { tick().catch(() => {}) }, tickIntervalMs) if (typeof timer.unref === 'function') timer.unref() } function stopTimer() { if (!timer) return clearInterval(timer) timer = null } async function concealKnownSessions() { if (typeof sessionPort?.archiveSession !== 'function') return 0 const state = await snapshot() const ids = new Set() for (const id of state.hiddenSessionIds || []) if (id) ids.add(id) for (const run of state.runs || []) if (run?.sessionId) ids.add(run.sessionId) let hidden = 0 for (const id of ids) { try { await sessionPort.archiveSession(id) hidden += 1 } catch { // Session may already be gone. } } return hidden } async function migrateOwnersIfNeeded() { const uds = getUdsAuth() await withState((current) => { const { state, changed } = migrateJobOwners(current, { getSessionOwner: (sessionId) => uds?.getSessionOwner?.(sessionId) || null, }) return changed ? state : current }) } async function recover() { await mkdir(join(dshHome(), 'ops-cron'), { recursive: true }) await migrateOwnersIfNeeded() const t = now() await withState((current) => { let next = interruptActiveRuns(current, t) next = { ...next, jobs: (next.jobs || []).map((job) => { try { if (job.nextRunAt === null) return job if (Number.isFinite(job.nextRunAt) && job.nextRunAt > t) return job const nextRunAt = nextFire(job.schedule, t, job.schedule?.timezone || next.settings.timezone) return { ...job, nextRunAt } } catch { return job } }), } return next }) await concealKnownSessions() } async function handleRequest(req, res) { const url = parseUrl(req) const path = url.pathname.replace(/\/+$/, '') || '/' const method = (req.method || 'GET').toUpperCase() const write = (status, body) => json(res, status, body) try { if (path === `${API_PREFIX}/health` && method === 'GET') { const state = await snapshot() write(200, { ok: true, plugin: PLUGIN_NAME, enabled: state.settings.enabled !== false, jobCount: (state.jobs || []).length, }) return } if (path === `${API_PREFIX}/settings` && method === 'GET') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const state = await snapshot() write(200, { ok: true, settings: state.settings }) return } if (path === `${API_PREFIX}/settings` && method === 'PUT') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const body = await readJsonBody(req) const settings = await updateSettings(body) write(200, { ok: true, settings }) return } if (path === `${API_PREFIX}/models` && method === 'GET') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const catalog = typeof sessionPort?.listModels === 'function' ? await sessionPort.listModels() : { groups: [], current: null } write(200, { ok: true, ...catalog }) return } if (path === `${API_PREFIX}/presets` && method === 'GET') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const catalog = typeof sessionPort?.listPresets === 'function' ? await sessionPort.listPresets() : { items: [], current: null } write(200, { ok: true, ...catalog }) return } if (path === `${API_PREFIX}/workspaces` && method === 'GET') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return let workspaces = typeof sessionPort?.listWorkspaces === 'function' ? await sessionPort.listWorkspaces() : [] if (!Array.isArray(workspaces)) workspaces = [] const uds = getUdsAuth() if (!canViewAllJobs(identity) && uds?.isUserPath) { workspaces = workspaces.filter((row) => uds.isUserPath(identity.empNo, row?.path)) } if (!canViewAllJobs(identity) && workspaces.length === 0) { const pathOnly = identity.workspacePath || uds?.getProvisionedWorkspacePath?.(identity.empNo) if (pathOnly) workspaces = [{ id: null, path: pathOnly, name: identity.displayName || identity.empNo }] } write(200, { ok: true, workspaces, viewer: viewerPayload(identity) }) return } if (path === `${API_PREFIX}/im-catalog` && method === 'GET') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const dshIm = getDshIm() if (!dshIm || typeof dshIm.listDeliveryCatalog !== 'function') { write(200, { ok: true, available: false, options: [], hint: 'dsh-im-ops missing or outdated — install ≥ops.24 for delivery picker', }) return } try { const options = await dshIm.listDeliveryCatalog() write(200, { ok: true, available: true, options: Array.isArray(options) ? options : [], }) } catch (error) { write(200, { ok: true, available: false, options: [], hint: error instanceof Error ? error.message : String(error), }) } return } if (path === `${API_PREFIX}/jobs` && method === 'GET') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const uds = getUdsAuth() const state = await withState((current) => { const { state: next, changed } = claimUnassignedForViewer(current, identity, { getSessionOwner: typeof uds?.getSessionOwner === 'function' ? (sessionId) => uds.getSessionOwner(sessionId) : undefined, }) return changed ? next : current }) const jobs = filterJobsForIdentity(listJobs(state), identity).map(jobView) write(200, { ok: true, jobs, viewer: viewerPayload(identity), }) return } if (path === `${API_PREFIX}/jobs` && method === 'POST') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const body = await readJsonBody(req) const job = await createJob(body, identity) write(200, { ok: true, job }) return } const jobMatch = path.match(new RegExp(`^${API_PREFIX}/jobs/([^/]+)(/run|/pause|/resume)?$`)) if (jobMatch) { const jobId = decodeURIComponent(jobMatch[1]) const rest = jobMatch[2] || '' const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const state = await snapshot() const existing = getJob(state, jobId) assertCanAccessJob(existing, identity) if (method === 'GET' && !rest) { write(200, { ok: true, job: jobView(existing) }) return } if (method === 'PATCH' && !rest) { const body = await readJsonBody(req) const job = await updateJob(jobId, body, identity) write(200, { ok: true, job }) return } if (method === 'DELETE' && !rest) { await deleteJob(jobId) write(200, { ok: true, id: jobId }) return } if (method === 'POST' && rest === '/run') { const result = await dispatchRun(jobId, 'run-now') write(200, { ok: true, ...result }) return } if (method === 'POST' && rest === '/pause') { const job = await pauseJob(jobId, false, identity) write(200, { ok: true, job }) return } if (method === 'POST' && rest === '/resume') { const job = await pauseJob(jobId, true, identity) write(200, { ok: true, job }) return } } const openMatch = path.match(new RegExp(`^${API_PREFIX}/runs/([^/]+)/open$`)) if (openMatch && method === 'POST') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const runId = decodeURIComponent(openMatch[1]) const state = await snapshot() const run = (state.runs || []).find((row) => row.id === runId) if (!run?.sessionId) return write(404, { ok: false, error: 'run not found' }) assertCanAccessJob(getJob(state, run.jobId), identity) let unarchived = false if (typeof sessionPort?.revealSession === 'function') { unarchived = await sessionPort.revealSession(run.sessionId) } write(200, { ok: true, sessionId: run.sessionId, runId: run.id, unarchived: unarchived !== false }) return } const adoptMatch = path.match(new RegExp(`^${API_PREFIX}/sessions/([^/]+)/adopt$`)) if (adoptMatch && method === 'POST') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const sessionId = decodeURIComponent(adoptMatch[1]) if (!sessionId) return write(400, { ok: false, error: 'sessionId required' }) if (typeof sessionPort?.adoptSession !== 'function') { return write(200, { ok: false, attached: false, sessionId }) } const result = await sessionPort.adoptSession(sessionId) write(200, { ok: true, sessionId, ...result }) return } if (path === `${API_PREFIX}/conceal` && method === 'POST') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const hidden = await concealKnownSessions() write(200, { ok: true, hidden }) return } if (path === `${API_PREFIX}/history` && method === 'GET') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const state = await snapshot() const jobId = url.searchParams.get('jobId') || undefined if (jobId) { assertCanAccessJob(getJob(state, jobId), identity) } const visibleJobs = filterJobsForIdentity(listJobs(state), identity) const runs = filterRunsForJobs(listHistory(state, jobId), visibleJobs).map(runView) write(200, { ok: true, runs, viewer: viewerPayload(identity) }) return } if (path === `${API_PREFIX}/preview` && method === 'POST') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const body = await readJsonBody(req) const settings = (await snapshot()).settings const schedule = validateSchedule(body.schedule || body, body.timezone || settings.timezone) const nextRunAt = nextFire(schedule, now(), schedule.timezone) write(200, { ok: true, nextRunAt, schedule }) return } if (path === `${API_PREFIX}/workspace-visible` && method === 'GET') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return const state = await snapshot() const listed = url.searchParams.getAll('id') write(200, { ok: true, hiddenSessionIds: state.hiddenSessionIds || [], visible: workspaceVisibleIds(listed, state.hiddenSessionIds), }) return } write(404, { ok: false, error: 'not found' }) } catch (error) { const code = error && error.code if (code === 'NOT_FOUND') return write(404, { ok: false, error: error.message }) if (code === 'LOGIN_REQUIRED') return write(401, { ok: false, error: 'login_required', message: error.message }) if (code === 'INVALID_CRON' || code === 'INVALID_AT' || code === 'INVALID_SCHEDULE' || code === 'INVALID_JOB' || code === 'INVALID_TIMEZONE' || code === 'INVALID_CWD') { return write(400, { ok: false, error: error.message, code }) } if (code === 'PAYLOAD_TOO_LARGE') return write(413, { ok: false, error: 'payload too large' }) write(500, { ok: false, error: error instanceof Error ? error.message : 'internal error' }) } } return { store, now, createJob, updateJob, pauseJob, deleteJob, updateSettings, dispatchRun, tick, recover, startTimer, stopTimer, handleRequest, snapshot, getUdsAuth, async listJobs(identity = null) { const jobs = listJobs(await snapshot()) const filtered = identity ? filterJobsForIdentity(jobs, identity) : jobs return filtered.map(jobView) }, async getJob(jobId, identity = null) { const job = getJob(await snapshot(), jobId) if (!job) return null if (identity) assertCanAccessJob(job, identity) return jobView(job) }, async listHistory(jobId, identity = null) { const state = await snapshot() if (jobId && identity) assertCanAccessJob(getJob(state, jobId), identity) const runs = listHistory(state, jobId) if (!identity) return runs.map(runView) const visible = filterJobsForIdentity(listJobs(state), identity) return filterRunsForJobs(runs, visible).map(runView) }, workspaceVisibleIds(allIds) { const hidden = store.snapshot().hiddenSessionIds return workspaceVisibleIds(allIds, hidden) }, } } function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)) } async function withTimeout(promise, ms, message) { let timer try { return await Promise.race([ promise, new Promise((_, reject) => { timer = setTimeout(() => { const error = new Error(message) error.code = 'RUN_TIMEOUT' reject(error) }, ms) }), ]) } finally { clearTimeout(timer) } } /** * Wait until a live agent finishes the turn started by followup. * `whenIdle` alone can resolve before the driver wakes; poll for `running` first. */ export async function waitForAgentTurn(agent, options = {}) { if (!agent) { return { status: 'failed', error: 'agent missing', summary: 'Agent disappeared after dispatch' } } const startTimeoutMs = Number.isFinite(options.startTimeoutMs) ? options.startTimeoutMs : 30_000 const turnTimeoutMs = Number.isFinite(options.turnTimeoutMs) ? options.turnTimeoutMs : 10 * 60 * 1000 const started = Date.now() let sawRunning = agent.status === 'running' while (!sawRunning && Date.now() - started < startTimeoutMs) { if (agent.status === 'running') { sawRunning = true break } await sleep(20) } if (agent.status === 'running') sawRunning = true if (sawRunning && typeof agent.whenIdle === 'function') { await withTimeout(agent.whenIdle(), turnTimeoutMs, 'run timed out') return { status: 'succeeded', summary: 'Turn finished' } } return { status: 'failed', error: 'agent_never_started', summary: 'Scheduled prompt was queued but the agent never started a turn', } } function workspacePathOf(workspace) { return String(workspace?.path || workspace?.record?.path || '').trim() } function listWorkspaces(ctx) { const registry = tryGet(ctx, 'workspaceRegistry') const list = typeof registry?.list === 'function' ? registry.list() : [] return Array.isArray(list) ? list : [] } export function listWorkspaceChoices(ctx) { return listWorkspaces(ctx).map((workspace) => { const path = workspacePathOf(workspace) const title = String(workspace?.title || workspace?.record?.title || '').trim() || (path ? basename(path) : '') return { id: String(workspace?.id || path), title: title || path, path, } }).filter((row) => row.path) } function sessionCwdOf(ctx, sessionId) { const live = tryGet(ctx, 'sessions')?.get?.(sessionId) const agent = tryGet(ctx, 'agents')?.get?.(sessionId) return String(live?.header?.cwd || live?.cwd || agent?.session?.header?.cwd || agent?.session?.cwd || '').trim() } /** * Prefer the job owner's provisioned workspace (uds-auth), then explicit cwd, * then recent workspace, then shared ops-cron fallback. * Official attachSession requires header.cwd === workspace.path. */ export function resolveSessionPlacement(ctx, job = {}, deps = {}) { const workspaces = listWorkspaces(ctx) const match = (path) => (path && workspaces.find((row) => workspacePathOf(row) === path)) || null const uds = deps.udsAuth || (typeof deps.getUdsAuth === 'function' ? deps.getUdsAuth() : tryGet(ctx, 'udsAuth')) const ownerEmpNo = String(job?.ownerEmpNo || '').trim() const ownerPath = ownerEmpNo && !ownerEmpNo.startsWith('__') && uds?.getProvisionedWorkspacePath ? uds.getProvisionedWorkspacePath(ownerEmpNo) : null let requested = String(job?.cwd || '').trim() if (requested && ownerEmpNo && !ownerEmpNo.startsWith('__') && uds?.isUserPath) { if (!uds.isUserPath(ownerEmpNo, requested) && ownerPath) { requested = ownerPath } } if (requested) return { cwd: requested, workspace: match(requested) } if (ownerPath) return { cwd: ownerPath, workspace: match(ownerPath) } const recent = workspaces[0] const recentPath = workspacePathOf(recent) if (recentPath) return { cwd: recentPath, workspace: recent } const isolated = defaultCwd() return { cwd: isolated, workspace: match(isolated) } } export async function attachLiveSessionToWorkspace(workspace, sessionId) { if (!sessionId || !workspace || typeof workspace.attachSession !== 'function') return false try { await workspace.attachSession(sessionId) return true } catch { return false } } /** * Put a forked (or still-loose) session into the workspace whose path matches * its cwd. Does not create a new workspace: a cwd mismatch cannot join a * different project. */ export async function adoptSessionIntoWorkspace(ctx, sessionId) { if (!sessionId) return { ok: false, attached: false } const cwd = sessionCwdOf(ctx, sessionId) const workspaces = listWorkspaces(ctx) const workspace = (cwd && workspaces.find((row) => workspacePathOf(row) === cwd)) || null const attached = await attachLiveSessionToWorkspace(workspace, sessionId) if (attached) { try { await unarchiveSession(ctx, sessionId) } catch { /* listing membership is enough */ } } return { ok: attached, attached, sessionId, cwd: cwd || null, workspaceId: workspace?.id || null, } } export async function archiveLiveSession(ctx, sessionId) { if (!sessionId) return false const registry = tryGet(ctx, 'workspaceRegistry') if (!registry || typeof registry.archiveSession !== 'function') return false try { await registry.archiveSession(sessionId) return true } catch { return false } } export async function unarchiveSession(ctx, sessionId) { if (!sessionId) return false let registry = tryGet(ctx, 'workspaceRegistry') if (!registry) { const started = Date.now() while (!registry && Date.now() - started < 1500) { await sleep(50) registry = tryGet(ctx, 'workspaceRegistry') } } if (!registry) return false if (typeof registry.enqueueOperation === 'function' && typeof registry.requireState === 'function' && typeof registry.setState === 'function') { await registry.enqueueOperation(async () => { const state = registry.requireState() const archived = state.archivedSessionIds || [] if (!archived.includes(sessionId)) return await registry.setState({ ...state, archivedSessionIds: archived.filter((id) => id !== sessionId), }) }) return true } const ids = typeof registry.archivedSessionIds === 'function' ? registry.archivedSessionIds() : registry.archivedSessionIds if (Array.isArray(ids) && typeof registry.setState === 'function') { if (!ids.includes(sessionId)) return true const state = typeof registry.requireState === 'function' ? registry.requireState() : { archivedSessionIds: ids } await registry.setState({ ...state, archivedSessionIds: ids.filter((id) => id !== sessionId), }) return true } return false } /** * Snapshot the same default model a New Session uses. * Persona templates interpolate `{{model}}`; an empty value fails assembly. */ export function currentDefaultModel(ctx) { try { const selection = tryGet(ctx, 'agentDefaultModel')?.currentSelection?.() const provider = typeof selection?.provider === 'string' ? selection.provider.trim() : '' const model = typeof selection?.model === 'string' ? selection.model.trim() : '' if (!provider || !model) return null return { provider, model, ...selection.reasoningEffort === undefined ? {} : { reasoningEffort: selection.reasoningEffort }, } } catch { return null } } export async function listModelChoices(ctx) { const llm = tryGet(ctx, 'llm') const providers = typeof llm?.listProviders === 'function' ? llm.listProviders() : [] const list = Array.isArray(providers) ? providers : [] const groups = [] for (const row of list) { const provider = String(row?.id || row?.provider || '').trim() if (!provider) continue let models = [] try { models = typeof llm.listModels === 'function' ? await llm.listModels(provider) : [] } catch { models = [] } groups.push({ provider, displayName: String(row?.name || row?.displayName || provider), models: (Array.isArray(models) ? models : []).map((entry) => ({ id: String(entry?.id || entry?.model || '').trim(), name: String(entry?.name || entry?.displayName || entry?.id || '').trim(), })).filter((entry) => entry.id), }) } return { groups, current: currentDefaultModel(ctx) } } export async function listPresetChoices(ctx) { const presets = tryGet(ctx, 'agentPresets') if (!presets || typeof presets.list !== 'function') { return { items: [], current: null } } let items = [] try { const listed = await presets.list() const rows = Array.isArray(listed) ? listed : (Array.isArray(listed?.items) ? listed.items : []) items = rows.map((row) => ({ id: String(row?.id || '').trim(), name: String(row?.name || row?.displayName || row?.id || '').trim(), })).filter((row) => row.id) } catch { items = [] } let current = null if (typeof presets.resolve === 'function') { try { const resolved = await presets.resolve() const id = typeof resolved?.id === 'string' ? resolved.id.trim() : '' if (id) current = { id, name: String(resolved?.name || id) } } catch { current = null } } return { items, current } } export async function resolveJobModel(ctx, job) { const provider = typeof job?.provider === 'string' ? job.provider.trim() : '' const model = typeof job?.model === 'string' ? job.model.trim() : '' if (provider && model) { return { provider, model, ...job.reasoningEffort ? { reasoningEffort: job.reasoningEffort } : {}, } } return resolveDefaultModel(ctx) } export async function resolveDefaultModel(ctx, options = {}) { const waitMs = Number.isFinite(options.waitMs) ? Math.max(0, options.waitMs) : 1500 let service = tryGet(ctx, 'agentDefaultModel') if (!service?.currentSelection && waitMs > 0) { const started = Date.now() while (!service?.currentSelection && Date.now() - started < waitMs) { await sleep(50) service = tryGet(ctx, 'agentDefaultModel') } } const selection = service?.currentSelection?.() const provider = typeof selection?.provider === 'string' ? selection.provider.trim() : '' const model = typeof selection?.model === 'string' ? selection.model.trim() : '' if (!provider || !model) { throw new Error('no default model is configured; pick a model in Models before running scheduled tasks') } return { provider, model, ...selection.reasoningEffort === undefined ? {} : { reasoningEffort: selection.reasoningEffort }, } } /** * Fill `{{provider}}` / `{{model}}` and route LLM requests, without importing * `@deepseek-ai/dsh-agent`. Same contract as official installModelSelection. */ export function bindModelSelection(agentCtx, selection) { if (!agentCtx || typeof agentCtx.on !== 'function' || !selection) return const snapshot = { provider: selection.provider, model: selection.model, ...selection.reasoningEffort === undefined ? {} : { reasoningEffort: selection.reasoningEffort }, } agentCtx.on('system-prompt/assemble', async (...args) => { const next = args.find((arg) => typeof arg === 'function') const assembled = next ? await next() : (args[0] || {}) return { ...assembled, variables: { ...assembled.variables, provider: snapshot.provider, model: snapshot.model, }, } }) agentCtx.on('agent/request', async (...args) => { const next = args.find((arg) => typeof arg === 'function') const resolved = next ? await next() : (args[0] || {}) const { reasoningEffort: _inherited, ...rest } = resolved || {} return { ...rest, provider: snapshot.provider, model: snapshot.model, ...snapshot.reasoningEffort === undefined ? {} : { reasoningEffort: snapshot.reasoningEffort }, } }) } async function composeCronAgent(ctx, selection, job) { const presets = tryGet(ctx, 'agentPresets') if (!presets || typeof presets.resolve !== 'function' || typeof presets.mount !== 'function') { return { setup: (agentCtx) => { bindModelSelection(agentCtx, selection) }, } } const wanted = typeof job?.agentPreset === 'string' ? job.agentPreset.trim() : '' const resolved = wanted ? await presets.resolve(wanted) : await presets.resolve() const presetId = resolved?.id return { agentPreset: presetId, setup: async (agentCtx) => { bindModelSelection(agentCtx, selection) if (presetId) await presets.mount(agentCtx, presetId) }, } } export function makeLiveSessionPort(ctx) { const handles = new Map() return { async createAndPrompt({ job, run, text, source }) { const agents = tryGet(ctx, 'agents') if (!agents || typeof agents.create !== 'function') { throw new Error('ctx.agents.create is unavailable') } const sessionId = randomUUID() const udsAuth = tryGet(ctx, 'udsAuth') const placement = resolveSessionPlacement(ctx, job, { udsAuth }) const cwd = placement.cwd || defaultCwd() await mkdir(cwd, { recursive: true }) const selection = await resolveJobModel(ctx, job) const composition = await composeCronAgent(ctx, selection, job) const handle = await agents.create({ sessionId, agentOptions: { provider: selection.provider, model: selection.model, ...selection.reasoningEffort === undefined ? {} : { reasoningEffort: selection.reasoningEffort }, }, meta: { cwd, ...composition.agentPreset ? { agentPreset: composition.agentPreset } : {}, }, setup: composition.setup, }) const agent = handle?.agent if (!agent || typeof agent.followup !== 'function') { throw new Error('created agent has no followup') } const message = { id: randomUUID(), role: 'user', content: [{ type: 'text', text }], source: source || { kind: 'plugin', plugin: PLUGIN_NAME }, } agent.followup(message) const timeoutMs = Math.max(60_000, (Number(job?.timeoutMinutes) || 10) * 60_000) handles.set(sessionId, { handle, timeoutMs, agent }) try { await attachLiveSessionToWorkspace(placement.workspace, sessionId) } catch { // Forks of unattached runs stay loose; listing hide is separate. } try { const owner = String(job?.ownerEmpNo || '').trim() if (owner && !owner.startsWith('__')) { udsAuth?.stampSessionOwner?.(sessionId, owner) } } catch { // Ownership stamp is best-effort. } try { if (typeof agent.session?.append === 'function') { agent.session.append('session/title', { title: `${TITLE_PREFIX}${job.name}`, source: 'plugin', }) } } catch { // Title is best-effort. } return { sessionId, handle, status: 'running', summary: `Started session for ${job.name}`, } }, async archiveSession(sessionId) { return archiveLiveSession(ctx, sessionId) }, async revealSession(sessionId) { return unarchiveSession(ctx, sessionId) }, async adoptSession(sessionId) { return adoptSessionIntoWorkspace(ctx, sessionId) }, async listWorkspaces() { return listWorkspaceChoices(ctx) }, async listModels() { return listModelChoices(ctx) }, async listPresets() { return listPresetChoices(ctx) }, async waitForTurn(sessionId) { const entry = handles.get(sessionId) const agents = tryGet(ctx, 'agents') const agent = entry?.agent || entry?.handle?.agent || agents?.get?.(sessionId) try { const finished = await waitForAgentTurn(agent, { turnTimeoutMs: entry?.timeoutMs, }) const summary = extractAssistantText(agent?.session?.deriveMessages?.() || []) return { ...finished, summary: summary || finished.summary, } } catch (error) { if (error && error.code === 'RUN_TIMEOUT' && agent && typeof agent.cancel === 'function') { try { agent.cancel({ kind: 'timeout' }) } catch { /* ignore */ } } throw error } finally { handles.delete(sessionId) } }, } } function tryGet(ctx, name) { try { return ctx.get(name) } catch { return undefined } } export { DEFAULT_SETTINGS }