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);