diff --git a/scripts/src/cli.ts b/scripts/src/cli.ts index 36884bf..ac19c7e 100644 --- a/scripts/src/cli.ts +++ b/scripts/src/cli.ts @@ -4,7 +4,6 @@ import { execFileSync, spawn } from "node:child_process"; import { closeSync, existsSync, openSync } from "node:fs"; import path from "node:path"; import { startManagerServer } from "../../src/mgr/server.js"; -import { MemoryAgentRunStore } from "../../src/mgr/store.js"; import { ManagerClient } from "../../src/mgr/client.js"; import { runOnce } from "../../src/runner/run-once.js"; import { renderRunnerJobDryRun } from "../../src/runner/k8s-job.js"; @@ -2032,8 +2031,9 @@ async function startServer(args: ParsedArgs): Promise { async function startServerForeground(args: ParsedArgs): Promise { const port = Number(flag(args, "port", "8080")); const host = flag(args, "host", "0.0.0.0"); - const storeMode = optionalFlag(args, "store") ?? process.env.AGENTRUN_STORE ?? process.env.AGENTRUN_MGR_STORE; - const started = await startManagerServer({ port, host, ...(storeMode === "memory" ? { store: new MemoryAgentRunStore() } : {}) }); + const storeMode = optionalFlag(args, "store"); + if (storeMode) process.env.AGENTRUN_STORE = storeMode; + const started = await startManagerServer({ port, host }); const database = await started.store.health(); return { serviceId: "agentrun-mgr", baseUrl: started.baseUrl, pid: process.pid, database, mode: "foreground", note: "foreground process; use server start without --foreground for local background mode" }; } diff --git a/src/common/kafka-events.ts b/src/common/kafka-events.ts index 912633f..db61466 100644 --- a/src/common/kafka-events.ts +++ b/src/common/kafka-events.ts @@ -1,5 +1,9 @@ -import { Kafka, logLevel, type Consumer, type EachMessagePayload, type Producer } from "kafkajs"; +import { Kafka, logLevel, type Consumer, type EachMessagePayload, type ISocketFactory, type Producer } from "kafkajs"; import { createHash } from "node:crypto"; +import { Resolver } from "node:dns"; +import * as net from "node:net"; +import { isIP, type LookupFunction } from "node:net"; +import * as tls from "node:tls"; import type { JsonRecord, JsonValue } from "./types.js"; import { redactJson, redactText } from "./redaction.js"; @@ -19,6 +23,8 @@ export interface AgentRunKafkaConfig { clientId: string; agentrunEventTopic: string; codexStdioTopic: string; + dnsServers: string[]; + dnsSearchDomains: string[]; enabled: boolean; valuesPrinted: false; } @@ -59,6 +65,8 @@ export function agentRunKafkaConfig(env: NodeJS.ProcessEnv = process.env, overri clientId, agentrunEventTopic: topicOverride && overrides.stream !== "stdio" ? topicOverride : agentrunEventTopic, codexStdioTopic: topicOverride && overrides.stream === "stdio" ? topicOverride : codexStdioTopic, + dnsServers: csv(env.AGENTRUN_KAFKA_DNS_SERVERS), + dnsSearchDomains: csv(env.AGENTRUN_KAFKA_DNS_SEARCH_DOMAINS), enabled: truthy(env.AGENTRUN_KAFKA_ENABLED) || truthy(env.AGENTRUN_KAFKA_SHADOW_PRODUCE_ENABLED) || truthy(env.AGENTRUN_KAFKA_STDIO_PRODUCE_ENABLED), valuesPrinted: false, }; @@ -156,7 +164,7 @@ export async function tailAgentRunKafka(options: KafkaTailOptions): Promise } function defaultProducerFactory(config: AgentRunKafkaConfig): Producer { - const kafka = new Kafka({ clientId: config.clientId, brokers: config.brokers, logLevel: logLevel.NOTHING }); + const kafka = createAgentRunKafka(config); return kafka.producer({ allowAutoTopicCreation: false }); } +export function createAgentRunKafka(config: AgentRunKafkaConfig, clientId = config.clientId): Kafka { + const socketFactory = kafkaSocketFactory(config.dnsServers, config.dnsSearchDomains); + return new Kafka({ clientId, brokers: config.brokers, logLevel: logLevel.NOTHING, ...(socketFactory ? { socketFactory } : {}) }); +} + +export function kafkaSocketFactory(dnsServers: string[], dnsSearchDomains: string[] = []): ISocketFactory | undefined { + const lookup = kafkaDnsLookup(dnsServers, dnsSearchDomains); + if (!lookup) return undefined; + return ({ host, port, ssl, onConnect }) => ssl + ? tls.connect({ host, port, ...ssl, lookup }, onConnect) + : net.connect({ host, port, lookup }, onConnect); +} + +export function kafkaDnsLookup(dnsServers: string[], dnsSearchDomains: string[] = []): LookupFunction | undefined { + if (dnsServers.length === 0) return undefined; + const resolver = new Resolver(); + resolver.setServers(dnsServers); + const lookup: LookupFunction = (hostname, options, callback) => { + const family = typeof options === "number" ? options : options.family; + const all = typeof options === "object" && options.all === true; + const literalFamily = isIP(hostname); + if (literalFamily > 0) { + callback(null, all ? [{ address: hostname, family: literalFamily }] : hostname, literalFamily); + return; + } + resolveKafkaHostname(resolver, hostname, dnsSearchDomains, family === 6 ? 6 : 4, (error, addresses) => { + if (error) { + callback(error, undefined as never); + return; + } + if (addresses.length === 0) { + const emptyError = new Error(`Kafka DNS lookup returned no IPv${family === 6 ? 6 : 4} addresses for ${hostname}.`) as NodeJS.ErrnoException; + emptyError.code = "agentrun_kafka_dns_empty"; + callback(emptyError, undefined as never); + return; + } + const resolvedFamily = family === 6 ? 6 : 4; + callback(null, all ? addresses.map((address) => ({ address, family: resolvedFamily })) : addresses[0]!, resolvedFamily); + }); + }; + return lookup; +} + +function resolveKafkaHostname(resolver: Resolver, hostname: string, searchDomains: string[], family: 4 | 6, callback: (error: NodeJS.ErrnoException | null, addresses: string[]) => void): void { + const candidates = [ + hostname, + ...searchDomains + .map((domain) => domain.replace(/^\.+|\.+$/gu, "")) + .filter(Boolean) + .filter((domain) => !hostname.endsWith(`.${domain}`)) + .map((domain) => `${hostname}.${domain}`), + ]; + const resolveNext = (index: number, previousError: NodeJS.ErrnoException | null = null): void => { + const candidate = candidates[index]; + if (!candidate) { + callback(previousError ?? new Error(`Kafka DNS lookup failed for ${hostname}`), []); + return; + } + const done = (error: NodeJS.ErrnoException | null, addresses: string[]): void => { + if (!error && addresses.length > 0) callback(null, addresses); + else resolveNext(index + 1, error ?? previousError); + }; + if (family === 6) resolver.resolve6(candidate, done); + else resolver.resolve4(candidate, done); + }; + resolveNext(0); +} + async function sendAgentRunKafkaMessage(producer: Producer, message: AgentRunKafkaMessage): Promise { await sendAgentRunKafkaMessageBatch(producer, [message]); } diff --git a/src/mgr/runner-dispatcher.ts b/src/mgr/runner-dispatcher.ts index 158465e..7c215ef 100644 --- a/src/mgr/runner-dispatcher.ts +++ b/src/mgr/runner-dispatcher.ts @@ -49,8 +49,9 @@ export async function dispatchRunnerIntentsOnce(input: { items.push(itemSummary(intent, "stale", null)); continue; } - const terminal = intent.attemptCount >= input.options.maxAttempts; const failure = dispatchFailure(error, intent); + const terminal = isDeterministicDispatchFailure(failure.failureKind) + || intent.attemptCount >= input.options.maxAttempts; const nextAttemptAt = new Date(Date.now() + input.options.retryBackoffMs).toISOString(); try { if (terminal) { @@ -115,6 +116,14 @@ function dispatchFailure(error: unknown, intent: RunnerDispatchIntentRecord): Js }; } +function isDeterministicDispatchFailure(failureKind: unknown): boolean { + return failureKind === "auth-missing" + || failureKind === "auth-failed" + || failureKind === "schema-invalid" + || failureKind === "tenant-policy-denied" + || failureKind === "secret-unavailable"; +} + function itemSummary(intent: RunnerDispatchIntentRecord, state: string, error: JsonRecord | null, completion?: RunnerDispatchCompletion): JsonRecord { return { dispatchIntentId: intent.id, runId: intent.runId, commandId: intent.commandId, runnerJobId: intent.runnerJobId, actualRunnerJobId: completion?.actualRunnerJobId ?? intent.actualRunnerJobId, activeRunnerId: completion?.activeRunnerId ?? intent.activeRunnerId, dispatchOutcome: completion?.dispatchOutcome ?? intent.dispatchOutcome, attemptCount: intent.attemptCount, state, failureKind: error?.failureKind ?? null, valuesPrinted: false }; } diff --git a/src/mgr/server.ts b/src/mgr/server.ts index c344d26..f70bbe7 100644 --- a/src/mgr/server.ts +++ b/src/mgr/server.ts @@ -68,7 +68,11 @@ function runnerJobDefaultsForRequest(defaults: ManagerServerOptions["runnerJobDe ...(defaults?.backendRetryMaxAttempts !== undefined ? { backendRetryMaxAttempts: defaults.backendRetryMaxAttempts } : optionalPositiveIntegerRecord("backendRetryMaxAttempts", process.env.AGENTRUN_BACKEND_RETRY_MAX_ATTEMPTS)), ...(defaults?.backendRetryInitialBackoffMs !== undefined ? { backendRetryInitialBackoffMs: defaults.backendRetryInitialBackoffMs } : optionalPositiveIntegerRecord("backendRetryInitialBackoffMs", process.env.AGENTRUN_BACKEND_RETRY_INITIAL_BACKOFF_MS)), ...(defaults?.backendRetryMaxBackoffMs !== undefined ? { backendRetryMaxBackoffMs: defaults.backendRetryMaxBackoffMs } : optionalPositiveIntegerRecord("backendRetryMaxBackoffMs", process.env.AGENTRUN_BACKEND_RETRY_MAX_BACKOFF_MS)), - ...(defaults?.kubectlCommand ? { kubectlCommand: defaults.kubectlCommand } : {}), + ...(defaults?.kubectlCommand + ? { kubectlCommand: defaults.kubectlCommand } + : process.env.AGENTRUN_KUBECTL_COMMAND + ? { kubectlCommand: process.env.AGENTRUN_KUBECTL_COMMAND } + : {}), ...(defaults?.unideskSshEndpointEnv ? { unideskSshEndpointEnv: defaults.unideskSshEndpointEnv } : {}), ...(retention ? { retention } : {}), }; diff --git a/src/runner/k8s-job.ts b/src/runner/k8s-job.ts index 3831286..9398720 100644 --- a/src/runner/k8s-job.ts +++ b/src/runner/k8s-job.ts @@ -479,6 +479,8 @@ function runnerKafkaEnvVars(env: NodeJS.ProcessEnv): JsonRecord[] { ...optionalEnvVar("AGENTRUN_KAFKA_STDIO_TOPIC", env.AGENTRUN_KAFKA_STDIO_TOPIC), ...optionalEnvVar("AGENTRUN_KAFKA_STDIO_CLIENT_ID", env.AGENTRUN_KAFKA_STDIO_CLIENT_ID), ...optionalEnvVar("AGENTRUN_KAFKA_STDIO_PRODUCE_ENABLED", env.AGENTRUN_KAFKA_STDIO_PRODUCE_ENABLED), + ...optionalEnvVar("AGENTRUN_KAFKA_DNS_SERVERS", env.AGENTRUN_KAFKA_DNS_SERVERS), + ...optionalEnvVar("AGENTRUN_KAFKA_DNS_SEARCH_DOMAINS", env.AGENTRUN_KAFKA_DNS_SEARCH_DOMAINS), ]; } diff --git a/src/selftest/cases/21-runner-durable-dispatch.ts b/src/selftest/cases/21-runner-durable-dispatch.ts index 3d4ef33..bdaff80 100644 --- a/src/selftest/cases/21-runner-durable-dispatch.ts +++ b/src/selftest/cases/21-runner-durable-dispatch.ts @@ -150,7 +150,11 @@ async function assertHungKubectlTimesOutAndRecovers(kubectlCommand: string): Pro const firstResult = await firstCycle; assert.equal(firstResult.retryCount, 1, JSON.stringify({ firstResult, intent: store.getRunnerDispatchIntent(command.id) })); assert.equal(store.getRunnerDispatchIntent(command.id)?.state, "retry"); - await new Promise((resolve) => setTimeout(resolve, 15)); + await waitFor( + () => Date.parse(store.getRunnerDispatchIntent(command.id)?.nextAttemptAt ?? "") <= Date.now(), + 1_000, + "runner dispatch retry did not become due", + ); assert.equal((await dispatchRunnerIntentsOnce({ store, defaults: defaults(kubectlCommand), options })).dispatchedCount, 1); assert.equal(store.getRunnerDispatchIntent(command.id)?.attemptCount, 2); assert.equal(store.getRunnerDispatchIntent(command.id)?.dispatchOutcome, "kubernetes-job"); @@ -292,6 +296,17 @@ async function assertTypedDispatchFailureFacts(kubectlCommand: string): Promise< assert.equal(event.payload.reason, "runner-dispatch-source-timeout"); } + const deterministicStore = new RunnerCreateFailureStore(new AgentRunError("schema-invalid", "runner assembly is invalid", { httpStatus: 400, details: { reason: "runner-assembly-invalid", valuesPrinted: false } })); + const deterministicRun = deterministicStore.createRun(runInput()); + const deterministicCommand = deterministicStore.createCommand(deterministicRun.id, validateCreateCommand({ type: "turn", payload: {}, dispatch: { kind: "kubernetes-job", input: {} } })); + const deterministic = await dispatchRunnerIntentsOnce({ store: deterministicStore, defaults: defaults(kubectlCommand), options: dispatcherOptions(5) }); + assert.equal(deterministic.failedCount, 1); + assert.equal(deterministic.retryCount, 0); + assert.equal(deterministicStore.getRunnerDispatchIntent(deterministicCommand.id)?.attemptCount, 1); + assert.equal(deterministicStore.getRunnerDispatchIntent(deterministicCommand.id)?.state, "failed"); + assert.equal(deterministicStore.getRun(deterministicRun.id).failureKind, "schema-invalid"); + assert.equal(deterministicStore.listEvents(deterministicRun.id, 0, 100).some((event) => event.payload.phase === "runner-dispatch-retry"), false); + const untypedStore = new RunnerCreateFailureStore(new Error("kubectl get jobs then kubectl create runner job failed")); const untypedRun = untypedStore.createRun(runInput()); const untypedCommand = untypedStore.createCommand(untypedRun.id, validateCreateCommand({ type: "turn", payload: {}, dispatch: { kind: "kubernetes-job", input: {} } })); diff --git a/src/selftest/cases/35-kafka-durable-outbox.ts b/src/selftest/cases/35-kafka-durable-outbox.ts index 6bf42d2..b33d53e 100644 --- a/src/selftest/cases/35-kafka-durable-outbox.ts +++ b/src/selftest/cases/35-kafka-durable-outbox.ts @@ -1,5 +1,5 @@ import assert from "node:assert/strict"; -import { closeAgentRunKafkaProducer, kafkaMessageMetadata, publishAgentRunKafkaMessage, publishAgentRunKafkaMessageBatch, rawFrameValue, setAgentRunKafkaProducerFactoryForSelfTest, type AgentRunKafkaConfig } from "../../common/kafka-events.js"; +import { agentRunKafkaConfig, closeAgentRunKafkaProducer, kafkaMessageMetadata, publishAgentRunKafkaMessage, publishAgentRunKafkaMessageBatch, rawFrameValue, setAgentRunKafkaProducerFactoryForSelfTest, type AgentRunKafkaConfig } from "../../common/kafka-events.js"; import { AgentRunError } from "../../common/errors.js"; import { agentRunOtelTraceContext } from "../../common/otel-trace.js"; import type { CreateRunInput, JsonRecord, KafkaEventOutboxRecord } from "../../common/types.js"; @@ -15,6 +15,8 @@ const config: AgentRunKafkaConfig = { clientId: "agentrun-selftest", agentrunEventTopic: "agentrun.event.v1", codexStdioTopic: "codex-stdio.raw.v1", + dnsServers: [], + dnsSearchDomains: [], enabled: true, valuesPrinted: false, }; @@ -30,6 +32,9 @@ const relayOptions: KafkaOutboxRelayOptions = { }; const selfTest: SelfTestCase = async () => { + const nativeDns = agentRunKafkaConfig({ AGENTRUN_KAFKA_DNS_SERVERS: "10.43.0.10,10.43.0.11", AGENTRUN_KAFKA_DNS_SEARCH_DOMAINS: "cluster.local" }); + assert.deepEqual(nativeDns.dnsServers, ["10.43.0.10", "10.43.0.11"]); + assert.deepEqual(nativeDns.dnsSearchDomains, ["cluster.local"]); await assert.rejects( () => openAgentRunStoreFromEnv({ AGENTRUN_STORE: "memory", AGENTRUN_KAFKA_ENABLED: "true" }), (error) => error instanceof AgentRunError && error.message.includes("configuration is incomplete"), @@ -180,7 +185,7 @@ const selfTest: SelfTestCase = async () => { await assertStaleClaimsCannotOverwrite(); await assertProducerRecreatedAfterDisconnected(); - return { name: "kafka-durable-outbox", tests: ["explicit-enable", "canonical-event", "warm-run-command-trace-authority", "run-level-trace-fallback", "stdio-command-trace-context", "kafka-tail-correlation-summary", "durable-key-batch", "native-broker-batch", "ack-before-delivered", "claim-duration-visible", "monotonic-backlog-drain", "settle-stale-classification", "cross-key-parallelism", "batch-bounded-parallelism", "failed-key-isolation", "atomic-dispatch-terminal", "terminal-command-idempotent-replay", "terminal-command-replay-no-new-events-or-outbox", "terminal-command-same-key-payload-conflict", "terminal-command-different-key-fails-closed", "stale-claim-fencing", "producer-reconnect"] }; + return { name: "kafka-durable-outbox", tests: ["explicit-enable", "native-dns-config", "canonical-event", "warm-run-command-trace-authority", "run-level-trace-fallback", "stdio-command-trace-context", "kafka-tail-correlation-summary", "durable-key-batch", "native-broker-batch", "ack-before-delivered", "claim-duration-visible", "monotonic-backlog-drain", "settle-stale-classification", "cross-key-parallelism", "batch-bounded-parallelism", "failed-key-isolation", "atomic-dispatch-terminal", "terminal-command-idempotent-replay", "terminal-command-replay-no-new-events-or-outbox", "terminal-command-same-key-payload-conflict", "terminal-command-different-key-fails-closed", "stale-claim-fencing", "producer-reconnect"] }; }; function assertWarmRunCommandTraceAuthority(): void { diff --git a/src/selftest/cases/36-stdio-event-reconstruction.ts b/src/selftest/cases/36-stdio-event-reconstruction.ts index 9323d6a..4a68e8c 100644 --- a/src/selftest/cases/36-stdio-event-reconstruction.ts +++ b/src/selftest/cases/36-stdio-event-reconstruction.ts @@ -225,6 +225,8 @@ async function assertOneShotPublisherOwnership(frames: ReturnType