Skip to content

Commit 85cf428

Browse files
authored
Cascade integration removal to every member and stop serving orphaned rows (#1991)
* Cascade integration removal to every member and stop serving orphaned rows * Register the MCP server before starting OAuth in the local app test
1 parent e29c360 commit 85cf428

10 files changed

Lines changed: 367 additions & 11 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@executor-js/sdk": patch
3+
---
4+
5+
Removing an integration now drops every member's connections and tools under it, not only the remover's own. Tool and connection listings no longer serve rows whose integration is gone from the catalog, invoking such a tool reports the missing integration, and `oauth.start` refuses an unknown integration before creating a session.

‎apps/local/src/mcp-oauth.test.ts‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,10 @@ const TEST_BASE_URL = "http://local.test";
7373

7474
interface Harness {
7575
readonly fetch: typeof globalThis.fetch;
76+
readonly registerRemoteServer: (input: {
77+
readonly slug: string;
78+
readonly endpoint: string;
79+
}) => Effect.Effect<void, unknown>;
7680
readonly dispose: () => Promise<void>;
7781
}
7882

@@ -132,6 +136,18 @@ const startHarness = async (tmpDir: string): Promise<Harness> => {
132136
webHandler(
133137
input instanceof Request ? input : new Request(input, init),
134138
)) as typeof globalThis.fetch,
139+
// `oauth.start` refuses an integration that is not in the catalog, so the
140+
// flow under test needs a registered MCP server to mint against.
141+
registerRemoteServer: ({ slug, endpoint }) =>
142+
executor.mcp
143+
.addServer({
144+
transport: "remote",
145+
name: slug,
146+
slug,
147+
endpoint,
148+
authenticationTemplate: [{ kind: "oauth2", slug: "oauth" }],
149+
})
150+
.pipe(Effect.asVoid),
135151
dispose: async () => {
136152
await Effect.runPromise(Effect.ignore(Effect.tryPromise(() => disposeHandler())));
137153
await Effect.runPromise(
@@ -191,6 +207,12 @@ describe("local oauth (real OAuth discovery + stubbed start)", () => {
191207
expect(probed.authorizationUrl).toBe(oauth.authorizationEndpoint);
192208
expect(probed.tokenUrl).toBe(oauth.tokenEndpoint);
193209

210+
// The catalog row `start` mints against.
211+
yield* harness.registerRemoteServer({
212+
slug: "mcp_remote",
213+
endpoint: oauth.mcpResourceUrl,
214+
});
215+
194216
// createClient — register an owner-scoped OAuth app for the start flow.
195217
const slug = `mcp-oauth2-${randomBytes(4).toString("hex")}`;
196218
const created = yield* run((client) =>

‎packages/core/sdk/src/errors.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -290,6 +290,9 @@ export type ExecuteError =
290290
| PluginNotLoadedError
291291
| NoHandlerError
292292
| ConnectionNotFoundError
293+
/** The tool row outlived its integration (an orphan the catalog no longer
294+
* lists), so there is no plugin config to invoke it against. */
295+
| IntegrationNotFoundError
293296
| CredentialProviderNotRegisteredError
294297
| CredentialResolutionError
295298
| ElicitationDeclinedError
@@ -298,6 +301,5 @@ export type ExecuteError =
298301
/** Convenience union spanning every typed error the SDK raises. */
299302
export type ExecutorError =
300303
| ExecuteError
301-
| IntegrationNotFoundError
302304
| IntegrationRemovalNotAllowedError
303305
| ArtifactNotFoundError;

‎packages/core/sdk/src/executor.ts‎

Lines changed: 57 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1986,6 +1986,22 @@ export const createExecutor = <const TPlugins extends readonly AnyPlugin[] = rea
19861986
const healthProbeInFlight = healthProbeGateFor(rootDbUntyped);
19871987
const fuma = makeFumaClient(rootDb);
19881988
const core = makeCoreDb(fuma);
1989+
// The ONE tenant-wide mutating handle: delete-only, tenant reach. Used
1990+
// solely by the integration-removal cascade, which must drop EVERY
1991+
// member's connections and tools under the removed slug — a bound admin
1992+
// can only reach its own rows, and the rest would survive as orphans that
1993+
// still list and invoke. The context rebinds inside the removal
1994+
// transaction, so the cascade commits or rolls back with the catalog row.
1995+
// Never exposed to plugins or request surfaces.
1996+
const cascadeCore = makeCoreDb(
1997+
makeFumaClient(rootDb, {
1998+
context: {
1999+
...ownerContext,
2000+
reach: "tenant",
2001+
writes: "delete-only",
2002+
} satisfies ExecutorOwnerPolicyContext,
2003+
}),
2004+
);
19892005
const blobs = config.blobs ?? makeFumaBlobStore(fuma);
19902006
const transaction = <A, E>(effect: Effect.Effect<A, E>) => fuma.transaction(effect);
19912007

@@ -3071,6 +3087,13 @@ export const createExecutor = <const TPlugins extends readonly AnyPlugin[] = rea
30713087
where: (b: AnyCb) => b("slug", "=", String(slug)),
30723088
});
30733089

3090+
/** Every slug in the tenant's catalog — the set an owned row's
3091+
* `integration` must belong to for the row to be servable. */
3092+
const listCatalogSlugs = (): Effect.Effect<ReadonlySet<string>, StorageFailure> =>
3093+
core
3094+
.findMany("integration", { select: ["slug"] })
3095+
.pipe(Effect.map((rows) => new Set(rows.map((row) => String(row.slug)))));
3096+
30743097
// Project a row's stored config into declared auth methods via the owning
30753098
// plugin's `describeAuthMethods` hook. The hook is plugin-authored, so a
30763099
// throw (malformed config it didn't guard) degrades to `[]` rather than
@@ -3337,14 +3360,21 @@ export const createExecutor = <const TPlugins extends readonly AnyPlugin[] = rea
33373360
),
33383361
);
33393362
}
3340-
// Drop owned connections / tools / definitions for this integration.
3341-
const where = (b: AnyCb) => b("integration", "=", String(slug));
3342-
yield* core.deleteMany("tool", { where });
3343-
yield* core.deleteMany("definition", { where });
3344-
yield* core.deleteMany("connection", { where });
3363+
// The catalog row goes first through the bound handle: a read-only
3364+
// (platform-view) context is refused here, before the widened
3365+
// cascade below could touch anything.
33453366
yield* core.deleteMany("integration", {
33463367
where: (b: AnyCb) => b("slug", "=", String(slug)),
33473368
});
3369+
// Drop connections / tools / definitions for this integration across
3370+
// EVERY subject in the tenant, not just the remover's own rows. A
3371+
// removed integration has no reason to keep anyone's rows, and rows
3372+
// left behind become orphans: invisible in the catalog, yet still
3373+
// listed to agents and still targetable by reconnect.
3374+
const where = (b: AnyCb) => b("integration", "=", String(slug));
3375+
yield* cascadeCore.deleteMany("tool", { where });
3376+
yield* cascadeCore.deleteMany("definition", { where });
3377+
yield* cascadeCore.deleteMany("connection", { where });
33483378
return existing.plugin_id;
33493379
}),
33503380
).pipe(
@@ -4650,7 +4680,13 @@ export const createExecutor = <const TPlugins extends readonly AnyPlugin[] = rea
46504680
filter?.owner === undefined ? true : b("owner", "=", filter.owner),
46514681
),
46524682
});
4653-
const connections = rows.map(rowToConnection);
4683+
// Same catalog gate as `toolsList`: a connection whose integration was
4684+
// removed is an orphan, and offering it (in the accounts list, or to
4685+
// an agent as a reconnect target) leads into flows that cannot mint.
4686+
const catalogSlugs = yield* listCatalogSlugs();
4687+
const connections = rows
4688+
.filter((row) => catalogSlugs.has(String(row.integration)))
4689+
.map(rowToConnection);
46544690
if (!activeToolPolicyProvider) return connections;
46554691

46564692
const visibleTools = yield* toolsList({ includeAnnotations: false });
@@ -5587,8 +5623,14 @@ export const createExecutor = <const TPlugins extends readonly AnyPlugin[] = rea
55875623
});
55885624
const includeBlocked = filter?.includeBlocked ?? false;
55895625
const policyRules = yield* listActivePolicyRuleSet();
5626+
// Only tools whose integration is still in the catalog. A tool row
5627+
// whose integration was removed is an orphan (a removal that could
5628+
// not reach this subject's rows): listing it invites an invoke that
5629+
// cannot resolve its config and a reconnect that cannot mint.
5630+
const catalogSlugs = yield* listCatalogSlugs();
55905631
const tools: Tool[] = [];
55915632
for (const row of rows) {
5633+
if (!catalogSlugs.has(String(row.integration))) continue;
55925634
const tool = rowToTool(row);
55935635
if (!matchesToolFilter(tool, filter)) continue;
55945636
if (!includeBlocked) {
@@ -6489,6 +6531,13 @@ export const createExecutor = <const TPlugins extends readonly AnyPlugin[] = rea
64896531
Effect.onError(() => Fiber.interrupt(integrationRowFiber)),
64906532
);
64916533
const integrationRow = yield* Fiber.join(integrationRowFiber);
6534+
// A tool row that outlived its integration (an orphan the catalog no
6535+
// longer lists) is not invokable: its plugin config is gone, and
6536+
// the auth-recovery hints would steer the caller into an OAuth
6537+
// flow that cannot mint. Report it as the missing integration it is.
6538+
if (!integrationRow) {
6539+
return yield* new IntegrationNotFoundError({ slug: parsed.integration });
6540+
}
64926541
const grantedScopes = grantedScopesFromRow(connectionRow);
64936542
const invokeTool = runtime.plugin.invokeTool;
64946543
const invokeWith = (
@@ -6501,7 +6550,7 @@ export const createExecutor = <const TPlugins extends readonly AnyPlugin[] = rea
65016550
template: AuthTemplateSlug.make(connectionRow.template),
65026551
value: resolved[PRIMARY_INPUT_VARIABLE] ?? null,
65036552
values: resolved,
6504-
config: integrationRow ? decodeJsonColumn(integrationRow.config) : undefined,
6553+
config: decodeJsonColumn(integrationRow.config),
65056554
...(grantedScopes ? { grantedScopes } : {}),
65066555
};
65076556
return wrapInvocationError(
@@ -6609,6 +6658,7 @@ export const createExecutor = <const TPlugins extends readonly AnyPlugin[] = rea
66096658
guardOrgWrite: (owner: Owner) => guardOrgWrite(owner),
66106659
defaultWritableProvider,
66116660
mintOAuthConnection: (input: MintOAuthConnectionInput) => mintOAuthConnection(input),
6661+
integrationExists: (slug) => findIntegrationRow(slug).pipe(Effect.map((row) => row !== null)),
66126662
connectionNameTaken: (ref) => findConnectionRow(ref).pipe(Effect.map((row) => row !== null)),
66136663
// One integration-row read + one projector run. Resolve the method this
66146664
// template selects exactly as the runtime's `selectAuthMethod` does —

‎packages/core/sdk/src/fuma-runtime.ts‎

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import { Cause, Context, Data, Effect, Exit, Layer, Predicate } from "effect";
2-
import type { AbstractQuery } from "@executor-js/fumadb/query";
2+
import { withQueryContext, type AbstractQuery } from "@executor-js/fumadb/query";
33
import type { AnySchema, AnyTable, Schema as FumaSchema } from "@executor-js/fumadb/schema";
44

55
export class StorageError extends Data.TaggedError("StorageError")<{
@@ -306,6 +306,12 @@ export type IFumaClient<TSchema extends AnySchema = AnySchema> = Readonly<{
306306

307307
export interface MakeFumaClientOptions {
308308
readonly tables?: ReadonlySet<string>;
309+
/** Owner-policy context to rebind EVERY query to, including queries issued
310+
* inside an enclosing transaction (whose handle otherwise carries the
311+
* context of whoever opened it). Lets a narrowly-scoped client (the
312+
* integration-removal cascade) join a bound transaction without inheriting
313+
* the bound reach. */
314+
readonly context?: unknown;
309315
}
310316

311317
const isAllowedTable = (tables: ReadonlySet<string> | undefined, table: PropertyKey): boolean =>
@@ -347,9 +353,11 @@ const makeSafeFumaQuery = <TSchema extends AnySchema>(
347353
};
348354

349355
export const makeFumaClient = (db: FumaDb, options: MakeFumaClientOptions = {}): IFumaClient => {
356+
const rebind = (handle: FumaDb): FumaDb =>
357+
options.context === undefined ? handle : withQueryContext(handle, options.context);
350358
const use: IFumaClient["use"] = (label, fn) =>
351359
Effect.flatMap(Effect.service(activeFumaDbRef), (active) =>
352-
fumaEffect(label, () => fn(makeSafeFumaQuery(active ?? db, options))),
360+
fumaEffect(label, () => fn(makeSafeFumaQuery(rebind(active ?? db), options))),
353361
).pipe(Effect.withSpan(`fumadb.${label}`));
354362

355363
const transaction = <A, E>(effect: Effect.Effect<A, E>): Effect.Effect<A, E | StorageFailure> =>

0 commit comments

Comments
 (0)