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
4 changes: 4 additions & 0 deletions listener/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,10 @@ RETRY_SCHEDULER_PROCESSOR_ID=
# Number of retry jobs to process per scheduler tick.
RETRY_SCHEDULER_BATCH_SIZE=10

# Request timeout for outbound webhook delivery (ms). Timed-out requests are
# treated as failures and are subject to normal retry backoff.
# WEBHOOK_DELIVERY_TIMEOUT_MS=10000

# -----------------------------------------------------------------------------
# Retry Policy
# -----------------------------------------------------------------------------
Expand Down
4 changes: 3 additions & 1 deletion listener/src/api/events-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,10 @@ import { handleTemplateRoutes } from './template-routes';
import { sendOk, sendErr, sendJson, ErrorCode } from '../utils/response';
import { normalizePaginationParams } from '../utils/pagination';
import { handleApiError, ApiError } from './error-handler';
import { applyRequestContext } from '../utils/request-id';
import { applyRequestIdMiddleware } from '../middleware/request-id';
import { TemplateService } from '../services/template-service';
import { handleTemplateRoutes } from './template-routes';
import { sendOk, sendErr, sendJson, ErrorCode } from '../utils/response';
import { validateContentType, getMimeType } from '../middleware/content-type';
import { NotificationHistoryService } from '../services/notification-history';
import { SearchSuggestionService } from '../services/search-suggestion';
Expand Down
1 change: 1 addition & 0 deletions listener/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -262,6 +262,7 @@ function loadAnalyticsConfig(): AnalyticsConfig {
function loadRetrySchedulerConfig(policy: RetryPolicyOptions): RetrySchedulerOptions {
return {
enabled: trimEnv('RETRY_SCHEDULER_ENABLED') !== 'false',
webhookTimeoutMs: parseIntegerEnv('WEBHOOK_DELIVERY_TIMEOUT_MS', '10000'),
pollIntervalMs: parseIntegerEnv('RETRY_SCHEDULER_POLL_INTERVAL_MS', '15000'),
lockTimeoutMs: parseIntegerEnv('RETRY_SCHEDULER_LOCK_TIMEOUT_MS', '60000'),
processorId: trimEnv('RETRY_SCHEDULER_PROCESSOR_ID'),
Expand Down
4 changes: 2 additions & 2 deletions listener/src/database/migration-system.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ export class MigrationRunner {
const rows = await this.db.all<{ id: string }>(
'SELECT id FROM migrations ORDER BY applied_at'
);
return rows.map((row) => row.id);
return (rows as unknown as { id: string }[]).map((row) => row.id);
}

async applyMigration(migration: Migration): Promise<void> {
Expand All @@ -67,7 +67,7 @@ export class MigrationRunner {
logger.info(`Migration ${migration.id} (${migration.name}) applied successfully`);
} catch (error) {
await this.db.run('ROLLBACK');
logger.error(`Migration ${migration.id} failed, rolling back:`, error);
logger.error(`Migration ${migration.id} failed, rolling back: ${(error as Error)?.message ?? String(error)}`);
throw error;
}
});
Expand Down
7 changes: 6 additions & 1 deletion listener/src/database/schema.sql
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,8 @@ CREATE TABLE IF NOT EXISTS scheduled_notifications (
contract_address TEXT, -- Stellar contract address (if applicable)
priority INTEGER NOT NULL DEFAULT 5, -- 1-10, lower = higher priority
metadata TEXT, -- Additional JSON metadata
next_retry_at DATETIME -- When the next retry should be attempted
next_retry_at DATETIME, -- When the next retry should be attempted
deduplication_key TEXT -- Caller-supplied key; duplicate inserts with the same key are silently skipped
);

-- Indexes for performance optimization
Expand All @@ -62,6 +63,10 @@ CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_created_at
CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_event_id
ON scheduled_notifications(event_id);

CREATE UNIQUE INDEX IF NOT EXISTS idx_scheduled_notifications_dedup_key
ON scheduled_notifications(deduplication_key)
WHERE deduplication_key IS NOT NULL;

CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_target
ON scheduled_notifications(target_recipient, status);

Expand Down
1 change: 1 addition & 0 deletions listener/src/middleware/security-headers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
* See: https://owasp.org/www-project-secure-headers/
*/

import type { ServerResponse } from 'http';
import type http from 'http';
import type { ServerResponse } from 'http';

Expand Down
34 changes: 34 additions & 0 deletions listener/src/migrations/003-notification-deduplication-key.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
import * as sqlite3 from 'sqlite3';

const migration = {
id: '003',
name: 'notification-deduplication-key',
up: async (db: sqlite3.Database) => {
await db.run(`
ALTER TABLE scheduled_notifications
ADD COLUMN deduplication_key TEXT
`);
await db.run(`
CREATE UNIQUE INDEX IF NOT EXISTS idx_scheduled_notifications_dedup_key
ON scheduled_notifications(deduplication_key)
WHERE deduplication_key IS NOT NULL
`);
},
down: async (db: sqlite3.Database) => {
await db.run('DROP INDEX IF EXISTS idx_scheduled_notifications_dedup_key');
// SQLite does not support DROP COLUMN before 3.35; recreate without the column.
await db.run(`
CREATE TABLE scheduled_notifications_backup AS
SELECT id, payload, payload_hash, notification_type, target_recipient,
execute_at, created_at, updated_at, status, retry_count, max_retries,
processing_started_at, processing_completed_at, processor_id,
lock_expires_at, last_error, error_details, event_id,
contract_address, priority, metadata, next_retry_at
FROM scheduled_notifications
`);
await db.run('DROP TABLE scheduled_notifications');
await db.run('ALTER TABLE scheduled_notifications_backup RENAME TO scheduled_notifications');
},
};

export default migration;
22 changes: 22 additions & 0 deletions listener/src/schema/compatibility-check.ts
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,14 @@ export function checkSchemaCompatibility(
breakingChanges.push('Consumer event schema is empty or malformed');
}

// Check 4: Topic structure changes
// Topics are appended as trailing topics - existing consumers ignore them
// Breaking change only if the core topic (event name) changes position
if (consumer.expectedTopics && contractEvent.topics.length < consumer.expectedTopics.length) {
const missingTopics = consumer.expectedTopics.filter(
(t) => !contractEvent.topics.includes(t)
);
if (missingTopics.length > 0) {
for (const consumerEvent of consumerSchema.events) {
const contractEvent = contractEvents.get(consumerEvent.eventName);
if (!contractEvent) {
Expand All @@ -185,6 +193,20 @@ export function checkSchemaCompatibility(
safeAdditions.push(...result.safeAdditions);
}

// Check 5: New data fields are safe (backward compatible)
// Any new data fields in the contract that weren't expected by the consumer
// are simply ignored - this is the Soroban trailing-topic pattern
const newDataFields = contractEvent.dataFields.filter(
(f) => !consumer.expectedFields.includes(f.name)
);
safeAdditions.push(
...newDataFields.map(
(f) => `Safe addition: New data field '${f.name}' in '${contractEvent.name}' (ignored by existing consumers)`
)
);

// Determine compatibility
const compatible = breakingChanges.length === 0;
for (const eventName of contractEvents.keys()) {
if (!consumerNames.has(eventName)) {
safeAdditions.push(`Safe addition: New event '${eventName}' is ignored by existing consumers`);
Expand Down
3 changes: 3 additions & 0 deletions listener/src/services/discord-notification.ts
Original file line number Diff line number Diff line change
Expand Up @@ -407,6 +407,7 @@ export class DiscordNotificationService implements NotificationProvider {
const originalLength = embed.title?.length ?? title.length;
title = title.slice(0, MAX_DISCORD_EMBED_TITLE_LENGTH - 3) + '...';
logger.warn('Discord embed title truncated', {
originalLength: title.length + 3,
originalLength,
maxLength: MAX_DISCORD_EMBED_TITLE_LENGTH,
});
Expand All @@ -427,8 +428,10 @@ export class DiscordNotificationService implements NotificationProvider {

let footer = embed.footer;
if (footer?.text && footer.text.length > MAX_DISCORD_FOOTER_TEXT_LENGTH) {
const originalLength = footer.text.length;
footer = { text: footer.text.slice(0, MAX_DISCORD_FOOTER_TEXT_LENGTH - 3) + '...' };
logger.warn('Discord footer text truncated', {
originalLength,
originalLength: footer.text.length,
maxLength: MAX_DISCORD_FOOTER_TEXT_LENGTH,
});
Expand Down
7 changes: 7 additions & 0 deletions listener/src/services/event-subscriber.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ export class EventSubscriber {
private eventQueue: EventProcessingQueue | null = null;
private expirationService: NotificationExpirationService | null = null;
private lastSuccessfulPollAt: number | null = null;
private backfillStartLedger: number | null = null;
private circuitBreaker: CircuitBreaker | null = null;
private backfillStartLedger: number | null = null;
/** Cold-start ledger resolved once per session by resolveBackfillStartLedger(). */
Expand Down Expand Up @@ -389,6 +390,9 @@ export class EventSubscriber {
if (lastCursor) {
// Normal real-time polling: continue from the last known cursor.
request = {
filters: [{ contractIds: [contractConfig.address], type: 'contract' }],
cursor: lastCursor,
limit: 100,
filters: [
{
contractIds: [contractConfig.address],
Expand All @@ -405,6 +409,9 @@ export class EventSubscriber {
// Cold start: apply the backfill safety limit.
const startLedger = await this.resolveBackfillStartLedger();
request = {
filters: [{ contractIds: [contractConfig.address], type: 'contract' }],
startLedger,
limit: 100,
filters: [
{
contractIds: [contractConfig.address],
Expand Down
2 changes: 1 addition & 1 deletion listener/src/services/notification-retry-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -367,6 +367,6 @@ function buildRetryFingerprint(
contractAddress: string
): string {
const eventName =
getEventName(event.topic) ?? event.topic.map((entry) => entry.toString()).join('|');
getEventName(event.topic) ?? event.topic.map((entry: { toString(): string }) => entry.toString()).join('|');
return `${contractAddress}:${event.id}:${eventName}:${event.txHash ?? ''}`;
}
1 change: 1 addition & 0 deletions listener/src/services/notification-stats-cache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,7 @@ export class NotificationStatsCache {
* @notice Check if stats are currently cached
*/
has(): boolean {
return this.cache.get(this.CACHE_KEY) !== undefined;
return (this.cache as any).has(this.CACHE_KEY);
}

Expand Down
5 changes: 5 additions & 0 deletions listener/src/services/retry-scheduler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,8 @@ export interface RetrySchedulerConfig {
maxDelayMs: number;
/** Add ±25 % random jitter to prevent thundering herd. Default: true. */
jitter: boolean;
/** Request timeout for outbound webhook delivery (ms). Default: 10 000. */
webhookTimeoutMs: number;
/**
* Hard ceiling on total delivery attempts, including the first one.
* `undefined` (default) leaves each notification's own `maxRetries` in
Expand All @@ -59,6 +61,7 @@ export const RETRY_SCHEDULER_DEFAULTS: RetrySchedulerConfig = {
multiplier: 2,
maxDelayMs: 60 * 60 * 1_000,
jitter: true,
webhookTimeoutMs: 10_000,
};

/**
Expand Down Expand Up @@ -125,6 +128,8 @@ export class RetryScheduler {
this.processorId = this.config.processorId ?? `retry-${uuidv4()}`;
this.repository = repository;
this.discordService = discordService ?? null;
this.webhookDeliveryService =
webhookDeliveryService ?? new WebhookDeliveryService({ timeoutMs: this.config.webhookTimeoutMs });
this.webhookDeliveryService = webhookDeliveryService ?? new WebhookDeliveryService();
this.deliveryReceiptRepository = deliveryReceiptRepository;
}
Expand Down
45 changes: 41 additions & 4 deletions listener/src/services/scheduled-notification-repository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,9 @@ export class ScheduledNotificationRepository {
}

/**
* Create a new scheduled notification
* Create a new scheduled notification.
* If a deduplicationKey is provided and a notification with that key already
* exists, the existing notification's id is returned without creating a duplicate.
*/
async create(input: CreateScheduledNotificationInput, requestId?: string): Promise<number> {
const payloadJson = JSON.stringify(input.payload);
Expand All @@ -42,8 +44,8 @@ export class ScheduledNotificationRepository {
const sql = `
INSERT INTO scheduled_notifications (
payload, payload_hash, notification_type, target_recipient, execute_at,
max_retries, event_id, contract_address, priority, metadata
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
max_retries, event_id, contract_address, priority, metadata, deduplication_key
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`;

const serializedPayload = compressPayload(input.payload);
Expand All @@ -59,8 +61,20 @@ export class ScheduledNotificationRepository {
input.contractAddress ?? null,
input.priority ?? 5,
input.metadata ? JSON.stringify(input.metadata) : null,
input.deduplicationKey ?? null,
];

try {
const result = await this.db.run(sql, params);

this.statsCache.invalidate();

logger.info('Scheduled notification created', {
requestId,
id: result.lastID,
executeAt: input.executeAt,
type: input.notificationType,
});
const result = await this.db.run(sql, params);

// Invalidate stats cache after creation
Expand All @@ -73,7 +87,29 @@ export class ScheduledNotificationRepository {
type: input.notificationType,
});

return result.lastID;
return result.lastID;
} catch (err) {
if (
input.deduplicationKey &&
(err as any)?.message?.includes('UNIQUE constraint failed')
) {
const existing = await this.db.get<{ id: number }>(
'SELECT id FROM scheduled_notifications WHERE deduplication_key = ?',
[input.deduplicationKey],
);

if (existing) {
logger.info('Duplicate notification skipped — deduplication key already exists', {
requestId,
deduplicationKey: input.deduplicationKey,
existingId: existing.id,
});
return existing.id;
}
}

throw err;
}
}

/**
Expand Down Expand Up @@ -902,6 +938,7 @@ export class ScheduledNotificationRepository {
priority: row.priority,
metadata: row.metadata,
nextRetryAt: row.next_retry_at ? new Date(row.next_retry_at) : null,
deduplicationKey: row.deduplication_key ?? null,
};
}
}
2 changes: 2 additions & 0 deletions listener/src/types/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,8 @@ export interface RetrySchedulerOptions {
multiplier: number;
maxDelayMs: number;
jitter: boolean;
/** Request timeout for outbound webhook delivery (ms). Default: 10 000. */
webhookTimeoutMs: number;
/**
* Retry-policy ceiling on total attempts. Mirrors `RetrySchedulerConfig`;
* `undefined` leaves each notification's own `maxRetries` in control.
Expand Down
4 changes: 4 additions & 0 deletions listener/src/types/scheduled-notification.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,8 @@ export interface ScheduledNotification {
cancellationReason?: string | null;
/** When the next retry should be attempted (null if not retrying). */
nextRetryAt?: Date | null;
/** Stable logical key used to prevent duplicate delivery of the same notification. */
deduplicationKey?: string | null;
}

export interface CreateScheduledNotificationInput {
Expand All @@ -51,6 +53,7 @@ export interface CreateScheduledNotificationInput {
targetRecipient: string;
executeAt: Date;
maxRetries?: number;
deduplicationKey?: string;
eventId?: string;
contractAddress?: string;
priority?: number;
Expand Down Expand Up @@ -80,6 +83,7 @@ export interface ScheduledNotificationRow {
priority: number;
metadata: string | null;
next_retry_at: string | null;
deduplication_key: string | null;
}

export interface NotificationExecutionLog {
Expand Down
2 changes: 1 addition & 1 deletion listener/src/utils/benchmark-utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -235,7 +235,7 @@ export function exportMetricsToJSON(metrics: BenchmarkMetrics[], outputPath: str
}

export class ResourceMonitor {
private intervalId?: NodeJS.Timer;
private intervalId?: ReturnType<typeof setInterval>;
private samples: Array<{ timestamp: number; memory: NodeJS.MemoryUsage; cpu: NodeJS.CpuUsage }> = [];

start(intervalMs: number = 1000): void {
Expand Down
Loading