Isolate scheduled tasks by empNo, aligned with uds-auth workspaces.

Super/fallback admins see all jobs grouped by user folder; admin/user only see their own. Stamp ownership on create, filter HTTP/tools, and prefer the owner's provisioned cwd on fire.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-09-08 15:49:54 +08:00
parent 0ffc48b43a
commit 7d27179e51
12 changed files with 785 additions and 141 deletions

View file

@ -23,6 +23,16 @@ import {
storePath,
upsertJob,
} from './store.js'
import {
assertCanAccessJob,
canViewAllJobs,
filterJobsForIdentity,
filterRunsForJobs,
migrateJobOwners,
normalizeOwnerEmpNo,
UNASSIGNED_OWNER,
viewerPayload,
} from './ownership.js'
export const PLUGIN_NAME = 'dsh-ops-cron'
export const API_PREFIX = '/dsh-ops-cron'
@ -68,24 +78,68 @@ function parseCookieHeader(header, name) {
return null
}
/** Browser UDS / fallback login cookies (validated more strictly by uds-auth when present). */
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')
}
function requireBrowserLogin(request, write) {
/**
* 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<object|null>} identity or null after writing error response
*/
async function requireIdentity(request, write, getUdsAuth) {
if (!isTrustedApiRequest(request)) {
write(403, { ok: false, error: 'forbidden' })
return false
return null
}
if (!browserEmpNo(request)) {
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 false
return null
}
return true
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
if (owner && uds?.isUserPath?.(owner, raw)) return raw
if (ownerPath) return ownerPath
return ''
}
function isTrustedApiRequest(request) {
@ -162,6 +216,7 @@ export function createHostService(options = {}) {
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
@ -215,10 +270,19 @@ export function createHostService(options = {}) {
return { job: jobView(job), run: runView(run), decision: claimedDecision }
}
async function createJob(input) {
async function createJob(input, identity = null) {
const t = now()
// Ignore client-supplied ids on create — otherwise POST/tools can overwrite.
const { id: _ignoredId, ...safeInput } = input && typeof input === 'object' ? input : {}
const { id: _ignoredId, ownerEmpNo: _ignoreOwner, ...safeInput } = input && typeof input === 'object' ? input : {}
if (identity?.empNo) {
safeInput.ownerEmpNo = identity.empNo
safeInput.ownerDisplayName = identity.displayName || identity.empNo
safeInput.cwd = sanitizeJobCwd(safeInput.cwd, identity, getUdsAuth, { forOwnerEmpNo: identity.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)
@ -227,7 +291,7 @@ export function createHostService(options = {}) {
return jobView(created)
}
async function updateJob(jobId, patch) {
async function updateJob(jobId, patch, identity = null) {
const t = now()
const state = await withState((current) => {
const job = getJob(current, jobId)
@ -236,6 +300,7 @@ export function createHostService(options = {}) {
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,
@ -254,6 +319,26 @@ export function createHostService(options = {}) {
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
@ -268,8 +353,8 @@ export function createHostService(options = {}) {
return jobView(getJob(state, jobId))
}
async function pauseJob(jobId, enabled) {
return updateJob(jobId, { enabled })
async function pauseJob(jobId, enabled, identity = null) {
return updateJob(jobId, { enabled }, identity)
}
async function deleteJob(jobId) {
@ -388,8 +473,20 @@ export function createHostService(options = {}) {
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)
@ -430,14 +527,16 @@ export function createHostService(options = {}) {
}
if (path === `${API_PREFIX}/settings` && method === 'GET') {
if (!requireBrowserLogin(req, write)) return
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') {
if (!requireBrowserLogin(req, write)) return
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 })
@ -445,7 +544,8 @@ export function createHostService(options = {}) {
}
if (path === `${API_PREFIX}/models` && method === 'GET') {
if (!requireBrowserLogin(req, write)) return
const identity = await requireIdentity(req, write, getUdsAuth)
if (!identity) return
const catalog = typeof sessionPort?.listModels === 'function'
? await sessionPort.listModels()
: { groups: [], current: null }
@ -454,7 +554,8 @@ export function createHostService(options = {}) {
}
if (path === `${API_PREFIX}/presets` && method === 'GET') {
if (!requireBrowserLogin(req, write)) return
const identity = await requireIdentity(req, write, getUdsAuth)
if (!identity) return
const catalog = typeof sessionPort?.listPresets === 'function'
? await sessionPort.listPresets()
: { items: [], current: null }
@ -463,16 +564,27 @@ export function createHostService(options = {}) {
}
if (path === `${API_PREFIX}/workspaces` && method === 'GET') {
if (!requireBrowserLogin(req, write)) return
const workspaces = typeof sessionPort?.listWorkspaces === 'function'
const identity = await requireIdentity(req, write, getUdsAuth)
if (!identity) return
let workspaces = typeof sessionPort?.listWorkspaces === 'function'
? await sessionPort.listWorkspaces()
: []
write(200, { ok: true, workspaces: Array.isArray(workspaces) ? workspaces : [] })
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') {
if (!requireBrowserLogin(req, write)) return
const identity = await requireIdentity(req, write, getUdsAuth)
if (!identity) return
const dshIm = getDshIm()
if (!dshIm || typeof dshIm.listDeliveryCatalog !== 'function') {
write(200, {
@ -502,16 +614,23 @@ export function createHostService(options = {}) {
}
if (path === `${API_PREFIX}/jobs` && method === 'GET') {
if (!requireBrowserLogin(req, write)) return
const identity = await requireIdentity(req, write, getUdsAuth)
if (!identity) return
const state = await snapshot()
write(200, { ok: true, jobs: listJobs(state).map(jobView) })
const jobs = filterJobsForIdentity(listJobs(state), identity).map(jobView)
write(200, {
ok: true,
jobs,
viewer: viewerPayload(identity),
})
return
}
if (path === `${API_PREFIX}/jobs` && method === 'POST') {
if (!requireBrowserLogin(req, write)) return
const identity = await requireIdentity(req, write, getUdsAuth)
if (!identity) return
const body = await readJsonBody(req)
const job = await createJob(body)
const job = await createJob(body, identity)
write(200, { ok: true, job })
return
}
@ -520,18 +639,18 @@ export function createHostService(options = {}) {
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) {
if (!requireBrowserLogin(req, write)) return
const state = await snapshot()
const job = getJob(state, jobId)
if (!job) return write(404, { ok: false, error: 'job not found' })
write(200, { ok: true, job: jobView(job) })
write(200, { ok: true, job: jobView(existing) })
return
}
if (!requireBrowserLogin(req, write)) return
if (method === 'PATCH' && !rest) {
const body = await readJsonBody(req)
const job = await updateJob(jobId, body)
const job = await updateJob(jobId, body, identity)
write(200, { ok: true, job })
return
}
@ -546,12 +665,12 @@ export function createHostService(options = {}) {
return
}
if (method === 'POST' && rest === '/pause') {
const job = await pauseJob(jobId, false)
const job = await pauseJob(jobId, false, identity)
write(200, { ok: true, job })
return
}
if (method === 'POST' && rest === '/resume') {
const job = await pauseJob(jobId, true)
const job = await pauseJob(jobId, true, identity)
write(200, { ok: true, job })
return
}
@ -559,11 +678,13 @@ export function createHostService(options = {}) {
const openMatch = path.match(new RegExp(`^${API_PREFIX}/runs/([^/]+)/open$`))
if (openMatch && method === 'POST') {
if (!requireBrowserLogin(req, write)) return
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)
@ -574,7 +695,8 @@ export function createHostService(options = {}) {
const adoptMatch = path.match(new RegExp(`^${API_PREFIX}/sessions/([^/]+)/adopt$`))
if (adoptMatch && method === 'POST') {
if (!requireBrowserLogin(req, write)) return
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') {
@ -586,22 +708,30 @@ export function createHostService(options = {}) {
}
if (path === `${API_PREFIX}/conceal` && method === 'POST') {
if (!requireBrowserLogin(req, write)) return
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') {
if (!requireBrowserLogin(req, write)) return
const identity = await requireIdentity(req, write, getUdsAuth)
if (!identity) return
const state = await snapshot()
const jobId = url.searchParams.get('jobId') || undefined
write(200, { ok: true, runs: listHistory(state, jobId).map(runView) })
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') {
if (!requireBrowserLogin(req, write)) return
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)
@ -611,7 +741,8 @@ export function createHostService(options = {}) {
}
if (path === `${API_PREFIX}/workspace-visible` && method === 'GET') {
if (!requireBrowserLogin(req, write)) return
const identity = await requireIdentity(req, write, getUdsAuth)
if (!identity) return
const state = await snapshot()
const listed = url.searchParams.getAll('id')
write(200, {
@ -626,7 +757,8 @@ export function createHostService(options = {}) {
} catch (error) {
const code = error && error.code
if (code === 'NOT_FOUND') return write(404, { ok: false, error: error.message })
if (code === 'INVALID_CRON' || code === 'INVALID_AT' || code === 'INVALID_SCHEDULE' || code === 'INVALID_JOB' || code === 'INVALID_TIMEZONE') {
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' })
@ -649,15 +781,25 @@ export function createHostService(options = {}) {
stopTimer,
handleRequest,
snapshot,
async listJobs() {
return listJobs(await snapshot()).map(jobView)
getUdsAuth,
async listJobs(identity = null) {
const jobs = listJobs(await snapshot())
const filtered = identity ? filterJobsForIdentity(jobs, identity) : jobs
return filtered.map(jobView)
},
async getJob(jobId) {
async getJob(jobId, identity = null) {
const job = getJob(await snapshot(), jobId)
return job ? jobView(job) : null
if (!job) return null
if (identity) assertCanAccessJob(job, identity)
return jobView(job)
},
async listHistory(jobId) {
return listHistory(await snapshot(), jobId).map(runView)
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
@ -748,15 +890,27 @@ function sessionCwdOf(ctx, sessionId) {
}
/**
* Prefer the user's most recently ordered Workspace so a scheduled session
* (and a later native fork) can occupy a real sidebar group. Official
* attachSession requires header.cwd === workspace.path.
* 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 = {}) {
const requested = String(job?.cwd || '').trim()
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 }
@ -1026,7 +1180,8 @@ export function makeLiveSessionPort(ctx) {
throw new Error('ctx.agents.create is unavailable')
}
const sessionId = randomUUID()
const placement = resolveSessionPlacement(ctx, job)
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)
@ -1062,6 +1217,14 @@ export function makeLiveSessionPort(ctx) {
} 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', {