From bd2e7298c131594dd040915ad30a415208084618 Mon Sep 17 00:00:00 2001 From: Ability Date: Thu, 25 Jun 2026 20:48:28 +0100 Subject: [PATCH 1/5] Add integration coverage for missing deposit fields, empty idempotency key, and auth challenge rate limits --- tests/mvp-express.integration.test.ts | 88 +++++++++++++++++++++++++++ 1 file changed, 88 insertions(+) diff --git a/tests/mvp-express.integration.test.ts b/tests/mvp-express.integration.test.ts index cc0f2f8c..a2e90fef 100644 --- a/tests/mvp-express.integration.test.ts +++ b/tests/mvp-express.integration.test.ts @@ -389,6 +389,32 @@ describe('MVP Express-mounted integration', () => { } }); + it('3a) auth challenge route returns 429 when authChallengeMax is exceeded', async () => { + const account = Keypair.random().publicKey(); + const headers = { 'x-forwarded-for': '10.0.0.99' }; + + const firstResponse = await invoke({ + path: `/auth/challenge?account=${account}`, + headers, + }); + expect(firstResponse.status).toBe(200); + + const secondResponse = await invoke({ + path: `/auth/challenge?account=${account}`, + headers, + }); + expect(secondResponse.status).toBe(200); + + const thirdResponse = await invoke({ + path: `/auth/challenge?account=${account}`, + headers, + }); + + expect(thirdResponse.status).toBe(429); + expect(thirdResponse.body.error).toBe('rate_limited'); + expect(thirdResponse.headers['retry-after']).toBeDefined(); + }); + it('4) unauthorized deposit interactive rejected', async () => { const response = await invoke({ method: 'POST', @@ -449,6 +475,38 @@ describe('MVP Express-mounted integration', () => { expect(response.body.id).toBeUndefined(); }); + it('5f) deposit missing asset_code returns invalid_request', async () => { + const response = await invoke({ + method: 'POST', + path: '/transactions/deposit/interactive', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${accessToken}`, + }, + body: { amount: '10' }, + }); + + expect(response.status).toBe(400); + expect(response.body.error).toBe('invalid_request'); + expect(response.body.message).toContain('asset_code and amount'); + }); + + it('5g) deposit missing amount returns invalid_request', async () => { + const response = await invoke({ + method: 'POST', + path: '/transactions/deposit/interactive', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${accessToken}`, + }, + body: { asset_code: 'USDC' }, + }); + + expect(response.status).toBe(400); + expect(response.body.error).toBe('invalid_request'); + expect(response.body.message).toContain('asset_code and amount'); + }); + it('5e) deposit with deposits_enabled: false asset is rejected', async () => { const disabledDbUrl = makeSqliteDbUrlForTests(); const disabledAnchor = createAnchor({ @@ -592,6 +650,36 @@ describe('MVP Express-mounted integration', () => { expect(response.body.idempotency_replay).toBe(true); }); + it('6d) empty Idempotency-Key header is treated as no key and creates a new deposit', async () => { + const firstResponse = await invoke({ + method: 'POST', + path: '/transactions/deposit/interactive', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${accessToken}`, + 'idempotency-key': ' ', + }, + body: { asset_code: 'USDC', amount: '12' }, + }); + + const secondResponse = await invoke({ + method: 'POST', + path: '/transactions/deposit/interactive', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${accessToken}`, + 'idempotency-key': ' ', + }, + body: { asset_code: 'USDC', amount: '12' }, + }); + + expect(firstResponse.status).toBe(201); + expect(secondResponse.status).toBe(201); + expect(firstResponse.body.id).not.toBe(secondResponse.body.id); + expect(firstResponse.body.idempotency_replay).toBeUndefined(); + expect(secondResponse.body.idempotency_replay).toBeUndefined(); + }); + it('7) transaction lookup fetches persisted data', async () => { const response = await invoke({ method: 'GET', From 79e3eeef5933a55832be305d8358db4c8c99d057 Mon Sep 17 00:00:00 2001 From: Ability Date: Tue, 28 Jul 2026 16:08:38 +0100 Subject: [PATCH 2/5] Normalize webhook provider headers --- PR_DESCRIPTION.md | 23 +++++++ src/runtime/http/express-router.ts | 25 ++++++-- tests/mvp-express.integration.test.ts | 56 ++++++++++++++++- tests/webhook-fallback.test.ts | 88 ++++++++++++++++++++++++++- 4 files changed, 184 insertions(+), 8 deletions(-) create mode 100644 PR_DESCRIPTION.md diff --git a/PR_DESCRIPTION.md b/PR_DESCRIPTION.md new file mode 100644 index 00000000..3a72e02c --- /dev/null +++ b/PR_DESCRIPTION.md @@ -0,0 +1,23 @@ +# Normalize webhook provider values and add webhook rate-limit coverage + +## What does this PR do? +- Trims whitespace-only webhook provider values before evaluating fallback logic. +- Falls back from a blank header to the body provider, then to `generic`. +- Accepts array-style `x-webhook-provider` and `x-anchor-signature` headers by reading the first non-empty value. +- Adds focused regression coverage for webhook provider fallback behavior and webhook route rate limiting. + +## How to test? +- Run `bun test tests/webhook-fallback.test.ts tests/mvp-express.integration.test.ts` +- Confirm the new webhook fallback and webhook rate-limit cases pass. + +## Checklist +- [ ] My code follows the code style of this project. +- [x] I have added tests for my changes. +- [ ] I have updated the documentation accordingly. +- [x] I have run `bun test` locally. + +## Issue Reference +Closes #355 +Closes #356 +Closes #359 +Closes #357 diff --git a/src/runtime/http/express-router.ts b/src/runtime/http/express-router.ts index 6b1b3dab..f1d56a0c 100644 --- a/src/runtime/http/express-router.ts +++ b/src/runtime/http/express-router.ts @@ -134,6 +134,23 @@ function extractClientIdentifier(req: IncomingMessage): string { return leftMost || socketIp || 'unknown'; } +function extractFirstHeaderValue(value: string | string[] | undefined): string | undefined { + if (Array.isArray(value)) { + for (const item of value) { + if (typeof item === 'string' && item.trim().length > 0) { + return item; + } + } + return undefined; + } + + if (typeof value === 'string') { + return value.trim().length > 0 ? value : undefined; + } + + return undefined; +} + function hasValidSignature(transaction: Transaction, publicKey: string): boolean { const keypair = Keypair.fromPublicKey(publicKey); const hash = transaction.hash(); @@ -617,15 +634,15 @@ export class AnchorExpressRouter { const eventIdField = payload.id; const eventId = typeof eventIdField === 'string' && eventIdField.length > 0 ? eventIdField : randomUUID(); - const providerHeader = req.headers['x-webhook-provider']; + const providerHeader = extractFirstHeaderValue(req.headers['x-webhook-provider']); const providerBody = payload.provider; const provider = - typeof providerHeader === 'string' && providerHeader.length > 0 + typeof providerHeader === 'string' && providerHeader.trim().length > 0 ? providerHeader - : typeof providerBody === 'string' && providerBody.length > 0 + : typeof providerBody === 'string' && providerBody.trim().length > 0 ? providerBody : 'generic'; - const signatureHeader = req.headers['x-anchor-signature']; + const signatureHeader = extractFirstHeaderValue(req.headers['x-anchor-signature']); const signature = typeof signatureHeader === 'string' ? signatureHeader : undefined; try { diff --git a/tests/mvp-express.integration.test.ts b/tests/mvp-express.integration.test.ts index a2e90fef..413e68c2 100644 --- a/tests/mvp-express.integration.test.ts +++ b/tests/mvp-express.integration.test.ts @@ -17,7 +17,7 @@ interface TestResponse { interface TestRequestOptions { method?: string; path: string; - headers?: Record; + headers?: Record; body?: Record; } @@ -30,7 +30,7 @@ function createMountedInvoker(anchor: AnchorInstance) { const req = Readable.from(serializedBody ? [serializedBody] : []) as IncomingMessage & { method: string; url: string; - headers: Record; + headers: Record; body?: Record; }; @@ -861,6 +861,58 @@ describe('MVP Express-mounted integration', () => { expect(response.body.provider).toBe('generic'); // Should default to 'generic' }); + it('8c) webhook route returns 429 after webhookMax is exceeded', async () => { + const headers = { 'content-type': 'application/json', 'x-forwarded-for': '10.0.0.200' }; + + for (let index = 0; index < 21; index += 1) { + const payload = { + id: `evt_rate_limit_${index}`, + type: 'deposit.completed', + transaction_id: transactionId, + }; + + const signature = createHmac('sha256', 'webhook-test-secret') + .update(JSON.stringify(payload)) + .digest('hex'); + + const response = await invoke({ + method: 'POST', + path: '/webhooks/events', + headers: { + ...headers, + 'x-webhook-provider': 'generic', + 'x-anchor-signature': signature, + }, + body: payload, + }); + + if (index < 20) { + expect(response.status).toBe(200); + } + } + + const rateLimitedResponse = await invoke({ + method: 'POST', + path: '/webhooks/events', + headers: { + ...headers, + 'x-webhook-provider': 'generic', + 'x-anchor-signature': createHmac('sha256', 'webhook-test-secret') + .update(JSON.stringify({ id: 'evt_rate_limit_21', type: 'deposit.completed', transaction_id: transactionId })) + .digest('hex'), + }, + body: { + id: 'evt_rate_limit_21', + type: 'deposit.completed', + transaction_id: transactionId, + }, + }); + + expect(rateLimitedResponse.status).toBe(429); + expect(rateLimitedResponse.body.error).toBe('rate_limited'); + expect(rateLimitedResponse.headers['retry-after']).toBeDefined(); + }); + it('8d) webhook without id field returns a generated event_id', async () => { const payload = { type: 'deposit.completed', diff --git a/tests/webhook-fallback.test.ts b/tests/webhook-fallback.test.ts index 2f8aca94..5c765d89 100644 --- a/tests/webhook-fallback.test.ts +++ b/tests/webhook-fallback.test.ts @@ -16,7 +16,7 @@ interface TestResponse { interface TestRequestOptions { method?: string; path: string; - headers?: Record; + headers?: Record; body?: Record; } @@ -29,7 +29,7 @@ function createMountedInvoker(anchor: AnchorInstance) { const req = Readable.from(serializedBody ? [serializedBody] : []) as IncomingMessage & { method: string; url: string; - headers: Record; + headers: Record; body?: Record; }; @@ -240,4 +240,88 @@ describe('Webhook Provider Fallback', () => { expect(response.status).toBe(200); expect(lastProvider).toBe('body-provider'); }); + + it('Whitespace-only header falls back to body provider', async () => { + const payload = { id: 'evt_whitespace_header', type: 'test', provider: 'body-provider' }; + const signature = createHmac('sha256', 'webhook-test-secret') + .update(JSON.stringify(payload)) + .digest('hex'); + + const response = await invoke({ + method: 'POST', + path: '/webhooks/events', + headers: { + 'content-type': 'application/json', + 'x-webhook-provider': ' ', + 'x-anchor-signature': signature, + }, + body: payload, + }); + + expect(response.status).toBe(200); + expect(lastProvider).toBe('body-provider'); + }); + + it('Whitespace-only header and body values fall back to generic', async () => { + const payload = { id: 'evt_blank_both', type: 'test', provider: ' ' }; + const signature = createHmac('sha256', 'webhook-test-secret') + .update(JSON.stringify(payload)) + .digest('hex'); + + const response = await invoke({ + method: 'POST', + path: '/webhooks/events', + headers: { + 'content-type': 'application/json', + 'x-webhook-provider': ' ', + 'x-anchor-signature': signature, + }, + body: payload, + }); + + expect(response.status).toBe(200); + expect(lastProvider).toBe('generic'); + }); + + it('Array-style provider headers use the first non-empty value', async () => { + const payload = { id: 'evt_array_provider', type: 'test' }; + const signature = createHmac('sha256', 'webhook-test-secret') + .update(JSON.stringify(payload)) + .digest('hex'); + + const response = await invoke({ + method: 'POST', + path: '/webhooks/events', + headers: { + 'content-type': 'application/json', + 'x-webhook-provider': ['', ' ', 'array-provider'], + 'x-anchor-signature': signature, + }, + body: payload, + }); + + expect(response.status).toBe(200); + expect(lastProvider).toBe('array-provider'); + }); + + it('Array-style signature headers are accepted when they contain a value', async () => { + const payload = { id: 'evt_array_signature', type: 'test' }; + const signature = createHmac('sha256', 'webhook-test-secret') + .update(JSON.stringify(payload)) + .digest('hex'); + + const response = await invoke({ + method: 'POST', + path: '/webhooks/events', + headers: { + 'content-type': 'application/json', + 'x-webhook-provider': 'generic', + 'x-anchor-signature': [signature], + }, + body: payload, + }); + + expect(response.status).toBe(200); + expect(lastProvider).toBe('generic'); + }); }); From 927a752c29aafe197206269beb8741c863da03a6 Mon Sep 17 00:00:00 2001 From: Shayam Prasad Sah Date: Sat, 29 Aug 2026 18:19:06 +0100 Subject: [PATCH 3/5] fix: deterministic pending ordering and adapter lifecycle safety --- src/runtime/database/sql-database-adapter.ts | 119 +++++++++--- src/runtime/interfaces.ts | 4 +- .../sql-adapter-interactive-tx.test.ts | 175 +++++++++++++++++- 3 files changed, 266 insertions(+), 32 deletions(-) diff --git a/src/runtime/database/sql-database-adapter.ts b/src/runtime/database/sql-database-adapter.ts index 4ff0a7af..28b17798 100644 --- a/src/runtime/database/sql-database-adapter.ts +++ b/src/runtime/database/sql-database-adapter.ts @@ -47,6 +47,8 @@ export class SqlDatabaseAdapter implements DatabaseAdapter { private readonly url: string; private sqlite: SqliteLike | null = null; private postgres: PostgresClient | null = null; + private connectPromise: Promise | null = null; + private disconnectPromise: Promise | null = null; constructor(databaseConfig: FrameworkConfig['database']) { this.provider = databaseConfig.provider; @@ -54,34 +56,84 @@ export class SqlDatabaseAdapter implements DatabaseAdapter { } public async connect(): Promise { - if (this.provider === 'sqlite') { - this.sqlite = new Database(toSqlitePath(this.url)); + if (this.sqlite || this.postgres) { return; } - if (this.provider === 'postgres') { - const moduleName = 'pg'; - const pgModuleUnknown: unknown = await import(moduleName); - const pgModule = pgModuleUnknown as { - Client: new (config: { connectionString: string }) => PostgresClient; - }; - this.postgres = new pgModule.Client({ connectionString: this.url }); - await this.postgres.connect(); + if (this.connectPromise) { + await this.connectPromise; return; } - throw new ConfigError(`Unsupported database provider: ${this.provider}`); + this.connectPromise = (async () => { + try { + if (this.provider === 'sqlite') { + this.sqlite = new Database(toSqlitePath(this.url)); + return; + } + + if (this.provider === 'postgres') { + const moduleName = 'pg'; + const pgModuleUnknown: unknown = await import(moduleName); + const pgModule = pgModuleUnknown as { + Client: new (config: { connectionString: string }) => PostgresClient; + }; + const client = new pgModule.Client({ connectionString: this.url }); + this.postgres = client; + await this.postgres.connect(); + return; + } + + throw new ConfigError(`Unsupported database provider: ${this.provider}`); + } catch (error) { + this.sqlite = null; + this.postgres = null; + throw error; + } + })(); + + try { + await this.connectPromise; + } finally { + this.connectPromise = null; + } } public async disconnect(): Promise { - if (this.sqlite) { - this.sqlite.close(); - this.sqlite = null; + if (!this.sqlite && !this.postgres) { + return; } - if (this.postgres) { - await this.postgres.end(); + if (this.disconnectPromise) { + await this.disconnectPromise; + return; + } + + this.disconnectPromise = (async () => { + const sqlite = this.sqlite; + const postgres = this.postgres; + this.sqlite = null; this.postgres = null; + + try { + if (sqlite) { + sqlite.close(); + } + + if (postgres) { + await postgres.end(); + } + } catch (error) { + this.sqlite = sqlite; + this.postgres = postgres; + throw error; + } + })(); + + try { + await this.disconnectPromise; + } finally { + this.disconnectPromise = null; } } @@ -263,19 +315,27 @@ export class SqlDatabaseAdapter implements DatabaseAdapter { }; } - public async markAuthChallengeConsumed(id: string): Promise { + public async markAuthChallengeConsumed(id: string): Promise { const consumedAt = nowIso(); if (this.sqlite) { - this.sqlite - .prepare('UPDATE auth_challenges SET consumed_at = ? WHERE id = ?') + const existing = this.sqlite + .prepare('SELECT consumed_at FROM auth_challenges WHERE id = ? LIMIT 1') + .get(id) as { consumed_at: string | null } | null; + if (!existing || existing.consumed_at) { + return false; + } + + const result = this.sqlite + .prepare('UPDATE auth_challenges SET consumed_at = ? WHERE id = ? AND consumed_at IS NULL') .run(consumedAt, id); - return; + return result.changes > 0; } - await this.requirePostgres().query( - 'UPDATE auth_challenges SET consumed_at = $1 WHERE id = $2', + const response = await this.requirePostgres().query<{ id: string }>( + 'UPDATE auth_challenges SET consumed_at = $1 WHERE id = $2 AND consumed_at IS NULL RETURNING id', [consumedAt, id], ); + return response.rows.length > 0; } public async insertInteractiveTransaction(input: { @@ -360,32 +420,33 @@ export class SqlDatabaseAdapter implements DatabaseAdapter { if (this.sqlite) { const rows = this.sqlite .prepare( - "SELECT * FROM interactive_transactions WHERE status = 'pending_user_transfer_start' AND created_at < ?", + "SELECT * FROM interactive_transactions WHERE status = 'pending_user_transfer_start' AND created_at < ? ORDER BY created_at ASC, id ASC", ) .all(cutoffIso) as Record[]; return rows.map((row) => this.mapTransactionRow(row)); } const response = await this.requirePostgres().query>( - "SELECT * FROM interactive_transactions WHERE status = 'pending_user_transfer_start' AND created_at < $1", + "SELECT * FROM interactive_transactions WHERE status = 'pending_user_transfer_start' AND created_at < $1 ORDER BY created_at ASC, id ASC", [cutoffIso], ); return response.rows.map((row) => this.mapTransactionRow(row)); } - public async updateTransactionStatus(id: string, status: string): Promise { + public async updateTransactionStatus(id: string, status: string): Promise { const updatedAt = nowIso(); if (this.sqlite) { - this.sqlite + const result = this.sqlite .prepare('UPDATE interactive_transactions SET status = ?, updated_at = ? WHERE id = ?') .run(status, updatedAt, id); - return; + return result.changes > 0; } - await this.requirePostgres().query( - 'UPDATE interactive_transactions SET status = $1, updated_at = $2 WHERE id = $3', + const response = await this.requirePostgres().query<{ id: string }>( + 'UPDATE interactive_transactions SET status = $1, updated_at = $2 WHERE id = $3 RETURNING id', [status, updatedAt, id], ); + return response.rows.length > 0; } public async getIdempotencyRecord( diff --git a/src/runtime/interfaces.ts b/src/runtime/interfaces.ts index 718e3ffc..a5aae048 100644 --- a/src/runtime/interfaces.ts +++ b/src/runtime/interfaces.ts @@ -61,7 +61,7 @@ export interface DatabaseAdapter { expiresAt: string; }): Promise; getAuthChallengeByChallenge(challenge: string): Promise; - markAuthChallengeConsumed(id: string): Promise; + markAuthChallengeConsumed(id: string): Promise; insertInteractiveTransaction(input: { id: string; @@ -73,7 +73,7 @@ export interface DatabaseAdapter { }): Promise; getInteractiveTransactionById(id: string): Promise; listPendingTransactionsBefore(cutoffIso: string): Promise; - updateTransactionStatus(id: string, status: string): Promise; + updateTransactionStatus(id: string, status: string): Promise; getIdempotencyRecord(scope: string, idempotencyKey: string): Promise; insertIdempotencyRecord(input: { diff --git a/tests/runtime/sql-adapter-interactive-tx.test.ts b/tests/runtime/sql-adapter-interactive-tx.test.ts index 48f0723b..64101e35 100644 --- a/tests/runtime/sql-adapter-interactive-tx.test.ts +++ b/tests/runtime/sql-adapter-interactive-tx.test.ts @@ -3,6 +3,7 @@ import { createSqlDatabaseAdapter } from '@/runtime/database/sql-database-adapte import { randomUUID } from 'node:crypto'; import { unlinkSync } from 'node:fs'; import { afterAll, beforeAll, describe, expect, it } from 'vitest'; +import { mock } from 'bun:test'; import type { DatabaseAdapter } from '@/runtime/interfaces.ts'; describe('SqlDatabaseAdapter – interactive transaction status updates', () => { @@ -55,7 +56,7 @@ describe('SqlDatabaseAdapter – interactive transaction status updates', () => expect(inserted.status).toBe('pending_user_transfer_start'); currentTime = new RealDate('2026-01-01T00:00:01.000Z').getTime(); - await db.updateTransactionStatus(txId, 'completed'); + await expect(db.updateTransactionStatus(txId, 'completed')).resolves.toBe(true); const fetched = await db.getInteractiveTransactionById(txId); expect(fetched).not.toBeNull(); @@ -65,4 +66,176 @@ describe('SqlDatabaseAdapter – interactive transaction status updates', () => globalThis.Date = RealDate; } }); + + it('orders pending transactions by created_at then id for tied timestamps', async () => { + const RealDate = Date; + let currentTime = new RealDate('2026-01-02T00:00:00.000Z').getTime(); + + class MockDate extends RealDate { + constructor(value?: string | number | Date) { + super(value === undefined ? currentTime : value); + } + + static override now(): number { + return currentTime; + } + } + + globalThis.Date = MockDate as DateConstructor; + + try { + await db.insertInteractiveTransaction({ + id: 'b-tx', + account: 'GTEST1234', + kind: 'deposit', + assetCode: 'USDC', + amount: '10.00', + status: 'pending_user_transfer_start', + }); + await db.insertInteractiveTransaction({ + id: 'a-tx', + account: 'GTEST1234', + kind: 'deposit', + assetCode: 'USDC', + amount: '20.00', + status: 'pending_user_transfer_start', + }); + + const pending = await db.listPendingTransactionsBefore('2026-01-03T00:00:00.000Z'); + expect(pending.map((tx) => tx.id)).toEqual(['a-tx', 'b-tx']); + } finally { + globalThis.Date = RealDate; + } + }); + + it('reports missing transaction IDs as failed updates and keeps existing timestamps stable', async () => { + const txId = randomUUID(); + const RealDate = Date; + let currentTime = new RealDate('2026-01-05T00:00:00.000Z').getTime(); + + class MockDate extends RealDate { + constructor(value?: string | number | Date) { + super(value === undefined ? currentTime : value); + } + + static override now(): number { + return currentTime; + } + } + + globalThis.Date = MockDate as DateConstructor; + + try { + const inserted = await db.insertInteractiveTransaction({ + id: txId, + account: 'GTEST1234', + kind: 'deposit', + assetCode: 'USDC', + amount: '75.00', + status: 'pending_user_transfer_start', + }); + + await expect(db.updateTransactionStatus('missing-id', 'completed')).resolves.toBe(false); + + currentTime = new RealDate('2026-01-05T00:00:01.000Z').getTime(); + await expect(db.updateTransactionStatus(txId, 'completed')).resolves.toBe(true); + + const fetched = await db.getInteractiveTransactionById(txId); + expect(fetched).not.toBeNull(); + expect(fetched!.status).toBe('completed'); + expect(fetched!.updatedAt).not.toBe(inserted.updatedAt); + } finally { + globalThis.Date = RealDate; + } + }); + + it('marks auth challenges consumed exactly once and leaves repeated consumption as a no-op', async () => { + const challengeId = randomUUID(); + const challenge = `challenge-${randomUUID()}`; + const RealDate = Date; + let currentTime = new RealDate('2026-01-06T00:00:00.000Z').getTime(); + + class MockDate extends RealDate { + constructor(value?: string | number | Date) { + super(value === undefined ? currentTime : value); + } + + static override now(): number { + return currentTime; + } + } + + globalThis.Date = MockDate as DateConstructor; + + try { + await db.insertAuthChallenge({ + id: challengeId, + account: 'GTEST1234', + challenge, + expiresAt: '2026-01-07T00:00:00.000Z', + }); + + await expect(db.markAuthChallengeConsumed(challengeId)).resolves.toBe(true); + + const firstState = await db.getAuthChallengeByChallenge(challenge); + expect(firstState).not.toBeNull(); + expect(firstState!.consumedAt).toBe('2026-01-06T00:00:00.000Z'); + + currentTime = new RealDate('2026-01-06T00:00:05.000Z').getTime(); + await expect(db.markAuthChallengeConsumed(challengeId)).resolves.toBe(false); + await expect(db.markAuthChallengeConsumed('missing-challenge-id')).resolves.toBe(false); + + const secondState = await db.getAuthChallengeByChallenge(challenge); + expect(secondState).not.toBeNull(); + expect(secondState!.consumedAt).toBe('2026-01-06T00:00:00.000Z'); + } finally { + globalThis.Date = RealDate; + } + }); + + it('serializes concurrent postgres connect attempts and retries after a failed connect', async () => { + const connectCalls: Array = []; + const endCalls: Array = []; + + mock.module('pg', () => ({ + Client: class { + public readonly connectionString: string; + + constructor(config: { connectionString: string }) { + this.connectionString = config.connectionString; + } + + async connect(): Promise { + connectCalls.push(this.connectionString); + if (connectCalls.length === 1) { + await new Promise((resolve) => setTimeout(resolve, 25)); + throw new Error('connect failed'); + } + await new Promise((resolve) => setTimeout(resolve, 25)); + } + + async end(): Promise { + endCalls.push(this.connectionString); + await new Promise((resolve) => setTimeout(resolve, 25)); + } + }, + })); + + try { + const { SqlDatabaseAdapter } = await import('@/runtime/database/sql-database-adapter.ts'); + const adapter = new SqlDatabaseAdapter({ provider: 'postgres', url: 'postgres://example' }); + + await expect(Promise.all([adapter.connect(), adapter.connect()])).rejects.toThrow('connect failed'); + expect(connectCalls).toHaveLength(1); + + await expect(adapter.connect()).resolves.toBeUndefined(); + expect(connectCalls).toHaveLength(2); + + const disconnectCalls = Promise.all([adapter.disconnect(), adapter.disconnect()]); + await expect(disconnectCalls).resolves.toEqual([undefined, undefined]); + expect(endCalls).toHaveLength(1); + } finally { + mock.restore(); + } + }); }); From ec972e2826df1181e99d44b5b9b47fa3a2347cc0 Mon Sep 17 00:00:00 2001 From: Shayam Prasad Sah Date: Tue, 29 Sep 2026 22:19:23 +0100 Subject: [PATCH 4/5] Handle webhook timeouts and queue lifecycle --- src/core/config.ts | 2 + src/core/factory.ts | 152 +++++++++++++----- src/runtime/queue/in-memory-queue.ts | 23 ++- .../webhooks/default-webhook-processor.ts | 53 ++++-- src/types/config.ts | 11 ++ src/utils/validation.ts | 12 ++ tests/core/factory-lifecycle.test.ts | 89 ++++++++++ tests/runtime/queue.unit.test.ts | 97 +++++++++++ tests/runtime/webhook-processor.unit.test.ts | 47 +++++- 9 files changed, 429 insertions(+), 57 deletions(-) create mode 100644 tests/core/factory-lifecycle.test.ts diff --git a/src/core/config.ts b/src/core/config.ts index a97e3c7c..27d77129 100644 --- a/src/core/config.ts +++ b/src/core/config.ts @@ -74,6 +74,8 @@ export class AnchorConfig { queue: { backend: input.framework.queue?.backend ?? 'memory', concurrency: input.framework.queue?.concurrency ?? 1, + maxPendingJobs: input.framework.queue?.maxPendingJobs, + onError: input.framework.queue?.onError, }, watchers: { enabled: input.framework.watchers?.enabled ?? true, diff --git a/src/core/factory.ts b/src/core/factory.ts index 9e4dd264..6901a18c 100644 --- a/src/core/factory.ts +++ b/src/core/factory.ts @@ -34,6 +34,8 @@ export class AnchorInstance { private initialized = false; private backgroundJobsRunning = false; + private initPromise: Promise | null = null; + private shutdownPromise: Promise | null = null; constructor(config: Partial) { this.config = new AnchorConfig(config); @@ -52,49 +54,82 @@ export class AnchorInstance { } /** - * Initialize registered plugins and all runtime services. + * Initialize registered plugins and runtime services. Concurrent calls share + * one initialization attempt; calls made during shutdown run afterward. */ public async init(): Promise { + if (this.shutdownPromise) { + await this.shutdownPromise; + return this.init(); + } if (this.initialized) return; + if (this.initPromise) return this.initPromise; + + const initPromise = this.initializeResources(); + this.initPromise = initPromise; + try { + await initPromise; + } finally { + if (this.initPromise === initPromise) this.initPromise = null; + } + } + private async initializeResources(): Promise { const frameworkConfig = this.config.get('framework'); - this.database = createSqlDatabaseAdapter(frameworkConfig.database); - await this.database.connect(); - await this.database.migrate(); - - const queueConcurrency = frameworkConfig.queue?.concurrency ?? 1; - this.queue = new InMemoryQueueAdapter({ concurrency: queueConcurrency }); - - this.webhookProcessor = new DefaultWebhookProcessor({ - config: this.config.getConfig(), - database: this.database, - }); - - const watchersEnabled = frameworkConfig.watchers?.enabled ?? true; - if (watchersEnabled) { - this.watchers = [ - new TransactionWatcher(this.database, this.queue, { - pollIntervalMs: frameworkConfig.watchers?.pollIntervalMs ?? 15000, - transactionTimeoutMs: frameworkConfig.watchers?.transactionTimeoutMs ?? 300000, - retentionDays: frameworkConfig.watchers?.retentionDays ?? 90, - }), - ]; - } - - this.expressRouter = new AnchorExpressRouter({ - config: this.config, - database: this.database, - webhookProcessor: this.webhookProcessor, - }).getMiddleware(); - - for (const plugin of this.plugins.values()) { - if (plugin.init) { - await plugin.init(this); + try { + this.database = this.createDatabaseAdapter(); + await this.database.connect(); + await this.database.migrate(); + + const queueConcurrency = frameworkConfig.queue?.concurrency ?? 1; + this.queue = new InMemoryQueueAdapter({ + concurrency: queueConcurrency, + maxPendingJobs: frameworkConfig.queue?.maxPendingJobs, + onError: frameworkConfig.queue?.onError, + }); + + this.webhookProcessor = new DefaultWebhookProcessor({ + config: this.config.getConfig(), + database: this.database, + }); + + const watchersEnabled = frameworkConfig.watchers?.enabled ?? true; + if (watchersEnabled) { + this.watchers = [ + new TransactionWatcher(this.database, this.queue, { + pollIntervalMs: frameworkConfig.watchers?.pollIntervalMs ?? 15000, + transactionTimeoutMs: frameworkConfig.watchers?.transactionTimeoutMs ?? 300000, + retentionDays: frameworkConfig.watchers?.retentionDays ?? 90, + }), + ]; + } + + this.expressRouter = new AnchorExpressRouter({ + config: this.config, + database: this.database, + webhookProcessor: this.webhookProcessor, + }).getMiddleware(); + + for (const plugin of this.plugins.values()) { + if (plugin.init) { + await plugin.init(this); + } } + + this.initialized = true; + } catch (error) { + try { + await this.releaseResources(); + } catch { + // Preserve the initialization error; shutdown can retry resource cleanup. + } + throw error; } + } - this.initialized = true; + protected createDatabaseAdapter(): DatabaseAdapter { + return createSqlDatabaseAdapter(this.config.get('framework').database); } /** @@ -128,12 +163,53 @@ export class AnchorInstance { } /** - * Cleanly shutdown all services. + * Cleanly shut down all services. If initialization is pending, shutdown waits + * for it to settle before releasing resources. The instance can then be initialized again. */ public async shutdown(): Promise { - if (!this.initialized) return; - await this.stopBackgroundJobs(); - await this.requireDatabase().disconnect(); + if (this.shutdownPromise) return this.shutdownPromise; + if ( + !this.initialized && + !this.initPromise && + !this.database && + !this.queue && + this.watchers.length === 0 + ) { + return; + } + + const shutdownPromise = this.shutdownResources(); + this.shutdownPromise = shutdownPromise; + try { + await shutdownPromise; + } finally { + if (this.shutdownPromise === shutdownPromise) this.shutdownPromise = null; + } + } + + private async shutdownResources(): Promise { + await this.initPromise?.catch(() => undefined); + await this.releaseResources(); + } + + private async releaseResources(): Promise { + if (this.backgroundJobsRunning) { + for (const watcher of this.watchers) { + await watcher.stop(); + } + await this.queue?.stop(); + this.backgroundJobsRunning = false; + } + + if (this.database) { + await this.database.disconnect(); + } + + this.database = null; + this.queue = null; + this.webhookProcessor = null; + this.watchers = []; + this.expressRouter = null; this.initialized = false; } diff --git a/src/runtime/queue/in-memory-queue.ts b/src/runtime/queue/in-memory-queue.ts index 6fa977bb..f4bc4986 100644 --- a/src/runtime/queue/in-memory-queue.ts +++ b/src/runtime/queue/in-memory-queue.ts @@ -2,10 +2,14 @@ import type { QueueAdapter, QueueJob } from '@/runtime/interfaces.ts'; interface InMemoryQueueOptions { concurrency: number; + maxPendingJobs?: number; + onError?: (job: QueueJob, error: unknown) => void | Promise; } export class InMemoryQueueAdapter implements QueueAdapter { private readonly concurrency: number; + private readonly maxPendingJobs: number | undefined; + private readonly onError: InMemoryQueueOptions['onError']; private readonly jobs: QueueJob[] = []; private running = false; private activeWorkers = 0; @@ -15,9 +19,20 @@ export class InMemoryQueueAdapter implements QueueAdapter { constructor(options: InMemoryQueueOptions) { this.concurrency = options.concurrency; + this.maxPendingJobs = options.maxPendingJobs; + this.onError = options.onError; } public async enqueue(job: QueueJob): Promise { + const canStartImmediately = this.running && this.worker && this.activeWorkers < this.concurrency; + if ( + this.maxPendingJobs !== undefined && + this.jobs.length >= this.maxPendingJobs && + !canStartImmediately + ) { + throw new Error('In-memory queue is full'); + } + this.jobs.push(job); this.kick(); } @@ -68,8 +83,12 @@ export class InMemoryQueueAdapter implements QueueAdapter { (async () => { try { await worker(job); - } catch { - // Best-effort queue for MVP: job errors are handled by worker logic. + } catch (error) { + try { + await this.onError?.(job, error); + } catch { + // Reporting failures must not prevent other queued jobs from running. + } } finally { this.activeWorkers -= 1; diff --git a/src/runtime/webhooks/default-webhook-processor.ts b/src/runtime/webhooks/default-webhook-processor.ts index 149abc0d..de496ecd 100644 --- a/src/runtime/webhooks/default-webhook-processor.ts +++ b/src/runtime/webhooks/default-webhook-processor.ts @@ -7,6 +7,10 @@ interface DefaultWebhookProcessorOptions { database: DatabaseAdapter; } +const DEFAULT_CALLBACK_TIMEOUT_MS = 30_000; +const CALLBACK_FAILURE_MESSAGE = 'Webhook callback failed'; +const CALLBACK_TIMEOUT_MESSAGE = 'Webhook callback timed out'; + function toComparableBuffer(value: string): Buffer { return Buffer.from(value, 'utf8'); } @@ -52,31 +56,52 @@ export class DefaultWebhookProcessor implements WebhookProcessor { } try { - await this.config.webhooks?.onEvent?.( - { - id: insertion.record.id, - eventId: insertion.record.eventId, - provider: insertion.record.provider, - payload: insertion.record.payload, - }, - { - receivedAt: insertion.record.createdAt, - signature: input.signature, - }, - ); + const callback = this.config.webhooks?.onEvent; + if (callback) { + let timeout: ReturnType | undefined; + try { + await Promise.race([ + Promise.resolve().then(() => + callback( + { + id: insertion.record.id, + eventId: insertion.record.eventId, + provider: insertion.record.provider, + payload: insertion.record.payload, + }, + { + receivedAt: insertion.record.createdAt, + signature: input.signature, + }, + ), + ), + new Promise((_, reject) => { + timeout = setTimeout( + () => reject(new Error(CALLBACK_TIMEOUT_MESSAGE)), + this.config.webhooks?.callbackTimeoutMs ?? DEFAULT_CALLBACK_TIMEOUT_MS, + ); + }), + ]); + } finally { + if (timeout) clearTimeout(timeout); + } + } await this.database.updateWebhookEventStatus({ id: insertion.record.id, status: 'processed', }); } catch (error) { - const message = error instanceof Error ? error.message : 'Unknown webhook callback error'; + const message = + error instanceof Error && error.message === CALLBACK_TIMEOUT_MESSAGE + ? CALLBACK_TIMEOUT_MESSAGE + : CALLBACK_FAILURE_MESSAGE; await this.database.updateWebhookEventStatus({ id: insertion.record.id, status: 'failed', errorMessage: message, }); - throw error; + throw new Error('Webhook processing failed'); } return { duplicate: false, eventId: insertion.record.eventId }; diff --git a/src/types/config.ts b/src/types/config.ts index a3262772..30233ee3 100644 --- a/src/types/config.ts +++ b/src/types/config.ts @@ -1,3 +1,5 @@ +import type { QueueJob } from '@/runtime/interfaces.ts'; + /** * Configuration Types for Anchor-Kit * Defines the complete configuration interface required to initialize an Anchor-Kit instance @@ -420,6 +422,12 @@ export interface FrameworkConfig { * @optional - defaults to 1 */ concurrency?: number; + + /** Maximum number of jobs waiting for a worker; in-flight jobs are excluded. Defaults to unlimited. */ + maxPendingJobs?: number; + + /** Called when a queue worker fails. Errors from this callback are ignored. */ + onError?: (job: QueueJob, error: unknown) => void | Promise; }; /** @@ -623,6 +631,9 @@ export interface AnchorKitConfig { * Webhook integration configuration. */ webhooks?: { + /** Maximum time to wait for onEvent, in milliseconds. Defaults to 30000. */ + callbackTimeoutMs?: number; + /** * Called after webhook event verification and persistence. */ diff --git a/src/utils/validation.ts b/src/utils/validation.ts index 26252c61..b912dbc4 100644 --- a/src/utils/validation.ts +++ b/src/utils/validation.ts @@ -250,6 +250,18 @@ export const AnchorKitConfigSchema = { if (framework.queue?.concurrency !== undefined && framework.queue.concurrency < 1) { throw new Error('framework.queue.concurrency must be >= 1'); } + if ( + framework.queue?.maxPendingJobs !== undefined && + (!Number.isInteger(framework.queue.maxPendingJobs) || framework.queue.maxPendingJobs < 0) + ) { + throw new Error('framework.queue.maxPendingJobs must be a non-negative integer'); + } + if ( + config.webhooks?.callbackTimeoutMs !== undefined && + (!Number.isFinite(config.webhooks.callbackTimeoutMs) || config.webhooks.callbackTimeoutMs <= 0) + ) { + throw new Error('webhooks.callbackTimeoutMs must be > 0'); + } if ( framework.watchers?.pollIntervalMs !== undefined && framework.watchers.pollIntervalMs < 10 diff --git a/tests/core/factory-lifecycle.test.ts b/tests/core/factory-lifecycle.test.ts new file mode 100644 index 00000000..ccb5cd62 --- /dev/null +++ b/tests/core/factory-lifecycle.test.ts @@ -0,0 +1,89 @@ +import { AnchorInstance } from '@/core/factory.ts'; +import type { AnchorKitConfig } from '@/types/config.ts'; +import type { DatabaseAdapter } from '@/runtime/interfaces.ts'; +import { Keypair } from '@stellar/stellar-sdk'; +import { describe, expect, it, vi } from 'vitest'; + +describe('AnchorInstance lifecycle', () => { + it('waits for pending init before shutdown and supports later reinitialization', async () => { + let signalConnectStarted: (() => void) | undefined; + const connectStarted = new Promise((resolve) => { + signalConnectStarted = resolve; + }); + let releaseConnect: (() => void) | undefined; + const connectGate = new Promise((resolve) => { + releaseConnect = resolve; + }); + const signingKey = Keypair.random().secret(); + const anchor = new DelayedConnectAnchor( + { + network: { network: 'testnet' }, + server: { port: 3000 }, + security: { + sep10SigningKey: signingKey, + interactiveJwtSecret: 'jwt-secret', + distributionAccountSecret: 'dist-secret', + }, + assets: { + assets: [ + { + code: 'USDC', + issuer: 'GBBD47IF6LWK7P7MDEVSCWR7DPUWV3NY3DTQEVFL4NAT4AQH3ZLLFLA5', + }, + ], + }, + framework: { + database: { provider: 'sqlite', url: 'file:unused' }, + watchers: { enabled: false }, + }, + }, + connectGate, + () => signalConnectStarted?.(), + ); + + const firstInit = anchor.init(); + const concurrentInit = anchor.init(); + await connectStarted; + + const shutdown = anchor.shutdown(); + const initDuringShutdown = anchor.init(); + releaseConnect?.(); + + await Promise.all([firstInit, concurrentInit, shutdown, initDuringShutdown]); + await anchor.init(); + expect(anchor.databaseAdapters).toHaveLength(2); + expect(anchor.databaseAdapters[0]?.disconnect).toHaveBeenCalledTimes(1); + expect(() => anchor.getExpressRouter()).not.toThrow(); + await anchor.shutdown(); + expect(anchor.databaseAdapters[1]?.disconnect).toHaveBeenCalledTimes(1); + }); +}); + +class DelayedConnectAnchor extends AnchorInstance { + public readonly databaseAdapters: DatabaseAdapter[] = []; + private connectCount = 0; + + constructor( + config: Partial, + private readonly connectGate: Promise, + private readonly onConnectStart: () => void, + ) { + super(config); + } + + protected override createDatabaseAdapter(): DatabaseAdapter { + const adapter = { + connect: vi.fn(async () => { + this.connectCount += 1; + if (this.connectCount === 1) { + this.onConnectStart(); + await this.connectGate; + } + }), + migrate: vi.fn().mockResolvedValue(undefined), + disconnect: vi.fn().mockResolvedValue(undefined), + } as unknown as DatabaseAdapter; + this.databaseAdapters.push(adapter); + return adapter; + } +} diff --git a/tests/runtime/queue.unit.test.ts b/tests/runtime/queue.unit.test.ts index 571adbf1..8f3bc932 100644 --- a/tests/runtime/queue.unit.test.ts +++ b/tests/runtime/queue.unit.test.ts @@ -309,4 +309,101 @@ describe('InMemoryQueueAdapter', () => { await stopPromise; expect(jobFinished).toBe(true); // stop() should only resolve after job is finished }); + + it('reports a failed job and continues processing later jobs', async () => { + const failures: Array<{ job: QueueJob; error: unknown }> = []; + const completed: string[] = []; + let resolveCompleted: (() => void) | undefined; + const completedPromise = new Promise((resolve) => { + resolveCompleted = resolve; + }); + const queue = new InMemoryQueueAdapter({ + concurrency: 1, + onError: (job, error) => { + failures.push({ job, error }); + }, + }); + + await queue.start(async (job) => { + const id = job.payload.id as string; + if (id === 'fail') throw new Error('worker failed'); + completed.push(id); + resolveCompleted?.(); + }); + await queue.enqueue({ type: 'process_watcher_task', payload: { id: 'fail' } }); + await queue.enqueue({ type: 'cleanup_records', payload: { id: 'succeed' } }); + await completedPromise; + await queue.stop(); + + expect(failures).toHaveLength(1); + expect(failures[0]?.job.type).toBe('process_watcher_task'); + expect(failures[0]?.error).toEqual(new Error('worker failed')); + expect(completed).toEqual(['succeed']); + }); + + it('limits pending jobs and accepts more after the queue drains', async () => { + const queue = new InMemoryQueueAdapter({ concurrency: 1, maxPendingJobs: 1 }); + let completedCount = 0; + let resolveDrained: (() => void) | undefined; + const drained = new Promise((resolve) => { + resolveDrained = resolve; + }); + + await queue.enqueue({ type: 'cleanup_records', payload: { id: 1 } }); + await expect( + queue.enqueue({ type: 'cleanup_records', payload: { id: 2 } }), + ).rejects.toThrow('In-memory queue is full'); + + await queue.start(async () => { + completedCount += 1; + if (completedCount === 2) resolveDrained?.(); + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + await queue.enqueue({ type: 'cleanup_records', payload: { id: 3 } }); + await drained; + await queue.stop(); + + expect(completedCount).toBe(2); + }); + + it('excludes in-flight work from capacity and rejects concurrent overflow', async () => { + const queue = new InMemoryQueueAdapter({ concurrency: 1, maxPendingJobs: 1 }); + let signalWorkerStarted: (() => void) | undefined; + const workerStarted = new Promise((resolve) => { + signalWorkerStarted = resolve; + }); + let releaseWorker: (() => void) | undefined; + const workerGate = new Promise((resolve) => { + releaseWorker = resolve; + }); + let signalPendingComplete: (() => void) | undefined; + const pendingComplete = new Promise((resolve) => { + signalPendingComplete = resolve; + }); + const completed: string[] = []; + + await queue.start(async (job) => { + const id = job.payload.id as string; + if (id === 'in-flight') { + signalWorkerStarted?.(); + await workerGate; + } + completed.push(id); + if (id === 'pending-1') signalPendingComplete?.(); + }); + await queue.enqueue({ type: 'cleanup_records', payload: { id: 'in-flight' } }); + await workerStarted; + + const enqueueResults = await Promise.allSettled([ + queue.enqueue({ type: 'cleanup_records', payload: { id: 'pending-1' } }), + queue.enqueue({ type: 'cleanup_records', payload: { id: 'pending-2' } }), + ]); + expect(enqueueResults.map((result) => result.status)).toEqual(['fulfilled', 'rejected']); + + releaseWorker?.(); + await pendingComplete; + await queue.stop(); + + expect(completed).toEqual(['in-flight', 'pending-1']); + }); }); diff --git a/tests/runtime/webhook-processor.unit.test.ts b/tests/runtime/webhook-processor.unit.test.ts index 0c1da6cc..07842ee7 100644 --- a/tests/runtime/webhook-processor.unit.test.ts +++ b/tests/runtime/webhook-processor.unit.test.ts @@ -40,14 +40,13 @@ describe('DefaultWebhookProcessor Unit Tests', () => { rawBody: '{}', }; - // Should rethrow the error - await expect(processor.process(input)).rejects.toThrow('Callback failed'); + await expect(processor.process(input)).rejects.toThrow('Webhook processing failed'); // Should have updated status to failed with error message expect(mockDatabase.updateWebhookEventStatus).toHaveBeenCalledWith({ id: 'internal-id', status: 'failed', - errorMessage: 'Callback failed', + errorMessage: 'Webhook callback failed', }); }); @@ -96,4 +95,46 @@ describe('DefaultWebhookProcessor Unit Tests', () => { status: 'processed', }); }); + + it('marks a stalled callback failed when its configured timeout elapses', async () => { + const mockDatabase = { + insertWebhookEvent: vi.fn().mockResolvedValue({ + inserted: true, + record: { + id: 'timed-out-id', + eventId: 'stalled-event', + provider: 'generic', + payload: {}, + createdAt: new Date().toISOString(), + }, + }), + updateWebhookEventStatus: vi.fn().mockResolvedValue(undefined), + } as unknown as DatabaseAdapter; + const mockConfig = { + security: { verifyWebhookSignatures: false }, + webhooks: { + callbackTimeoutMs: 5, + onEvent: vi.fn(() => new Promise(() => undefined)), + }, + } as unknown as AnchorKitConfig; + const processor = new DefaultWebhookProcessor({ + config: mockConfig, + database: mockDatabase, + }); + + await expect( + processor.process({ + eventId: 'stalled-event', + provider: 'generic', + payload: {}, + rawBody: '{}', + }), + ).rejects.toThrow('Webhook processing failed'); + + expect(mockDatabase.updateWebhookEventStatus).toHaveBeenCalledWith({ + id: 'timed-out-id', + status: 'failed', + errorMessage: 'Webhook callback timed out', + }); + }); }); From 29b800070d938ba9e1cebbe8dfe9ffd295f48182 Mon Sep 17 00:00:00 2001 From: Ikemhood <91676649+ikemHood@users.noreply.github.com> Date: Wed, 30 Sep 2026 10:25:27 +0100 Subject: [PATCH 5/5] fix: complete queue and shutdown safeguards --- src/core/factory.ts | 13 ++++-- src/runtime/queue/in-memory-queue.ts | 9 +++- .../webhooks/default-webhook-processor.ts | 9 +++- src/utils/validation-helpers.ts | 7 ++++ tests/core/config.test.ts | 11 +++++ tests/runtime/queue.unit.test.ts | 42 +++++++++++-------- .../sql-adapter-interactive-tx.test.ts | 1 - 7 files changed, 66 insertions(+), 26 deletions(-) diff --git a/src/core/factory.ts b/src/core/factory.ts index fea50033..29a36d11 100644 --- a/src/core/factory.ts +++ b/src/core/factory.ts @@ -80,7 +80,7 @@ export class AnchorInstance { try { validatePluginRoutes(pluginRoutes); - this.database = createSqlDatabaseAdapter(frameworkConfig.database); + this.database = this.createDatabaseAdapter(); await this.database.connect(); await this.database.migrate(); @@ -175,15 +175,16 @@ export class AnchorInstance { } /** - * Cleanly shut down all services. If initialization is pending, shutdown waits - * for it to settle before releasing resources. The instance can then be initialized again. + * Cleanly shut down all services. If initialization is pending, shutdown waits + * for it to settle before releasing resources. The instance can then be initialized again. */ public async shutdown(): Promise { - if (!this.initialized && !this.shutdownPromise) return; + if (!this.initialized && !this.initPromise && !this.shutdownPromise) return; if (this.shutdownPromise) return this.shutdownPromise; this.shutdownPromise = (async () => { try { + await this.initPromise?.catch(() => undefined); if (!this.initialized) return; await this.stopBackgroundJobs(); @@ -203,6 +204,10 @@ export class AnchorInstance { return this.shutdownPromise; } + protected createDatabaseAdapter(): DatabaseAdapter { + return createSqlDatabaseAdapter(this.config.get('framework').database); + } + /** * Return middleware compatible with Express router mounting. * diff --git a/src/runtime/queue/in-memory-queue.ts b/src/runtime/queue/in-memory-queue.ts index 4115645a..a26b3747 100644 --- a/src/runtime/queue/in-memory-queue.ts +++ b/src/runtime/queue/in-memory-queue.ts @@ -22,6 +22,12 @@ export class InMemoryQueueAdapter implements QueueAdapter { if (!Number.isSafeInteger(options.concurrency) || options.concurrency < 1) { throw new ConfigError('InMemoryQueueAdapter concurrency must be a positive safe integer'); } + if ( + options.maxPendingJobs !== undefined && + (!Number.isSafeInteger(options.maxPendingJobs) || options.maxPendingJobs < 1) + ) { + throw new ConfigError('InMemoryQueueAdapter maxPendingJobs must be a positive safe integer'); + } this.concurrency = options.concurrency; this.maxPendingJobs = options.maxPendingJobs; @@ -29,7 +35,8 @@ export class InMemoryQueueAdapter implements QueueAdapter { } public async enqueue(job: QueueJob): Promise { - const canStartImmediately = this.running && this.worker && this.activeWorkers < this.concurrency; + const canStartImmediately = + this.running && this.worker && this.activeWorkers < this.concurrency; if ( this.maxPendingJobs !== undefined && this.jobs.length >= this.maxPendingJobs && diff --git a/src/runtime/webhooks/default-webhook-processor.ts b/src/runtime/webhooks/default-webhook-processor.ts index 12a45e83..72171370 100644 --- a/src/runtime/webhooks/default-webhook-processor.ts +++ b/src/runtime/webhooks/default-webhook-processor.ts @@ -41,7 +41,12 @@ export class DefaultWebhookProcessor implements WebhookProcessor { payload: Record; rawBody: string | Buffer | Uint8Array; signature?: string; - }): Promise<{ duplicate: boolean; eventId: string; provider: string }> { + }): Promise<{ + duplicate: boolean; + eventId: string; + provider: string; + status?: 'pending' | 'processed' | 'failed'; + }> { this.verifySignatureIfEnabled(input); const insertion = await this.database.insertOrGetWebhookEvent({ @@ -105,7 +110,7 @@ export class DefaultWebhookProcessor implements WebhookProcessor { status: 'failed', errorMessage: message, }); - throw new Error('Webhook processing failed'); + throw new Error('Webhook processing failed', { cause: error }); } return { diff --git a/src/utils/validation-helpers.ts b/src/utils/validation-helpers.ts index 5cf715b2..f2a0facc 100644 --- a/src/utils/validation-helpers.ts +++ b/src/utils/validation-helpers.ts @@ -578,6 +578,13 @@ function validateAnchorKitConfig(config: AnchorKitConfig): boolean { if (!assets) throw new Error('Missing required top-level field: assets'); if (!framework) throw new Error('Missing required top-level field: framework'); + if ( + config.webhooks?.callbackTimeoutMs !== undefined && + !isPositiveSafeInteger(config.webhooks.callbackTimeoutMs) + ) { + throw new Error('webhooks.callbackTimeoutMs must be a positive safe integer'); + } + NetworkConfigSchema.validate(network); SecurityConfigSchema.validate(security); if (security.enableClientAttribution) { diff --git a/tests/core/config.test.ts b/tests/core/config.test.ts index b7d0f176..95be819c 100644 --- a/tests/core/config.test.ts +++ b/tests/core/config.test.ts @@ -30,6 +30,17 @@ describe('AnchorConfig', () => { }, }; + it('requires a positive safe webhook callback timeout', () => { + for (const callbackTimeoutMs of [0, -1, 1.5, Number.NaN, Number.POSITIVE_INFINITY]) { + expect(() => + new AnchorConfig({ + ...validBaseConfig, + webhooks: { callbackTimeoutMs }, + }).validate(), + ).toThrow(/webhooks\.callbackTimeoutMs must be a positive safe integer/); + } + }); + function configWithPlugins( plugins: NonNullable, ): AnchorKitConfig { diff --git a/tests/runtime/queue.unit.test.ts b/tests/runtime/queue.unit.test.ts index 07eb5c23..7ffe4613 100644 --- a/tests/runtime/queue.unit.test.ts +++ b/tests/runtime/queue.unit.test.ts @@ -499,20 +499,19 @@ describe('InMemoryQueueAdapter', () => { }); await queue.start(async (job) => { - const id = job.payload.id as string; - if (id === 'fail') throw new Error('worker failed'); - completed.push(id); + if (job.type === 'process_watcher_task') throw new Error('worker failed'); + completed.push(job.type); resolveCompleted?.(); }); - await queue.enqueue({ type: 'process_watcher_task', payload: { id: 'fail' } }); - await queue.enqueue({ type: 'cleanup_records', payload: { id: 'succeed' } }); + await queue.enqueue({ type: 'process_watcher_task', payload: { watcherTaskId: 'fail' } }); + await queue.enqueue({ type: 'cleanup_records', payload: { retentionDays: 90 } }); await completedPromise; await queue.stop(); expect(failures).toHaveLength(1); expect(failures[0]?.job.type).toBe('process_watcher_task'); expect(failures[0]?.error).toEqual(new Error('worker failed')); - expect(completed).toEqual(['succeed']); + expect(completed).toEqual(['cleanup_records']); }); it('limits pending jobs and accepts more after the queue drains', async () => { @@ -523,9 +522,9 @@ describe('InMemoryQueueAdapter', () => { resolveDrained = resolve; }); - await queue.enqueue({ type: 'cleanup_records', payload: { id: 1 } }); + await queue.enqueue({ type: 'cleanup_records', payload: { retentionDays: 1 } }); await expect( - queue.enqueue({ type: 'cleanup_records', payload: { id: 2 } }), + queue.enqueue({ type: 'cleanup_records', payload: { retentionDays: 2 } }), ).rejects.toThrow('In-memory queue is full'); await queue.start(async () => { @@ -533,7 +532,7 @@ describe('InMemoryQueueAdapter', () => { if (completedCount === 2) resolveDrained?.(); }); await new Promise((resolve) => setTimeout(resolve, 0)); - await queue.enqueue({ type: 'cleanup_records', payload: { id: 3 } }); + await queue.enqueue({ type: 'cleanup_records', payload: { retentionDays: 3 } }); await drained; await queue.stop(); @@ -557,27 +556,34 @@ describe('InMemoryQueueAdapter', () => { const completed: string[] = []; await queue.start(async (job) => { - const id = job.payload.id as string; - if (id === 'in-flight') { + if (job.type === 'process_watcher_task') { signalWorkerStarted?.(); await workerGate; } - completed.push(id); - if (id === 'pending-1') signalPendingComplete?.(); + completed.push( + job.type === 'process_watcher_task' + ? 'in-flight' + : job.type === 'cleanup_records' + ? `${job.payload.retentionDays}` + : 'unexpected', + ); + if (job.type === 'cleanup_records' && job.payload.retentionDays === 1) { + signalPendingComplete?.(); + } }); - await queue.enqueue({ type: 'cleanup_records', payload: { id: 'in-flight' } }); + await queue.enqueue({ type: 'process_watcher_task', payload: { watcherTaskId: 'in-flight' } }); await workerStarted; const enqueueResults = await Promise.allSettled([ - queue.enqueue({ type: 'cleanup_records', payload: { id: 'pending-1' } }), - queue.enqueue({ type: 'cleanup_records', payload: { id: 'pending-2' } }), + queue.enqueue({ type: 'cleanup_records', payload: { retentionDays: 1 } }), + queue.enqueue({ type: 'cleanup_records', payload: { retentionDays: 2 } }), ]); expect(enqueueResults.map((result) => result.status)).toEqual(['fulfilled', 'rejected']); releaseWorker?.(); - await pendingComplete; + await pendingComplete; await queue.stop(); - expect(completed).toEqual(['in-flight', 'pending-1']); + expect(completed).toEqual(['in-flight', '1']); }); }); diff --git a/tests/runtime/sql-adapter-interactive-tx.test.ts b/tests/runtime/sql-adapter-interactive-tx.test.ts index 4e217ce0..7a50209f 100644 --- a/tests/runtime/sql-adapter-interactive-tx.test.ts +++ b/tests/runtime/sql-adapter-interactive-tx.test.ts @@ -3,7 +3,6 @@ import { createSqlDatabaseAdapter } from '@/runtime/database/sql-database-adapte import { randomUUID } from 'node:crypto'; import { unlinkSync } from 'node:fs'; import { afterAll, beforeAll, describe, expect, it } from 'vitest'; -import { mock } from 'bun:test'; import type { DatabaseAdapter } from '@/runtime/interfaces.ts'; describe('SqlDatabaseAdapter – interactive transaction status updates', () => {