Files
john 6df82818c5 feat(orchestrator): harden zero-impact shadow rollout
Gate and bound Portal shadow observations while preserving Native execution. Add fail-closed service boundaries, terminal retention controls, Canary readiness telemetry, ops visibility, and isolated regression coverage.
2026-07-25 07:28:37 +08:00

249 lines
7.7 KiB
JavaScript

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,
};