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
6 changes: 6 additions & 0 deletions .changeset/openapi-transport-unreachable.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
---
"executor": patch
"@executor-js/plugin-openapi": patch
---

OpenAPI tools that cannot reach the upstream server now return an `upstream_unreachable` error instead of `Internal tool error [id]`. The message names the integration and origin that could not be reached, `details` carries the sanitized `host` and errno-style `code` (`ECONNREFUSED`, `ENOTFOUND`, …), and the failure is logged with the same classification.
11 changes: 10 additions & 1 deletion apps/cloud/src/mcp/session-build-semaphore.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { describe, expect, it, beforeEach } from "@effect/vitest";
import { describe, expect, it, beforeEach, afterEach, vi } from "@effect/vitest";

import {
acquireBuildSlot,
Expand All @@ -13,6 +13,10 @@ describe("session-build-semaphore", () => {
resetBuildSlotsForTest();
});

afterEach(() => {
vi.useRealTimers();
});

it("grants up to the cap immediately, with no wait", async () => {
const results = await Promise.all([
acquireBuildSlot().promise,
Expand Down Expand Up @@ -214,6 +218,7 @@ describe("session-build-semaphore", () => {
});

it("proceeds without a slot when the queue wait exceeds the timeout, and does not count it as active", async () => {
vi.useFakeTimers();
await Promise.all([
acquireBuildSlot().promise,
acquireBuildSlot().promise,
Expand All @@ -223,6 +228,10 @@ describe("session-build-semaphore", () => {
expect(currentActiveBuildsForTest()).toBe(4);

const timedOutHandle = acquireBuildSlot(10);
await vi.advanceTimersByTimeAsync(9);
expect(currentQueueLengthForTest()).toBe(1);
expect(currentActiveBuildsForTest()).toBe(4);
await vi.advanceTimersByTimeAsync(1);
const result = await timedOutHandle.promise;

expect(result).toEqual({ acquired: false, waitMs: expect.any(Number), timedOut: true });
Expand Down
226 changes: 226 additions & 0 deletions e2e/scenarios/openapi-unreachable-artifact.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,226 @@
// Cross-target: an artifact whose OpenAPI query cannot reach its upstream gets
// an actionable network error, not the opaque defect mask. This walks the real
// path from a saved artifact through the nested shell, execute-action, sandbox,
// OpenAPI transport, and back into ArtifactError.
import { randomBytes } from "node:crypto";
import { createServer } from "node:http";

import { expect } from "@effect/vitest";
import { Effect } from "effect";
import type { Page } from "playwright";
import { composePluginApi } from "@executor-js/api/server";
import { openApiHttpPlugin } from "@executor-js/plugin-openapi/api";
import { ConnectionName, IntegrationSlug, type ArtifactId } from "@executor-js/sdk/shared";

import { scenario } from "../src/scenario";
import { Api, Browser, Mcp, Target } from "../src/services";
import { visit } from "../src/surfaces/browser";
import type { McpSession } from "../src/surfaces/mcp";

const api = composePluginApi([openApiHttpPlugin()] as const);

const unique = (prefix: string) => `${prefix}_${randomBytes(4).toString("hex")}`;

type DroppingUpstream = {
readonly url: string;
readonly requests: () => number;
readonly close: () => void;
};

// Accept the request, then drop the socket before sending response headers.
// This produces a real transport failure without relying on a hardcoded or
// temporarily-unused port.
const serveDroppingUpstream = () =>
Effect.acquireRelease(
Effect.callback<DroppingUpstream>((resume) => {
let hits = 0;
const server = createServer((_request, response) => {
hits += 1;
response.destroy();
});
server.listen(0, "127.0.0.1", () => {
const address = server.address();
const port = typeof address === "object" && address ? address.port : 0;
resume(
Effect.succeed({
url: `http://127.0.0.1:${port}`,
requests: () => hits,
close: () => {
server.close();
server.closeAllConnections();
},
}),
);
});
}),
(server) => Effect.sync(server.close),
);

const unreachableSpec = (baseUrl: string): string =>
JSON.stringify({
openapi: "3.0.3",
info: { title: "Unreachable API", version: "1.0.0" },
servers: [{ url: baseUrl }],
paths: {
"/things": {
get: {
tags: ["things"],
operationId: "listThings",
summary: "List things",
responses: {
"200": {
description: "Things",
content: {
"application/json": {
schema: { type: "array", items: { type: "object" } },
},
},
},
},
},
},
},
});

const createConnectionCode = (slug: string) => `
const created = await tools.executor.coreTools.connections.create({
owner: "org",
name: "public",
integration: ${JSON.stringify(slug)},
template: "none",
});
return JSON.stringify(created.ok ? { ok: true } : { ok: false, error: created.error });
`;

const executeApproved = (session: McpSession, code: string) =>
Effect.gen(function* () {
let result = yield* session.call("execute", { code });
let guard = 0;
while (result.text.includes("executionId:") && guard < 10) {
result = yield* session.approvePaused(result.text);
guard += 1;
}
expect(result.ok, `execute completed (got: ${result.text.slice(0, 400)})`).toBe(true);
return result.text;
});

const artifactSource = (slug: string) => `
function App() {
const query = useQuery(tools.${slug}.things.listThings.queryOptions({}));
const result = query.data;
return (
<div className="flex h-full flex-col gap-4">
<h2>Upstream status</h2>
<div data-testid="upstream-state" className="min-h-0 flex-1">
{query.isLoading ? (
<ArtifactLoading />
) : query.error ? (
<ArtifactError error={query.error} onRetry={query.refetch} />
) : result?.ok === false ? (
<ArtifactError error={result.error} onRetry={query.refetch} />
) : (
<p>Unexpected upstream success</p>
)}
</div>
</div>
);
}
`;

const structuredOf = (result: { readonly raw: unknown }): Record<string, unknown> =>
((result.raw as { structuredContent?: Record<string, unknown> }).structuredContent ??
{}) as Record<string, unknown>;

const artifactContent = (page: Page) =>
page.frameLocator('[data-testid="artifact-shell-frame"]').frameLocator("iframe");

scenario(
"Artifacts · an unreachable OpenAPI host shows actionable retry guidance instead of an internal error",
{ timeout: 180_000 },
Effect.scoped(
Effect.gen(function* () {
const target = yield* Target;
const browser = yield* Browser;
const mcp = yield* Mcp;
const { client: makeClient } = yield* Api;

const identity = yield* target.newIdentity();
const client = yield* makeClient(api, identity);
const session = mcp.session(identity);
const upstream = yield* serveDroppingUpstream();
const slug = unique("unreachable");
const title = `Unreachable upstream ${randomBytes(4).toString("hex")}`;
let artifactId: ArtifactId | undefined;

yield* Effect.ensuring(
Effect.gen(function* () {
yield* client.openapi.addSpec({
payload: {
spec: { kind: "blob", value: unreachableSpec(upstream.url) },
slug,
baseUrl: upstream.url,
},
});

const created = yield* executeApproved(session, createConnectionCode(slug));
expect(created, `the no-auth connection was created: ${created}`).toContain('"ok":true');

const rendered = yield* session.call("create-artifact", {
code: artifactSource(slug),
title,
description: "Shows whether the upstream API is reachable",
connections: { [slug]: `${slug}.org.public` },
});
expect(rendered.ok, `create-artifact succeeded: ${rendered.text}`).toBe(true);

const structured = structuredOf(rendered);
artifactId = structured.artifactId as ArtifactId;
expect(artifactId, "the artifact was persisted").toBeTruthy();

yield* browser.session(identity, async ({ page, step }) => {
await step("Open the artifact that reads from the unreachable API", async () => {
await visit(page, String(structured.url));
await page.getByRole("heading", { name: title }).waitFor({ timeout: 20_000 });
});

await step(
"The artifact explains that the upstream host could not be reached",
async () => {
const state = artifactContent(page).getByTestId("upstream-state");
await state.locator('[data-slot="artifact-error"]').waitFor({ timeout: 30_000 });
const message = await state.innerText();

expect(message, "the user gets actionable network guidance").toContain(
`Could not reach the upstream server for "${slug}"`,
);
expect(message, "the opaque defect mask never reaches the artifact").not.toContain(
"Internal tool error",
);
expect(message, "the request path is not leaked").not.toContain("/things");
},
);
});

expect(upstream.requests(), "the artifact made a real upstream request").toBeGreaterThan(
0,
);
}),
Effect.gen(function* () {
if (artifactId !== undefined) {
yield* client.artifacts.remove({ params: { artifactId } }).pipe(Effect.ignore);
}
yield* client.connections
.remove({
params: {
owner: "org",
integration: IntegrationSlug.make(slug),
name: ConnectionName.make("public"),
},
})
.pipe(Effect.ignore);
yield* client.openapi.removeSpec({ params: { slug } }).pipe(Effect.ignore);
}),
);
}),
),
);
39 changes: 38 additions & 1 deletion packages/plugins/openapi/src/sdk/backing.ts
Original file line number Diff line number Diff line change
Expand Up @@ -629,6 +629,22 @@ export const resolveOpenApiBackedTools = ({
};
});

// Transport failures used to escape as defects, which the hosts log with a
// correlation id. As a typed tool failure nothing else records them, so log
// and annotate the span with the sanitized classification operators need to
// tell DNS from refused from TLS.
const recordUpstreamUnreachable = (integration: string, error: OpenApiInvocationError) => {
const annotations = {
"plugin.openapi.integration": integration,
"plugin.openapi.upstream.host": error.upstreamHost ?? "unknown",
"plugin.openapi.upstream.transport_code": error.transportCode ?? "unknown",
};
return Effect.logWarning("OpenAPI upstream unreachable").pipe(
Effect.annotateLogs(annotations),
Effect.andThen(Effect.annotateCurrentSpan(annotations)),
);
};

export const invokeOpenApiBackedTool = (input: {
readonly ctx: PluginCtx<OpenapiStore>;
readonly toolRow: { readonly integration: string; readonly name: string };
Expand Down Expand Up @@ -727,7 +743,28 @@ export const invokeOpenApiBackedTool = (input: {
details: error.cause ?? error,
}),
})
: Effect.fail(error),
: error.reason === "transport_error"
? recordUpstreamUnreachable(integration, error).pipe(
Effect.as({
ok: false as const,
failure: ToolResult.fail({
code: "upstream_unreachable",
// Executor sends the request, not the user's browser, so
// point at what the user can act on: the configured
// origin and the service behind it.
message: `Could not reach the upstream server for "${integration}"${error.upstreamHost ? ` at ${error.upstreamHost}` : ""}. Verify the integration's base URL and that the service is online, then try again.`,
// Unlike the timeout branches, `error.cause` is withheld:
// the TransportError carries the whole request, including
// resolved auth headers. Absent fields are dropped, not
// `undefined`: the result must stay a JSON value.
details: {
...(error.upstreamHost !== undefined ? { host: error.upstreamHost } : {}),
...(error.transportCode !== undefined ? { code: error.transportCode } : {}),
},
}),
}),
)
: Effect.fail(error),
),
);

Expand Down
14 changes: 13 additions & 1 deletion packages/plugins/openapi/src/sdk/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,19 @@ export class OpenApiSpecOverrideError extends Schema.TaggedErrorClass<OpenApiSpe
export class OpenApiInvocationError extends Data.TaggedError("OpenApiInvocationError")<{
readonly message: string;
readonly statusCode: Option.Option<number>;
readonly reason?: "response_headers_timeout" | "response_body_timeout" | "unknown_arguments";
readonly reason?:
| "response_headers_timeout"
| "response_body_timeout"
| "unknown_arguments"
| "transport_error";
// `host[:port]` of a request that failed at the transport layer. It is the
// integration's configured origin, so it is safe to show; the path, query,
// and headers stay on `cause`.
readonly upstreamHost?: string | undefined;
// Errno-style code behind a transport failure (`ECONNREFUSED`, `ENOTFOUND`,
// `UND_ERR_SOCKET`, …) when the runtime exposes one. Tells DNS from refused
// from TLS without exposing the request.
readonly transportCode?: string | undefined;
readonly cause?: unknown;
}> {}

Expand Down
Loading
Loading