Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 15 additions & 4 deletions src/core/factory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,9 +61,14 @@ 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<void> {
if (this.shutdownPromise) {
await this.shutdownPromise;
return this.init();
}
if (this.initialized) return;
if (this.initPromise) return this.initPromise;

Expand All @@ -75,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();

Expand Down Expand Up @@ -170,14 +175,16 @@ 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<void> {
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();
Expand All @@ -197,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.
*
Expand Down
30 changes: 28 additions & 2 deletions src/runtime/queue/in-memory-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,14 @@ import type { QueueAdapter, QueueJob } from '@/runtime/interfaces.ts';

interface InMemoryQueueOptions {
concurrency: number;
maxPendingJobs?: number;
onError?: (job: QueueJob, error: unknown) => void | Promise<void>;
}

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;
Expand All @@ -18,11 +22,29 @@ 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;
this.onError = options.onError;
}

public async enqueue(job: QueueJob): Promise<void> {
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();
}
Expand Down Expand Up @@ -80,8 +102,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;

Expand Down
60 changes: 45 additions & 15 deletions src/runtime/webhooks/default-webhook-processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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');
}
Expand Down Expand Up @@ -37,7 +41,12 @@ export class DefaultWebhookProcessor implements WebhookProcessor {
payload: Record<string, unknown>;
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({
Expand All @@ -56,31 +65,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<typeof setTimeout> | 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<never>((_, 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', { cause: error });
}

return {
Expand Down
11 changes: 11 additions & 0 deletions src/types/config.ts
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -422,6 +424,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<void>;
};

/**
Expand Down Expand Up @@ -633,6 +641,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.
*/
Expand Down
7 changes: 7 additions & 0 deletions src/utils/validation-helpers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
11 changes: 11 additions & 0 deletions tests/core/config.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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['framework']['plugins']>,
): AnchorKitConfig {
Expand Down
89 changes: 89 additions & 0 deletions tests/core/factory-lifecycle.test.ts
Original file line number Diff line number Diff line change
@@ -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<void>((resolve) => {
signalConnectStarted = resolve;
});
let releaseConnect: (() => void) | undefined;
const connectGate = new Promise<void>((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<AnchorKitConfig>,
private readonly connectGate: Promise<void>,
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;
}
}
Loading
Loading