#!/usr/bin/env node const assert = require("assert"); const cp = require("child_process"); const crypto = require("crypto"); const fs = require("fs"); const net = require("net"); const path = require("path"); const { coordinatorWireRequest } = require("./coordinator-wire"); const { nodeIdentity, signedNodeRequest } = require("./node-signing"); const repo = path.resolve(__dirname, ".."); const digest = (value) => `sha256:${crypto.createHash("sha256").update(value).digest("hex")}`; const environmentDigest = digest("scheduler-linux-container"); const dependencyDigest = digest("scheduler-toolchain-dependencies"); const sourceDigest = digest("scheduler-source-tree"); function buildFlagshipBundle() { const output = cp.execFileSync( "cargo", [ "run", "-q", "-p", "clusterflux-cli", "--bin", "clusterflux", "--", "build", "--project", "examples/launch-build-demo", "--json", ], { cwd: repo, encoding: "utf8" } ); const report = JSON.parse(output); const directory = path.resolve(repo, report.bundle_artifact.directory); const manifest = JSON.parse(fs.readFileSync(path.join(directory, "manifest.json"), "utf8")); const entrypoints = JSON.parse( fs.readFileSync(path.join(directory, manifest.entrypoints), "utf8") ); const entrypoint = entrypoints.find((candidate) => candidate.name === "build"); assert(entrypoint, "flagship bundle omitted build entrypoint"); const taskDescriptors = JSON.parse( fs.readFileSync(path.join(directory, manifest.task_descriptors), "utf8") ); const prepareSource = taskDescriptors.find( (candidate) => candidate.name === "prepare_source" ); assert(prepareSource, "flagship bundle omitted prepare_source task"); return { digest: manifest.bundle_digest, taskExport: prepareSource.export, wasmModuleBase64: fs.readFileSync(path.join(directory, "module.wasm")).toString("base64"), }; } function waitForJsonLine(child) { return new Promise((resolve, reject) => { let buffer = ""; child.stdout.on("data", (chunk) => { buffer += chunk.toString(); const newline = buffer.indexOf("\n"); if (newline < 0) return; try { resolve(JSON.parse(buffer.slice(0, newline).trim())); } catch (error) { reject(error); } }); child.once("exit", (code) => { reject(new Error(`process exited before JSON line with code ${code}`)); }); }); } function send(addr, message) { return new Promise((resolve, reject) => { const socket = net.connect(addr.port, addr.host, () => { socket.write(`${JSON.stringify(coordinatorWireRequest(message))}\n`); }); let buffer = ""; socket.on("data", (chunk) => { buffer += chunk.toString(); const newline = buffer.indexOf("\n"); if (newline < 0) return; socket.end(); try { resolve(JSON.parse(buffer.slice(0, newline))); } catch (error) { reject(error); } }); socket.on("error", reject); }); } function linuxCapabilities() { return { os: "Linux", arch: "x86_64", capabilities: [ "Command", "Containers", "RootlessPodman", "SourceFilesystem", "VfsArtifacts" ], environment_backends: ["Container"], source_providers: ["filesystem"] }; } function gitCapabilities() { const capabilities = linuxCapabilities(); capabilities.capabilities = [...capabilities.capabilities, "SourceGit"].sort(); capabilities.source_providers = [...capabilities.source_providers, "git"].sort(); return capabilities; } async function attachNode(addr, node) { const identity = nodeIdentity("scheduler-placement-smoke", node); const attached = await send(addr, { type: "attach_node", tenant: "tenant", project: "project", node, public_key: identity.publicKey }); assert.strictEqual(attached.type, "node_attached"); assert.strictEqual(attached.node, node); return identity; } async function reportNode(addr, node, identity, locality) { const recorded = await send(addr, signedNodeRequest(node, identity, "report_node_capabilities", { type: "report_node_capabilities", tenant: "tenant", project: "project", node, capabilities: locality.capabilities || linuxCapabilities(), cached_environment_digests: locality.cached_environment_digests, dependency_cache_digests: locality.dependency_cache_digests, source_snapshots: locality.source_snapshots, artifact_locations: locality.artifact_locations, direct_connectivity: locality.direct_connectivity !== false, online: true })); assert.strictEqual( recorded.type, "node_capabilities_recorded", JSON.stringify(recorded) ); assert.strictEqual(recorded.node, node); return recorded; } (async () => { const bundle = buildFlagshipBundle(); const coordinator = cp.spawn( "cargo", [ "run", "-q", "-p", "clusterflux-coordinator", "--bin", "clusterflux-coordinator", "--", "--listen", "127.0.0.1:0", "--allow-local-trusted-loopback" ], { cwd: repo } ); let coordinatorStderr = ""; coordinator.stderr.on("data", (chunk) => { coordinatorStderr += chunk.toString(); }); try { const ready = await waitForJsonLine(coordinator); const [host, portText] = ready.listen.split(":"); const addr = { host, port: Number(portText) }; assert.strictEqual((await send(addr, { type: "ping" })).type, "pong"); const coldNode = await attachNode(addr, "cold-node"); const warmNode = await attachNode(addr, "warm-node"); const cold = await reportNode(addr, "cold-node", coldNode, { cached_environment_digests: [], dependency_cache_digests: [], source_snapshots: [], artifact_locations: [], direct_connectivity: false }); assert.strictEqual(cold.node_descriptors, 1); const warm = await reportNode(addr, "warm-node", warmNode, { cached_environment_digests: [environmentDigest], dependency_cache_digests: [dependencyDigest], source_snapshots: [sourceDigest], artifact_locations: ["toolchain-cache"] }); assert.strictEqual(warm.node_descriptors, 2); const reportedNodes = new Set([cold.node, warm.node]); assert.strictEqual(reportedNodes.size, 2); assert(reportedNodes.has("cold-node")); assert(reportedNodes.has("warm-node")); const inspected = await send(addr, { type: "list_node_descriptors", tenant: "tenant", project: "project", actor_user: "operator" }); assert.strictEqual(inspected.type, "node_descriptors"); assert.strictEqual(inspected.actor, "operator"); assert.strictEqual(inspected.descriptors.length, 2); const warmDescriptor = inspected.descriptors.find( (descriptor) => descriptor.id === "warm-node" ); assert(warmDescriptor, "warm node descriptor must be visible to inspector state"); assert(warmDescriptor.capabilities.capabilities.includes("Command")); assert(warmDescriptor.capabilities.capabilities.includes("RootlessPodman")); assert(warmDescriptor.cached_environments.includes(environmentDigest)); assert(warmDescriptor.dependency_caches.includes(dependencyDigest)); assert(warmDescriptor.source_snapshots.includes(sourceDigest)); assert(warmDescriptor.artifact_locations.includes("toolchain-cache")); const crossScopeInspection = await send(addr, { type: "list_node_descriptors", tenant: "other-tenant", project: "project", actor_user: "operator" }); assert.strictEqual(crossScopeInspection.type, "node_descriptors"); assert.strictEqual(crossScopeInspection.descriptors.length, 0); const crossTenantReport = await send(addr, signedNodeRequest("warm-node", warmNode, "report_node_capabilities", { type: "report_node_capabilities", tenant: "other-tenant", project: "project", node: "warm-node", capabilities: linuxCapabilities(), cached_environment_digests: [], dependency_cache_digests: [], source_snapshots: [], artifact_locations: [], direct_connectivity: true, online: true })); assert.strictEqual(crossTenantReport.type, "error"); assert.match(crossTenantReport.message, /tenant\/project scope/); const placement = await send(addr, { type: "schedule_task", tenant: "tenant", project: "project", environment: { os: "Linux", arch: null, capabilities: ["Containers", "RootlessPodman"] }, environment_digest: environmentDigest, required_capabilities: ["Command"], dependency_cache: dependencyDigest, source_snapshot: sourceDigest, required_artifacts: ["toolchain-cache"], prefer_node: null }); assert.strictEqual(placement.type, "task_placement"); assert.strictEqual(placement.placement.node, "warm-node"); assert.ok(placement.placement.score > 0); assert.ok(placement.placement.reasons.includes("warm environment cache")); assert.ok(placement.placement.reasons.includes("warm dependency cache")); assert.ok(placement.placement.reasons.includes("source snapshot already local")); assert.ok( placement.placement.reasons.includes("1 required artifact(s) already local") ); const impossible = await send(addr, { type: "schedule_task", tenant: "tenant", project: "project", environment: null, environment_digest: null, required_capabilities: ["WindowsCommandDev"], dependency_cache: null, source_snapshot: null, required_artifacts: [], prefer_node: null }); assert.strictEqual(impossible.type, "error"); assert.match(impossible.message, /WindowsCommandDev/); const started = await send(addr, { type: "start_process", tenant: "tenant", project: "project", actor_user: "operator", process: "vp-wait-for-git" }); assert.strictEqual(started.type, "process_started"); const queued = await send(addr, { type: "launch_task", tenant: "tenant", project: "project", actor_user: "operator", task_spec: { tenant: "tenant", project: "project", process: "vp-wait-for-git", task_definition: "prepare_source", task_instance: "prepare_source-1", dispatch: { kind: "coordinator_node_wasm", export: bundle.taskExport, abi: "task_v1", }, environment_id: null, environment: null, environment_digest: null, required_capabilities: ["SourceFilesystem", "SourceGit"], dependency_cache: null, source_snapshot: null, required_artifacts: [], args: [], vfs_epoch: started.epoch, bundle_digest: bundle.digest, }, wait_for_node: true, artifact_path: "/vfs/artifacts/git-status.txt", wasm_module_base64: bundle.wasmModuleBase64, }); assert.strictEqual(queued.type, "task_queued"); assert.strictEqual(queued.process, "vp-wait-for-git"); assert.strictEqual(queued.task, "prepare_source-1"); assert.match(queued.reason, /SourceGit/); assert.strictEqual(queued.queued_tasks, 1); const gitNode = await attachNode(addr, "git-node"); const gitRecorded = await reportNode(addr, "git-node", gitNode, { capabilities: gitCapabilities(), cached_environment_digests: [], dependency_cache_digests: [], source_snapshots: [], artifact_locations: [], direct_connectivity: false }); assert.strictEqual(gitRecorded.type, "node_capabilities_recorded"); const pendingAssignment = await send(addr, signedNodeRequest("git-node", gitNode, "poll_task_assignment", { type: "poll_task_assignment", tenant: "tenant", project: "project", node: "git-node" })); assert.strictEqual(pendingAssignment.type, "task_assignment"); assert(pendingAssignment.assignment, "late capable node should receive queued assignment"); assert.strictEqual(pendingAssignment.assignment.process, "vp-wait-for-git"); assert.strictEqual(pendingAssignment.assignment.task, "prepare_source-1"); assert.strictEqual(pendingAssignment.assignment.node, "git-node"); assert.strictEqual( pendingAssignment.assignment.task_spec.dispatch.export, bundle.taskExport ); const emptyAssignment = await send(addr, signedNodeRequest("git-node", gitNode, "poll_task_assignment", { type: "poll_task_assignment", tenant: "tenant", project: "project", node: "git-node" })); assert.strictEqual(emptyAssignment.type, "task_assignment"); assert.strictEqual(emptyAssignment.assignment, null); await reportNode(addr, "warm-node", warmNode, { cached_environment_digests: [], dependency_cache_digests: [], source_snapshots: [], artifact_locations: [], direct_connectivity: false }); const disconnectedTransfer = await send(addr, { type: "schedule_task", tenant: "tenant", project: "project", environment: null, environment_digest: null, required_capabilities: ["Command"], dependency_cache: null, source_snapshot: sourceDigest, required_artifacts: ["toolchain-cache"], prefer_node: null }); assert.strictEqual(disconnectedTransfer.type, "error"); assert.match(disconnectedTransfer.message, /source snapshot unavailable/); assert.match(disconnectedTransfer.message, /required artifact\(s\) unavailable/); assert.match(disconnectedTransfer.message, /direct connectivity unavailable/); } catch (error) { if (coordinatorStderr) { error.message = `${error.message}\ncoordinator stderr:\n${coordinatorStderr}`; } throw error; } finally { coordinator.kill("SIGTERM"); } console.log("Scheduler placement smoke passed"); })().catch((error) => { console.error(error.stack || error.message); process.exit(1); });