6df82818c5
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.
249 lines
7.7 KiB
JavaScript
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,
|
|
};
|