import * as plugins from '../plugins.js'; import type { TSha256Digest, TImmutableContainerPlatform } from './immutableimage.js'; import { canonicalizeStrictJson, createCanonicalJsonSha256Hex, } from '../private/canonicaljson.js'; export interface IClusterNodeMetrics { cpuUsagePercent: number; memoryUsedMB: number; memoryAvailableMB: number; diskUsedGB: number; diskAvailableGB: number; containerCount: number; timestamp: number; } /** * Lightweight high-frequency node metrics sample (1s cadence). * Streamed spark -> cloudly -> admin UIs; never persisted. */ export interface INodeMetricsSample { timestamp: number; cpuUsagePercent: number; memoryUsedMB: number; memoryAvailableMB: number; /** Host-level network receive rate across non-loopback interfaces. */ networkRxBytesPerSec?: number; /** Host-level network transmit rate across non-loopback interfaces. */ networkTxBytesPerSec?: number; } export type TNodeActionName = 'systemUpgrade' | 'reboot' | 'serveZoneServiceUpdate'; export type TNodeActionStatus = 'pending' | 'running' | 'succeeded' | 'failed'; /** * An operator-requested action executed by spark on a node. * Delivered to spark in metrics-sample responses; results posted back. */ export interface INodeAction { id: string; nodeId: string; actionName: TNodeActionName; payload?: { /** Reboot only: gracefully stop workloads before rebooting. */ drainFirst?: boolean; }; status: TNodeActionStatus; requestedAt: number; requestedBy?: string; startedAt?: number; finishedAt?: number; resultText?: string; } export interface IServeZoneServiceRuntimeInfo { name: string; serviceId?: string; image?: string; desiredImage?: string; imageVersion?: string; runningImageId?: string; runningTaskCount?: number; running: boolean; checkedAt: number; updatedAt?: number; error?: string; } /** Spark -> Cloudly HTTP body for POST /spark/v1/nodes/metrics-sample */ export interface ISparkMetricsSampleRequest { nodeId: string; nodeToken: string; sample: INodeMetricsSample; } export interface ISparkMetricsSampleResponse { accepted: boolean; message?: string; /** Actions queued for this node; spark executes them and posts results. */ pendingActions?: INodeAction[]; } /** Spark -> Cloudly HTTP body for POST /spark/v1/nodes/action-result */ export interface ISparkActionResultRequest { nodeId: string; nodeToken: string; actionId: string; status: Extract; resultText?: string; } export interface ISparkActionResultResponse { accepted: boolean; message?: string; } export interface ISwarmManagerObservedNodeV1 { swarmNodeId: string; nodeName: string; platform: TImmutableContainerPlatform; role: 'manager' | 'worker'; availability: 'active' | 'pause' | 'drain'; state: 'unknown' | 'down' | 'ready' | 'disconnected'; managerStatus?: { leader: boolean; reachability: 'unknown' | 'unreachable' | 'reachable'; }; } export interface ISwarmManagerSnapshotV1 { snapshotVersion: 1; snapshotDigest: TSha256Digest; nodeCount: number; managerCount: number; nodes: ISwarmManagerObservedNodeV1[]; } interface ISparkSwarmObservationBaseV1 { schemaVersion: 1; reporterSessionId: string; observationSequence: number; observedAt: number; } export type TSparkSwarmObservationV1 = | (ISparkSwarmObservationBaseV1 & { state: 'unknown'; reason: 'docker-unavailable' | 'query-failed' | 'not-initialized'; }) | (ISparkSwarmObservationBaseV1 & { state: 'not-member'; }) | (ISparkSwarmObservationBaseV1 & { state: 'member'; swarmClusterId: string; localSwarmNodeId: string; controlAvailable: boolean; managerSnapshot?: ISwarmManagerSnapshotV1; }); interface ISparkSwarmObservationBaseV2 { schemaVersion: 2; reporterSessionId: string; observationSequence: number; observedAt: number; } export type TSparkSwarmObservationV2 = | (ISparkSwarmObservationBaseV2 & { state: 'unknown'; reason: 'docker-unavailable' | 'query-failed' | 'not-initialized'; }) | (ISparkSwarmObservationBaseV2 & { state: 'not-member'; }) | (ISparkSwarmObservationBaseV2 & { state: 'member'; localSwarmNodeId: string; controlAvailable: false; }) | (ISparkSwarmObservationBaseV2 & { state: 'member'; swarmClusterId: string; localSwarmNodeId: string; controlAvailable: true; managerSnapshot?: ISwarmManagerSnapshotV1; }); export type TSparkSwarmObservationV2RejectionCode = | 'sequence-conflict' | 'sequence-stale' | 'sequence-gap' | 'session-sequence-invalid' | 'session-retired' | 'observation-too-old' | 'observation-from-future' | 'observation-time-not-advancing' | 'acceptance-race'; export const sparkSwarmObservationV2Contract = Object.freeze({ endpoint: '/spark/v2/nodes/swarm-observation', digestDomain: 'serve.zone/spark-swarm-observation-v2', maximumPastAgeMs: 300_000, maximumFutureSkewMs: 30_000, retainedRetiredSessionIds: 32, rejectionRetryability: Object.freeze({ 'sequence-conflict': false, 'sequence-stale': false, 'sequence-gap': false, 'session-sequence-invalid': false, 'session-retired': false, 'observation-too-old': false, 'observation-from-future': false, 'observation-time-not-advancing': false, 'acceptance-race': true, } satisfies Record), } as const); /** Cloudly derives node and cluster scope from this authenticated Spark request. */ export interface ISparkSwarmObservationV2Request { nodeId: string; nodeToken: string; observation: TSparkSwarmObservationV2; } export interface ISparkSwarmObservationV2AcceptedReceipt { schemaVersion: 2; state: 'accepted'; nodeId: string; reporterSessionId: string; observationSequence: number; observationDigest: TSha256Digest; acceptedAt: number; /** Runtime target generation observed immediately after this acceptance. */ targetGenerationAtAcceptance: number; } export interface ISparkSwarmObservationV2RejectedReceipt { schemaVersion: 2; state: 'rejected'; nodeId: string; reporterSessionId: string; observationSequence: number; observationDigest: TSha256Digest; rejectedAt: number; code: TSparkSwarmObservationV2RejectionCode; retryable: boolean; } /** * Bound receipts are returned only after request validation and node * authentication. Malformed and authentication failures use generic HTTP * errors so they cannot expose credential or persisted replay state. */ export type TSparkSwarmObservationV2Response = | ISparkSwarmObservationV2AcceptedReceipt | ISparkSwarmObservationV2RejectedReceipt; /** Cloudly derives node and cluster scope from this authenticated Spark request. */ export interface ISparkSwarmObservationRequest { nodeId: string; nodeToken: string; observation: TSparkSwarmObservationV1; } export interface ISparkSwarmObservationResponse { accepted: boolean; acceptedSequence?: number; targetGeneration?: number; message?: string; } export interface IClusterRuntimeTargetV1 { cloudlyNodeId: string; swarmClusterId: string; swarmNodeId: string; nodeName: string; platform: TImmutableContainerPlatform; } export interface IClusterRuntimeTargetSetReadyV1 { schemaVersion: 1; state: 'ready'; cloudlyClusterId: string; swarmClusterId: string; generation: number; targetSetDigest: TSha256Digest; acceptedAt: number; freshUntil: number; targets: IClusterRuntimeTargetV1[]; } export interface IClusterRuntimeTargetSetUnavailableV1 { schemaVersion: 1; state: 'unavailable'; cloudlyClusterId: string; generation: number; changedAt: number; reason: | 'never-observed' | 'stale' | 'no-manager-consensus' | 'identity-conflict' | 'no-schedulable-targets'; } export type TClusterRuntimeTargetSetV1 = | IClusterRuntimeTargetSetReadyV1 | IClusterRuntimeTargetSetUnavailableV1; const runtimeIdentifierRegex = /^[A-Za-z0-9][A-Za-z0-9:._-]{0,199}$/; const runtimeDigestRegex = /^sha256:[a-f0-9]{64}$/; const runtimePlatforms = new Set([ 'linux/amd64', 'linux/arm64', ]); const maximumSwarmNodes = 1024; const isRuntimeRecord = (valueArg: unknown): valueArg is Record => ( Boolean(valueArg) && typeof valueArg === 'object' && !Array.isArray(valueArg) ); const hasExactRuntimeKeys = ( valueArg: Record, keysArg: string[], ): boolean => JSON.stringify(Object.keys(valueArg).sort()) === JSON.stringify([...keysArg].sort()); const isRuntimeIdentifier = (valueArg: unknown): valueArg is string => ( typeof valueArg === 'string' && runtimeIdentifierRegex.test(valueArg) ); const isPositiveSafeInteger = (valueArg: unknown): valueArg is number => ( Number.isSafeInteger(valueArg) && (valueArg as number) > 0 ); const computeSha256 = async (inputArg: string): Promise => { const digest = new Uint8Array(await globalThis.crypto.subtle.digest( 'SHA-256', new TextEncoder().encode(inputArg), )); return `sha256:${[...digest] .map((byteArg) => byteArg.toString(16).padStart(2, '0')) .join('')}` as TSha256Digest; }; const failSparkSwarmObservationV2Canonicalization = (reasonArg: string): never => { throw new Error(`Spark Swarm observation v2 cannot be canonicalized: ${reasonArg}`); }; export const createSparkSwarmObservationV2DigestInput = ( observationArg: TSparkSwarmObservationV2, ): string => canonicalizeStrictJson( { domain: sparkSwarmObservationV2Contract.digestDomain, observation: observationArg, }, failSparkSwarmObservationV2Canonicalization, 'Spark Swarm observation v2 digest input', ); export const computeSparkSwarmObservationV2Digest = async ( observationArg: TSparkSwarmObservationV2, ): Promise => ( `sha256:${await createCanonicalJsonSha256Hex( createSparkSwarmObservationV2DigestInput(observationArg), failSparkSwarmObservationV2Canonicalization, )}` as TSha256Digest ); const canonicalSwarmNode = ( nodeArg: ISwarmManagerObservedNodeV1, ): Record => ({ swarmNodeId: nodeArg.swarmNodeId, nodeName: nodeArg.nodeName, platform: nodeArg.platform, role: nodeArg.role, availability: nodeArg.availability, state: nodeArg.state, managerStatus: nodeArg.managerStatus ? { leader: nodeArg.managerStatus.leader, reachability: nodeArg.managerStatus.reachability, } : null, }); export const createSwarmManagerSnapshotDigestInput = ( swarmClusterIdArg: string, snapshotArg: ISwarmManagerSnapshotV1, ): string => JSON.stringify({ schemaVersion: 1, purpose: 'serve.zone/swarm-manager-snapshot', swarmClusterId: swarmClusterIdArg, snapshotVersion: snapshotArg.snapshotVersion, nodeCount: snapshotArg.nodeCount, managerCount: snapshotArg.managerCount, nodes: snapshotArg.nodes.map(canonicalSwarmNode), }); export const computeSwarmManagerSnapshotDigest = async ( swarmClusterIdArg: string, snapshotArg: ISwarmManagerSnapshotV1, ): Promise => computeSha256( createSwarmManagerSnapshotDigestInput(swarmClusterIdArg, snapshotArg), ); const validateSwarmManagerSnapshot = async ( swarmClusterIdArg: string, localSwarmNodeIdArg: string, snapshotArg: unknown, ): Promise => { if (!isRuntimeRecord(snapshotArg) || !hasExactRuntimeKeys(snapshotArg, [ 'snapshotVersion', 'snapshotDigest', 'nodeCount', 'managerCount', 'nodes', ])) { return ['Spark Swarm manager snapshot must use its exact schema']; } const errors: string[] = []; if (snapshotArg.snapshotVersion !== 1) { errors.push('Spark Swarm manager snapshotVersion must be 1'); } if (typeof snapshotArg.snapshotDigest !== 'string' || !runtimeDigestRegex.test(snapshotArg.snapshotDigest)) { errors.push('Spark Swarm manager snapshot digest must be canonical'); } if (!Array.isArray(snapshotArg.nodes) || snapshotArg.nodes.length === 0 || snapshotArg.nodes.length > maximumSwarmNodes) { errors.push('Spark Swarm manager snapshot nodes must be a non-empty bounded array'); return errors; } if (snapshotArg.nodeCount !== snapshotArg.nodes.length || !isPositiveSafeInteger(snapshotArg.nodeCount)) { errors.push('Spark Swarm manager snapshot nodeCount must match nodes'); } const nodeIds = new Set(); const nodeNames = new Set(); let managerCount = 0; let leaderCount = 0; let previousSortKey: string | undefined; for (const [index, nodeArg] of snapshotArg.nodes.entries()) { if (!isRuntimeRecord(nodeArg)) { errors.push(`Spark Swarm manager snapshot nodes[${index}] must be an object`); continue; } const isManager = nodeArg.role === 'manager'; const expectedKeys = [ 'swarmNodeId', 'nodeName', 'platform', 'role', 'availability', 'state', ...(isManager ? ['managerStatus'] : []), ]; if (!hasExactRuntimeKeys(nodeArg, expectedKeys)) { errors.push(`Spark Swarm manager snapshot nodes[${index}] must use its exact role schema`); } if (!isRuntimeIdentifier(nodeArg.swarmNodeId) || !isRuntimeIdentifier(nodeArg.nodeName)) { errors.push(`Spark Swarm manager snapshot nodes[${index}] identity must be canonical`); } if (!runtimePlatforms.has(nodeArg.platform as TImmutableContainerPlatform)) { errors.push(`Spark Swarm manager snapshot nodes[${index}] platform must be supported`); } if (!['manager', 'worker'].includes(nodeArg.role as string) || !['active', 'pause', 'drain'].includes(nodeArg.availability as string) || !['unknown', 'down', 'ready', 'disconnected'].includes(nodeArg.state as string)) { errors.push(`Spark Swarm manager snapshot nodes[${index}] scheduler state must be canonical`); } if (isManager) { managerCount++; if (!isRuntimeRecord(nodeArg.managerStatus) || !hasExactRuntimeKeys(nodeArg.managerStatus, ['leader', 'reachability']) || typeof nodeArg.managerStatus.leader !== 'boolean' || !['unknown', 'unreachable', 'reachable'] .includes(nodeArg.managerStatus.reachability as string)) { errors.push(`Spark Swarm manager snapshot nodes[${index}] managerStatus must be canonical`); } else if (nodeArg.managerStatus.leader) { leaderCount++; } } if (isRuntimeIdentifier(nodeArg.swarmNodeId)) { if (nodeIds.has(nodeArg.swarmNodeId)) { errors.push('Spark Swarm manager snapshot node IDs must be unique'); } nodeIds.add(nodeArg.swarmNodeId); } if (isRuntimeIdentifier(nodeArg.nodeName)) { if (nodeNames.has(nodeArg.nodeName)) { errors.push('Spark Swarm manager snapshot node names must be unique'); } nodeNames.add(nodeArg.nodeName); } if (isRuntimeIdentifier(nodeArg.swarmNodeId) && isRuntimeIdentifier(nodeArg.nodeName)) { const sortKey = `${nodeArg.swarmNodeId}\0${nodeArg.nodeName}\0${nodeArg.platform}`; if (previousSortKey !== undefined && sortKey <= previousSortKey) { errors.push('Spark Swarm manager snapshot nodes must be uniquely sorted'); } previousSortKey = sortKey; } } if (snapshotArg.managerCount !== managerCount || !isPositiveSafeInteger(snapshotArg.managerCount)) { errors.push('Spark Swarm manager snapshot managerCount must match nodes'); } if (leaderCount > 1) errors.push('Spark Swarm manager snapshot may contain at most one leader'); const sourceNode = snapshotArg.nodes.find((nodeArg) => ( isRuntimeRecord(nodeArg) && nodeArg.swarmNodeId === localSwarmNodeIdArg )); if (!isRuntimeRecord(sourceNode) || sourceNode.role !== 'manager' || sourceNode.availability !== 'active' || sourceNode.state !== 'ready' || !isRuntimeRecord(sourceNode.managerStatus) || sourceNode.managerStatus.reachability !== 'reachable') { errors.push('Spark Swarm manager snapshot source must be an active reachable manager'); } if (errors.length === 0 && snapshotArg.snapshotDigest !== await computeSwarmManagerSnapshotDigest( swarmClusterIdArg, snapshotArg as unknown as ISwarmManagerSnapshotV1, )) { errors.push('Spark Swarm manager snapshot digest does not match'); } return errors; }; export const validateSparkSwarmObservation = async ( observationArg: unknown, ): Promise => { try { if (!isRuntimeRecord(observationArg)) return ['Spark Swarm observation must be an object']; const observation = observationArg as Record; const commonKeys = [ 'schemaVersion', 'reporterSessionId', 'observationSequence', 'observedAt', 'state', ]; const expectedKeys = observation.state === 'unknown' ? [...commonKeys, 'reason'] : observation.state === 'not-member' ? commonKeys : observation.state === 'member' ? [ ...commonKeys, 'swarmClusterId', 'localSwarmNodeId', 'controlAvailable', ...(observation.managerSnapshot === undefined ? [] : ['managerSnapshot']), ] : []; if (expectedKeys.length === 0 || !hasExactRuntimeKeys(observation, expectedKeys)) { return ['Spark Swarm observation must use its exact state schema']; } const errors: string[] = []; if (observation.schemaVersion !== 1) errors.push('Spark Swarm observation schemaVersion must be 1'); if (!isRuntimeIdentifier(observation.reporterSessionId) || !isPositiveSafeInteger(observation.observationSequence) || !isPositiveSafeInteger(observation.observedAt)) { errors.push('Spark Swarm observation identity, sequence, and timestamp must be canonical'); } if (observation.state === 'unknown' && !['docker-unavailable', 'query-failed', 'not-initialized'] .includes(observation.reason as string)) { errors.push('Spark Swarm observation unknown reason must be canonical'); } if (observation.state === 'member') { if (!isRuntimeIdentifier(observation.swarmClusterId) || !isRuntimeIdentifier(observation.localSwarmNodeId) || typeof observation.controlAvailable !== 'boolean') { errors.push('Spark Swarm member observation must have canonical local authority'); } if (observation.managerSnapshot !== undefined) { if (observation.controlAvailable !== true) { errors.push('Spark Swarm manager snapshot requires controlAvailable'); } else if (isRuntimeIdentifier(observation.swarmClusterId) && isRuntimeIdentifier(observation.localSwarmNodeId)) { errors.push(...await validateSwarmManagerSnapshot( observation.swarmClusterId, observation.localSwarmNodeId, observation.managerSnapshot, )); } } } return errors; } catch { return ['Spark Swarm observation must be safely inspectable']; } }; export const validateSparkSwarmObservationV2 = async ( observationArg: unknown, ): Promise => { try { if (!isRuntimeRecord(observationArg)) return ['Spark Swarm observation v2 must be an object']; const observation = observationArg as Record; const commonKeys = [ 'schemaVersion', 'reporterSessionId', 'observationSequence', 'observedAt', 'state', ]; const expectedKeys = observation.state === 'unknown' ? [...commonKeys, 'reason'] : observation.state === 'not-member' ? commonKeys : observation.state === 'member' && observation.controlAvailable === false ? [...commonKeys, 'localSwarmNodeId', 'controlAvailable'] : observation.state === 'member' && observation.controlAvailable === true ? [ ...commonKeys, 'swarmClusterId', 'localSwarmNodeId', 'controlAvailable', ...(observation.managerSnapshot === undefined ? [] : ['managerSnapshot']), ] : []; if (expectedKeys.length === 0 || !hasExactRuntimeKeys(observation, expectedKeys)) { return ['Spark Swarm observation v2 must use its exact state schema']; } const errors: string[] = []; if (observation.schemaVersion !== 2) { errors.push('Spark Swarm observation v2 schemaVersion must be 2'); } if (!isRuntimeIdentifier(observation.reporterSessionId) || !isPositiveSafeInteger(observation.observationSequence) || !isPositiveSafeInteger(observation.observedAt)) { errors.push('Spark Swarm observation v2 identity, sequence, and timestamp must be canonical'); } if (observation.state === 'unknown' && !['docker-unavailable', 'query-failed', 'not-initialized'] .includes(observation.reason as string)) { errors.push('Spark Swarm observation v2 unknown reason must be canonical'); } if (observation.state === 'member') { if (!isRuntimeIdentifier(observation.localSwarmNodeId) || typeof observation.controlAvailable !== 'boolean') { errors.push('Spark Swarm member observation v2 must have canonical local authority'); } if (observation.controlAvailable === true) { if (!isRuntimeIdentifier(observation.swarmClusterId)) { errors.push('Spark Swarm manager observation v2 must have canonical cluster authority'); } else if (observation.managerSnapshot !== undefined && isRuntimeIdentifier(observation.localSwarmNodeId)) { errors.push(...await validateSwarmManagerSnapshot( observation.swarmClusterId, observation.localSwarmNodeId, observation.managerSnapshot, )); } } } return errors; } catch { return ['Spark Swarm observation v2 must be safely inspectable']; } }; export const validateSparkSwarmObservationV2Request = async ( requestArg: unknown, ): Promise => { try { if (!isRuntimeRecord(requestArg) || !hasExactRuntimeKeys(requestArg, ['nodeId', 'nodeToken', 'observation'])) { return ['Spark Swarm observation v2 request must use its exact schema']; } try { canonicalizeStrictJson( requestArg, failSparkSwarmObservationV2Canonicalization, 'Spark Swarm observation v2 request', ); } catch { return ['Spark Swarm observation v2 request must contain strict canonical JSON data']; } const errors = await validateSparkSwarmObservationV2(requestArg.observation); if (!isRuntimeIdentifier(requestArg.nodeId)) { errors.push('Spark Swarm observation v2 request nodeId must be canonical'); } if (typeof requestArg.nodeToken !== 'string' || requestArg.nodeToken.length < 1 || requestArg.nodeToken.length > 512 || requestArg.nodeToken.trim() !== requestArg.nodeToken) { errors.push('Spark Swarm observation v2 request nodeToken must be a bounded opaque token'); } return errors; } catch { return ['Spark Swarm observation v2 request must be safely inspectable']; } }; export const validateSparkSwarmObservationV2Response = async ( responseArg: unknown, expectedRequestArg: ISparkSwarmObservationV2Request, ): Promise => { try { const requestErrors = await validateSparkSwarmObservationV2Request( expectedRequestArg, ); if (requestErrors.length > 0) { return ['expected Spark Swarm observation v2 request must be valid']; } if (!isRuntimeRecord(responseArg)) { return ['Spark Swarm observation v2 response must be an object']; } const acceptedKeys = [ 'schemaVersion', 'state', 'nodeId', 'reporterSessionId', 'observationSequence', 'observationDigest', 'acceptedAt', 'targetGenerationAtAcceptance', ]; const rejectedKeys = [ 'schemaVersion', 'state', 'nodeId', 'reporterSessionId', 'observationSequence', 'observationDigest', 'rejectedAt', 'code', 'retryable', ]; const expectedKeys = responseArg.state === 'accepted' ? acceptedKeys : responseArg.state === 'rejected' ? rejectedKeys : []; if (expectedKeys.length === 0 || !hasExactRuntimeKeys(responseArg, expectedKeys)) { return ['Spark Swarm observation v2 response must use its exact state schema']; } try { canonicalizeStrictJson( responseArg, failSparkSwarmObservationV2Canonicalization, 'Spark Swarm observation v2 response', ); } catch { return ['Spark Swarm observation v2 response must contain strict canonical JSON data']; } const errors: string[] = []; const observationDigest = await computeSparkSwarmObservationV2Digest( expectedRequestArg.observation, ); if (responseArg.schemaVersion !== 2 || responseArg.nodeId !== expectedRequestArg.nodeId || responseArg.reporterSessionId !== expectedRequestArg.observation.reporterSessionId || responseArg.observationSequence !== expectedRequestArg.observation.observationSequence || responseArg.observationDigest !== observationDigest) { errors.push('Spark Swarm observation v2 response must bind the expected request'); } if (responseArg.state === 'accepted') { if (!isPositiveSafeInteger(responseArg.acceptedAt) || !Number.isSafeInteger(responseArg.targetGenerationAtAcceptance) || (responseArg.targetGenerationAtAcceptance as number) < 0) { errors.push('Spark Swarm observation v2 accepted receipt fields must be canonical'); } } else { const retryability = sparkSwarmObservationV2Contract.rejectionRetryability; if (!isPositiveSafeInteger(responseArg.rejectedAt) || typeof responseArg.code !== 'string' || !Object.hasOwn(retryability, responseArg.code) || responseArg.retryable !== retryability[ responseArg.code as TSparkSwarmObservationV2RejectionCode ]) { errors.push('Spark Swarm observation v2 rejected receipt fields must be canonical'); } } return errors; } catch { return ['Spark Swarm observation v2 response must be safely inspectable']; } }; const canonicalRuntimeTarget = ( targetArg: IClusterRuntimeTargetV1, ): IClusterRuntimeTargetV1 => ({ cloudlyNodeId: targetArg.cloudlyNodeId, swarmClusterId: targetArg.swarmClusterId, swarmNodeId: targetArg.swarmNodeId, nodeName: targetArg.nodeName, platform: targetArg.platform, }); export const createClusterRuntimeTargetSetDigestInput = ( targetSetArg: IClusterRuntimeTargetSetReadyV1, ): string => JSON.stringify({ schemaVersion: targetSetArg.schemaVersion, purpose: 'serve.zone/cluster-runtime-target-set', cloudlyClusterId: targetSetArg.cloudlyClusterId, swarmClusterId: targetSetArg.swarmClusterId, generation: targetSetArg.generation, targets: targetSetArg.targets.map(canonicalRuntimeTarget), }); export const computeClusterRuntimeTargetSetDigest = async ( targetSetArg: IClusterRuntimeTargetSetReadyV1, ): Promise => computeSha256(createClusterRuntimeTargetSetDigestInput(targetSetArg)); export const validateClusterRuntimeTargetSet = async ( targetSetArg: unknown, ): Promise => { try { if (!isRuntimeRecord(targetSetArg)) return ['cluster runtime target set must be an object']; if (targetSetArg.state === 'unavailable') { if (!hasExactRuntimeKeys(targetSetArg, [ 'schemaVersion', 'state', 'cloudlyClusterId', 'generation', 'changedAt', 'reason', ])) { return ['unavailable cluster runtime target set must use its exact schema']; } const errors: string[] = []; if (targetSetArg.schemaVersion !== 1 || !isRuntimeIdentifier(targetSetArg.cloudlyClusterId) || !Number.isSafeInteger(targetSetArg.generation) || (targetSetArg.generation as number) < 0 || !isPositiveSafeInteger(targetSetArg.changedAt) || ![ 'never-observed', 'stale', 'no-manager-consensus', 'identity-conflict', 'no-schedulable-targets', ].includes(targetSetArg.reason as string)) { errors.push('unavailable cluster runtime target set fields must be canonical'); } if (targetSetArg.reason === 'never-observed' && targetSetArg.generation !== 0) { errors.push('never-observed target state must use generation zero'); } return errors; } if (targetSetArg.state !== 'ready' || !hasExactRuntimeKeys(targetSetArg, [ 'schemaVersion', 'state', 'cloudlyClusterId', 'swarmClusterId', 'generation', 'targetSetDigest', 'acceptedAt', 'freshUntil', 'targets', ])) { return ['ready cluster runtime target set must use its exact schema']; } const errors: string[] = []; if (targetSetArg.schemaVersion !== 1 || !isRuntimeIdentifier(targetSetArg.cloudlyClusterId) || !isRuntimeIdentifier(targetSetArg.swarmClusterId) || !isPositiveSafeInteger(targetSetArg.generation) || !isPositiveSafeInteger(targetSetArg.acceptedAt) || !isPositiveSafeInteger(targetSetArg.freshUntil) || (targetSetArg.freshUntil as number) <= (targetSetArg.acceptedAt as number) || typeof targetSetArg.targetSetDigest !== 'string' || !runtimeDigestRegex.test(targetSetArg.targetSetDigest)) { errors.push('ready cluster runtime target set authority fields must be canonical'); } if (!Array.isArray(targetSetArg.targets) || targetSetArg.targets.length === 0 || targetSetArg.targets.length > maximumSwarmNodes) { errors.push('ready cluster runtime target set targets must be a non-empty bounded array'); return errors; } const cloudlyNodeIds = new Set(); const swarmNodeScopes = new Set(); const nodeNames = new Set(); let previousSortKey: string | undefined; for (const [index, targetArg] of targetSetArg.targets.entries()) { if (!isRuntimeRecord(targetArg) || !hasExactRuntimeKeys(targetArg, [ 'cloudlyNodeId', 'swarmClusterId', 'swarmNodeId', 'nodeName', 'platform', ]) || !isRuntimeIdentifier(targetArg.cloudlyNodeId) || !isRuntimeIdentifier(targetArg.swarmClusterId) || !isRuntimeIdentifier(targetArg.swarmNodeId) || !isRuntimeIdentifier(targetArg.nodeName) || targetArg.swarmClusterId !== targetSetArg.swarmClusterId || !runtimePlatforms.has(targetArg.platform as TImmutableContainerPlatform)) { errors.push(`cluster runtime target set targets[${index}] must be canonical`); continue; } const sortKey = `${targetArg.cloudlyNodeId}\0${targetArg.swarmClusterId}\0${ targetArg.swarmNodeId }\0${targetArg.nodeName}\0${targetArg.platform}`; if (previousSortKey !== undefined && sortKey <= previousSortKey) { errors.push('cluster runtime target set targets must be uniquely sorted'); } previousSortKey = sortKey; const swarmScope = `${targetArg.swarmClusterId}\0${targetArg.swarmNodeId}`; if (cloudlyNodeIds.has(targetArg.cloudlyNodeId) || swarmNodeScopes.has(swarmScope) || nodeNames.has(targetArg.nodeName)) { errors.push('cluster runtime target identities must be unique'); } cloudlyNodeIds.add(targetArg.cloudlyNodeId); swarmNodeScopes.add(swarmScope); nodeNames.add(targetArg.nodeName); } if (errors.length === 0 && targetSetArg.targetSetDigest !== await computeClusterRuntimeTargetSetDigest( targetSetArg as unknown as IClusterRuntimeTargetSetReadyV1, )) { errors.push('cluster runtime target set digest does not match'); } return errors; } catch { return ['cluster runtime target set must be safely inspectable']; } }; export const isClusterRuntimeTargetSetFresh = ( targetSetArg: TClusterRuntimeTargetSetV1, nowArg: number, ): targetSetArg is IClusterRuntimeTargetSetReadyV1 => ( targetSetArg.state === 'ready' && Number.isSafeInteger(nowArg) && nowArg >= targetSetArg.acceptedAt && nowArg < targetSetArg.freshUntil ); export type TSparkNodeMode = 'cloudly' | 'coreflow-node'; export type TSparkCloudlyConnectionStatus = | 'not-configured' | 'connecting' | 'connected' | 'failed'; export interface ISparkNodeRuntimeInfo { runtime: 'spark'; nodeId: string; mode?: TSparkNodeMode; hostname?: string; platform: string; arch: string; osRelease?: string; sparkVersion: string; cloudlyUrl?: string; cloudlyConnectionStatus: TSparkCloudlyConnectionStatus; dockerAvailable: boolean; swarmNodeId?: string; serveZoneServices?: IServeZoneServiceRuntimeInfo[]; checkedAt: number; lastError?: string; } export interface IClusterNode { id: string; data: { /** * Reference to the cluster this node belongs to */ clusterId: string; /** * Reference to the physical server (if applicable) */ baremetalId?: string; /** * Type of node */ nodeType: 'baremetal' | 'vm' | 'container'; /** * Current status of the node */ status: 'initializing' | 'online' | 'offline' | 'maintenance'; /** * Role of the node in the cluster */ role: 'master' | 'worker'; /** * Timestamp when node joined the cluster */ joinedAt: number; /** * Last health check timestamp */ lastHealthCheck: number; /** * Current metrics for the node */ metrics?: IClusterNodeMetrics; /** * Runtime status reported by the Spark node agent. */ sparkRuntimeInfo?: ISparkNodeRuntimeInfo; /** * Docker swarm node ID if part of swarm */ swarmNodeId?: string; /** * SSH keys deployed to this node */ sshKeys: plugins.tsclass.network.ISshKey[]; /** * Debian packages installed on this node */ requiredDebianPackages: string[]; }; }