From 0423820bb68122a5e6d805a4b7ec9001a20ba460 Mon Sep 17 00:00:00 2001 From: fforres Date: Sat, 26 Sep 2026 13:25:53 -0700 Subject: [PATCH] fix(openapi): refreshing one connection no longer loads every operation in the org listOperations read every stored OpenAPI operation row for all integrations and schema-decoded each binding just to read its integration. With a large catalog such as the Cloudflare API (about 2,700 operations) that allocated over 200 MB per refresh, and the second refresh in a session isolate crossed the Durable Object memory limit, surfacing as "Execution lost: the session was reset". Plugin storage list now pushes keyPrefix into SQL, and the store lists only the current and legacy key prefixes for the integration, decoding bindings of matching rows once. Claude-Session: https://claude.ai/code/session_01VMqJkxznzTQaFVHHcttxpJ --- packages/core/sdk/src/executor.ts | 13 +- .../plugins/openapi/src/sdk/store.test.ts | 124 ++++++++++++++++++ packages/plugins/openapi/src/sdk/store.ts | 22 +++- 3 files changed, 154 insertions(+), 5 deletions(-) diff --git a/packages/core/sdk/src/executor.ts b/packages/core/sdk/src/executor.ts index a7249e627e..eb27d6e5d3 100644 --- a/packages/core/sdk/src/executor.ts +++ b/packages/core/sdk/src/executor.ts @@ -1606,6 +1606,15 @@ const makePluginStorageFacade = (input: { key === undefined ? true : b("key", "=", key), ); + const whereForPrefix = + (collection: string, keyPrefix: string | undefined): CoreWhere => + (b: AnyCb) => + b.and( + b("plugin_id", "=", input.pluginId), + b("collection", "=", collection), + keyPrefix === undefined ? true : b("key", "starts with", keyPrefix), + ); + const whereOwner = (owner: Owner, collection: string, key: string): CoreWhere => { const os = ownerSubject(owner); return (b: AnyCb) => @@ -1813,7 +1822,7 @@ const makePluginStorageFacade = (input: { if (validationError) return yield* validationError; const rows = yield* input.core.findMany("plugin_storage", { - where: whereFor(definition.name), + where: whereForPrefix(definition.name, queryInput?.keyPrefix), }); const filtered = sortByOwnerPrecedence(rows) .filter((row) => @@ -1889,7 +1898,7 @@ const makePluginStorageFacade = (input: { list: (storageInput) => Effect.gen(function* () { const rows = yield* input.core.findMany("plugin_storage", { - where: whereFor(storageInput.collection), + where: whereForPrefix(storageInput.collection, storageInput.keyPrefix), }); return sortByOwnerPrecedence(rows) .filter((row) => diff --git a/packages/plugins/openapi/src/sdk/store.test.ts b/packages/plugins/openapi/src/sdk/store.test.ts index ef4a4ed2cf..a9177fa78c 100644 --- a/packages/plugins/openapi/src/sdk/store.test.ts +++ b/packages/plugins/openapi/src/sdk/store.test.ts @@ -13,6 +13,103 @@ import { import { makeDefaultOpenapiStore } from "./store"; import { OperationBinding } from "./types"; +function makeRecordingStore() { + const rows = new Map(); + const listPrefixes: Array = []; + const now = new Date(); + const storageKey = (collection: string, key: string) => `${collection}\0${key}`; + const seed = (owner: "org" | "user", collection: string, key: string, data: unknown) => { + const entry: PluginStorageEntry = { + id: storageKey(collection, key), + owner, + pluginId: "openapi", + collection, + key, + data, + createdAt: now, + updatedAt: now, + }; + rows.set(entry.id, entry); + return entry; + }; + const unusedCollection: ReturnType = { + get: () => Effect.succeed(null), + getForOwner: () => Effect.succeed(null), + list: () => Effect.succeed([]), + put: (input) => Effect.succeed(seed(input.owner, "unused", input.key, input.data)), + query: () => Effect.succeed([]), + count: () => Effect.succeed(0), + remove: () => Effect.void, + }; + const pluginStorage: PluginStorageFacade = { + collection: () => unusedCollection, + get: (input: { readonly collection: string; readonly key: string }) => + Effect.succeed( + (rows.get(storageKey(input.collection, input.key)) as PluginStorageEntry | undefined) ?? + null, + ), + getForOwner: (input: { readonly collection: string; readonly key: string }) => + Effect.succeed( + (rows.get(storageKey(input.collection, input.key)) as PluginStorageEntry | undefined) ?? + null, + ), + list: (input: { readonly collection: string; readonly keyPrefix?: string }) => + Effect.sync(() => { + listPrefixes.push(input.keyPrefix); + return [...rows.values()].filter( + (row) => + row.collection === input.collection && + (input.keyPrefix === undefined || row.key.startsWith(input.keyPrefix)), + ) as PluginStorageEntry[]; + }), + put: (input: { + readonly owner: "org" | "user"; + readonly collection: string; + readonly key: string; + readonly data: unknown; + }) => + Effect.succeed( + seed(input.owner, input.collection, input.key, input.data) as PluginStorageEntry, + ), + putMany: (input) => + Effect.sync(() => { + for (const entry of input.entries) + seed(input.owner, entry.collection, entry.key, entry.data); + }), + remove: (input) => + Effect.sync(() => { + rows.delete(storageKey(input.collection, input.key)); + }), + removeMany: (input) => + Effect.sync(() => { + for (const entry of input.entries) rows.delete(storageKey(entry.collection, entry.key)); + }), + }; + const blobs: PluginBlobStore = { + get: () => Effect.succeed(null), + put: () => Effect.void, + delete: () => Effect.void, + has: () => Effect.succeed(false), + }; + const store = makeDefaultOpenapiStore({ + owner: { tenant: Tenant.make("tenant"), subject: Subject.make("subject") }, + blobs, + pluginStorage, + } satisfies StorageDeps); + return { store, seed, listPrefixes }; +} + +function getBinding(pathTemplate: string) { + return OperationBinding.make({ + method: "get", + servers: [], + pathTemplate, + parameters: [], + requestBody: Option.none(), + responseBody: Option.none(), + }); +} + describe("OpenAPI operation store", () => { it.effect("bounds operation storage keys while preserving tool-name lookup", () => Effect.gen(function* () { @@ -146,4 +243,31 @@ describe("OpenAPI operation store", () => { expect(operation?.binding.pathTemplate).toBe("/users/{userId}/messages"); }), ); + + it.effect("lists one integration's operations without reading other integrations' rows", () => + Effect.gen(function* () { + const { store, seed, listPrefixes } = makeRecordingStore(); + yield* store.putOperations("google_sheets", [ + { integration: "google_sheets", toolName: "values.get", binding: getBinding("/values") }, + ]); + seed("org", "operation", "google_sheets.sheets.copyTo", { + integration: "google_sheets", + toolName: "sheets.copyTo", + binding: { method: "post", servers: [], pathTemplate: "/copyTo", parameters: [] }, + }); + seed("org", "operation", "op.someotherhash.tool", { + integration: "cloudflare_api", + toolName: "zones.list", + binding: { notABinding: true }, + }); + + const operations = yield* store.listOperations("google_sheets"); + + expect(operations.map((operation) => operation.toolName).sort()).toEqual([ + "sheets.copyTo", + "values.get", + ]); + expect(listPrefixes).not.toContain(undefined); + }), + ); }); diff --git a/packages/plugins/openapi/src/sdk/store.ts b/packages/plugins/openapi/src/sdk/store.ts index 548a521031..8028c997e4 100644 --- a/packages/plugins/openapi/src/sdk/store.ts +++ b/packages/plugins/openapi/src/sdk/store.ts @@ -147,15 +147,31 @@ export const makeDefaultOpenapiStore = ({ pluginStorage, blobs }: StorageDeps): ...(operation.description !== undefined ? { description: operation.description } : {}), }); - const listRows = (integration: string) => + const belongsTo = (row: PluginStorageEntry, integration: string): boolean => { + const decoded = decodeOperationStorage(row.data); + return Option.isSome(decoded) && decoded.value.integration === integration; + }; + + const listRowsWithPrefix = (integration: string, keyPrefix: string) => pluginStorage - .list({ collection: OPERATION_COLLECTION }) + .list({ collection: OPERATION_COLLECTION, keyPrefix }) .pipe( Effect.map((rows: readonly PluginStorageEntry[]) => - rows.filter((row) => rowToOperation(row)?.integration === integration), + rows.filter((row) => belongsTo(row, integration)), ), ); + const listRows = (integration: string) => + Effect.all([ + listRowsWithPrefix(integration, `${OPERATION_KEY_VERSION}.${stableKeyHash(integration)}.`), + listRowsWithPrefix(integration, `${integration}.`), + ]).pipe( + Effect.map(([current, legacy]) => { + const seen = new Set(current.map((row) => row.key)); + return [...current, ...legacy.filter((row) => !seen.has(row.key))]; + }), + ); + const removeOperations = (integration: string) => Effect.gen(function* () { const rows = yield* listRows(integration);