Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 37 additions & 17 deletions open-sse/executors/antigravity.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,10 @@ import {
} from "./antigravity/sseCollect.ts";
// processAntigravitySSEPayload re-exported for external importers (tests).
export { processAntigravitySSEPayload } from "./antigravity/sseCollect.ts";
import {
handleAntigravityFallbackChainError,
handleAntigravityFallback400,
} from "./antigravity/proFallbackChain.ts";
import {
applyAntigravityClientProfileHeaders,
removeHeaderCaseInsensitive,
Expand Down Expand Up @@ -1070,31 +1074,47 @@ export class AntigravityExecutor extends BaseExecutor {
let firstResult: Awaited<ReturnType<AntigravityExecutor["executeOnce"]>> | null = null;
for (let i = 0; i < chain.length; i++) {
const candidate = chain[i];
const result = await this.executeOnce(input, candidate);
let result: Awaited<ReturnType<AntigravityExecutor["executeOnce"]>>;
try {
result = await this.executeOnce(input, candidate);
} catch (error) {
const outcome = handleAntigravityFallbackChainError(
input,
error,
candidate,
i,
chain,
firstResult,
resolvedUpstreamId
);
switch (outcome.action) {
case "throw":
throw outcome.error;
case "return":
return outcome.result;
default:
continue;
}
}

// Success (or any non-400) on a candidate → return immediately.
if (result.response.status !== HTTP_STATUS.BAD_REQUEST) {
return result;
}

// Remember the FIRST 400 so the exhausted-chain case surfaces the original error.
if (i === 0) firstResult = result;

const isLast = i === chain.length - 1;
if (!isLast) {
input.log?.debug?.(
"AG_PRO_FALLBACK",
`400 on "${candidate}" — retrying with next Pro candidate "${chain[i + 1]}"`
);
continue;
}

// Chain exhausted: surface the FIRST candidate's sanitized 400.
input.log?.warn?.(
"AG_PRO_FALLBACK",
`Pro fallback chain exhausted (all ${chain.length} candidates 400'd) for "${resolvedUpstreamId}"`
if (!firstResult) firstResult = result;

const outcome400 = handleAntigravityFallback400(
input,
result,
firstResult,
candidate,
i,
chain,
resolvedUpstreamId
);
return firstResult ?? result;
if (outcome400.action === "return") return outcome400.result;
}

// Unreachable (loop always returns), but keeps the type checker happy.
Expand Down
104 changes: 104 additions & 0 deletions open-sse/executors/antigravity/proFallbackChain.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
// Pure Pro-family fallback-chain decision helpers for the Antigravity executor (#7290):
// decide what execute()'s per-candidate loop does after executeOnce() throws or
// returns a 400, without depending on executor instance state (no `this`).
// Extracted from antigravity.ts (file-size cap) -- mirrors the existing
// antigravity/sseCollect.ts submodule pattern.
import type { ExecuteInput } from "../base.ts";

/** Shape of one execute()/executeOnce() result (kept local to avoid importing the class). */
export type AntigravityExecuteResult = {
response: Response;
url: string;
headers: Record<string, string>;
transformedBody: unknown;
};

/** True for an aborted request (caller disconnect) — never retried across candidates. */
export function isAntigravityAbortError(input: ExecuteInput, error: unknown): boolean {
return Boolean(
input.signal?.aborted ||
(error instanceof DOMException && error.name === "AbortError") ||
(error instanceof Error && error.name === "AbortError")
);
}

export type AntigravityFallbackChainErrorOutcome =
| { action: "throw"; error: unknown }
| { action: "return"; result: AntigravityExecuteResult }
| { action: "continue" };

/**
* Decide what execute()'s Pro-fallback loop does after executeOnce() THROWS for one
* candidate: propagate an abort immediately, retry the next candidate, surface the
* first 400 if the chain is exhausted, or throw a chain-exhausted error.
*/
export function handleAntigravityFallbackChainError(
input: ExecuteInput,
error: unknown,
candidate: string,
i: number,
chain: readonly string[],
firstResult: AntigravityExecuteResult | null,
resolvedUpstreamId: string
): AntigravityFallbackChainErrorOutcome {
// Abort signal (user disconnect) — propagate immediately, do not retry.
if (isAntigravityAbortError(input, error)) {
return { action: "throw", error };
}
if (i < chain.length - 1) {
input.log?.debug?.(
"AG_PRO_FALLBACK",
`Exception on "${candidate}" (${error instanceof Error ? error.message : String(error)}) -- retrying with next Pro candidate "${chain[i + 1]}"`
);
return { action: "continue" };
}
// Last candidate also threw -- return original 400 if available, otherwise throw.
if (firstResult) {
input.log?.warn?.(
"AG_PRO_FALLBACK",
`Pro fallback chain exhausted (last candidate threw, but first candidate returned 400) for "${resolvedUpstreamId}". Returning original 400.`
);
return { action: "return", result: firstResult };
}
return {
action: "throw",
error: new Error(
`Pro fallback chain exhausted (all ${chain.length} candidates failed). Last error: ${error instanceof Error ? error.message : String(error)}`
),
};
}

export type AntigravityFallback400Outcome =
| { action: "return"; result: AntigravityExecuteResult }
| { action: "continue" };

/**
* Decide what execute()'s Pro-fallback loop does after one candidate returns a 400:
* retry the next candidate, or (chain exhausted) surface the first candidate's
* sanitized 400.
*/
export function handleAntigravityFallback400(
input: ExecuteInput,
result: AntigravityExecuteResult,
firstResult: AntigravityExecuteResult | null,
candidate: string,
i: number,
chain: readonly string[],
resolvedUpstreamId: string
): AntigravityFallback400Outcome {
const isLast = i === chain.length - 1;
if (!isLast) {
input.log?.debug?.(
"AG_PRO_FALLBACK",
`400 on "${candidate}" — retrying with next Pro candidate "${chain[i + 1]}"`
);
return { action: "continue" };
}

// Chain exhausted: surface the FIRST candidate's sanitized 400.
input.log?.warn?.(
"AG_PRO_FALLBACK",
`Pro fallback chain exhausted (all ${chain.length} candidates 400'd) for "${resolvedUpstreamId}"`
);
return { action: "return", result: firstResult ?? result };
}
128 changes: 90 additions & 38 deletions src/mitm/dns/provision.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,8 @@
* is guarded and unit-testable without spawning the MITM server (#6127 / #6198).
*/

import { addDNSEntry, addDNSEntries } from "./dnsConfig.ts";
import { addDNSEntry, addDNSEntries, isSudoAvailable } from "./dnsConfig.ts";
import { isRoot } from "../systemCommands.ts";
import { ALL_TARGETS } from "../targets/index.ts";
import { getAllAgentBridgeStates } from "@/lib/db/agentBridgeState.ts";
import { listCustomHosts } from "@/lib/db/inspectorCustomHosts.ts";
Expand All @@ -23,45 +24,39 @@ export interface DnsProvisionDeps {
addHostsDns?: (hosts: string[], sudoPassword: string) => Promise<void>;
getAgentStates?: () => ReturnType<typeof getAllAgentBridgeStates>;
listEnabledCustomHosts?: () => ReturnType<typeof listCustomHosts>;
/** Return true if privileged host-file writes are possible (sudo or root). */
canElevate?: () => boolean;
logger?: DnsProvisionLogger;
}

/**
* Provision every AgentBridge DNS entry (Antigravity defaults + agents with
* `dns_enabled=true` + enabled custom hosts). **Every step is best-effort**: a failure
* is logged with the full `err` — which carries the privileged command's stderr
* (`systemCommands.ts` folds stderr into the Error message) — and never aborts the
* bridge start.
*
* Previously the default step (`addDNSEntry`) was called unguarded while the two
* sibling steps and cert install were wrapped, so in containers/headless (Docker
* `USER node`, no `sudo`, read-only /etc/hosts) it threw out of `startMitmInternal`
* and killed the whole start (#6127); its stderr also never reached app.log — only a
* bare exit code hit the toast (#6198). Extracting + guarding all three steps here
* restores the symmetry and makes the behavior unit-testable without spawning the
* MITM server.
*/
export async function provisionDnsEntries(
/** Fully-resolved dependency set used by the per-step provisioning helpers below. */
type ResolvedDnsProvisionDeps = {
addDefaultDns: (sudoPassword: string) => Promise<void>;
addHostsDns: (hosts: string[], sudoPassword: string) => Promise<void>;
getAgentStates: () => ReturnType<typeof getAllAgentBridgeStates>;
listEnabledCustomHosts: () => ReturnType<typeof listCustomHosts>;
logger: DnsProvisionLogger;
};

/** Antigravity default hosts — best-effort, never throws. */
async function provisionDefaultDns(
sudoPassword: string,
deps: DnsProvisionDeps = {}
deps: ResolvedDnsProvisionDeps
): Promise<void> {
const addDefaultDns = deps.addDefaultDns ?? addDNSEntry;
const addHostsDns = deps.addHostsDns ?? addDNSEntries;
const getAgentStates = deps.getAgentStates ?? getAllAgentBridgeStates;
const listEnabledCustomHosts =
deps.listEnabledCustomHosts ?? (() => listCustomHosts({ enabledOnly: true }));
const logger = deps.logger ?? defaultLog;

// Antigravity default hosts.
try {
await addDefaultDns(sudoPassword);
await deps.addDefaultDns(sudoPassword);
} catch (err) {
logger.error({ err }, "Failed to add default DNS entries (continuing)");
deps.logger.error({ err }, "Failed to add default DNS entries (continuing)");
}
}

// Collect hosts from agents that have dns_enabled=true in the DB.
/** Hosts for agents with `dns_enabled=true` in the DB — best-effort, never throws. */
async function provisionAgentDns(
sudoPassword: string,
deps: ResolvedDnsProvisionDeps
): Promise<void> {
try {
const agentStates = getAgentStates();
const agentStates = deps.getAgentStates();
const agentHostsToAdd: string[] = [];
for (const state of agentStates) {
if (!state.dns_enabled) continue;
Expand All @@ -71,22 +66,79 @@ export async function provisionDnsEntries(
}
}
if (agentHostsToAdd.length > 0) {
logger.info({ count: agentHostsToAdd.length }, "Adding DNS for agent host(s)...");
await addHostsDns(agentHostsToAdd, sudoPassword);
deps.logger.info({ count: agentHostsToAdd.length }, "Adding DNS for agent host(s)...");
await deps.addHostsDns(agentHostsToAdd, sudoPassword);
}
} catch (err) {
logger.error({ err }, "Failed to add agent DNS entries (continuing)");
deps.logger.error({ err }, "Failed to add agent DNS entries (continuing)");
}
}

// Collect enabled custom hosts.
/** Enabled custom hosts — best-effort, never throws. */
async function provisionCustomHostsDns(
sudoPassword: string,
deps: ResolvedDnsProvisionDeps
): Promise<void> {
try {
const customHosts = listEnabledCustomHosts();
const customHosts = deps.listEnabledCustomHosts();
const customHostNames = customHosts.map((h) => h.host);
if (customHostNames.length > 0) {
logger.info({ count: customHostNames.length }, "Adding DNS for custom host(s)...");
await addHostsDns(customHostNames, sudoPassword);
deps.logger.info({ count: customHostNames.length }, "Adding DNS for custom host(s)...");
await deps.addHostsDns(customHostNames, sudoPassword);
}
} catch (err) {
logger.error({ err }, "Failed to add custom host DNS entries (continuing)");
deps.logger.error({ err }, "Failed to add custom host DNS entries (continuing)");
}
}

/**
* Provision every AgentBridge DNS entry (Antigravity defaults + agents with
* `dns_enabled=true` + enabled custom hosts). **Every step is best-effort**: a failure
* is logged with the full `err` — which carries the privileged command's stderr
* (`systemCommands.ts` folds stderr into the Error message) — and never aborts the
* bridge start.
*
* Previously the default step (`addDNSEntry`) was called unguarded while the two
* sibling steps and cert install were wrapped, so in containers/headless (Docker
* `USER node`, no `sudo`, read-only /etc/hosts) it threw out of `startMitmInternal`
* and killed the whole start (#6127); its stderr also never reached app.log — only a
* bare exit code hit the toast (#6198). Extracting + guarding all three steps here
* restores the symmetry and makes the behavior unit-testable without spawning the
* MITM server.
*/
export async function provisionDnsEntries(
sudoPassword: string,
deps: DnsProvisionDeps = {}
): Promise<void> {
const canElevate = deps.canElevate ?? (() => isSudoAvailable() || isRoot());
const logger = deps.logger ?? defaultLog;

// Explicit opt-out: skip all DNS modification when the env var is set.
if (process.env.SKIP_ANTIGRAVITY_DNS === "true") {
logger.info("Skipping DNS entries - SKIP_ANTIGRAVITY_DNS=true");
return;
}

// In containers (USER node, no sudo installed, not root) we cannot write
// to /etc/hosts. Rather than attempting sudo and swallowing the error,
// detect the condition up-front and bail out with a clear message.
if (!canElevate()) {
logger.info(
"Skipping DNS entries - sudo not available and not running as root (likely a container)"
);
return;
}

const resolvedDeps: ResolvedDnsProvisionDeps = {
addDefaultDns: deps.addDefaultDns ?? addDNSEntry,
addHostsDns: deps.addHostsDns ?? addDNSEntries,
getAgentStates: deps.getAgentStates ?? getAllAgentBridgeStates,
listEnabledCustomHosts:
deps.listEnabledCustomHosts ?? (() => listCustomHosts({ enabledOnly: true })),
logger,
};

await provisionDefaultDns(sudoPassword, resolvedDeps);
await provisionAgentDns(sudoPassword, resolvedDeps);
await provisionCustomHostsDns(sudoPassword, resolvedDeps);
}
Loading