diff --git a/.changeset/clean-mcp-streams.md b/.changeset/clean-mcp-streams.md new file mode 100644 index 0000000..ede3cae --- /dev/null +++ b/.changeset/clean-mcp-streams.md @@ -0,0 +1,5 @@ +--- +'@inflowpayai/inflow': patch +--- + +Report errors returned by streaming MCP commands and close interrupted command streams. diff --git a/.changeset/safe-secret-lifecycle-recovery.md b/.changeset/safe-secret-lifecycle-recovery.md new file mode 100644 index 0000000..710cc36 --- /dev/null +++ b/.changeset/safe-secret-lifecycle-recovery.md @@ -0,0 +1,6 @@ +--- +'@inflowpayai/inflow-core': patch +'@inflowpayai/inflow': patch +--- + +Preserve recoverable credential lifecycle state when secure storage deletion is temporarily unavailable. diff --git a/packages/cli/test/unit/mcp-stream-errors.test.ts b/packages/cli/test/unit/mcp-stream-errors.test.ts new file mode 100644 index 0000000..497ccb4 --- /dev/null +++ b/packages/cli/test/unit/mcp-stream-errors.test.ts @@ -0,0 +1,75 @@ +import { Mcp } from 'incur'; +import { describe, expect, it, vi } from 'vitest'; + +function streamingTool(run: (context: { error: (input: { code: string; message: string }) => never }) => unknown) { + return { + name: 'streaming_test', + description: 'Streaming MCP test command', + inputSchema: { type: 'object' as const, properties: {} }, + command: { run }, + }; +} + +describe('MCP streaming command errors', () => { + it('surfaces a generator return c.error as an MCP tool error', async () => { + const result = await Mcp.callTool( + streamingTool(async function* (context) { + await Promise.resolve(); + yield { phase: 'started' }; + return context.error({ code: 'VAULT_LOCKED', message: 'The InFlow vault is locked.' }); + }), + {}, + ); + + expect(result).toMatchObject({ + content: [{ type: 'text', text: 'The InFlow vault is locked.' }], + isError: true, + }); + }); + + it('continues to buffer successful generator chunks', async () => { + const result = await Mcp.callTool( + streamingTool(async function* () { + await Promise.resolve(); + yield { phase: 'started' }; + yield { phase: 'complete' }; + }), + {}, + ); + + expect(result).toEqual({ + content: [ + { + type: 'text', + text: '[{"phase":"started"},{"phase":"complete"}]', + }, + ], + }); + }); + + it('closes the generator when a progress notification fails', async () => { + let finalized = false; + const result = await Mcp.callTool( + streamingTool(async function* () { + try { + await Promise.resolve(); + yield { phase: 'started' }; + yield { phase: 'complete' }; + } finally { + finalized = true; + } + }), + {}, + { + extra: { mcpReq: { _meta: { progressToken: 'progress-1' } } }, + sendNotification: vi.fn().mockRejectedValue(new Error('progress delivery failed')), + }, + ); + + expect(result).toMatchObject({ + content: [{ type: 'text', text: 'progress delivery failed' }], + isError: true, + }); + expect(finalized).toBe(true); + }); +}); diff --git a/packages/core/src/secure-storage/lifecycle.ts b/packages/core/src/secure-storage/lifecycle.ts index 4648450..78ffe76 100644 --- a/packages/core/src/secure-storage/lifecycle.ts +++ b/packages/core/src/secure-storage/lifecycle.ts @@ -1,3 +1,4 @@ +import { SecureStorageError } from './errors.js'; import type { SecretReference, SecretReferenceManifest, @@ -5,7 +6,6 @@ import type { SyncSecretReferenceManifest, SyncSecureSecretStore, } from './secret-store.js'; -import { SecureStorageError } from './errors.js'; import type { SecureSqliteRepository } from './sqlite.js'; function errorFromCause(cause: unknown): Error { diff --git a/packages/core/test/unit/secure-storage/lifecycle.test.ts b/packages/core/test/unit/secure-storage/lifecycle.test.ts index e648d39..3b485cc 100644 --- a/packages/core/test/unit/secure-storage/lifecycle.test.ts +++ b/packages/core/test/unit/secure-storage/lifecycle.test.ts @@ -139,6 +139,17 @@ describe('SecureSecretLifecycleCoordinator', () => { expect(repository.listSecretLifecycle('deleting')).toEqual([]); }); + it('completes a retried asynchronous delete when the secret is already absent', async () => { + const reference = { purpose: 'api-key', reference: 'missing-delete' }; + await coordinator.create(reference, Buffer.from('secret'), null); + await store.delete(reference); + + await coordinator.delete(reference); + + expect(repository.listSecretLifecycle('deleting')).toEqual([]); + expect(await manifest.read()).toEqual([]); + }); + it('preserves interrupted work when the asynchronous secret store is temporarily unavailable', async () => { const pending = { purpose: 'api-key', reference: 'locked-pending' }; repository.beginSecretLifecycle(pending, null); @@ -152,7 +163,6 @@ describe('SecureSecretLifecycleCoordinator', () => { await expect(lockedCoordinator.recoverInterruptedWork()).rejects.toMatchObject({ secureStorageCode: 'vault_locked', }); - expect(repository.listSecretLifecycle('pending')).toEqual([pending]); }); @@ -225,6 +235,20 @@ describe('SecureSecretLifecycleCoordinator', () => { expect(repository.listSecretLifecycle('deleting')).toEqual([]); }); + it('completes a retried synchronous delete when the secret is already absent', () => { + const syncStore = new SyncMemorySecretStore(); + const syncManifest = new SyncSecretReferenceManifestStore(syncStore); + const syncCoordinator = new SyncSecureSecretLifecycleCoordinator(repository, syncStore, syncManifest); + const reference = { purpose: 'api-key', reference: 'sync-missing-delete' }; + syncCoordinator.create(reference, Buffer.from('secret'), null); + syncStore.delete(reference); + + syncCoordinator.delete(reference); + + expect(repository.listSecretLifecycle('deleting')).toEqual([]); + expect(syncManifest.read()).toEqual([]); + }); + it('preserves interrupted work when the synchronous secret store is temporarily unavailable', () => { const pending = { purpose: 'api-key', reference: 'sync-locked-pending' }; repository.beginSecretLifecycle(pending, null); diff --git a/patches/incur.patch b/patches/incur.patch index d1e43fe..e803b47 100644 --- a/patches/incur.patch +++ b/patches/incur.patch @@ -85,7 +85,55 @@ : []), --- a/dist/Mcp.js +++ b/dist/Mcp.js -@@ -130,6 +130,7 @@ +@@ -69,18 +69,46 @@ + const progressToken = options.extra?.mcpReq?._meta?.progressToken; + let i = 0; ++ const iterator = result.stream[Symbol.asyncIterator](); ++ let completed = false; + try { +- for await (const chunk of result.stream) { ++ while (true) { ++ const item = await iterator.next(); ++ if (item.done) { ++ completed = true; ++ const returned = item.value; ++ if (returned !== null && ++ typeof returned === 'object' && ++ returned[Symbol.for('incur.sentinel')] === 'error') { ++ const cta = formatCtaBlock(options.name ?? tool.name, returned.cta); ++ const text = returned.message ?? 'Command failed'; ++ return { ++ content: [{ type: 'text', text: cta ? `${text}\n\n${renderCtaText(cta)}` : text }], ++ ...(cta ? { _meta: { cta } } : undefined), ++ isError: true, ++ }; ++ } ++ break; ++ } ++ const chunk = item.value; + chunks.push(chunk); + if (progressToken !== undefined && options.sendNotification) + await options.sendNotification({ + method: 'notifications/progress', + params: { progressToken, progress: ++i, message: Json.stringify(chunk) }, + }); + } + } + catch (err) { + return { + content: [{ type: 'text', text: err instanceof Error ? err.message : String(err) }], + isError: true, + }; + } ++ finally { ++ if (!completed && iterator.return) { ++ try { ++ await iterator.return(); ++ } ++ catch { } ++ } ++ } +@@ -130,6 +158,7 @@ const hasInput = Object.keys(mergedShape).length > 0; server.registerTool(tool.name, { ...(tool.description ? { description: tool.description } : undefined), @@ -93,7 +141,7 @@ ...(hasInput ? { inputSchema: z.object(mergedShape) } : undefined), ...(tool.outputSchema ? { outputSchema: options.fromJsonSchema(tool.outputSchema) } -@@ -159,6 +160,7 @@ +@@ -159,6 +188,7 @@ return toolResult({ tools: page.map((tool) => ({ name: tool.name, @@ -101,7 +149,7 @@ ...(tool.description ? { description: tool.description } : undefined), ...(tool.annotations ? { annotations: tool.annotations } : undefined), })), -@@ -177,6 +179,7 @@ +@@ -177,6 +207,7 @@ return toolError(`Unknown tool: ${params.name}`); return toolResult({ name: tool.name, @@ -109,7 +157,7 @@ ...(tool.description ? { description: tool.description } : undefined), inputSchema: tool.inputSchema, ...(tool.outputSchema ? { outputSchema: tool.outputSchema } : undefined), -@@ -284,7 +287,8 @@ +@@ -284,7 +315,8 @@ export function collectTools(commands, prefix, parentMiddlewares = [], filter) { const tools = filterTools(collectToolEntries(commands, prefix, parentMiddlewares), filter); assertUniqueToolNames(tools); @@ -119,7 +167,7 @@ } function collectToolEntries(commands, prefix, parentMiddlewares = []) { const result = []; -@@ -306,6 +310,7 @@ +@@ -306,6 +338,7 @@ const outputSchema = entry.output ? mcpOutputSchema(entry.output) : undefined; result.push({ name: mcp?.name ?? path.join('_'), diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 3c15665..f6b39cd 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -5,7 +5,7 @@ settings: excludeLinksFromLockfile: false patchedDependencies: - incur: 811c32d70903670aacd7e84800e2dac041ec445a5b47f9e4c34416077c2c793b + incur: aeae32ab87b3e46a5c59987af280484b29ad31f08a42eb18234f04f432e71dde importers: @@ -103,7 +103,7 @@ importers: version: 2.19.0 incur: specifier: ^0.4.5 - version: 0.4.19(patch_hash=811c32d70903670aacd7e84800e2dac041ec445a5b47f9e4c34416077c2c793b) + version: 0.4.19(patch_hash=aeae32ab87b3e46a5c59987af280484b29ad31f08a42eb18234f04f432e71dde) ink: specifier: ^5.2.1 version: 5.2.1(@types/react@18.3.31)(react@18.3.1) @@ -4826,7 +4826,7 @@ snapshots: imurmurhash@0.1.4: {} - incur@0.4.19(patch_hash=811c32d70903670aacd7e84800e2dac041ec445a5b47f9e4c34416077c2c793b): + incur@0.4.19(patch_hash=aeae32ab87b3e46a5c59987af280484b29ad31f08a42eb18234f04f432e71dde): dependencies: '@cfworker/json-schema': 4.1.1 '@modelcontextprotocol/server': 2.0.0-alpha.4 @@ -5240,7 +5240,7 @@ snapshots: dependencies: '@stripe/stripe-js': 9.9.0 eventsource-parser: 3.1.0 - incur: 0.4.19(patch_hash=811c32d70903670aacd7e84800e2dac041ec445a5b47f9e4c34416077c2c793b) + incur: 0.4.19(patch_hash=aeae32ab87b3e46a5c59987af280484b29ad31f08a42eb18234f04f432e71dde) ox: 0.14.33(typescript@5.9.3)(zod@4.4.3) viem: 2.55.10(typescript@5.9.3)(zod@4.4.3) zod: 4.4.3