feat: bootstrap commercial AI drama platform
This commit is contained in:
@@ -0,0 +1,282 @@
|
||||
import { dbAll, dbGet, dbRun, withTransaction } from "./db.mjs";
|
||||
import { executeGenerationJob } from "./execution.mjs";
|
||||
|
||||
const workerId = process.env.AI_DRAMA_WORKER_ID || `local-worker-${process.pid}`;
|
||||
const pollMs = Math.max(250, Number(process.env.AI_DRAMA_WORKER_POLL_MS || 1200));
|
||||
const maxConcurrency = Math.max(1, Math.min(8, Number(process.env.AI_DRAMA_WORKER_CONCURRENCY || 2)));
|
||||
const leaseMs = Math.max(30_000, Number(process.env.AI_DRAMA_WORKER_LEASE_MS || 300_000));
|
||||
const staleAfterMs = Math.max(15_000, Number(process.env.AI_DRAMA_WORKER_STALE_MS || Math.max(30_000, pollMs * 5)));
|
||||
const enabled = process.env.AI_DRAMA_WORKER_ENABLED !== "0";
|
||||
|
||||
const state = {
|
||||
workerId,
|
||||
enabled,
|
||||
pollMs,
|
||||
maxConcurrency,
|
||||
inFlight: new Set(),
|
||||
startedAt: enabled ? new Date().toISOString() : null,
|
||||
lastPollAt: null,
|
||||
lastClaimAt: null,
|
||||
lastCompletedAt: null,
|
||||
lastFailedAt: null,
|
||||
lastHeartbeatAt: null,
|
||||
lastReclaimAt: null,
|
||||
lastReclaimedCount: 0,
|
||||
lastRetryAt: null,
|
||||
lastRetryCount: 0,
|
||||
lastError: "",
|
||||
timer: null,
|
||||
pumping: false
|
||||
};
|
||||
|
||||
function now() {
|
||||
return new Date().toISOString();
|
||||
}
|
||||
|
||||
function workerContext(job) {
|
||||
const organization = dbGet("SELECT * FROM organizations WHERE id = ?", [job.organization_id]);
|
||||
const workspace = dbGet("SELECT * FROM workspaces WHERE id = ?", [job.workspace_id]);
|
||||
const project = dbGet("SELECT * FROM projects WHERE id = ?", [job.project_id]);
|
||||
const user = dbGet("SELECT * FROM users WHERE id = 'u-local-worker'") || dbGet("SELECT * FROM users WHERE id = ?", [job.created_by]);
|
||||
return {
|
||||
user,
|
||||
organization,
|
||||
workspace,
|
||||
project,
|
||||
permissions: ["job:create", "queue:manage", "model:manage", "usage:view", "audit:view"],
|
||||
roles: [{ key: "local_worker", name: "本地 Worker", scope: "system" }],
|
||||
systemAdmin: false,
|
||||
orgElevated: true,
|
||||
workspaces: [],
|
||||
projects: []
|
||||
};
|
||||
}
|
||||
|
||||
function localQueueDepth() {
|
||||
return Number(dbGet(
|
||||
`SELECT COUNT(*) AS count
|
||||
FROM generation_jobs j
|
||||
JOIN model_connectors m ON m.id = j.adapter_id
|
||||
WHERE j.status = 'queued' AND m.cost_mode = 'local' AND m.status = 'ready'`,
|
||||
[]
|
||||
)?.count || 0);
|
||||
}
|
||||
|
||||
function queueAlert() {
|
||||
const oldest = dbGet("SELECT MIN(created_at) AS oldest FROM generation_jobs WHERE status = 'queued'", []);
|
||||
const queueDepth = localQueueDepth();
|
||||
const oldestAt = oldest?.oldest || null;
|
||||
const oldestAgeMs = oldestAt ? Math.max(0, Date.now() - Date.parse(oldestAt)) : 0;
|
||||
return {
|
||||
level: queueDepth === 0 ? "none" : oldestAgeMs >= 10 * 60 * 1000 ? "critical" : oldestAgeMs >= 3 * 60 * 1000 ? "warning" : "normal",
|
||||
queueDepth,
|
||||
oldestAt,
|
||||
oldestAgeMs
|
||||
};
|
||||
}
|
||||
|
||||
function updateWorkerHealth(status = enabled ? "ready" : "paused") {
|
||||
const timestamp = now();
|
||||
state.lastHeartbeatAt = timestamp;
|
||||
const alert = queueAlert();
|
||||
const metadata = {
|
||||
workerId,
|
||||
enabled,
|
||||
concurrency: maxConcurrency,
|
||||
inFlight: state.inFlight.size,
|
||||
pollMs,
|
||||
lastPollAt: state.lastPollAt,
|
||||
lastClaimAt: state.lastClaimAt,
|
||||
lastCompletedAt: state.lastCompletedAt,
|
||||
lastFailedAt: state.lastFailedAt,
|
||||
lastHeartbeatAt: state.lastHeartbeatAt,
|
||||
staleAfterMs,
|
||||
heartbeatAgeMs: 0,
|
||||
lastReclaimAt: state.lastReclaimAt,
|
||||
lastReclaimedCount: state.lastReclaimedCount,
|
||||
lastRetryAt: state.lastRetryAt,
|
||||
lastRetryCount: state.lastRetryCount,
|
||||
queueAlert: alert,
|
||||
lastError: state.lastError
|
||||
};
|
||||
dbRun(
|
||||
`UPDATE service_health
|
||||
SET status = ?, queue_depth = ?, last_heartbeat = ?, version = ?, metadata_json = ?, updated_at = ?
|
||||
WHERE service_key = 'local-worker'`,
|
||||
[status, localQueueDepth(), timestamp, "node-24-local-worker", JSON.stringify(metadata), timestamp]
|
||||
);
|
||||
}
|
||||
|
||||
function reclaimStaleLeases() {
|
||||
const timestamp = now();
|
||||
const staleBefore = new Date(Date.now() - leaseMs).toISOString();
|
||||
const staleJobs = dbAll("SELECT id, max_attempts FROM generation_jobs WHERE status = 'running' AND leased_at IS NOT NULL AND leased_at < ?", [staleBefore]);
|
||||
if (!staleJobs.length) return 0;
|
||||
withTransaction(() => {
|
||||
for (const job of staleJobs) {
|
||||
const attempt = Number(dbGet("SELECT MAX(attempt_number) AS attempt_number FROM job_attempts WHERE job_id = ?", [job.id])?.attempt_number || 0);
|
||||
const message = `Worker 租约在 ${leaseMs}ms 后失效,已回收第 ${attempt} 次执行`;
|
||||
const exhausted = attempt >= Number(job.max_attempts || 3);
|
||||
dbRun("UPDATE job_attempts SET status = 'failed', error_message = ?, finished_at = ? WHERE job_id = ? AND attempt_number = ? AND status = 'running'", [message, timestamp, job.id, attempt]);
|
||||
dbRun("UPDATE generation_jobs SET status = ?, error_message = ?, next_run_at = ?, leased_by = NULL, leased_at = NULL, finished_at = CASE WHEN ? THEN ? ELSE finished_at END, updated_at = ? WHERE id = ?", [exhausted ? "failed" : "queued", message, exhausted ? null : timestamp, exhausted ? 1 : 0, exhausted ? timestamp : null, timestamp, job.id]);
|
||||
}
|
||||
});
|
||||
state.lastReclaimAt = timestamp;
|
||||
state.lastReclaimedCount = staleJobs.length;
|
||||
state.lastError = `${staleJobs.length} 个过期任务租约已回收`;
|
||||
return staleJobs.length;
|
||||
}
|
||||
|
||||
function scheduleRetry(jobId, message) {
|
||||
const job = dbGet("SELECT max_attempts FROM generation_jobs WHERE id = ?", [jobId]);
|
||||
const attempt = Number(dbGet("SELECT MAX(attempt_number) AS attempt_number FROM job_attempts WHERE job_id = ?", [jobId])?.attempt_number || 0);
|
||||
const maxAttempts = Number(job?.max_attempts || 3);
|
||||
if (!job || attempt >= maxAttempts) return false;
|
||||
const delayMs = Math.min(120_000, 2 ** Math.max(0, attempt - 1) * 2_000);
|
||||
const nextRunAt = new Date(Date.now() + delayMs).toISOString();
|
||||
dbRun("UPDATE generation_jobs SET status = 'queued', next_run_at = ?, error_message = ?, leased_by = NULL, leased_at = NULL, updated_at = ? WHERE id = ? AND status = 'failed'", [`${message};将在 ${Math.ceil(delayMs / 1000)} 秒后自动重试(${attempt}/${maxAttempts})`, nextRunAt, now(), jobId]);
|
||||
state.lastRetryAt = now();
|
||||
state.lastRetryCount += 1;
|
||||
return true;
|
||||
}
|
||||
|
||||
function claimNextJob() {
|
||||
const timestamp = now();
|
||||
const staleBefore = new Date(Date.now() - leaseMs).toISOString();
|
||||
return withTransaction(() => {
|
||||
const row = dbGet(
|
||||
`SELECT j.*
|
||||
FROM generation_jobs j
|
||||
JOIN model_connectors m ON m.id = j.adapter_id
|
||||
WHERE j.status = 'queued'
|
||||
AND (j.next_run_at IS NULL OR j.next_run_at <= ?)
|
||||
AND (j.leased_by IS NULL OR j.leased_at < ?)
|
||||
AND m.cost_mode = 'local'
|
||||
AND m.status = 'ready'
|
||||
AND NOT EXISTS (
|
||||
SELECT 1
|
||||
FROM job_dependencies d
|
||||
JOIN generation_jobs dependency ON dependency.id = d.depends_on_job_id
|
||||
WHERE d.job_id = j.id AND dependency.status <> 'completed'
|
||||
)
|
||||
ORDER BY j.priority DESC, j.created_at ASC
|
||||
LIMIT 1`,
|
||||
[timestamp, staleBefore]
|
||||
);
|
||||
if (!row) return null;
|
||||
const result = dbRun(
|
||||
`UPDATE generation_jobs
|
||||
SET leased_by = ?, leased_at = ?, updated_at = ?
|
||||
WHERE id = ? AND status = 'queued' AND (leased_by IS NULL OR leased_at < ?)`,
|
||||
[workerId, timestamp, timestamp, row.id, staleBefore]
|
||||
);
|
||||
if (!Number(result?.changes || 0)) return null;
|
||||
state.lastClaimAt = timestamp;
|
||||
return { ...row, leased_by: workerId, leased_at: timestamp };
|
||||
});
|
||||
}
|
||||
|
||||
function clearLease(jobId) {
|
||||
dbRun(
|
||||
"UPDATE generation_jobs SET leased_by = NULL, leased_at = NULL, updated_at = ? WHERE id = ? AND leased_by = ?",
|
||||
[now(), jobId, workerId]
|
||||
);
|
||||
}
|
||||
|
||||
async function processClaimedJob(job) {
|
||||
try {
|
||||
const result = await executeGenerationJob(workerContext(job), job.id, { workerId, source: "local-worker" });
|
||||
state.lastCompletedAt = now();
|
||||
state.lastError = "";
|
||||
return { jobId: job.id, status: result.job?.status || "completed" };
|
||||
} catch (error) {
|
||||
state.lastFailedAt = now();
|
||||
state.lastError = String(error.message || error).slice(0, 500);
|
||||
const retryScheduled = scheduleRetry(job.id, state.lastError);
|
||||
return { jobId: job.id, status: retryScheduled ? "queued" : "failed", retryScheduled, error: state.lastError };
|
||||
} finally {
|
||||
clearLease(job.id);
|
||||
}
|
||||
}
|
||||
|
||||
async function pump() {
|
||||
if (!enabled || state.pumping) return;
|
||||
state.pumping = true;
|
||||
state.lastPollAt = now();
|
||||
try {
|
||||
reclaimStaleLeases();
|
||||
while (state.inFlight.size < maxConcurrency) {
|
||||
const job = claimNextJob();
|
||||
if (!job) break;
|
||||
state.inFlight.add(job.id);
|
||||
const task = processClaimedJob(job);
|
||||
void task.finally(() => {
|
||||
state.inFlight.delete(job.id);
|
||||
updateWorkerHealth();
|
||||
});
|
||||
}
|
||||
updateWorkerHealth();
|
||||
} catch (error) {
|
||||
// SQLite can briefly reject BEGIN IMMEDIATE while another local process commits.
|
||||
// Keep the worker alive and let the next poll retry instead of taking down the API.
|
||||
state.lastError = String(error?.message || error).slice(0, 500);
|
||||
try {
|
||||
updateWorkerHealth("degraded");
|
||||
} catch {
|
||||
// A second lock while reporting health is still transient; the next poll retries.
|
||||
}
|
||||
} finally {
|
||||
state.pumping = false;
|
||||
}
|
||||
}
|
||||
|
||||
export function workerStatus() {
|
||||
const heartbeatAt = state.lastHeartbeatAt || state.lastPollAt || state.startedAt;
|
||||
const heartbeatAgeMs = heartbeatAt ? Math.max(0, Date.now() - Date.parse(heartbeatAt)) : null;
|
||||
const stale = enabled && (heartbeatAgeMs === null || heartbeatAgeMs > staleAfterMs);
|
||||
return {
|
||||
...state,
|
||||
inFlight: state.inFlight.size,
|
||||
queueDepth: localQueueDepth(),
|
||||
staleAfterMs,
|
||||
lastHeartbeatAt: state.lastHeartbeatAt,
|
||||
heartbeatAgeMs,
|
||||
healthStatus: stale ? "stale" : enabled ? "ready" : "paused",
|
||||
stale,
|
||||
queueAlert: queueAlert(),
|
||||
timer: undefined,
|
||||
pumping: undefined
|
||||
};
|
||||
}
|
||||
|
||||
export async function runWorkerOnce() {
|
||||
if (!enabled) return { ...workerStatus(), dispatched: 0, disabled: true };
|
||||
reclaimStaleLeases();
|
||||
const job = claimNextJob();
|
||||
if (!job) {
|
||||
updateWorkerHealth();
|
||||
return { ...workerStatus(), dispatched: 0 };
|
||||
}
|
||||
state.inFlight.add(job.id);
|
||||
const result = await processClaimedJob(job);
|
||||
state.inFlight.delete(job.id);
|
||||
updateWorkerHealth();
|
||||
return { ...workerStatus(), dispatched: 1, result };
|
||||
}
|
||||
|
||||
export function startLocalWorker() {
|
||||
if (!enabled || state.timer) {
|
||||
updateWorkerHealth(enabled ? "ready" : "paused");
|
||||
return workerStatus();
|
||||
}
|
||||
updateWorkerHealth("ready");
|
||||
state.timer = setInterval(() => { void pump(); }, pollMs);
|
||||
void pump();
|
||||
return workerStatus();
|
||||
}
|
||||
|
||||
export function stopLocalWorker() {
|
||||
if (state.timer) clearInterval(state.timer);
|
||||
state.timer = null;
|
||||
updateWorkerHealth("paused");
|
||||
}
|
||||
Reference in New Issue
Block a user