import { pathToFileURL } from 'node:url'; import { createOrchestratorApp } from './app.mjs'; import { createOrchestratorCheckpoint } from './checkpoint.mjs'; import { createExecutorJobPersistence } from './executor-job-store.mjs'; import { createLangGraphOrchestratorRuntime } from './runtime.mjs'; function positivePort(value, fallback = 8093) { const parsed = Number(value); if (!Number.isInteger(parsed) || parsed < 1 || parsed > 65535) return fallback; return parsed; } function envFlag(value, fallback = false) { const normalized = String(value ?? '').trim().toLowerCase(); if (!normalized) return fallback; return ['1', 'true', 'yes', 'on'].includes(normalized); } function csvList(value) { return [...new Set( String(value ?? '') .split(',') .map((item) => item.trim().toLowerCase()) .filter(Boolean), )]; } function isLoopbackHost(value) { const host = String(value ?? '').trim().toLowerCase(); return host === 'localhost' || host === '::1' || host === '[::1]' || /^127(?:\.\d{1,3}){3}$/.test(host); } function assertSecureServerConfig({ host, executionEnabled, serviceToken, workerToken, } = {}) { const hasServiceToken = Boolean(String(serviceToken ?? '').trim()); const hasWorkerToken = Boolean(String(workerToken ?? '').trim()); if (!isLoopbackHost(host) && !hasServiceToken) { const error = new Error( 'MEMIND_ORCHESTRATOR_SERVICE_TOKEN is required for non-loopback binding', ); error.code = 'ORCHESTRATOR_SERVICE_TOKEN_REQUIRED'; throw error; } if (executionEnabled && (!hasServiceToken || !hasWorkerToken)) { const error = new Error( 'Execution requires both MEMIND_ORCHESTRATOR_SERVICE_TOKEN and MEMIND_ORCHESTRATOR_WORKER_TOKEN', ); error.code = 'ORCHESTRATOR_EXECUTION_TOKENS_REQUIRED'; throw error; } if ( executionEnabled && String(serviceToken).trim() === String(workerToken).trim() ) { const error = new Error('Execution service and worker tokens must be distinct'); error.code = 'ORCHESTRATOR_EXECUTION_TOKENS_NOT_DISTINCT'; throw error; } } function startRetentionSweep(runtime, { retentionDays, intervalMs = 24 * 60 * 60 * 1000, limit = 100, nowMs = Date.now, setIntervalFn = setInterval, clearIntervalFn = clearInterval, logger = console, } = {}) { const days = Number(retentionDays); if (!Number.isFinite(days) || days <= 0) return null; const boundedDays = Math.min(3650, Math.max(1, Math.floor(days))); const configuredInterval = Number(intervalMs); const boundedInterval = Number.isFinite(configuredInterval) ? Math.max(60 * 60 * 1000, configuredInterval) : 24 * 60 * 60 * 1000; let active = false; const sweep = async () => { if (active) return null; active = true; try { const result = await runtime.purgeTerminalRuns({ before: nowMs() - boundedDays * 24 * 60 * 60 * 1000, limit, dryRun: false, }); if (result.deleted.length || result.failures.length) { logger.log( `[orchestrator-retention] deleted=${result.deleted.length} failures=${result.failures.length}`, ); } return result; } catch (error) { logger.warn( '[orchestrator-retention] sweep failed:', error instanceof Error ? error.message : error, ); return null; } finally { active = false; } }; const timer = setIntervalFn(() => void sweep(), boundedInterval); timer?.unref?.(); return { retentionDays: boundedDays, intervalMs: boundedInterval, sweep, stop() { clearIntervalFn(timer); }, }; } export async function startOrchestratorServer({ env = process.env, logger = console, } = {}) { const host = String(env.MEMIND_ORCHESTRATOR_HOST ?? '127.0.0.1').trim() || '127.0.0.1'; const executionEnabled = envFlag(env.MEMIND_ORCHESTRATOR_EXECUTION_ENABLED, false); assertSecureServerConfig({ host, executionEnabled, serviceToken: env.MEMIND_ORCHESTRATOR_SERVICE_TOKEN, workerToken: env.MEMIND_ORCHESTRATOR_WORKER_TOKEN, }); const checkpoint = await createOrchestratorCheckpoint({ mode: env.MEMIND_ORCHESTRATOR_CHECKPOINT_MODE, connectionString: env.MEMIND_ORCHESTRATOR_DATABASE_URL, schema: env.MEMIND_ORCHESTRATOR_DATABASE_SCHEMA, }); let executorJobs; try { executorJobs = await createExecutorJobPersistence({ mode: env.MEMIND_ORCHESTRATOR_CHECKPOINT_MODE, connectionString: env.MEMIND_ORCHESTRATOR_DATABASE_URL, schema: env.MEMIND_ORCHESTRATOR_DATABASE_SCHEMA, }); } catch (error) { await Promise.allSettled([checkpoint.close()]); throw error; } const runtime = createLangGraphOrchestratorRuntime({ checkpointer: checkpoint.checkpointer, checkpointKind: checkpoint.kind, durable: checkpoint.durable, checkpointProbe: checkpoint.probe, executorJobStore: executorJobs.store, executorJobProbe: executorJobs.probe, executionEnabled, enabledExecutors: csvList(env.MEMIND_ORCHESTRATOR_ENABLED_EXECUTORS), admissionConfig: { tenantAllowlist: csvList(env.MEMIND_ORCHESTRATOR_TENANT_ALLOWLIST), userAllowlist: csvList(env.MEMIND_ORCHESTRATOR_USER_ALLOWLIST), workspaceKindAllowlist: csvList( env.MEMIND_ORCHESTRATOR_WORKSPACE_KIND_ALLOWLIST ?? 'workspace-alias', ), maxGlobalActive: env.MEMIND_ORCHESTRATOR_MAX_GLOBAL_ACTIVE, maxTenantActive: env.MEMIND_ORCHESTRATOR_MAX_TENANT_ACTIVE, maxUserActive: env.MEMIND_ORCHESTRATOR_MAX_USER_ACTIVE, maxTenantPerMinute: env.MEMIND_ORCHESTRATOR_MAX_TENANT_PER_MINUTE, maxUserPerMinute: env.MEMIND_ORCHESTRATOR_MAX_USER_PER_MINUTE, }, }); const app = createOrchestratorApp({ runtime, serviceToken: env.MEMIND_ORCHESTRATOR_SERVICE_TOKEN, workerToken: env.MEMIND_ORCHESTRATOR_WORKER_TOKEN, }); const port = positivePort(env.MEMIND_ORCHESTRATOR_PORT, 8093); let retentionSweep = null; let server; try { server = await new Promise((resolve, reject) => { const listening = app.listen(port, host, () => resolve(listening)); listening.once('error', reject); }); } catch (error) { await Promise.allSettled([executorJobs.close(), checkpoint.close()]); throw error; } retentionSweep = startRetentionSweep(runtime, { retentionDays: env.MEMIND_ORCHESTRATOR_RETENTION_DAYS, intervalMs: env.MEMIND_ORCHESTRATOR_RETENTION_SWEEP_INTERVAL_MS, limit: env.MEMIND_ORCHESTRATOR_RETENTION_SWEEP_LIMIT, logger, }); logger.log(`[orchestrator] listening on http://${host}:${port} (${checkpoint.kind})`); async function close() { retentionSweep?.stop(); let serverCloseError = null; try { await new Promise((resolve, reject) => { server.close((error) => (error ? reject(error) : resolve())); }); } catch (error) { serverCloseError = error; } const results = await Promise.allSettled([ executorJobs.close(), checkpoint.close(), ]); if (serverCloseError) throw serverCloseError; const failed = results.find((result) => result.status === 'rejected'); if (failed) throw failed.reason; } return { app, server, runtime, checkpoint, executorJobs, retentionSweep, close, }; } const isEntrypoint = process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href; if (isEntrypoint) { const running = await startOrchestratorServer(); const shutdown = async (signal) => { console.log(`[orchestrator] received ${signal}, shutting down`); await running.close(); process.exit(0); }; process.once('SIGINT', () => void shutdown('SIGINT')); process.once('SIGTERM', () => void shutdown('SIGTERM')); } export const orchestratorServerInternals = { assertSecureServerConfig, csvList, envFlag, isLoopbackHost, positivePort, startRetentionSweep, };