From aff1372f33ea8c868637dd9bb1066edf4fa141ca Mon Sep 17 00:00:00 2001 From: jazelly Date: Sun, 26 Jul 2026 10:51:46 +0930 Subject: [PATCH 1/2] fix(cli): make daemon entry executable --- packages/cli/package.json | 2 +- packages/cli/src/daemon/daemon-entry.mts | 55 + packages/cli/src/daemon/errors.ts | 224 +++++ packages/cli/src/daemon/job-store.ts | 1170 ++++++++++++++++++++++ packages/cli/src/daemon/lifecycle.ts | 42 + packages/cli/src/daemon/server.ts | 468 +++++++++ 6 files changed, 1960 insertions(+), 1 deletion(-) create mode 100644 packages/cli/src/daemon/daemon-entry.mts create mode 100644 packages/cli/src/daemon/errors.ts create mode 100644 packages/cli/src/daemon/job-store.ts create mode 100644 packages/cli/src/daemon/lifecycle.ts create mode 100644 packages/cli/src/daemon/server.ts diff --git a/packages/cli/package.json b/packages/cli/package.json index 26314b7..b509abc 100644 --- a/packages/cli/package.json +++ b/packages/cli/package.json @@ -24,7 +24,7 @@ }, "scripts": { "build": "npm run build:js && npm run build:native", - "build:js": "node -e \"require('node:fs').rmSync('dist',{recursive:true,force:true})\" && tsc -p tsconfig.json && node -e \"require('node:fs').chmodSync('dist/src/tokenless.mjs',0o755)\"", + "build:js": "node -e \"require('node:fs').rmSync('dist',{recursive:true,force:true})\" && tsc -p tsconfig.json && node -e \"const fs=require('node:fs'); fs.chmodSync('dist/src/tokenless.mjs',0o755); fs.chmodSync('dist/src/daemon/daemon-entry.mjs',0o755)\"", "build:native": "node scripts/build-rust-binaries.mjs", "lint": "tsc -p tsconfig.json --noEmit", "prepack": "npm run build:js" diff --git a/packages/cli/src/daemon/daemon-entry.mts b/packages/cli/src/daemon/daemon-entry.mts new file mode 100644 index 0000000..2ffde20 --- /dev/null +++ b/packages/cli/src/daemon/daemon-entry.mts @@ -0,0 +1,55 @@ +#!/usr/bin/env node +import { nativeBinaryBuildInfo } from './server.js' +import { startDaemon } from './lifecycle.js' + +async function main() { + const args = process.argv.slice(2) + if (args.length === 1 && args[0] === '--tokenless-build-info') { + console.log(JSON.stringify(nativeBinaryBuildInfo('tokenless-daemon'))) + return + } + const options = parseArgs(args) + await startDaemon(options) +} + +function parseArgs(args: string[]) { + let homeDir: string | undefined + let host = '127.0.0.1' + let port = 7331 + const positional: string[] = [] + for (let index = 0; index < args.length; index += 1) { + const arg = args[index] + if (arg === undefined) continue + if (arg === '--home') { + homeDir = requireValue(args, ++index, '--home') + } else if (arg === '--host') { + host = requireValue(args, ++index, '--host') + } else if (arg === '--port') { + port = parsePort(requireValue(args, ++index, '--port')) + } else { + positional.push(arg) + } + } + if (positional.length === 1 && positional[0] === 'serve') { + return { homeDir, host, port } + } + if (positional.length === 0) return { homeDir, host, port } + throw new Error('usage: daemon-entry.mjs [--home ] [serve] [--host ] [--port ]') +} + +function requireValue(args: string[], index: number, flag: string) { + const value = args[index] + if (!value) throw new Error(`${flag} requires a value`) + return value +} + +function parsePort(value: string) { + const port = Number(value) + if (!Number.isInteger(port) || port < 0 || port > 65535) throw new Error('--port must be a valid TCP port') + return port +} + +main().catch((error) => { + console.error(error instanceof Error && error.message ? error.message : String(error)) + process.exit(1) +}) diff --git a/packages/cli/src/daemon/errors.ts b/packages/cli/src/daemon/errors.ts new file mode 100644 index 0000000..02516fe --- /dev/null +++ b/packages/cli/src/daemon/errors.ts @@ -0,0 +1,224 @@ +import { DAEMON_ERROR_PROTOCOL } from '../generated/protocol-constants.js' + +export type JobStatus = + | 'queued' + | 'claimed' + | 'running' + | 'waiting_for_user' + | 'succeeded' + | 'failed' + | 'canceled' + | 'timed_out' + +type DaemonErrorKind = + | 'io' + | 'random' + | 'sqlite' + | 'json' + | 'missing_home' + | 'invalid_input' + | 'non_loopback_bind' + | 'invalid_status' + | 'job_not_found' + | 'claim_rejected' + | 'claim_expired' + | 'bridge_busy' + | 'control_auth_missing' + | 'control_auth_rejected' + | 'invalid_job_state' + +export class DaemonError extends Error { + readonly kind: DaemonErrorKind + readonly jobId?: string | undefined + readonly statusValue?: string | undefined + readonly expected?: string | undefined + readonly actual?: JobStatus | undefined + readonly host?: string | undefined + readonly cause?: unknown + + constructor(kind: DaemonErrorKind, message: string, options: { + jobId?: string | undefined + statusValue?: string | undefined + expected?: string | undefined + actual?: JobStatus | undefined + host?: string | undefined + cause?: unknown + } = {}) { + super(message) + this.name = 'DaemonError' + this.kind = kind + this.jobId = options.jobId + this.statusValue = options.statusValue + this.expected = options.expected + this.actual = options.actual + this.host = options.host + this.cause = options.cause + } +} + +export function ioError(error: unknown) { + return new DaemonError('io', `I/O error: ${errorText(error)}`, { cause: error }) +} + +export function sqliteError(error: unknown) { + return new DaemonError('sqlite', `SQLite error: ${errorText(error)}`, { cause: error }) +} + +export function jsonError(error: unknown) { + return new DaemonError('json', `JSON error: ${errorText(error)}`, { cause: error }) +} + +export function missingHomeError() { + return new DaemonError( + 'missing_home', + 'cannot resolve Tokenless home; pass --home or set TOKENLESS_HOME/HOME' + ) +} + +export function invalidInput(message: string) { + return new DaemonError('invalid_input', `invalid input: ${message}`) +} + +export function nonLoopbackBind(host: string) { + return new DaemonError( + 'non_loopback_bind', + `refusing to bind daemon to non-loopback host ${host}; Tokenless daemon is a local control plane`, + { host } + ) +} + +export function invalidStatus(status: string) { + return new DaemonError('invalid_status', `invalid job status: ${status}`, { statusValue: status }) +} + +export function jobNotFound(jobId: string) { + return new DaemonError('job_not_found', `job not found: ${jobId}`, { jobId }) +} + +export function claimRejected(jobId: string) { + return new DaemonError('claim_rejected', `claim rejected for job: ${jobId}`, { jobId }) +} + +export function claimExpired(jobId: string) { + return new DaemonError('claim_expired', `claim lease expired for job: ${jobId}`, { jobId }) +} + +export function controlAuthMissing() { + return new DaemonError('control_auth_missing', 'missing bearer token') +} + +export function controlAuthRejected() { + return new DaemonError('control_auth_rejected', 'invalid bearer token') +} + +export function invalidJobState(jobId: string, expected: string, actual: JobStatus) { + return new DaemonError( + 'invalid_job_state', + `invalid state for job ${jobId}: expected ${expected}, found ${actual}`, + { jobId, expected, actual } + ) +} + +export function toDaemonError(error: unknown) { + if (error instanceof DaemonError) return error + return sqliteError(error) +} + +export function daemonErrorStatus(error: DaemonError) { + switch (error.kind) { + case 'invalid_input': + case 'non_loopback_bind': + case 'invalid_status': + return 400 + case 'control_auth_missing': + return 401 + case 'job_not_found': + return 404 + case 'claim_rejected': + case 'control_auth_rejected': + return 403 + case 'claim_expired': + case 'invalid_job_state': + case 'bridge_busy': + return 409 + case 'io': + case 'random': + case 'sqlite': + case 'json': + case 'missing_home': + return 500 + } +} + +export function daemonErrorCodeRetryable(error: DaemonError) { + switch (error.kind) { + case 'io': + return { code: 'daemon_io_error', retryable: true } + case 'random': + return { code: 'daemon_random_error', retryable: true } + case 'sqlite': + return { code: 'daemon_store_error', retryable: true } + case 'json': + return { code: 'daemon_json_error', retryable: false } + case 'missing_home': + return { code: 'daemon_home_missing', retryable: false } + case 'invalid_input': + return { code: 'invalid_input', retryable: false } + case 'non_loopback_bind': + return { code: 'non_loopback_bind', retryable: false } + case 'invalid_status': + return { code: 'invalid_status', retryable: false } + case 'job_not_found': + return { code: 'job_not_found', retryable: false } + case 'claim_rejected': + return { code: 'claim_rejected', retryable: false } + case 'claim_expired': + return { code: 'claim_expired', retryable: false } + case 'bridge_busy': + return { code: 'bridge_busy', retryable: true } + case 'control_auth_missing': + return { code: 'control_auth_missing', retryable: false } + case 'control_auth_rejected': + return { code: 'control_auth_rejected', retryable: false } + case 'invalid_job_state': + return { code: 'invalid_job_state', retryable: false } + } +} + +export function daemonErrorBody(error: DaemonError) { + const { code, retryable } = daemonErrorCodeRetryable(error) + const envelope: Record = { + protocol: DAEMON_ERROR_PROTOCOL, + code, + message: error.message, + retryable, + } + const details = daemonErrorDetails(error) + if (details) envelope.details = details + return { error: envelope } +} + +function daemonErrorDetails(error: DaemonError) { + switch (error.kind) { + case 'non_loopback_bind': + return { host: error.host } + case 'invalid_status': + return { status: error.statusValue } + case 'job_not_found': + case 'claim_rejected': + case 'claim_expired': + return { job_id: error.jobId } + case 'invalid_job_state': + return { + job_id: error.jobId, + expected: error.expected, + actual: error.actual, + } + default: + return null + } +} + +function errorText(error: unknown) { + return error instanceof Error && error.message ? error.message : String(error) +} diff --git a/packages/cli/src/daemon/job-store.ts b/packages/cli/src/daemon/job-store.ts new file mode 100644 index 0000000..16934ec --- /dev/null +++ b/packages/cli/src/daemon/job-store.ts @@ -0,0 +1,1170 @@ +import { randomBytes, randomUUID } from 'node:crypto' +import fsSync from 'node:fs' +import fs from 'node:fs/promises' +import os from 'node:os' +import path from 'node:path' +import { DatabaseSync, type SQLInputValue } from 'node:sqlite' + +import { + removeStagedVisibleAttachmentBundle, + validateVisibleAttachmentDescriptor, +} from '../visible-attachments.js' +import { + claimExpired, + claimRejected, + controlAuthRejected, + invalidInput, + invalidJobState, + invalidStatus, + ioError, + jobNotFound, + jsonError, + missingHomeError, + sqliteError, + toDaemonError, + type JobStatus, +} from './errors.js' + +export type { JobStatus } from './errors.js' + +export type ExecutionBackend = 'legacy_extension' | 'playwright' + +export type Job = { + job_id: string + claim_token: string + execution_backend: ExecutionBackend + profile_id: string | null + provider: string + action: string + status: JobStatus + request_json: unknown + result_json: unknown | null + error_json: unknown | null + blocker_json: unknown | null + checkpoint_json: unknown | null + resume_json: unknown | null + created_at: string + updated_at: string + claim_expires_at_ms: number | null +} + +export type JobView = Omit +export type JobWithClaimToken = Omit + +export type CreateJobInput = { + provider: string + action: string + request_json: unknown + execution_backend?: ExecutionBackend | undefined + profile_id?: string | null | undefined + job_id?: string | undefined + claim_token?: string | undefined +} + +export type ListJobsInput = { + status?: JobStatus | undefined + execution_backend?: ExecutionBackend | undefined + profile_id?: string | undefined + provider?: string | undefined + task_id?: string | undefined + limit?: number | undefined +} + +export type ClaimNextInput = { + provider?: string | undefined + action?: string | undefined +} + +const DATABASE_FILE_NAME = 'tokenless.sqlite3' +const CONTROL_TOKEN_FILE_NAME = 'daemon.token' +const DEFAULT_CLAIM_LEASE_MS = 30_000 +const SECRET_TOKEN_BYTES = 32 +const SUMMARY_SCALAR_CHARS = 256 +const PROFILE_ID_CHARS = 128 +const MAX_VISIBLE_ATTACHMENTS = 100 +const MAX_VISIBLE_ATTACHMENT_REQUEST_BYTES = 512 * 1024 * 1024 +const ACTIVE_STATUSES = new Set(['claimed', 'running', 'waiting_for_user']) +const JOB_STATUSES = new Set([ + 'queued', + 'claimed', + 'running', + 'waiting_for_user', + 'succeeded', + 'failed', + 'canceled', + 'timed_out', +]) +const EXECUTION_BACKENDS = new Set(['legacy_extension', 'playwright']) + +type RequestSummaryMetadata = { + task_id: string | null + project_name: string | null + chat_name: string | null + idempotency_key: string | null + task_keys: string[] +} + +export class JobStore { + readonly homeDir: string + readonly databasePath: string + readonly controlTokenPath: string + readonly claimLeaseMs: number + + #db: DatabaseSync + #closed = false + + static async open(homeDir = defaultHomeDir(), claimLeaseMs = DEFAULT_CLAIM_LEASE_MS) { + await ensureTokenlessHome(homeDir) + const canonicalHome = await fs.realpath(homeDir) + const store = new JobStore(canonicalHome, Math.max(1, Math.floor(claimLeaseMs))) + await ensureControlToken(store.controlTokenPath) + store.initialize() + return store + } + + private constructor(homeDir: string, claimLeaseMs: number) { + this.homeDir = homeDir + this.databasePath = path.join(homeDir, DATABASE_FILE_NAME) + this.controlTokenPath = path.join(homeDir, CONTROL_TOKEN_FILE_NAME) + this.claimLeaseMs = claimLeaseMs + try { + this.#db = new DatabaseSync(this.databasePath) + this.#db.exec('PRAGMA foreign_keys = ON;') + this.#db.exec('PRAGMA busy_timeout = 5000;') + } catch (error) { + throw sqliteError(error) + } + } + + close() { + if (this.#closed) return + this.#closed = true + this.#db.close() + } + + controlToken() { + try { + return fsSync.readFileSync(this.controlTokenPath, 'utf8').trim() + } catch (error) { + throw ioError(error) + } + } + + requireControlToken(token: string) { + const expected = this.controlToken() + if (!constantTimeEqual(Buffer.from(token), Buffer.from(expected))) { + throw controlAuthRejected() + } + } + + createJob(input: CreateJobInput) { + const provider = normalizeNonempty(String(input.provider ?? ''), 'provider') + const action = normalizeNonempty(String(input.action ?? ''), 'action') + const executionBackend = input.execution_backend ?? 'legacy_extension' + assertExecutionBackend(executionBackend) + const profileId = validateJobBackendProfile(executionBackend, input.profile_id ?? null) + const jobId = input.job_id === undefined + ? randomUUID() + : normalizeNonempty(input.job_id, 'job_id') + const claimToken = input.claim_token === undefined + ? generateSecretToken() + : normalizeNonempty(input.claim_token, 'claim_token') + const now = nowRfc3339() + const summary = requestSummaryMetadata(input.request_json) + const requestJson = stringifyJson(input.request_json) + + this.transaction(() => { + this.run( + `INSERT INTO jobs ( + job_id, claim_token, execution_backend, profile_id, + provider, action, status, request_json, + result_json, error_json, blocker_json, created_at, updated_at, + summary_task_id, summary_project_name, summary_chat_name, + summary_idempotency_key + ) VALUES ( + ?, ?, ?, ?, ?, ?, ?, ?, NULL, NULL, NULL, ?, ?, + ?, ?, ?, ? + )`, + jobId, + claimToken, + executionBackend, + profileId, + provider, + action, + 'queued', + requestJson, + now, + now, + summary.task_id, + summary.project_name, + summary.chat_name, + summary.idempotency_key + ) + for (const taskKey of summary.task_keys) { + this.run('INSERT INTO job_task_keys (job_id, task_id) VALUES (?, ?)', jobId, taskKey) + } + }) + + return this.getJob(jobId) + } + + listJobs(query: ListJobsInput = {}) { + this.requeueExpiredClaims() + if (query.status !== undefined) assertJobStatus(query.status) + if (query.execution_backend !== undefined) assertExecutionBackend(query.execution_backend) + const profileId = query.profile_id === undefined ? undefined : normalizeProfileId(query.profile_id, 'profile_id') + validateFilterBackendProfile(query.execution_backend, profileId) + const provider = query.provider === undefined ? undefined : normalizeNonempty(query.provider, 'provider') + const taskId = query.task_id === undefined ? undefined : normalizeSummaryFilter(query.task_id, 'task_id') + const limit = clampLimit(query.limit, 100, 1000) + + let sql = `SELECT + jobs.job_id, jobs.claim_token, jobs.execution_backend, jobs.profile_id, + jobs.provider, jobs.action, jobs.status, jobs.request_json, jobs.result_json, + jobs.error_json, jobs.blocker_json, jobs.checkpoint_json, jobs.resume_json, + jobs.created_at, jobs.updated_at, jobs.claim_expires_at + FROM jobs` + const params: SQLInputValue[] = [] + if (taskId !== undefined) { + sql += ` INNER JOIN job_task_keys AS matched_task + ON matched_task.job_id = jobs.job_id + AND matched_task.task_id = ?` + params.push(taskId) + } + sql += ' WHERE 1 = 1' + if (query.status !== undefined) { + sql += ' AND jobs.status = ?' + params.push(query.status) + } + if (query.execution_backend !== undefined) { + sql += ' AND jobs.execution_backend = ?' + params.push(query.execution_backend) + } + if (profileId !== undefined) { + sql += ' AND jobs.profile_id = ?' + params.push(profileId) + } + if (provider !== undefined) { + sql += ' AND jobs.provider = ?' + params.push(provider) + } + sql += ' ORDER BY jobs.created_at DESC, jobs.job_id DESC LIMIT ?' + params.push(limit) + return this.all(sql, ...params).map(rowToJob) + } + + getJob(jobId: string) { + this.requeueExpiredClaims() + return this.getJobWithoutRecovery(jobId) + } + + claimJob(jobId: string, claimToken: string) { + const nowMs = nowUnixMillis() + this.requeueExpiredClaimsAt(nowMs) + const now = nowRfc3339() + const expiresAt = saturatingAdd(nowMs, this.claimLeaseMs) + const result = this.run( + `UPDATE jobs + SET status = ?, updated_at = ?, claim_expires_at = ? + WHERE job_id = ? AND claim_token = ? AND status = ?`, + 'claimed', + now, + expiresAt, + jobId, + claimToken, + 'queued' + ) + if (result.changes === 1) return this.getJobWithoutRecovery(jobId) + return this.explainClaimFailure(jobId, claimToken) + } + + claimNextJob( + query: ClaimNextInput = {}, + executionBackend: ExecutionBackend = 'legacy_extension', + profileIdInput: string | null = null + ) { + const nowMs = nowUnixMillis() + const profileId = profileIdInput === null ? null : normalizeProfileId(profileIdInput, 'profile_id') + validateClaimBackendProfile(executionBackend, profileId) + const provider = query.provider === undefined ? null : normalizeNonempty(query.provider, 'provider') + const action = query.action === undefined ? null : normalizeNonempty(query.action, 'action') + this.requeueExpiredClaimsAt(nowMs) + const now = nowRfc3339() + const expiresAt = saturatingAdd(nowMs, this.claimLeaseMs) + const nextClaimToken = generateSecretToken() + const row = this.get( + `UPDATE jobs + SET status = ?, updated_at = ?, claim_expires_at = ?, claim_token = ? + WHERE job_id = ( + SELECT job_id + FROM jobs + WHERE status = ? + AND execution_backend = ? + AND ((? IS NULL AND profile_id IS NULL) OR profile_id = ?) + AND (? IS NULL OR provider = ?) + AND (? IS NULL OR action = ?) + ORDER BY created_at ASC, job_id ASC + LIMIT 1 + ) + RETURNING + job_id, claim_token, execution_backend, profile_id, + provider, action, status, request_json, result_json, error_json, + blocker_json, checkpoint_json, resume_json, + created_at, updated_at, claim_expires_at`, + 'claimed', + now, + expiresAt, + nextClaimToken, + 'queued', + executionBackend, + profileId, + profileId, + provider, + provider, + action, + action + ) + return row ? rowToJob(row) : null + } + + renewClaim(jobId: string, claimToken: string) { + const nowMs = nowUnixMillis() + const now = nowRfc3339() + const expiresAt = saturatingAdd(nowMs, this.claimLeaseMs) + const result = this.run( + `UPDATE jobs + SET claim_expires_at = ?, updated_at = ? + WHERE job_id = ? + AND claim_token = ? + AND status IN ('claimed', 'running', 'waiting_for_user') + AND claim_expires_at > ?`, + expiresAt, + now, + jobId, + claimToken, + nowMs + ) + if (result.changes === 1) return this.getJobWithoutRecovery(jobId) + return this.explainActiveClaimFailure(jobId, claimToken, nowMs) + } + + markRunning(jobId: string, claimToken: string) { + const nowMs = nowUnixMillis() + const now = nowRfc3339() + const expiresAt = saturatingAdd(nowMs, this.claimLeaseMs) + const result = this.run( + `UPDATE jobs + SET status = ?, blocker_json = NULL, claim_expires_at = ?, updated_at = ? + WHERE job_id = ? + AND claim_token = ? + AND status IN ('claimed', 'waiting_for_user') + AND claim_expires_at > ?`, + 'running', + expiresAt, + now, + jobId, + claimToken, + nowMs + ) + if (result.changes === 1) return this.getJobWithoutRecovery(jobId) + return this.explainRunningFailure(jobId, claimToken, nowMs) + } + + markWaitingForUser(jobId: string, claimToken: string, blockerJson: unknown) { + const nowMs = nowUnixMillis() + const now = nowRfc3339() + const expiresAt = saturatingAdd(nowMs, this.claimLeaseMs) + const result = this.run( + `UPDATE jobs + SET status = ?, blocker_json = ?, claim_expires_at = ?, updated_at = ? + WHERE job_id = ? + AND claim_token = ? + AND status = 'running' + AND claim_expires_at > ?`, + 'waiting_for_user', + stringifyJson(blockerJson), + expiresAt, + now, + jobId, + claimToken, + nowMs + ) + if (result.changes === 1) return this.getJobWithoutRecovery(jobId) + return this.explainActiveClaimFailure(jobId, claimToken, nowMs) + } + + checkpointJob(jobId: string, claimToken: string, checkpointJson: unknown) { + const nowMs = nowUnixMillis() + const now = nowRfc3339() + const result = this.run( + `UPDATE jobs + SET checkpoint_json = ?, updated_at = ? + WHERE job_id = ? + AND claim_token = ? + AND execution_backend = 'playwright' + AND status IN ('claimed', 'running', 'waiting_for_user') + AND claim_expires_at > ?`, + stringifyJson(checkpointJson), + now, + jobId, + claimToken, + nowMs + ) + if (result.changes === 1) return this.getJobWithoutRecovery(jobId) + return this.explainPlaywrightActiveFailure( + jobId, + claimToken, + nowMs, + 'only playwright jobs can persist browser checkpoints', + 'claimed, running, or waiting_for_user' + ) + } + + parkJob(jobId: string, claimToken: string, blockerJson: unknown, checkpointJson: unknown) { + const nowMs = nowUnixMillis() + const now = nowRfc3339() + const replacementToken = generateSecretToken() + const result = this.run( + `UPDATE jobs + SET status = ?, claim_token = ?, blocker_json = ?, + checkpoint_json = ?, resume_json = NULL, + claim_expires_at = NULL, updated_at = ? + WHERE job_id = ? + AND claim_token = ? + AND execution_backend = 'playwright' + AND status IN ('claimed', 'running', 'waiting_for_user') + AND claim_expires_at > ?`, + 'waiting_for_user', + replacementToken, + stringifyJson(blockerJson), + stringifyJson(checkpointJson), + now, + jobId, + claimToken, + nowMs + ) + if (result.changes === 1) return this.getJobWithoutRecovery(jobId) + return this.explainPlaywrightActiveFailure( + jobId, + claimToken, + nowMs, + 'only playwright jobs can be parked for browser resume', + 'claimed, running, or waiting_for_user' + ) + } + + resumeJob(jobId: string, resumeJson: unknown) { + const nowMs = nowUnixMillis() + this.requeueExpiredClaimsAt(nowMs) + const now = nowRfc3339() + const replacementToken = generateSecretToken() + const serializedResume = stringifyJson(resumeJson) + return this.transaction(() => { + const job = this.getJobWithoutRecovery(jobId) + if (job.execution_backend !== 'playwright') { + throw invalidInput('only playwright jobs can be resumed with browser visibility') + } + if ( + (job.status === 'queued' || job.status === 'claimed' || job.status === 'running') && + job.resume_json !== null + ) { + return job + } + if (job.status !== 'waiting_for_user' || job.checkpoint_json === null || job.claim_expires_at_ms !== null) { + throw invalidJobState(job.job_id, 'parked playwright waiting_for_user', job.status) + } + const result = this.run( + `UPDATE jobs + SET status = 'queued', claim_token = ?, resume_json = ?, + claim_expires_at = NULL, updated_at = ? + WHERE job_id = ? + AND execution_backend = 'playwright' + AND status = 'waiting_for_user' + AND checkpoint_json IS NOT NULL + AND claim_expires_at IS NULL`, + replacementToken, + serializedResume, + now, + jobId + ) + if (result.changes !== 1) { + throw invalidJobState(jobId, 'parked playwright waiting_for_user', job.status) + } + return this.getJobWithoutRecovery(jobId) + }) + } + + completeJob(jobId: string, claimToken: string, completion: { result_json: unknown } | { error_json: unknown }) { + const nowMs = nowUnixMillis() + const now = nowRfc3339() + const status: JobStatus = 'result_json' in completion ? 'succeeded' : 'failed' + const resultJson = 'result_json' in completion ? stringifyJson(completion.result_json) : null + const errorJson = 'error_json' in completion ? stringifyJson(completion.error_json) : null + const result = this.run( + `UPDATE jobs + SET status = ?, result_json = ?, error_json = ?, blocker_json = NULL, + checkpoint_json = NULL, resume_json = NULL, + updated_at = ?, claim_expires_at = NULL + WHERE job_id = ? + AND claim_token = ? + AND status IN ('claimed', 'running', 'waiting_for_user') + AND claim_expires_at > ?`, + status, + resultJson, + errorJson, + now, + jobId, + claimToken, + nowMs + ) + if (result.changes === 1) return this.getJobWithoutRecovery(jobId) + return this.explainActiveClaimFailure(jobId, claimToken, nowMs) + } + + async cancelJob(jobId: string, reason: unknown | undefined) { + const now = nowRfc3339() + const errorJson = stringifyJson(reason === undefined || reason === null + ? { code: 'job_canceled' } + : { code: 'job_canceled', reason }) + const result = this.run( + `UPDATE jobs + SET status = ?, result_json = NULL, error_json = ?, blocker_json = NULL, + checkpoint_json = NULL, resume_json = NULL, + updated_at = ?, claim_expires_at = NULL + WHERE job_id = ? AND status IN ('queued', 'claimed', 'running', 'waiting_for_user')`, + 'canceled', + errorJson, + now, + jobId + ) + const job = this.getJobWithoutRecovery(jobId) + if (result.changes === 1) { + await cleanupVisibleAttachmentBundlesForRequest(this.homeDir, job.request_json).catch(() => undefined) + return job + } + throw invalidJobState(jobId, 'queued, claimed, running, or waiting_for_user', job.status) + } + + requeueExpiredClaims() { + return this.requeueExpiredClaimsAt(nowUnixMillis()) + } + + private requeueExpiredClaimsAt(nowMs: number) { + const now = nowRfc3339() + this.parkExpiredCheckpointedWaitingClaimsAt(nowMs) + this.failExpiredWaitingClaimsAt(nowMs) + const exists = this.exists( + `SELECT EXISTS ( + SELECT 1 FROM jobs + WHERE status IN ('claimed', 'running') + AND (claim_expires_at IS NULL OR claim_expires_at <= ?) + ) AS present`, + nowMs + ) + if (!exists) return 0 + return this.transaction(() => { + const expiredJobIds = this.all( + `SELECT job_id + FROM jobs + WHERE status IN ('claimed', 'running') + AND (claim_expires_at IS NULL OR claim_expires_at <= ?) + ORDER BY job_id ASC`, + nowMs + ).map((row) => String(row.job_id)) + let requeued = 0 + for (const jobId of expiredJobIds) { + const result = this.run( + `UPDATE jobs + SET status = 'queued', claim_token = ?, claim_expires_at = NULL, + blocker_json = NULL, resume_json = NULL, updated_at = ? + WHERE job_id = ? + AND status IN ('claimed', 'running') + AND (claim_expires_at IS NULL OR claim_expires_at <= ?)`, + generateSecretToken(), + now, + jobId, + nowMs + ) + requeued += Number(result.changes) + } + return requeued + }) + } + + private parkExpiredCheckpointedWaitingClaimsAt(nowMs: number) { + const now = nowRfc3339() + const exists = this.exists( + `SELECT EXISTS ( + SELECT 1 FROM jobs + WHERE status = 'waiting_for_user' + AND execution_backend = 'playwright' + AND checkpoint_json IS NOT NULL + AND claim_expires_at <= ? + ) AS present`, + nowMs + ) + if (!exists) return 0 + return this.transaction(() => { + const expiredJobs = this.all( + `SELECT job_id, blocker_json + FROM jobs + WHERE status = 'waiting_for_user' + AND execution_backend = 'playwright' + AND checkpoint_json IS NOT NULL + AND claim_expires_at <= ? + ORDER BY job_id ASC`, + nowMs + ) + let parked = 0 + for (const row of expiredJobs) { + const result = this.run( + `UPDATE jobs + SET claim_token = ?, blocker_json = ?, + claim_expires_at = NULL, resume_json = NULL, + updated_at = ? + WHERE job_id = ? + AND status = 'waiting_for_user' + AND execution_backend = 'playwright' + AND checkpoint_json IS NOT NULL + AND claim_expires_at <= ?`, + generateSecretToken(), + stringifyJson(parkedResumeBlockerJson(parseOptionalJson(row.blocker_json))), + now, + String(row.job_id), + nowMs + ) + parked += Number(result.changes) + } + return parked + }) + } + + private failExpiredWaitingClaimsAt(nowMs: number) { + const now = nowRfc3339() + const errorJson = stringifyJson({ + code: 'playwright_user_handover_lease_lost', + message: 'The managed Playwright job lost its lease while waiting for user handover; retry from the same task state instead of replaying partial page actions.', + retryable: true, + }) + return this.run( + `UPDATE jobs + SET status = 'failed', error_json = ?, result_json = NULL, + blocker_json = NULL, checkpoint_json = NULL, resume_json = NULL, + claim_expires_at = NULL, updated_at = ? + WHERE status = 'waiting_for_user' + AND (execution_backend != 'playwright' OR checkpoint_json IS NULL) + AND (claim_expires_at IS NULL OR claim_expires_at <= ?)`, + errorJson, + now, + nowMs + ).changes + } + + private initialize() { + this.exec(` + PRAGMA journal_mode = WAL; + PRAGMA foreign_keys = ON; + CREATE TABLE IF NOT EXISTS jobs ( + job_id TEXT PRIMARY KEY NOT NULL, + claim_token TEXT NOT NULL, + execution_backend TEXT NOT NULL DEFAULT 'legacy_extension' CHECK ( + execution_backend IN ('legacy_extension', 'playwright') + ), + profile_id TEXT CHECK ( + profile_id IS NULL OR length(profile_id) BETWEEN 1 AND 128 + ), + provider TEXT NOT NULL, + action TEXT NOT NULL, + status TEXT NOT NULL CHECK ( + status IN ( + 'queued', + 'claimed', + 'running', + 'waiting_for_user', + 'succeeded', + 'failed', + 'canceled', + 'timed_out' + ) + ), + request_json TEXT NOT NULL, + result_json TEXT, + error_json TEXT, + blocker_json TEXT, + checkpoint_json TEXT, + resume_json TEXT, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + claim_expires_at INTEGER, + summary_task_id TEXT CHECK ( + summary_task_id IS NULL OR length(summary_task_id) <= 256 + ), + summary_project_name TEXT CHECK ( + summary_project_name IS NULL OR length(summary_project_name) <= 256 + ), + summary_chat_name TEXT CHECK ( + summary_chat_name IS NULL OR length(summary_chat_name) <= 256 + ), + summary_idempotency_key TEXT CHECK ( + summary_idempotency_key IS NULL OR length(summary_idempotency_key) <= 256 + ) + ); + CREATE TABLE IF NOT EXISTS job_task_keys ( + job_id TEXT NOT NULL REFERENCES jobs(job_id) ON DELETE CASCADE, + task_id TEXT NOT NULL CHECK (length(task_id) BETWEEN 1 AND 256), + PRIMARY KEY (job_id, task_id) + ); + CREATE INDEX IF NOT EXISTS jobs_status_created_at_idx + ON jobs(status, created_at); + CREATE INDEX IF NOT EXISTS jobs_provider_action_idx + ON jobs(provider, action); + CREATE INDEX IF NOT EXISTS jobs_claim_expires_at_idx + ON jobs(claim_expires_at); + CREATE INDEX IF NOT EXISTS jobs_backend_profile_status_fifo_idx + ON jobs(execution_backend, profile_id, status, created_at, job_id); + CREATE INDEX IF NOT EXISTS job_task_keys_task_id_idx + ON job_task_keys(task_id, job_id); + `) + restrictFilePermissionsSync(this.databasePath) + } + + private getJobWithoutRecovery(jobId: string) { + const row = this.get( + `SELECT + job_id, claim_token, execution_backend, profile_id, + provider, action, status, request_json, result_json, error_json, + blocker_json, checkpoint_json, resume_json, + created_at, updated_at, claim_expires_at + FROM jobs + WHERE job_id = ?`, + jobId + ) + if (!row) throw jobNotFound(jobId) + return rowToJob(row) + } + + private explainClaimFailure(jobId: string, claimToken: string): never { + const job = this.getJobWithoutRecovery(jobId) + if (job.claim_token !== claimToken) throw claimRejected(jobId) + throw invalidJobState(jobId, 'queued', job.status) + } + + private explainActiveClaimFailure(jobId: string, claimToken: string, nowMs: number): never { + const job = this.getJobWithoutRecovery(jobId) + if (job.claim_token !== claimToken) throw claimRejected(jobId) + if (ACTIVE_STATUSES.has(job.status) && (job.claim_expires_at_ms === null || job.claim_expires_at_ms <= nowMs)) { + throw claimExpired(jobId) + } + throw invalidJobState(jobId, 'claimed, running, or waiting_for_user', job.status) + } + + private explainRunningFailure(jobId: string, claimToken: string, nowMs: number): never { + const job = this.getJobWithoutRecovery(jobId) + if (job.claim_token !== claimToken) throw claimRejected(jobId) + if (ACTIVE_STATUSES.has(job.status) && (job.claim_expires_at_ms === null || job.claim_expires_at_ms <= nowMs)) { + throw claimExpired(jobId) + } + throw invalidJobState(jobId, 'claimed or waiting_for_user', job.status) + } + + private explainPlaywrightActiveFailure( + jobId: string, + claimToken: string, + nowMs: number, + backendMessage: string, + expected: string + ): never { + const job = this.getJobWithoutRecovery(jobId) + if (job.claim_token !== claimToken) throw claimRejected(jobId) + if (ACTIVE_STATUSES.has(job.status) && (job.claim_expires_at_ms === null || job.claim_expires_at_ms <= nowMs)) { + throw claimExpired(jobId) + } + if (job.execution_backend !== 'playwright') throw invalidInput(backendMessage) + throw invalidJobState(job.job_id, expected, job.status) + } + + private exec(sql: string) { + try { + this.#db.exec(sql) + } catch (error) { + throw sqliteError(error) + } + } + + private run(sql: string, ...params: SQLInputValue[]) { + try { + return this.#db.prepare(sql).run(...params) + } catch (error) { + throw toDaemonError(error) + } + } + + private get(sql: string, ...params: SQLInputValue[]) { + try { + return this.#db.prepare(sql).get(...params) + } catch (error) { + throw toDaemonError(error) + } + } + + private all(sql: string, ...params: SQLInputValue[]) { + try { + return this.#db.prepare(sql).all(...params) + } catch (error) { + throw toDaemonError(error) + } + } + + private exists(sql: string, ...params: SQLInputValue[]) { + const row = this.get(sql, ...params) + return Boolean(row && Number(Object.values(row)[0]) !== 0) + } + + private transaction(callback: () => T) { + this.exec('BEGIN IMMEDIATE') + try { + const result = callback() + this.exec('COMMIT') + return result + } catch (error) { + try { + this.exec('ROLLBACK') + } catch { + // Keep the original failure; a rollback failure means the connection is already unusable. + } + throw error + } + } +} + +export function publicView(job: Job): JobView { + return { + job_id: job.job_id, + execution_backend: job.execution_backend, + profile_id: job.profile_id, + provider: job.provider, + action: job.action, + status: job.status, + request_json: job.request_json, + result_json: job.result_json, + error_json: job.error_json, + blocker_json: job.blocker_json, + created_at: job.created_at, + updated_at: job.updated_at, + } +} + +export function withClaimToken(job: Job): JobWithClaimToken { + return { + job_id: job.job_id, + claim_token: job.claim_token, + execution_backend: job.execution_backend, + profile_id: job.profile_id, + provider: job.provider, + action: job.action, + status: job.status, + request_json: job.request_json, + result_json: job.result_json, + error_json: job.error_json, + blocker_json: job.blocker_json, + checkpoint_json: job.checkpoint_json, + resume_json: job.resume_json, + created_at: job.created_at, + updated_at: job.updated_at, + } +} + +export function defaultHomeDir() { + if (process.env.TOKENLESS_HOME) return process.env.TOKENLESS_HOME + if (process.env.HOME) return path.join(process.env.HOME, '.tokenless') + const home = os.homedir() + if (home) return path.join(home, '.tokenless') + throw missingHomeError() +} + +async function ensureTokenlessHome(homeDir: string) { + try { + await fs.mkdir(homeDir, { recursive: true, mode: 0o700 }) + await fs.chmod(homeDir, 0o700).catch(() => undefined) + } catch (error) { + throw ioError(error) + } +} + +async function ensureControlToken(tokenPath: string) { + try { + const token = fsSync.existsSync(tokenPath) ? fsSync.readFileSync(tokenPath, 'utf8').trim() : '' + if (token) { + restrictFilePermissionsSync(tokenPath) + return + } + const descriptor = fsSync.openSync(tokenPath, fsSync.constants.O_CREAT | fsSync.constants.O_EXCL | fsSync.constants.O_WRONLY, 0o600) + try { + fsSync.writeFileSync(descriptor, `${generateSecretToken()}\n`) + fsSync.fsyncSync(descriptor) + } finally { + fsSync.closeSync(descriptor) + } + restrictFilePermissionsSync(tokenPath) + } catch (error) { + const code = (error as NodeJS.ErrnoException).code + if (code === 'EEXIST') { + restrictFilePermissionsSync(tokenPath) + const existing = fsSync.readFileSync(tokenPath, 'utf8').trim() + if (!existing) throw invalidInput(`${tokenPath} is empty`) + return + } + throw ioError(error) + } +} + +function rowToJob(row: Record): Job { + const executionBackend = String(row.execution_backend) + assertExecutionBackend(executionBackend) + const status = String(row.status) + assertJobStatus(status) + return { + job_id: String(row.job_id), + claim_token: String(row.claim_token), + execution_backend: executionBackend, + profile_id: nullableString(row.profile_id), + provider: String(row.provider), + action: String(row.action), + status, + request_json: parseJson(row.request_json), + result_json: parseOptionalJson(row.result_json), + error_json: parseOptionalJson(row.error_json), + blocker_json: parseOptionalJson(row.blocker_json), + checkpoint_json: parseOptionalJson(row.checkpoint_json), + resume_json: parseOptionalJson(row.resume_json), + created_at: String(row.created_at), + updated_at: String(row.updated_at), + claim_expires_at_ms: nullableNumber(row.claim_expires_at), + } +} + +function parseJson(value: unknown) { + if (typeof value !== 'string') throw jsonError(new Error('expected JSON text')) + try { + return JSON.parse(value) as unknown + } catch (error) { + throw jsonError(error) + } +} + +function parseOptionalJson(value: unknown) { + if (value === null || value === undefined) return null + return parseJson(value) +} + +function stringifyJson(value: unknown) { + try { + return JSON.stringify(value) + } catch (error) { + throw jsonError(error) + } +} + +async function cleanupVisibleAttachmentBundlesForRequest(homeDir: string, request: unknown) { + const bundleIds = visibleAttachmentBundleIdsForRequest(request) + if (bundleIds === null) return + for (const bundleId of bundleIds) { + await removeStagedVisibleAttachmentBundle({ homeDir, bundleId }).catch(() => undefined) + } +} + +function visibleAttachmentBundleIdsForRequest(request: unknown) { + const requestObject = jsonRecord(request) + if (!requestObject || !Object.hasOwn(requestObject, 'attachments')) return new Set() + const attachments = requestObject.attachments + if (!Array.isArray(attachments) || attachments.length === 0 || attachments.length > MAX_VISIBLE_ATTACHMENTS) { + return null + } + const bundleIds = new Set() + const attachmentIds = new Set() + let expectedBundleId: string | undefined + let totalBytes = 0 + for (const attachment of attachments) { + let descriptor + try { + descriptor = validateVisibleAttachmentDescriptor(attachment) + } catch { + return null + } + if (expectedBundleId !== undefined && descriptor.bundleId !== expectedBundleId) return null + expectedBundleId = descriptor.bundleId + if (attachmentIds.has(descriptor.attachmentId)) return null + attachmentIds.add(descriptor.attachmentId) + if (descriptor.size > MAX_VISIBLE_ATTACHMENT_REQUEST_BYTES) return null + totalBytes += descriptor.size + if (!Number.isSafeInteger(totalBytes) || totalBytes > MAX_VISIBLE_ATTACHMENT_REQUEST_BYTES) return null + bundleIds.add(descriptor.bundleId) + } + return bundleIds +} + +function requestSummaryMetadata(request: unknown): RequestSummaryMetadata { + const requestObject = jsonRecord(request) + const metadataObject = jsonRecord(requestObject?.metadata) + const requestValue = (key: string) => boundedNonemptySummaryValue(requestObject?.[key]) + const metadataValue = (key: string) => boundedNonemptySummaryValue(metadataObject?.[key]) + const requestTaskId = requestValue('taskId') + const requestIdempotencyKey = requestValue('idempotencyKey') + const requestId = requestValue('requestId') + const metadataTaskId = metadataValue('taskId') + const metadataIdempotencyKey = metadataValue('idempotencyKey') + const taskKeys: string[] = [] + for (const key of [requestTaskId, requestIdempotencyKey, requestId, metadataTaskId, metadataIdempotencyKey]) { + if (key && !taskKeys.includes(key)) taskKeys.push(key) + } + return { + task_id: metadataTaskId ?? requestTaskId ?? null, + project_name: metadataValue('projectName') ?? requestValue('projectName') ?? null, + chat_name: metadataValue('chatName') ?? requestValue('chatName') ?? null, + idempotency_key: metadataIdempotencyKey ?? requestIdempotencyKey ?? null, + task_keys: taskKeys, + } +} + +function parkedResumeBlockerJson(blockerJson: unknown) { + const browser = { + windowOpen: false, + resumeRequired: true, + } + if (blockerJson && typeof blockerJson === 'object' && !Array.isArray(blockerJson)) { + const object = { ...blockerJson as Record } + if (object.browser && typeof object.browser === 'object' && !Array.isArray(object.browser)) { + object.browser = { + ...object.browser as Record, + windowOpen: false, + resumeRequired: true, + } + } else { + object.browser = browser + } + return object + } + if (blockerJson !== null) { + return { blocker: blockerJson, browser } + } + return { browser } +} + +function validateJobBackendProfile(executionBackend: ExecutionBackend, profileId: string | null) { + const normalized = profileId === null ? null : normalizeProfileId(profileId, 'profile_id') + if (executionBackend === 'legacy_extension') { + if (normalized !== null) throw invalidInput('legacy_extension jobs must not set profile_id') + return null + } + if (normalized === null) throw invalidInput('playwright jobs require profile_id') + return normalized +} + +function validateFilterBackendProfile(executionBackend: ExecutionBackend | undefined, profileId: string | undefined) { + if (executionBackend === 'legacy_extension' && profileId !== undefined) { + throw invalidInput('legacy_extension filters must not set profile_id') + } +} + +function validateClaimBackendProfile(executionBackend: ExecutionBackend, profileId: string | null) { + assertExecutionBackend(executionBackend) + if (executionBackend === 'legacy_extension') { + if (profileId !== null) throw invalidInput('legacy_extension claims must not set profile_id') + return + } + if (profileId === null) throw invalidInput('playwright claims require profile_id') +} + +function assertJobStatus(value: string): asserts value is JobStatus { + if (!JOB_STATUSES.has(value as JobStatus)) throw invalidStatus(value) +} + +function assertExecutionBackend(value: string): asserts value is ExecutionBackend { + if (!EXECUTION_BACKENDS.has(value as ExecutionBackend)) { + throw invalidInput(`invalid execution_backend: ${value}`) + } +} + +function normalizeNonempty(value: string, field: string) { + const trimmed = value.trim() + if (!trimmed) throw invalidInput(`${field} must be a nonempty string`) + return trimmed +} + +function normalizeSummaryFilter(value: string, field: string) { + const normalized = normalizeNonempty(value, field) + if (Array.from(normalized).length > SUMMARY_SCALAR_CHARS) { + throw invalidInput(`${field} must be at most ${SUMMARY_SCALAR_CHARS} characters`) + } + return normalized +} + +function normalizeProfileId(value: string, field: string) { + const normalized = normalizeNonempty(value, field) + if (Array.from(normalized).length > PROFILE_ID_CHARS) { + throw invalidInput(`${field} must be at most ${PROFILE_ID_CHARS} characters`) + } + return normalized +} + +function boundedNonemptySummaryValue(value: unknown) { + if (typeof value !== 'string') return null + const trimmed = value.trim() + if (!trimmed) return null + return Array.from(trimmed).slice(0, SUMMARY_SCALAR_CHARS).join('') +} + +function clampLimit(value: unknown, fallback: number, max: number) { + const numeric = Number(value) + const finite = Number.isFinite(numeric) ? Math.floor(numeric) : fallback + return Math.max(1, Math.min(max, finite || fallback)) +} + +function generateSecretToken() { + return randomBytes(SECRET_TOKEN_BYTES).toString('base64url') +} + +function nowRfc3339() { + return new Date().toISOString() +} + +function nowUnixMillis() { + return Date.now() +} + +function saturatingAdd(left: number, right: number) { + const result = left + right + return Number.isSafeInteger(result) ? result : Number.MAX_SAFE_INTEGER +} + +function nullableString(value: unknown) { + return value === null || value === undefined ? null : String(value) +} + +function nullableNumber(value: unknown) { + return value === null || value === undefined ? null : Number(value) +} + +function jsonRecord(value: unknown): Record | null { + return value && typeof value === 'object' && !Array.isArray(value) + ? value as Record + : null +} + +function restrictFilePermissionsSync(filePath: string) { + if (process.platform === 'win32' || !fsSync.existsSync(filePath)) return + try { + fsSync.chmodSync(filePath, 0o600) + } catch (error) { + throw ioError(error) + } +} + +function constantTimeEqual(left: Buffer, right: Buffer) { + let diff = left.length ^ right.length + const maxLength = Math.max(left.length, right.length) + for (let index = 0; index < maxLength; index += 1) { + diff |= (left[index] ?? 0) ^ (right[index] ?? 0) + } + return diff === 0 +} diff --git a/packages/cli/src/daemon/lifecycle.ts b/packages/cli/src/daemon/lifecycle.ts new file mode 100644 index 0000000..b6c3166 --- /dev/null +++ b/packages/cli/src/daemon/lifecycle.ts @@ -0,0 +1,42 @@ +import process from 'node:process' + +import { JobStore, defaultHomeDir } from './job-store.js' +import { serveHttp, type DaemonServer } from './server.js' + +export type StartDaemonOptions = { + homeDir?: string | undefined + host?: string | undefined + port?: number | undefined +} + +export async function startDaemon({ + homeDir = defaultHomeDir(), + host = '127.0.0.1', + port = 7331, +}: StartDaemonOptions = {}) { + const store = await JobStore.open(homeDir) + let daemon: DaemonServer + try { + daemon = await serveHttp({ store, host, port }) + } catch (error) { + store.close() + throw error + } + installProcessShutdownHandlers(daemon) + return daemon +} + +function installProcessShutdownHandlers(daemon: DaemonServer) { + let closing = false + const close = async () => { + if (closing) return + closing = true + await daemon.close() + } + process.once('SIGINT', () => { + void close().finally(() => process.exit(130)) + }) + process.once('SIGTERM', () => { + void close().finally(() => process.exit(143)) + }) +} diff --git a/packages/cli/src/daemon/server.ts b/packages/cli/src/daemon/server.ts new file mode 100644 index 0000000..a8bef14 --- /dev/null +++ b/packages/cli/src/daemon/server.ts @@ -0,0 +1,468 @@ +import { createHmac } from 'node:crypto' +import http, { type IncomingMessage, type ServerResponse } from 'node:http' +import net from 'node:net' + +import { + DAEMON_PROTOCOL, + DAEMON_READY_PROOF_PROTOCOL, + MANAGED_PLAYWRIGHT_JOB_PROTOCOL_VERSION_V1, + MANAGED_PLAYWRIGHT_JOB_PROTOCOL_VERSION_V2, + NATIVE_BINARY_BUILD_INFO_PROTOCOL, + NATIVE_PROTOCOL, + VISIBLE_ACTION_PROTOCOL_VERSION_V1, + VISIBLE_ACTION_PROTOCOL_VERSION_V2, +} from '../generated/protocol-constants.js' +import { tokenlessPackageVersion } from '../platform-package.js' +import { + controlAuthMissing, + controlAuthRejected, + daemonErrorBody, + daemonErrorStatus, + invalidInput, + nonLoopbackBind, + toDaemonError, + type DaemonError, +} from './errors.js' +import { JobStore, publicView, withClaimToken, type ExecutionBackend, type JobStatus } from './job-store.js' + +export type DaemonServer = { + close(): Promise + server: http.Server + store: JobStore +} + +const READY_CHALLENGE_BYTES = 32 +const READY_CHALLENGE_BASE64URL_CHARS = 43 +const MAX_HTTP_BODY_BYTES = 2 * 1024 * 1024 +const BODY_LIMIT_EXCEEDED_MESSAGE = 'Failed to buffer the request body: length limit exceeded' + +type JsonRecord = Record + +class BodyLimitExceededError extends Error { + constructor() { + super(BODY_LIMIT_EXCEEDED_MESSAGE) + this.name = 'BodyLimitExceededError' + } +} + +export async function serveHttp({ + store, + host, + port, +}: { + store: JobStore + host: string + port: number +}) { + validateLoopbackHost(host) + const server = http.createServer((request, response) => { + void handleRequest(store, server, request, response) + }) + await new Promise((resolve, reject) => { + const onError = (error: Error) => { + server.off('listening', onListening) + reject(error) + } + const onListening = () => { + server.off('error', onError) + resolve() + } + server.once('error', onError) + server.once('listening', onListening) + server.listen(port, host) + }) + return { + server, + store, + close: () => closeServer(server, store), + } satisfies DaemonServer +} + +export function validateLoopbackHost(host: string) { + const ip = net.isIP(host) + const loopback = host === 'localhost' || + host === '::1' || + host === '[::1]' || + (ip === 4 && host.startsWith('127.')) || + (ip === 6 && host === '::1') + if (!loopback) throw nonLoopbackBind(host) +} + +export function nativeBinaryBuildInfo(binary: string) { + return { + protocol: NATIVE_BINARY_BUILD_INFO_PROTOCOL, + binary, + version: tokenlessPackageVersion(), + platform: process.platform === 'darwin' ? 'darwin' : process.platform, + arch: process.arch === 'x64' ? 'x64' : process.arch, + } +} + +async function handleRequest( + store: JobStore, + server: http.Server, + request: IncomingMessage, + response: ServerResponse +) { + try { + const url = new URL(request.url || '/', 'http://127.0.0.1') + const method = request.method || 'GET' + if (method === 'GET' && url.pathname === '/health') { + writeJson(response, 200, healthResponse(store)) + return + } + if (method === 'GET' && url.pathname === '/ready') { + const challenge = url.searchParams.get('challenge') ?? '' + validateReadyChallenge(challenge) + const health = healthResponse(store) + writeJson(response, 200, { + ...health, + ready_proof_protocol: DAEMON_READY_PROOF_PROTOCOL, + ready_challenge: challenge, + ready_proof: daemonReadyProof(store, challenge, health.home_dir), + supported_protocols: supportedProtocols(), + }) + return + } + + requireControlAuth(store, request) + + if (method === 'POST' && url.pathname === '/jobs') { + const body = await readJsonObject(request) + const job = store.createJob({ + provider: requiredString(body.provider, 'provider'), + action: requiredString(body.action, 'action'), + request_json: requireField(body, 'request_json'), + execution_backend: optionalExecutionBackend(body.execution_backend), + profile_id: optionalString(body.profile_id), + job_id: optionalString(body.job_id) ?? undefined, + claim_token: optionalString(body.claim_token) ?? undefined, + }) + writeJson(response, 200, withClaimToken(job)) + return + } + + if (method === 'GET' && url.pathname === '/jobs') { + const jobs = store.listJobs({ + status: optionalJobStatus(url.searchParams.get('status')), + execution_backend: optionalExecutionBackend(url.searchParams.get('execution_backend')), + profile_id: optionalQueryString(url.searchParams.get('profile_id')), + provider: optionalQueryString(url.searchParams.get('provider')), + task_id: optionalQueryString(url.searchParams.get('task_id')), + limit: optionalLimit(url.searchParams.get('limit')), + }) + writeJson(response, 200, jobs.map(publicView)) + return + } + + const jobRoute = matchJobRoute(url.pathname) + if (jobRoute && method === 'GET' && jobRoute.action === null) { + writeJson(response, 200, publicView(store.getJob(jobRoute.jobId))) + return + } + if (jobRoute && method === 'POST' && jobRoute.action === 'claim') { + const body = await readJsonObject(request) + writeJson(response, 200, publicView(store.claimJob(jobRoute.jobId, requiredString(body.claim_token, 'claim_token')))) + return + } + if (jobRoute && method === 'POST' && jobRoute.action === 'complete') { + const body = await readJsonObject(request) + const hasResult = body.result_json !== undefined && body.result_json !== null + const hasError = body.error_json !== undefined && body.error_json !== null + if (hasResult === hasError) throw invalidInput('pass exactly one of result_json or error_json') + const claimToken = requiredString(body.claim_token, 'claim_token') + const job = hasResult + ? store.completeJob(jobRoute.jobId, claimToken, { result_json: body.result_json }) + : store.completeJob(jobRoute.jobId, claimToken, { error_json: body.error_json }) + writeJson(response, 200, publicView(job)) + return + } + if (jobRoute && method === 'POST' && jobRoute.action === 'resume') { + const body = await readJsonObject(request) + if (Object.keys(body).some((key) => key !== 'browser_visibility')) { + throw invalidInput('request body must be valid JSON: unknown field') + } + if (body.browser_visibility !== 'headed') throw invalidInput('request body must be valid JSON: invalid browser_visibility') + writeJson(response, 200, publicView(store.resumeJob(jobRoute.jobId, { browser_visibility: 'headed' }))) + return + } + + if (method === 'POST' && url.pathname === '/control/jobs/claim-next') { + const executionBackend = optionalExecutionBackend(url.searchParams.get('execution_backend')) ?? 'legacy_extension' + const job = store.claimNextJob( + { + provider: optionalQueryString(url.searchParams.get('provider')), + action: optionalQueryString(url.searchParams.get('action')), + }, + executionBackend, + optionalQueryString(url.searchParams.get('profile_id')) ?? null + ) + writeJson(response, 200, { job: job ? withClaimToken(job) : null }) + return + } + + const controlJobRoute = matchControlJobRoute(url.pathname) + if (controlJobRoute && method === 'POST') { + const rawBody = await readBody(request) + const body = rawBody ? parseJsonObject(rawBody) : {} + if (controlJobRoute.action === 'running') { + writeJson(response, 200, publicView(store.markRunning(controlJobRoute.jobId, requiredString(body.claim_token, 'claim_token')))) + return + } + if (controlJobRoute.action === 'waiting-for-user') { + writeJson(response, 200, publicView(store.markWaitingForUser( + controlJobRoute.jobId, + requiredString(body.claim_token, 'claim_token'), + requireField(body, 'blocker_json') + ))) + return + } + if (controlJobRoute.action === 'checkpoint') { + writeJson(response, 200, publicView(store.checkpointJob( + controlJobRoute.jobId, + requiredString(body.claim_token, 'claim_token'), + requireField(body, 'checkpoint_json') + ))) + return + } + if (controlJobRoute.action === 'park') { + writeJson(response, 200, publicView(store.parkJob( + controlJobRoute.jobId, + requiredString(body.claim_token, 'claim_token'), + requireField(body, 'blocker_json'), + requireField(body, 'checkpoint_json') + ))) + return + } + if (controlJobRoute.action === 'renew') { + writeJson(response, 200, publicView(store.renewClaim(controlJobRoute.jobId, requiredString(body.claim_token, 'claim_token')))) + return + } + if (controlJobRoute.action === 'cancel') { + writeJson(response, 200, publicView(await store.cancelJob(controlJobRoute.jobId, body.reason))) + return + } + } + + if (method === 'POST' && url.pathname === '/control/shutdown') { + writeJson(response, 200, { ok: true, status: 'shutting_down', pid: process.pid }) + setImmediate(() => { + void closeServer(server, store) + }) + return + } + + writeJson(response, 404, { error: { message: 'not found' } }) + } catch (error) { + if (error instanceof BodyLimitExceededError) { + writeText(response, 413, BODY_LIMIT_EXCEEDED_MESSAGE) + return + } + writeDaemonError(response, toDaemonError(error)) + } +} + +function healthResponse(store: JobStore) { + return { + protocol: DAEMON_PROTOCOL, + daemon_protocol: DAEMON_PROTOCOL, + version: tokenlessPackageVersion(), + native_protocol: NATIVE_PROTOCOL, + status: 'ok', + ready: true, + home_dir: store.homeDir, + pid: process.pid, + } +} + +function supportedProtocols() { + return { + daemon: [DAEMON_PROTOCOL], + job: [MANAGED_PLAYWRIGHT_JOB_PROTOCOL_VERSION_V1, MANAGED_PLAYWRIGHT_JOB_PROTOCOL_VERSION_V2], + action: [VISIBLE_ACTION_PROTOCOL_VERSION_V1, VISIBLE_ACTION_PROTOCOL_VERSION_V2], + } +} + +function requireControlAuth(store: JobStore, request: IncomingMessage) { + const authorization = request.headers.authorization + if (authorization === undefined) throw controlAuthMissing() + const value = Array.isArray(authorization) ? authorization[0] : authorization + if (!value || !value.startsWith('Bearer ')) throw controlAuthRejected() + store.requireControlToken(value.slice('Bearer '.length)) +} + +async function readJsonObject(request: IncomingMessage) { + let raw + try { + raw = await readBody(request) + } catch (error) { + if (error instanceof BodyLimitExceededError) { + throw invalidInput('request body must be valid JSON: Failed to buffer the request body') + } + throw error + } + if (!raw) throw invalidInput('request body must be valid JSON: missing request body') + return parseJsonObject(raw) +} + +function parseJsonObject(raw: string) { + try { + const parsed = JSON.parse(raw) as unknown + if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) { + throw new Error('request body must be a JSON object') + } + return parsed as JsonRecord + } catch (error) { + throw invalidInput(`request body must be valid JSON: ${error instanceof Error ? error.message : String(error)}`) + } +} + +async function readBody(request: IncomingMessage) { + const contentLength = request.headers['content-length'] + const declaredLength = Array.isArray(contentLength) ? contentLength[0] : contentLength + if (declaredLength !== undefined && declaredLength !== '') { + const length = Number(declaredLength) + if (Number.isFinite(length) && length > MAX_HTTP_BODY_BYTES) throw new BodyLimitExceededError() + } + const chunks: Buffer[] = [] + let totalBytes = 0 + for await (const chunk of request) { + const buffer = Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk) + totalBytes += buffer.length + if (totalBytes > MAX_HTTP_BODY_BYTES) throw new BodyLimitExceededError() + chunks.push(buffer) + } + const body = Buffer.concat(chunks).toString('utf8') + return body +} + +function writeJson(response: ServerResponse, status: number, body: unknown) { + const payload = JSON.stringify(body) + response.writeHead(status, { + 'content-type': 'application/json', + 'content-length': Buffer.byteLength(payload), + }) + response.end(payload) +} + +function writeText(response: ServerResponse, status: number, body: string) { + response.writeHead(status, { + 'content-type': 'text/plain; charset=utf-8', + 'content-length': Buffer.byteLength(body), + }) + response.end(body) +} + +function writeDaemonError(response: ServerResponse, error: DaemonError) { + writeJson(response, daemonErrorStatus(error), daemonErrorBody(error)) +} + +function validateReadyChallenge(challenge: string) { + if (challenge.length !== READY_CHALLENGE_BASE64URL_CHARS) { + throw invalidInput(`challenge must be canonical unpadded base64url encoding of ${READY_CHALLENGE_BYTES} bytes`) + } + let decoded: Buffer + try { + decoded = Buffer.from(challenge, 'base64url') + } catch { + decoded = Buffer.alloc(0) + } + if (decoded.length !== READY_CHALLENGE_BYTES || decoded.toString('base64url') !== challenge) { + throw invalidInput(`challenge must be canonical unpadded base64url encoding of ${READY_CHALLENGE_BYTES} bytes`) + } +} + +function daemonReadyProof(store: JobStore, challenge: string, canonicalHome: string) { + validateReadyChallenge(challenge) + return createHmac('sha256', store.controlToken()) + .update(lengthPrefixedMessage([ + DAEMON_READY_PROOF_PROTOCOL, + challenge, + DAEMON_PROTOCOL, + NATIVE_PROTOCOL, + canonicalHome, + ])) + .digest('base64url') +} + +function lengthPrefixedMessage(fields: string[]) { + return Buffer.concat(fields.flatMap((field) => { + const value = Buffer.from(field, 'utf8') + const length = Buffer.allocUnsafe(4) + length.writeUInt32BE(value.length) + return [length, value] + })) +} + +function matchJobRoute(pathname: string) { + const match = /^\/jobs\/([^/]+)(?:\/([^/]+))?$/.exec(pathname) + if (!match) return null + const action = match[2] ?? null + if (action !== null && !['claim', 'complete', 'resume'].includes(action)) return null + return { jobId: decodeURIComponent(match[1] || ''), action } +} + +function matchControlJobRoute(pathname: string) { + const match = /^\/control\/jobs\/([^/]+)\/([^/]+)$/.exec(pathname) + if (!match) return null + const action = match[2] ?? '' + if (!['checkpoint', 'park', 'running', 'waiting-for-user', 'renew', 'cancel'].includes(action)) return null + return { jobId: decodeURIComponent(match[1] || ''), action } +} + +function requireField(body: JsonRecord, field: string) { + if (!Object.hasOwn(body, field)) throw invalidInput(`missing field ${field}`) + return body[field] +} + +function requiredString(value: unknown, field: string) { + if (typeof value !== 'string') throw invalidInput(`${field} must be a string`) + return value +} + +function optionalString(value: unknown) { + if (value === undefined || value === null) return null + if (typeof value !== 'string') throw invalidInput('optional string field must be a string') + return value +} + +function optionalQueryString(value: string | null) { + return value === null ? undefined : value +} + +function optionalLimit(value: string | null) { + if (value === null) return undefined + if (!/^(?:0|[1-9][0-9]*)$/.test(value)) { + throw invalidInput('query parameters are invalid: Failed to deserialize query string') + } + const parsed = BigInt(value) + if (parsed > 18_446_744_073_709_551_615n) { + throw invalidInput('query parameters are invalid: Failed to deserialize query string') + } + return parsed > BigInt(Number.MAX_SAFE_INTEGER) ? Number.MAX_SAFE_INTEGER : Number(parsed) +} + +function optionalJobStatus(value: string | null) { + if (value === null) return undefined + const statuses: JobStatus[] = ['queued', 'claimed', 'running', 'waiting_for_user', 'succeeded', 'failed', 'canceled', 'timed_out'] + if (!statuses.includes(value as JobStatus)) throw invalidInput(`invalid status: ${value}`) + return value as JobStatus +} + +function optionalExecutionBackend(value: unknown) { + if (value === undefined || value === null) return undefined + if (value !== 'legacy_extension' && value !== 'playwright') { + throw invalidInput(`invalid execution_backend: ${String(value)}`) + } + return value as ExecutionBackend +} + +async function closeServer(server: http.Server, store: JobStore) { + await new Promise((resolve, reject) => { + server.close((error) => { + if (error) reject(error) + else resolve() + }) + }).catch(() => undefined) + store.close() +} From 393e556582b44323e9fc6283074bc64ca0be0f55 Mon Sep 17 00:00:00 2001 From: jazelly Date: Sun, 26 Jul 2026 11:22:31 +0930 Subject: [PATCH 2/2] feat(cli): complete durable TypeScript daemon --- .github/workflows/ts-daemon-conformance.yml | 35 + package.json | 1 + packages/cli/src/daemon/job-store.ts | 6 +- packages/cli/src/daemon/server.ts | 25 +- test/ts-daemon-conformance.test.mjs | 677 ++++++++++++++++++++ 5 files changed, 735 insertions(+), 9 deletions(-) create mode 100644 .github/workflows/ts-daemon-conformance.yml create mode 100644 test/ts-daemon-conformance.test.mjs diff --git a/.github/workflows/ts-daemon-conformance.yml b/.github/workflows/ts-daemon-conformance.yml new file mode 100644 index 0000000..8693f37 --- /dev/null +++ b/.github/workflows/ts-daemon-conformance.yml @@ -0,0 +1,35 @@ +name: TS daemon conformance + +on: + pull_request: + branches: + - dev + workflow_dispatch: + +permissions: + contents: read + +concurrency: + group: ts-daemon-conformance-${{ github.ref }}-${{ matrix.node-version }} + cancel-in-progress: true + +jobs: + conformance: + name: Node ${{ matrix.node-version }} + runs-on: ubuntu-24.04 + timeout-minutes: 20 + strategy: + fail-fast: false + matrix: + node-version: + - 22.13.0 + - 24 + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-node@v4 + with: + node-version: ${{ matrix.node-version }} + cache: npm + - uses: dtolnay/rust-toolchain@stable + - run: npm ci + - run: npm run test:ts-daemon-conformance diff --git a/package.json b/package.json index e6f9006..eed0e3c 100644 --- a/package.json +++ b/package.json @@ -15,6 +15,7 @@ "release:version": "node scripts/release/version.mjs", "sync:skill": "npm run build:scripts && node dist/scripts/sync-tokenless-skill.mjs", "test": "npm run build && node dist/scripts/test-all.mjs", + "test:ts-daemon-conformance": "npm run build && node --test --test-concurrency=1 test/ts-daemon-conformance.test.mjs", "test:e2e:protocol-cross-version": "npm run build && node scripts/run-protocol-cross-version-e2e.mjs", "test:e2e": "npm run build && node --test --test-concurrency=1 test/live-setup-auth-report.e2e.mjs test/live-provider-guest-access.e2e.mjs test/live-managed-playwright-prompt-actions.e2e.mjs test/live-managed-playwright.e2e.mjs", "test:e2e:live-managed-playwright-m1": "npm run build && node --test --test-concurrency=1 test/live-managed-playwright-prompt-actions.e2e.mjs", diff --git a/packages/cli/src/daemon/job-store.ts b/packages/cli/src/daemon/job-store.ts index 16934ec..7065fd1 100644 --- a/packages/cli/src/daemon/job-store.ts +++ b/packages/cli/src/daemon/job-store.ts @@ -461,6 +461,10 @@ export class JobStore { const serializedResume = stringifyJson(resumeJson) return this.transaction(() => { const job = this.getJobWithoutRecovery(jobId) + const checkpointPresent = this.exists( + 'SELECT EXISTS (SELECT 1 FROM jobs WHERE job_id = ? AND checkpoint_json IS NOT NULL) AS present', + jobId + ) if (job.execution_backend !== 'playwright') { throw invalidInput('only playwright jobs can be resumed with browser visibility') } @@ -470,7 +474,7 @@ export class JobStore { ) { return job } - if (job.status !== 'waiting_for_user' || job.checkpoint_json === null || job.claim_expires_at_ms !== null) { + if (job.status !== 'waiting_for_user' || !checkpointPresent || job.claim_expires_at_ms !== null) { throw invalidJobState(job.job_id, 'parked playwright waiting_for_user', job.status) } const result = this.run( diff --git a/packages/cli/src/daemon/server.ts b/packages/cli/src/daemon/server.ts index a8bef14..86dc2df 100644 --- a/packages/cli/src/daemon/server.ts +++ b/packages/cli/src/daemon/server.ts @@ -145,7 +145,7 @@ async function handleRequest( if (method === 'GET' && url.pathname === '/jobs') { const jobs = store.listJobs({ status: optionalJobStatus(url.searchParams.get('status')), - execution_backend: optionalExecutionBackend(url.searchParams.get('execution_backend')), + execution_backend: optionalQueryExecutionBackend(url.searchParams.get('execution_backend')), profile_id: optionalQueryString(url.searchParams.get('profile_id')), provider: optionalQueryString(url.searchParams.get('provider')), task_id: optionalQueryString(url.searchParams.get('task_id')), @@ -188,7 +188,7 @@ async function handleRequest( } if (method === 'POST' && url.pathname === '/control/jobs/claim-next') { - const executionBackend = optionalExecutionBackend(url.searchParams.get('execution_backend')) ?? 'legacy_extension' + const executionBackend = optionalQueryExecutionBackend(url.searchParams.get('execution_backend')) ?? 'legacy_extension' const job = store.claimNextJob( { provider: optionalQueryString(url.searchParams.get('provider')), @@ -203,8 +203,13 @@ async function handleRequest( const controlJobRoute = matchControlJobRoute(url.pathname) if (controlJobRoute && method === 'POST') { - const rawBody = await readBody(request) - const body = rawBody ? parseJsonObject(rawBody) : {} + if (controlJobRoute.action === 'cancel') { + const rawBody = await readBody(request) + const body = rawBody ? parseJsonObject(rawBody) : {} + writeJson(response, 200, publicView(await store.cancelJob(controlJobRoute.jobId, body.reason))) + return + } + const body = await readJsonObject(request) if (controlJobRoute.action === 'running') { writeJson(response, 200, publicView(store.markRunning(controlJobRoute.jobId, requiredString(body.claim_token, 'claim_token')))) return @@ -238,10 +243,6 @@ async function handleRequest( writeJson(response, 200, publicView(store.renewClaim(controlJobRoute.jobId, requiredString(body.claim_token, 'claim_token')))) return } - if (controlJobRoute.action === 'cancel') { - writeJson(response, 200, publicView(await store.cancelJob(controlJobRoute.jobId, body.reason))) - return - } } if (method === 'POST' && url.pathname === '/control/shutdown') { @@ -457,6 +458,14 @@ function optionalExecutionBackend(value: unknown) { return value as ExecutionBackend } +function optionalQueryExecutionBackend(value: string | null) { + if (value === null) return undefined + if (value !== 'legacy_extension' && value !== 'playwright') { + throw invalidInput('query parameters are invalid: Failed to deserialize query string') + } + return value +} + async function closeServer(server: http.Server, store: JobStore) { await new Promise((resolve, reject) => { server.close((error) => { diff --git a/test/ts-daemon-conformance.test.mjs b/test/ts-daemon-conformance.test.mjs new file mode 100644 index 0000000..034c5bf --- /dev/null +++ b/test/ts-daemon-conformance.test.mjs @@ -0,0 +1,677 @@ +import assert from 'node:assert/strict' +import { spawn, spawnSync } from 'node:child_process' +import { createHmac, randomBytes, randomUUID } from 'node:crypto' +import fs from 'node:fs' +import net from 'node:net' +import os from 'node:os' +import path from 'node:path' +import test from 'node:test' +import { fileURLToPath } from 'node:url' + +const root = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '..') +const cliDir = path.join(root, 'packages/cli') +const tsDaemonEntry = path.join(cliDir, 'dist/src/daemon/daemon-entry.mjs') +const executableSuffix = process.platform === 'win32' ? '.exe' : '' +const nativeTuple = `${process.platform}-${process.arch}` +const rustDaemon = path.join(cliDir, 'npm', `tokenless-native-${nativeTuple}`, 'bin', `tokenless-daemon${executableSuffix}`) + +const createdChildren = new Set() + +test.after(async () => { + await Promise.all([...createdChildren].map((child) => terminateChild(child))) +}) + +test('Rust and opt-in TS daemons share durable jobs through the same SQLite store', { + timeout: 120_000, +}, async () => { + requireBuiltArtifacts() + const homeDir = tempHome('tokenless-ts-rust-interop-') + try { + await rustWritesTsAdvancesRustVerifies(homeDir) + await tsWritesRustAdvancesTsVerifies(homeDir) + assertUnixRestrictivePermissions(homeDir) + } finally { + await terminateChildrenForHome(homeDir) + fs.rmSync(homeDir, { recursive: true, force: true }) + } +}) + +test('TS daemon claim-next is atomic across independent real Node clients', { + timeout: 60_000, +}, async () => { + requireBuiltArtifacts() + const homeDir = tempHome('tokenless-ts-claim-next-') + const daemon = await startTsDaemon(homeDir) + try { + const token = readControlToken(homeDir) + const jobId = randomUUID() + await daemonRequest(daemon.url, token, 'POST', '/jobs', { + provider: 'chatgpt', + action: 'prompt.submit', + request_json: { prompt: 'claim exactly once' }, + job_id: jobId, + }) + + const barrierDir = fs.mkdtempSync(path.join(homeDir, 'claim-next-barrier-')) + const releaseMarker = path.join(barrierDir, 'release') + const readyMarkers = [ + path.join(barrierDir, 'client-a.ready'), + path.join(barrierDir, 'client-b.ready'), + ] + const clientPromises = readyMarkers.map((readyMarker) => runNodeClaimClient({ + homeDir, + daemonUrl: daemon.url, + readyMarker, + releaseMarker, + })) + let clients + try { + await Promise.all(readyMarkers.map((readyMarker) => waitForFile(readyMarker, 10_000))) + fs.writeFileSync(releaseMarker, `${Date.now()}\n`, { flag: 'wx' }) + clients = await Promise.all(clientPromises) + } catch (error) { + if (!fs.existsSync(releaseMarker)) { + fs.writeFileSync(releaseMarker, `${Date.now()}\n`, { flag: 'wx' }) + } + await Promise.allSettled(clientPromises) + throw error + } + const claimed = clients.filter((result) => result.job !== null) + const empty = clients.filter((result) => result.job === null) + assert.equal(claimed.length, 1, JSON.stringify(clients)) + assert.equal(empty.length, 1, JSON.stringify(clients)) + assert.equal(claimed[0].job.job_id, jobId) + assert.equal(claimed[0].job.status, 'claimed') + assert.equal(typeof claimed[0].job.claim_token, 'string') + } finally { + await shutdownDaemon(daemon).catch(() => undefined) + await terminateChildrenForHome(homeDir) + fs.rmSync(homeDir, { recursive: true, force: true }) + } +}) + +test('TS daemon preserves Playwright state, recovers leases, filters summaries, and authenticates control', { + timeout: 150_000, +}, async () => { + requireBuiltArtifacts() + const homeDir = tempHome('tokenless-ts-playwright-') + try { + await verifyPlaywrightRestartResumeCancelAndAuth(homeDir) + await verifyLeaseCrashRecovery(homeDir) + } finally { + await terminateChildrenForHome(homeDir) + fs.rmSync(homeDir, { recursive: true, force: true }) + } +}) + +async function rustWritesTsAdvancesRustVerifies(homeDir) { + let rust = await startRustDaemon(homeDir) + const token = readControlToken(homeDir) + const jobId = randomUUID() + const claimToken = `rust-created-${randomUUID()}` + await daemonRequest(rust.url, token, 'POST', '/jobs', { + provider: 'chatgpt', + action: 'prompt.submit', + request_json: { + prompt: 'created by rust', + metadata: { + taskId: 'rust-origin-task', + projectName: 'conformance', + }, + }, + job_id: jobId, + claim_token: claimToken, + }) + await shutdownDaemon(rust) + rust = null + + let ts = await startTsDaemon(homeDir) + try { + const claimed = await daemonRequest(ts.url, token, 'POST', `/jobs/${encodeURIComponent(jobId)}/claim`, { + claim_token: claimToken, + }) + assert.equal(claimed.status, 'claimed') + const completed = await daemonRequest(ts.url, token, 'POST', `/jobs/${encodeURIComponent(jobId)}/complete`, { + claim_token: claimToken, + result_json: { runtime: 'ts', step: 'advanced' }, + }) + assert.equal(completed.status, 'succeeded') + } finally { + await shutdownDaemon(ts).catch(() => undefined) + } + + const verified = rustCli(homeDir, ['get', jobId]) + assert.equal(verified.job_id, jobId) + assert.equal(verified.status, 'succeeded') + assert.deepEqual(verified.result_json, { runtime: 'ts', step: 'advanced' }) +} + +async function tsWritesRustAdvancesTsVerifies(homeDir) { + let ts = await startTsDaemon(homeDir) + const token = readControlToken(homeDir) + const jobId = randomUUID() + const claimToken = `ts-created-${randomUUID()}` + await daemonRequest(ts.url, token, 'POST', '/jobs', { + provider: 'claude', + action: 'prompt.submit', + request_json: { + prompt: 'created by ts', + taskId: 'ts-origin-task', + idempotencyKey: 'ts-origin-idempotency', + }, + job_id: jobId, + claim_token: claimToken, + }) + await shutdownDaemon(ts) + ts = null + + const claimed = rustCli(homeDir, ['claim', jobId, '--claim-token', claimToken]) + assert.equal(claimed.status, 'claimed') + const completed = rustCli(homeDir, [ + 'complete', + jobId, + '--claim-token', + claimToken, + '--result-json', + JSON.stringify({ runtime: 'rust', step: 'advanced' }), + ]) + assert.equal(completed.status, 'succeeded') + + ts = await startTsDaemon(homeDir) + try { + const verified = await daemonRequest(ts.url, token, 'GET', `/jobs/${encodeURIComponent(jobId)}`) + assert.equal(verified.job_id, jobId) + assert.equal(verified.status, 'succeeded') + assert.deepEqual(verified.result_json, { runtime: 'rust', step: 'advanced' }) + } finally { + await shutdownDaemon(ts).catch(() => undefined) + } +} + +async function verifyPlaywrightRestartResumeCancelAndAuth(homeDir) { + let daemon = await startTsDaemon(homeDir) + let token = readControlToken(homeDir) + const tokenBeforeRestart = token + const challenge = randomBytes(32).toString('base64url') + const readyBefore = await readyProbe(daemon.url, challenge) + assert.equal(readyBefore.home_dir, homeDir) + assert.equal(readyBefore.pid, daemon.child.pid) + assert.equal(readyBefore.ready_proof, readyProof(token, challenge, homeDir)) + + const missingShutdown = await fetch(`${daemon.url}/control/shutdown`, { method: 'POST' }) + assert.equal(missingShutdown.status, 401) + const rejectedShutdown = await fetch(`${daemon.url}/control/shutdown`, { + method: 'POST', + headers: { authorization: 'Bearer not-the-token' }, + }) + assert.equal(rejectedShutdown.status, 403) + assert.equal((await readyProbe(daemon.url)).status, 'ok') + + const profileId = 'default-profile' + const resumedJobId = randomUUID() + const otherJobId = randomUUID() + await daemonRequest(daemon.url, token, 'POST', '/jobs', { + provider: 'chatgpt', + action: 'prompt.submit', + execution_backend: 'playwright', + profile_id: profileId, + job_id: resumedJobId, + request_json: { + prompt: 'checkpoint and resume', + taskId: 'root-task-key', + idempotencyKey: 'root-idempotency-key', + metadata: { + taskId: 'metadata-task-key', + projectName: 'conformance-project', + chatName: 'checkpoint-chat', + idempotencyKey: 'metadata-idempotency-key', + }, + }, + }) + await daemonRequest(daemon.url, token, 'POST', '/jobs', { + provider: 'chatgpt', + action: 'prompt.submit', + execution_backend: 'playwright', + profile_id: profileId, + job_id: otherJobId, + request_json: { + prompt: 'not selected by summary filter', + metadata: { taskId: 'other-task-key' }, + }, + }) + + const filtered = await daemonRequest(daemon.url, token, 'GET', '/jobs?task_id=metadata-task-key&execution_backend=playwright&limit=10') + assert.deepEqual(filtered.map((job) => job.job_id), [resumedJobId]) + assert.equal(Object.hasOwn(filtered[0], 'checkpoint_json'), false) + assert.equal(Object.hasOwn(filtered[0], 'resume_json'), false) + + const claimed = await daemonRequest( + daemon.url, + token, + 'POST', + `/control/jobs/claim-next?execution_backend=playwright&profile_id=${encodeURIComponent(profileId)}` + ) + assert.equal(claimed.job.job_id, resumedJobId) + assert.equal(claimed.job.status, 'claimed') + const firstClaimToken = claimed.job.claim_token + const checkpoint = { + profileId, + provider: 'chatgpt', + actionCursor: 1, + phase: { name: 'after-submit' }, + page: { url: 'https://chatgpt.com/c/real-boundary' }, + } + const running = await daemonRequest(daemon.url, token, 'POST', `/control/jobs/${resumedJobId}/running`, { + claim_token: firstClaimToken, + }) + assert.equal(running.status, 'running') + await daemonRequest(daemon.url, token, 'POST', `/control/jobs/${resumedJobId}/checkpoint`, { + claim_token: firstClaimToken, + checkpoint_json: checkpoint, + }) + const parked = await daemonRequest(daemon.url, token, 'POST', `/control/jobs/${resumedJobId}/park`, { + claim_token: firstClaimToken, + blocker_json: { reason: 'provider_verification', browser: { windowOpen: false } }, + checkpoint_json: checkpoint, + }) + assert.equal(parked.status, 'waiting_for_user') + assert.equal(parked.blocker_json.browser.resumeRequired, undefined) + + const shutdown = await shutdownDaemon(daemon) + assert.equal(shutdown.status, 'shutting_down') + daemon = null + + daemon = await startTsDaemon(homeDir) + token = readControlToken(homeDir) + const readyAfter = await readyProbe(daemon.url, challenge) + assert.equal(token, tokenBeforeRestart) + assert.equal(readyAfter.home_dir, readyBefore.home_dir) + assert.equal(readyAfter.pid, daemon.child.pid) + assert.equal(readyAfter.ready_proof, readyBefore.ready_proof) + assert.equal(readyAfter.ready_proof, readyProof(token, challenge, homeDir)) + + const resumed = await daemonRequest(daemon.url, token, 'POST', `/jobs/${resumedJobId}/resume`, { + browser_visibility: 'headed', + }) + assert.equal(resumed.status, 'queued') + const reclaimed = await daemonRequest( + daemon.url, + token, + 'POST', + `/control/jobs/claim-next?execution_backend=playwright&profile_id=${encodeURIComponent(profileId)}` + ) + assert.equal(reclaimed.job.job_id, resumedJobId) + assert.notEqual(reclaimed.job.claim_token, firstClaimToken) + assert.deepEqual(reclaimed.job.resume_json, { browser_visibility: 'headed' }) + assert.deepEqual(reclaimed.job.checkpoint_json, checkpoint) + + const canceled = await daemonRequest(daemon.url, token, 'POST', `/control/jobs/${resumedJobId}/cancel`) + assert.equal(canceled.status, 'canceled') + assert.equal(canceled.error_json.code, 'job_canceled') + + const rootKeyFiltered = await daemonRequest(daemon.url, token, 'GET', '/jobs?task_id=root-idempotency-key&limit=10') + assert.deepEqual(rootKeyFiltered.map((job) => job.job_id), [resumedJobId]) + + await shutdownDaemon(daemon) + daemon = null +} + +async function verifyLeaseCrashRecovery(homeDir) { + let daemon = await startTsDaemon(homeDir) + const token = readControlToken(homeDir) + const jobId = randomUUID() + await daemonRequest(daemon.url, token, 'POST', '/jobs', { + provider: 'gemini', + action: 'prompt.submit', + request_json: { prompt: 'recover expired claim' }, + job_id: jobId, + }) + const firstClaim = await daemonRequest(daemon.url, token, 'POST', '/control/jobs/claim-next') + assert.equal(firstClaim.job.job_id, jobId) + const staleClaimToken = firstClaim.job.claim_token + await killChild(daemon.child) + daemon = null + + await delay(32_000) + + daemon = await startTsDaemon(homeDir) + try { + const reclaimed = await daemonRequest(daemon.url, token, 'POST', '/control/jobs/claim-next') + assert.equal(reclaimed.job.job_id, jobId) + assert.equal(reclaimed.job.status, 'claimed') + assert.notEqual(reclaimed.job.claim_token, staleClaimToken) + + const staleCompletion = await fetch(`${daemon.url}/jobs/${encodeURIComponent(jobId)}/complete`, { + method: 'POST', + headers: jsonHeaders(token), + body: JSON.stringify({ + claim_token: staleClaimToken, + result_json: { stale: true }, + }), + }) + assert.equal(staleCompletion.status, 403) + const staleBody = await staleCompletion.json() + assert.equal(staleBody.error.code, 'claim_rejected') + + const completed = await daemonRequest(daemon.url, token, 'POST', `/jobs/${encodeURIComponent(jobId)}/complete`, { + claim_token: reclaimed.job.claim_token, + result_json: { recovered: true }, + }) + assert.equal(completed.status, 'succeeded') + assert.deepEqual(completed.result_json, { recovered: true }) + } finally { + await shutdownDaemon(daemon).catch(() => undefined) + } +} + +async function startTsDaemon(homeDir) { + const port = await freePort() + const child = spawn(process.execPath, [ + tsDaemonEntry, + '--home', + homeDir, + 'serve', + '--host', + '127.0.0.1', + '--port', + String(port), + ], { + cwd: root, + env: { ...process.env, TOKENLESS_HOME: homeDir }, + stdio: ['ignore', 'pipe', 'pipe'], + }) + return waitForDaemon(child, `http://127.0.0.1:${port}`, homeDir, 'TS') +} + +async function startRustDaemon(homeDir) { + const port = await freePort() + const child = spawn(rustDaemon, [ + '--home', + homeDir, + 'serve', + '--host', + '127.0.0.1', + '--port', + String(port), + ], { + cwd: root, + env: { ...process.env, TOKENLESS_HOME: homeDir }, + stdio: ['ignore', 'pipe', 'pipe'], + }) + return waitForDaemon(child, `http://127.0.0.1:${port}`, homeDir, 'Rust') +} + +async function waitForDaemon(child, url, homeDir, label) { + trackChild(child, homeDir) + let stdout = '' + let stderr = '' + child.stdout?.on('data', (chunk) => { + stdout += chunk.toString('utf8') + }) + child.stderr?.on('data', (chunk) => { + stderr += chunk.toString('utf8') + }) + const exited = new Promise((resolve) => { + child.once('exit', (code, signal) => resolve({ code, signal })) + }) + const deadline = Date.now() + 10_000 + let lastError + while (Date.now() < deadline) { + const exit = await Promise.race([exited, delay(50).then(() => null)]) + if (exit) { + throw new Error(`${label} daemon exited before ready: ${JSON.stringify(exit)}\nstdout:\n${stdout}\nstderr:\n${stderr}`) + } + try { + const ready = await readyProbe(url) + assert.equal(ready.ready, true) + assert.equal(ready.home_dir, homeDir) + return { child, url, homeDir, label } + } catch (error) { + lastError = error + } + } + throw new Error(`${label} daemon did not become ready: ${lastError?.message ?? lastError}\nstdout:\n${stdout}\nstderr:\n${stderr}`) +} + +async function shutdownDaemon(daemon) { + if (!daemon) return null + const token = readControlToken(daemon.homeDir) + const body = await daemonRequest(daemon.url, token, 'POST', '/control/shutdown') + await waitForExit(daemon.child, 5_000) + createdChildren.delete(daemon.child) + return body +} + +async function daemonRequest(daemonUrl, token, method, requestPath, body) { + const init = { + method, + headers: jsonHeaders(token), + } + if (body !== undefined) init.body = JSON.stringify(body) + const response = await fetch(`${daemonUrl}${requestPath}`, init) + const text = await response.text() + let payload = null + if (text) payload = JSON.parse(text) + assert.equal(response.ok, true, `${method} ${requestPath} returned ${response.status}: ${text}`) + return payload +} + +function jsonHeaders(token) { + return { + accept: 'application/json', + authorization: `Bearer ${token}`, + 'content-type': 'application/json', + } +} + +async function readyProbe(daemonUrl, challenge = randomBytes(32).toString('base64url')) { + const response = await fetch(`${daemonUrl}/ready?challenge=${challenge}`) + const text = await response.text() + assert.equal(response.ok, true, `/ready returned ${response.status}: ${text}`) + return JSON.parse(text) +} + +function readyProof(token, challenge, homeDir) { + return createHmac('sha256', token) + .update(lengthPrefixedMessage([ + 'tokenless.daemon-ready-proof.v1', + challenge, + 'tokenless.daemon.v1', + 'tokenless.native.v1', + homeDir, + ])) + .digest('base64url') +} + +function lengthPrefixedMessage(fields) { + return Buffer.concat(fields.flatMap((field) => { + const value = Buffer.from(field, 'utf8') + const length = Buffer.allocUnsafe(4) + length.writeUInt32BE(value.length) + return [length, value] + })) +} + +function rustCli(homeDir, args) { + const result = spawnSync(rustDaemon, ['--home', homeDir, ...args], { + cwd: root, + env: { ...process.env, TOKENLESS_HOME: homeDir }, + encoding: 'utf8', + timeout: 20_000, + }) + assert.equal(result.status, 0, `rust daemon ${args.join(' ')} failed\nstdout:\n${result.stdout}\nstderr:\n${result.stderr}`) + return JSON.parse(result.stdout) +} + +async function runNodeClaimClient({ homeDir, daemonUrl, readyMarker, releaseMarker }) { + const script = ` +const { existsSync, readFileSync, writeFileSync } = await import('node:fs') +const path = await import('node:path') +const token = readFileSync(path.join(process.env.TOKENLESS_HOME, 'daemon.token'), 'utf8').trim() +writeFileSync(process.env.TOKENLESS_READY_MARKER, process.pid + '\\n', { flag: 'wx' }) +const releaseDeadline = Date.now() + 10000 +while (!existsSync(process.env.TOKENLESS_RELEASE_MARKER)) { + if (Date.now() > releaseDeadline) { + console.error('claim client timed out waiting for release marker') + process.exit(1) + } + await new Promise((resolve) => setTimeout(resolve, 10)) +} +const response = await fetch(process.env.TOKENLESS_DAEMON_URL + '/control/jobs/claim-next', { + method: 'POST', + headers: { accept: 'application/json', authorization: 'Bearer ' + token }, +}) +const body = await response.text() +if (!response.ok) { + console.error(body) + process.exit(1) +} +console.log(body) +` + const child = spawn(process.execPath, ['--input-type=module', '--eval', script], { + cwd: root, + env: { + ...process.env, + TOKENLESS_HOME: homeDir, + TOKENLESS_DAEMON_URL: daemonUrl, + TOKENLESS_READY_MARKER: readyMarker, + TOKENLESS_RELEASE_MARKER: releaseMarker, + }, + stdio: ['ignore', 'pipe', 'pipe'], + }) + trackChild(child, homeDir) + let stdout = '' + let stderr = '' + child.stdout?.on('data', (chunk) => { + stdout += chunk.toString('utf8') + }) + child.stderr?.on('data', (chunk) => { + stderr += chunk.toString('utf8') + }) + return new Promise((resolve, reject) => { + const timeout = setTimeout(() => { + child.kill('SIGKILL') + reject(new Error(`claim client timed out\nstdout:\n${stdout}\nstderr:\n${stderr}`)) + }, 10_000) + child.once('exit', (code, signal) => { + clearTimeout(timeout) + createdChildren.delete(child) + if (code !== 0) { + reject(new Error(`claim client failed with ${JSON.stringify({ code, signal })}\nstdout:\n${stdout}\nstderr:\n${stderr}`)) + return + } + try { + resolve(JSON.parse(stdout)) + } catch (error) { + reject(new Error(`claim client returned invalid JSON: ${error instanceof Error ? error.message : String(error)}\nstdout:\n${stdout}\nstderr:\n${stderr}`)) + } + }) + }) +} + +function trackChild(child, homeDir) { + child.tokenlessHome = homeDir + createdChildren.add(child) + child.once('exit', () => { + createdChildren.delete(child) + }) +} + +async function terminateChildrenForHome(homeDir) { + await Promise.all([...createdChildren] + .filter((child) => child.tokenlessHome === homeDir) + .map((child) => terminateChild(child))) +} + +async function killChild(child) { + if (!child || child.exitCode !== null || child.signalCode !== null) { + if (child) createdChildren.delete(child) + return + } + child.kill('SIGKILL') + await waitForExit(child, 5_000).catch(() => undefined) + createdChildren.delete(child) +} + +async function terminateChild(child) { + if (!child || child.exitCode !== null || child.signalCode !== null) { + if (child) createdChildren.delete(child) + return + } + child.kill('SIGTERM') + try { + await waitForExit(child, 2_000) + } catch { + child.kill('SIGKILL') + await waitForExit(child, 2_000).catch(() => undefined) + } finally { + createdChildren.delete(child) + } +} + +async function waitForFile(filePath, timeoutMs) { + const deadline = Date.now() + timeoutMs + while (Date.now() < deadline) { + if (fs.existsSync(filePath)) return + await delay(10) + } + throw new Error(`file did not appear within ${timeoutMs} ms: ${filePath}`) +} + +function waitForExit(child, timeoutMs) { + if (child.exitCode !== null || child.signalCode !== null) return Promise.resolve() + return new Promise((resolve, reject) => { + const timeout = setTimeout(() => { + child.off('exit', onExit) + reject(new Error(`process ${child.pid} did not exit within ${timeoutMs} ms`)) + }, timeoutMs) + const onExit = () => { + clearTimeout(timeout) + resolve() + } + child.once('exit', onExit) + }) +} + +function tempHome(prefix) { + return fs.realpathSync(fs.mkdtempSync(path.join(os.tmpdir(), prefix))) +} + +function readControlToken(homeDir) { + return fs.readFileSync(path.join(homeDir, 'daemon.token'), 'utf8').trim() +} + +function assertUnixRestrictivePermissions(homeDir) { + if (process.platform === 'win32') return + assert.equal(fileMode(homeDir), 0o700) + assert.equal(fileMode(path.join(homeDir, 'daemon.token')), 0o600) + assert.equal(fileMode(path.join(homeDir, 'tokenless.sqlite3')), 0o600) +} + +function fileMode(filePath) { + return fs.statSync(filePath).mode & 0o777 +} + +async function freePort() { + return new Promise((resolve, reject) => { + const server = net.createServer() + server.unref() + server.once('error', reject) + server.listen(0, '127.0.0.1', () => { + const address = server.address() + server.close(() => { + if (!address || typeof address === 'string') reject(new Error('failed to allocate a port')) + else resolve(address.port) + }) + }) + }) +} + +function delay(ms) { + return new Promise((resolve) => setTimeout(resolve, ms)) +} + +function requireBuiltArtifacts() { + assert.equal(fs.existsSync(tsDaemonEntry), true, `missing compiled TS daemon: ${tsDaemonEntry}`) + assert.equal(fs.existsSync(rustDaemon), true, `missing packaged Rust daemon: ${rustDaemon}`) +}