diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml
new file mode 100644
index 00000000..f2bfabde
--- /dev/null
+++ b/.github/workflows/ci.yml
@@ -0,0 +1,147 @@
+name: CI
+
+on:
+ pull_request:
+ branches: [ main ]
+
+jobs:
+ format-lint-typecheck:
+ name: Code Quality Checks
+ runs-on: ubuntu-latest
+
+ steps:
+ - name: Checkout code
+ uses: actions/checkout@v4
+ with:
+ fetch-depth: 0
+
+ - name: Setup Node.js
+ uses: actions/setup-node@v4
+ with:
+ node-version: '22'
+
+ - name: Install listener dependencies
+ working-directory: ./listener
+ run: npm install --ignore-scripts
+
+ - name: Install dashboard dependencies
+ working-directory: ./dashboard
+ run: npm install --ignore-scripts
+
+ - name: Identify changed source files
+ id: changed
+ shell: bash
+ run: |
+ git diff --name-only "${{ github.event.pull_request.base.sha }}" "${{ github.event.pull_request.head.sha }}" > changed-files.txt
+
+ listener_files=$(grep -E '^listener/src/.*\.ts$' changed-files.txt | sed 's#^listener/##' | tr '\n' ' ' || true)
+ dashboard_files=$(grep -E '^dashboard/src/.*\.(ts|tsx)$' changed-files.txt | sed 's#^dashboard/##' | tr '\n' ' ' || true)
+
+ echo "listener_files=$listener_files" >> "$GITHUB_OUTPUT"
+ echo "dashboard_files=$dashboard_files" >> "$GITHUB_OUTPUT"
+
+ - name: Format check listener changes
+ if: steps.changed.outputs.listener_files != ''
+ working-directory: ./listener
+ run: |
+ echo "Checking listener formatting:"
+ printf '%s\n' "${{ steps.changed.outputs.listener_files }}"
+ npx prettier --check ${{ steps.changed.outputs.listener_files }} --config ../.prettierrc
+
+ - name: Format check dashboard changes
+ if: steps.changed.outputs.dashboard_files != ''
+ working-directory: ./dashboard
+ run: |
+ echo "Checking dashboard formatting:"
+ printf '%s\n' "${{ steps.changed.outputs.dashboard_files }}"
+ npx prettier --check ${{ steps.changed.outputs.dashboard_files }} --config ../.prettierrc
+
+ - name: Lint listener changes
+ if: steps.changed.outputs.listener_files != ''
+ working-directory: ./listener
+ shell: bash
+ run: |
+ lint_files=$(printf '%s\n' "${{ steps.changed.outputs.listener_files }}" | tr ' ' '\n' | grep -v '^src/test-utils/' | grep -v '^src/__tests__/' | tr '\n' ' ' || true)
+
+ if [ -z "$lint_files" ]; then
+ echo "No changed listener files are included in the configured ESLint scope."
+ exit 0
+ fi
+
+ echo "Linting changed listener files:"
+ printf '%s\n' "$lint_files"
+ npx eslint $lint_files --max-warnings=0
+
+ - name: Lint dashboard changes
+ if: steps.changed.outputs.dashboard_files != ''
+ working-directory: ./dashboard
+ run: |
+ echo "Linting changed dashboard files:"
+ printf '%s\n' "${{ steps.changed.outputs.dashboard_files }}"
+ npx eslint ${{ steps.changed.outputs.dashboard_files }} --max-warnings=0
+
+ - name: Static analysis listener changes
+ if: steps.changed.outputs.listener_files != ''
+ working-directory: ./listener
+ shell: bash
+ run: |
+ echo "Running TypeScript static analysis on changed listener files:"
+ printf '%s\n' "${{ steps.changed.outputs.listener_files }}"
+
+ set +e
+ tsc_output=$(npx tsc --noEmit --skipLibCheck --target es2020 --module commonjs --moduleResolution node --esModuleInterop --types node,jest ${{ steps.changed.outputs.listener_files }} 2>&1)
+ tsc_status=$?
+ set -e
+
+ printf '%s\n' "$tsc_output"
+
+ changed_errors=""
+ for file in ${{ steps.changed.outputs.listener_files }}; do
+ file_errors=$(printf '%s\n' "$tsc_output" | grep -F "${file}(" || true)
+ if [ -n "$file_errors" ]; then
+ changed_errors="${changed_errors}${file_errors}"$'\n'
+ fi
+ done
+
+ if [ -n "$changed_errors" ]; then
+ echo "TypeScript errors found in changed listener files:"
+ printf '%s\n' "$changed_errors"
+ exit 1
+ fi
+
+ if [ "$tsc_status" -ne 0 ]; then
+ echo "TypeScript reported errors only in unchanged/imported files; those errors are outside this PR's changed-file scope."
+ fi
+
+ - name: Static analysis dashboard
+ if: steps.changed.outputs.dashboard_files != ''
+ working-directory: ./dashboard
+ shell: bash
+ run: |
+ echo "Running TypeScript static analysis for dashboard changes:"
+ printf '%s\n' "${{ steps.changed.outputs.dashboard_files }}"
+
+ set +e
+ tsc_output=$(npx tsc --noEmit -p tsconfig.json 2>&1)
+ tsc_status=$?
+ set -e
+
+ printf '%s\n' "$tsc_output"
+
+ changed_errors=""
+ for file in ${{ steps.changed.outputs.dashboard_files }}; do
+ file_errors=$(printf '%s\n' "$tsc_output" | grep -F "${file}(" || true)
+ if [ -n "$file_errors" ]; then
+ changed_errors="${changed_errors}${file_errors}"$'\n'
+ fi
+ done
+
+ if [ -n "$changed_errors" ]; then
+ echo "TypeScript errors found in changed dashboard files:"
+ printf '%s\n' "$changed_errors"
+ exit 1
+ fi
+
+ if [ "$tsc_status" -ne 0 ]; then
+ echo "TypeScript reported errors only in unchanged/imported files; those errors are outside this PR's changed-file scope."
+ fi
diff --git a/dashboard/src/pages/EventExplorerPage.test.tsx b/dashboard/src/pages/EventExplorerPage.test.tsx
new file mode 100644
index 00000000..90e80bbf
--- /dev/null
+++ b/dashboard/src/pages/EventExplorerPage.test.tsx
@@ -0,0 +1,114 @@
+import '@testing-library/jest-dom';
+import { render, screen, waitFor, act } from '@testing-library/react';
+import { EventExplorerPage } from './EventExplorerPage';
+import { useEventStore } from '../store/eventStore';
+import { generateMockEvents } from '../utils/eventData';
+import { fetchEvents } from '../services/eventsApi';
+
+jest.mock('../services/eventsApi', () => ({
+ fetchEvents: jest.fn(),
+ fetchStatus: jest.fn(() => Promise.resolve({ contracts: [] })),
+}));
+
+jest.mock('../services/wallet', () => ({
+ restoreWalletSession: jest.fn(() => Promise.resolve()),
+}));
+
+jest.mock('../components/WalletConnectButton', () => ({
+ WalletConnectButton: () =>
,
+}));
+
+const mockedFetchEvents = fetchEvents as jest.MockedFunction;
+
+describe('EventExplorerPage refresh states', () => {
+ beforeEach(() => {
+ jest.useFakeTimers();
+ useEventStore.setState({
+ events: [],
+ filters: {
+ search: '',
+ contractAddress: 'all',
+ eventType: 'all',
+ status: 'all',
+ dateFrom: '',
+ dateTo: '',
+ },
+ isLoading: false,
+ error: null,
+ lastFetchedAt: 0,
+ });
+ mockedFetchEvents.mockReset();
+ });
+
+ afterEach(() => {
+ jest.useRealTimers();
+ });
+
+ it('shows refresh indicator and preserves events during background refresh', async () => {
+ const initialEvents = generateMockEvents(3);
+ mockedFetchEvents.mockResolvedValueOnce(initialEvents);
+
+ render();
+
+ // Wait for initial load
+ await waitFor(() => {
+ expect(screen.queryByText(/Loading events/i)).not.toBeInTheDocument();
+ });
+ expect(screen.getAllByRole('row').length).toBeGreaterThan(1);
+
+ // Setup next fetch to be pending
+ let resolveRefresh!: (value: typeof initialEvents) => void;
+ mockedFetchEvents.mockReturnValueOnce(
+ new Promise((resolve) => {
+ resolveRefresh = resolve;
+ }),
+ );
+
+ // Advance time to trigger background refresh (15s)
+ act(() => {
+ jest.advanceTimersByTime(15000);
+ });
+
+ // Refresh indicator appears, events remain
+ await waitFor(() => {
+ expect(screen.getByText(/Refreshing events/i)).toBeInTheDocument();
+ });
+ expect(screen.getAllByRole('row').length).toBeGreaterThan(1); // Table still rendered
+
+ // Complete refresh
+ await act(async () => {
+ resolveRefresh(initialEvents);
+ });
+
+ // Indicator disappears
+ await waitFor(() => {
+ expect(screen.queryByText(/Refreshing events/i)).not.toBeInTheDocument();
+ });
+ });
+
+ it('shows appropriate error state when refresh fails, preserving existing events', async () => {
+ const initialEvents = generateMockEvents(3);
+ mockedFetchEvents.mockResolvedValueOnce(initialEvents);
+
+ render();
+
+ await waitFor(() => {
+ expect(screen.queryByText(/Loading events/i)).not.toBeInTheDocument();
+ });
+
+ // Setup refresh failure
+ mockedFetchEvents.mockRejectedValueOnce(new Error('Network error'));
+
+ // Trigger refresh
+ act(() => {
+ jest.advanceTimersByTime(15000);
+ });
+
+ // Error banner appears, events remain
+ await waitFor(() => {
+ expect(screen.getByText(/Refresh Error:/i)).toBeInTheDocument();
+ expect(screen.getByText(/Background refresh failed/i)).toBeInTheDocument();
+ });
+ expect(screen.getAllByRole('row').length).toBeGreaterThan(1);
+ });
+});
diff --git a/dashboard/src/pages/EventExplorerPage.tsx b/dashboard/src/pages/EventExplorerPage.tsx
index 615377d9..9219c359 100644
--- a/dashboard/src/pages/EventExplorerPage.tsx
+++ b/dashboard/src/pages/EventExplorerPage.tsx
@@ -9,7 +9,11 @@ import { NotificationDetailsDrawer } from '../components/NotificationDetailsDraw
import { IndexingHealthPanel } from '../components/IndexingHealthPanel';
import { NotificationHealthPanel } from '../components/NotificationHealthPanel';
import { EmptyState } from '../components/EmptyState';
-import { useEventFilters, useEventLoadingState, useFilteredEvents } from '../hooks/useEventSelectors';
+import {
+ useEventFilters,
+ useEventLoadingState,
+ useFilteredEvents,
+} from '../hooks/useEventSelectors';
import { useEventStore } from '../store/eventStore';
import { fetchEvents, fetchStatus, type ContractStatus } from '../services/eventsApi';
import { resolveIndexingHealthUrl } from '../services/indexingHealthApi';
@@ -47,6 +51,8 @@ export function EventExplorerPage() {
const [limit, setLimit] = useState(() => parseLimitParam(initialSearch));
const [selectedNotification, setSelectedNotification] = useState(null);
const [contractStatuses, setContractStatuses] = useState([]);
+ const [isRefreshing, setIsRefreshing] = useState(false);
+ const [refreshError, setRefreshError] = useState(null);
const setEvents = useEventStore((state) => state.setEvents);
const setLoading = useEventStore((state) => state.setLoading);
@@ -121,6 +127,8 @@ export function EventExplorerPage() {
// Poll for status updates so delivered/failed notifications are reflected
// without requiring a manual page refresh.
const intervalId = setInterval(async () => {
+ setIsRefreshing(true);
+ setRefreshError(null);
try {
const remoteEvents = await fetchEvents(API_URL);
if (!cancelled) {
@@ -129,8 +137,13 @@ export function EventExplorerPage() {
}
} catch {
if (!cancelled) {
+ setRefreshError('Background refresh failed');
markSyncFailure('Background refresh failed');
}
+ } finally {
+ if (!cancelled) {
+ setIsRefreshing(false);
+ }
}
}, POLL_INTERVAL_MS);
@@ -165,7 +178,7 @@ export function EventExplorerPage() {
const pageCount = useMemo(
() => Math.max(1, Math.ceil(filteredEvents.length / limit)),
- [filteredEvents.length, limit]
+ [filteredEvents.length, limit],
);
useEffect(() => {
@@ -176,7 +189,14 @@ export function EventExplorerPage() {
useEffect(() => {
setPage(1);
- }, [filters.search, filters.contractAddress, filters.eventType, filters.status, filters.dateFrom, filters.dateTo]);
+ }, [
+ filters.search,
+ filters.contractAddress,
+ filters.eventType,
+ filters.status,
+ filters.dateFrom,
+ filters.dateTo,
+ ]);
useEffect(() => {
if (typeof window === 'undefined') {
@@ -209,21 +229,21 @@ export function EventExplorerPage() {
}, [setSearch, setContractFilter, setEventTypeFilter, setStatusFilter, setDateFrom, setDateTo]);
const handleRetry = useCallback(async () => {
- setLoading(true);
- setError(null);
+ setIsRefreshing(true);
+ setRefreshError(null);
try {
const remoteEvents = await fetchEvents(API_URL);
setEvents(remoteEvents);
markSyncSuccess();
+ setError(null);
} catch {
- setEvents(generateMockEvents(DEFAULT_EVENT_COUNT));
- setError('Retry failed — still using demo event data.');
+ setRefreshError('Manual refresh failed');
markSyncFailure('Manual refresh failed');
} finally {
- setLoading(false);
+ setIsRefreshing(false);
}
- }, [markSyncFailure, markSyncSuccess, setError, setEvents, setLoading]);
+ }, [markSyncFailure, markSyncSuccess, setError, setEvents]);
const handleSelectEvent = useCallback((event: BlockchainEvent) => {
setSelectedNotification(event);
@@ -240,8 +260,8 @@ export function EventExplorerPage() {
Event Explorer
Smart Contract Event Log
- Browse Soroban contract events across registered contracts with filters,
- pagination, and copy-to-clipboard contract metadata.
+ Browse Soroban contract events across registered contracts with filters, pagination, and
+ copy-to-clipboard contract metadata.
@@ -254,13 +274,13 @@ export function EventExplorerPage() {
{contractStatuses.map((contract) => (
{contract.address}
-
+
{contract.paused ? 'PAUSED' : 'ACTIVE'}
{contract.error && (
-
- Error: {contract.error}
-
+
Error: {contract.error}
)}
))}
@@ -273,7 +293,7 @@ export function EventExplorerPage() {
- {error && (
+ {error && !filteredEvents.length && (
Error: {error}
@@ -284,15 +304,33 @@ export function EventExplorerPage() {
)}
+ {refreshError && filteredEvents.length > 0 && (
+
+
+ Refresh Error: {refreshError} — existing events are still displayed.
+
+
+
+ )}
+
Showing {fromIndex.toLocaleString()}–{toIndex.toLocaleString()} of{' '}
{filteredEvents.length.toLocaleString()} events
- {isLoading &&
Loading events…
}
+ {(isLoading || isRefreshing) && (
+
+ {isRefreshing ? 'Refreshing events…' : 'Loading events…'}
+
+ )}
- {isLoading ? (
+ {isLoading && filteredEvents.length === 0 ? (
) : currentPageEvents.length > 0 ? (
= 0
+ ) {
+ latencySum += r.deliveryLatencyMs;
+ latencyCount++;
+ }
}
const total = success + failure + retry + skipped;
const terminal = success + failure;
const successRate = terminal > 0 ? success / terminal : 0;
const averageDurationMs = durationCount > 0 ? durationSum / durationCount : 0;
+ const averageDeliveryLatencyMs = latencyCount > 0 ? latencySum / latencyCount : 0;
return {
total,
@@ -279,12 +294,11 @@ export class NotificationAnalyticsAggregator {
skipped,
successRate,
averageDurationMs,
+ averageDeliveryLatencyMs,
};
}
- private computeByType(
- visible: AnalyticsDeliveryRecord[],
- ): AnalyticsByTypeSnapshot[] {
+ private computeByType(visible: AnalyticsDeliveryRecord[]): AnalyticsByTypeSnapshot[] {
const map = new Map();
for (const r of visible) {
const entry = map.get(r.notificationType) ?? { total: 0, success: 0, failure: 0 };
@@ -309,9 +323,7 @@ export class NotificationAnalyticsAggregator {
return result;
}
- private computeByContract(
- visible: AnalyticsDeliveryRecord[],
- ): AnalyticsByContractSnapshot[] {
+ private computeByContract(visible: AnalyticsDeliveryRecord[]): AnalyticsByContractSnapshot[] {
const map = new Map();
for (const r of visible) {
if (!r.contractAddress) continue;
@@ -345,8 +357,7 @@ export class NotificationAnalyticsAggregator {
now: number,
): AnalyticsBucketSnapshot[] {
const newestBucketStart = Math.floor(now / this.bucketSizeMs) * this.bucketSizeMs;
- const oldestBucketStart =
- newestBucketStart - (this.maxBuckets - 1) * this.bucketSizeMs;
+ const oldestBucketStart = newestBucketStart - (this.maxBuckets - 1) * this.bucketSizeMs;
const buckets: AnalyticsBucketSnapshot[] = [];
const indexByStart = new Map();
@@ -360,6 +371,7 @@ export class NotificationAnalyticsAggregator {
retry: 0,
skipped: 0,
averageDurationMs: 0,
+ averageDeliveryLatencyMs: 0,
};
indexByStart.set(t, buckets.length);
buckets.push(snapshot);
@@ -369,9 +381,12 @@ export class NotificationAnalyticsAggregator {
let durationCountBucket = 0;
let durationBucketIdx = -1;
+ let latencySumBucket = 0;
+ let latencyCountBucket = 0;
+ let latencyBucketIdx = -1;
+
for (const r of visible) {
- const bucketStart =
- Math.floor(r.timestamp / this.bucketSizeMs) * this.bucketSizeMs;
+ const bucketStart = Math.floor(r.timestamp / this.bucketSizeMs) * this.bucketSizeMs;
const idx = indexByStart.get(bucketStart);
if (idx === undefined) continue;
@@ -398,14 +413,30 @@ export class NotificationAnalyticsAggregator {
averageDurationMs: durationSumBucket / durationCountBucket,
};
}
+
+ if (
+ r.outcome === 'success' &&
+ r.deliveryLatencyMs !== undefined &&
+ r.deliveryLatencyMs >= 0
+ ) {
+ if (idx !== latencyBucketIdx) {
+ latencySumBucket = 0;
+ latencyCountBucket = 0;
+ latencyBucketIdx = idx;
+ }
+ latencySumBucket += r.deliveryLatencyMs;
+ latencyCountBucket++;
+ buckets[idx] = {
+ ...buckets[idx],
+ averageDeliveryLatencyMs: latencySumBucket / latencyCountBucket,
+ };
+ }
}
return buckets;
}
- private computeErrorBreakdown(
- visible: AnalyticsDeliveryRecord[],
- ): Record {
+ private computeErrorBreakdown(visible: AnalyticsDeliveryRecord[]): Record {
const counts = new Map();
for (const r of visible) {
if (r.outcome !== 'failure') continue;
diff --git a/listener/src/services/notification-scheduler.ts b/listener/src/services/notification-scheduler.ts
index 8d0cda9b..b0c7b882 100644
--- a/listener/src/services/notification-scheduler.ts
+++ b/listener/src/services/notification-scheduler.ts
@@ -10,6 +10,10 @@ import { getWorkerManager } from './worker-manager';
import { getJobMonitor } from './job-monitor';
import { ProviderRegistry, getProviderRegistry } from './provider-registry';
import { verifyPayloadIntegrity } from '../utils/payload-integrity';
+import {
+ getNotificationAnalyticsAggregator,
+ NotificationAnalyticsAggregator,
+} from './notification-analytics-aggregator';
/**
* Background scheduler that processes scheduled notifications
@@ -32,13 +36,14 @@ export class NotificationScheduler {
* When not supplied the module-level singleton is used.
*/
private providerRegistry: ProviderRegistry;
+ private analytics: NotificationAnalyticsAggregator;
constructor(
repository: ScheduledNotificationRepository,
config: SchedulerConfig,
discordService?: DiscordNotificationService | null,
batchValidator?: BatchValidationService,
- providerRegistry?: ProviderRegistry
+ providerRegistry?: ProviderRegistry,
) {
this.repository = repository;
this.config = { retryDelayMs: 5_000, ...config };
@@ -46,6 +51,7 @@ export class NotificationScheduler {
this.processorId = config.processorId || uuidv4();
this.batchValidator = batchValidator ?? new BatchValidationService();
this.providerRegistry = providerRegistry ?? getProviderRegistry();
+ this.analytics = getNotificationAnalyticsAggregator();
}
/**
@@ -126,7 +132,7 @@ export class NotificationScheduler {
this.processorId,
this.config.lockTimeoutMs,
this.config.batchSize,
- requestId
+ requestId,
);
if (notifications.length === 0) {
@@ -140,7 +146,7 @@ export class NotificationScheduler {
}
const batchRejection = this.batchValidator.rejectIfInvalid(
- this.toValidationBatch(notifications)
+ this.toValidationBatch(notifications),
);
if (batchRejection) {
@@ -153,9 +159,11 @@ export class NotificationScheduler {
for (const notification of notifications) {
await this.repository.markAsFailedOrRetry(
notification.id!,
- new Error(`Batch validation failed: ${batchRejection.errors.map((e) => e.message).join('; ')}`),
+ new Error(
+ `Batch validation failed: ${batchRejection.errors.map((e) => e.message).join('; ')}`,
+ ),
notification.retryCount,
- notification.maxRetries
+ notification.maxRetries,
);
}
return;
@@ -180,7 +188,7 @@ export class NotificationScheduler {
notification.id!,
new Error('Scheduler shutting down'),
notification.retryCount,
- notification.maxRetries
+ notification.maxRetries,
);
}
return;
@@ -197,7 +205,7 @@ export class NotificationScheduler {
notification.id!,
new Error('Scheduler shutting down'),
notification.retryCount,
- notification.maxRetries
+ notification.maxRetries,
);
continue;
}
@@ -237,7 +245,7 @@ export class NotificationScheduler {
private async processNotification(
notification: ScheduledNotification,
requestId: string,
- jobId?: string
+ jobId?: string,
): Promise {
const startTime = Date.now();
const executionAttempt = notification.retryCount + 1;
@@ -268,7 +276,7 @@ export class NotificationScheduler {
notification.id!,
new Error('Not yet due for execution'),
notification.retryCount,
- notification.maxRetries
+ notification.maxRetries,
);
if (jobId) {
jobMonitor.failJob(jobId, 'Not yet due for execution', {
@@ -296,7 +304,9 @@ export class NotificationScheduler {
requestId,
id: notification.id,
});
- } else if (!verifyPayloadIntegrity(notification.payload, notification.payloadHash, secret)) {
+ } else if (
+ !verifyPayloadIntegrity(notification.payload, notification.payloadHash, secret)
+ ) {
logger.error('Payload integrity verification failed — rejecting notification', {
requestId,
id: notification.id,
@@ -306,7 +316,7 @@ export class NotificationScheduler {
notification.id!,
new Error('Payload integrity check failed: hash mismatch'),
notification.maxRetries, // exhaust retries — don't retry a tampered payload
- notification.maxRetries
+ notification.maxRetries,
);
if (jobId) {
jobMonitor.failJob(jobId, 'Payload integrity check failed: hash mismatch', {
@@ -339,11 +349,25 @@ export class NotificationScheduler {
});
}
+ const deliveryLatencyMs = notification.createdAt
+ ? Date.now() - new Date(notification.createdAt).getTime()
+ : undefined;
+
+ this.analytics.record({
+ notificationType: notification.notificationType,
+ contractAddress: notification.contractAddress ?? undefined,
+ outcome: 'success',
+ durationMs,
+ deliveryLatencyMs,
+ timestamp: Date.now(),
+ });
+
logger.info('Notification delivered successfully', {
requestId,
id: notification.id,
type: notification.notificationType,
durationMs,
+ deliveryLatencyMs,
});
} else {
throw new Error('Notification delivery returned false');
@@ -365,6 +389,15 @@ export class NotificationScheduler {
});
}
+ this.analytics.record({
+ notificationType: notification.notificationType,
+ contractAddress: notification.contractAddress ?? undefined,
+ outcome: 'failure',
+ durationMs,
+ errorReason: (error as Error).message,
+ timestamp: Date.now(),
+ });
+
const willRetry = notification.retryCount + 1 < notification.maxRetries;
const nextRetryAt = willRetry
? new Date(Date.now() + (this.config.retryDelayMs ?? 5_000))
@@ -375,7 +408,7 @@ export class NotificationScheduler {
error as Error,
notification.retryCount,
notification.maxRetries,
- nextRetryAt
+ nextRetryAt,
);
await this.repository.logExecution({
@@ -403,7 +436,7 @@ export class NotificationScheduler {
*/
private async executeNotification(
notification: ScheduledNotification,
- requestId: string
+ requestId: string,
): Promise {
const payload = JSON.parse(notification.payload);
const type = notification.notificationType;
@@ -442,33 +475,33 @@ export class NotificationScheduler {
case 'discord':
if (!this.discordService) {
throw new Error(
- 'Discord service not configured and no Discord provider registered in the registry'
+ 'Discord service not configured and no Discord provider registered in the registry',
);
}
return await this.discordService.sendEventNotification(
payload.event,
payload.contractConfig,
- `scheduler-${notification.id}-${requestId}`
+ `scheduler-${notification.id}-${requestId}`,
);
case 'webhook':
throw new Error(
- 'Webhook delivery not yet implemented. Register a WebhookNotificationProvider in the ProviderRegistry.'
+ 'Webhook delivery not yet implemented. Register a WebhookNotificationProvider in the ProviderRegistry.',
);
case 'email':
throw new Error(
- 'Email delivery not yet implemented. Register an email NotificationProvider in the ProviderRegistry.'
+ 'Email delivery not yet implemented. Register an email NotificationProvider in the ProviderRegistry.',
);
case 'sms':
throw new Error(
- 'SMS delivery not yet implemented. Register an SMS NotificationProvider in the ProviderRegistry.'
+ 'SMS delivery not yet implemented. Register an SMS NotificationProvider in the ProviderRegistry.',
);
default:
throw new Error(
- `Unsupported notification type: "${type}". Register a provider for this type in the ProviderRegistry.`
+ `Unsupported notification type: "${type}". Register a provider for this type in the ProviderRegistry.`,
);
}
}
diff --git a/listener/src/test-utils/event-fixtures.test.ts b/listener/src/test-utils/event-fixtures.test.ts
new file mode 100644
index 00000000..701835dc
--- /dev/null
+++ b/listener/src/test-utils/event-fixtures.test.ts
@@ -0,0 +1,38 @@
+import { EventFixtures } from './event-fixtures';
+import { xdr } from '@stellar/stellar-sdk';
+
+describe('EventFixtures', () => {
+ it('generates a valid event', () => {
+ const event = EventFixtures.valid();
+ expect(event.id).toBe('evt-valid-1');
+ expect(event.type).toBe('contract');
+ expect(event.topic).toBeDefined();
+ expect(event.value).toBeDefined();
+ });
+
+ it('generates duplicate events', () => {
+ const events = EventFixtures.duplicate();
+ expect(events.length).toBe(2);
+ expect(events[0].id).toBe(events[1].id);
+ expect(events[0].txHash).toBe(events[1].txHash);
+ });
+
+ it('generates missing fields event', () => {
+ const event = EventFixtures.missingFields();
+ expect(event.type).toBe('contract');
+ expect(event.id).toBeUndefined();
+ });
+
+ it('generates unsupported version event', () => {
+ const event = EventFixtures.unsupportedVersion();
+ const val = event.value.map();
+ expect(val).toBeDefined();
+ const versionEntry = val?.find((entry) => entry.key().sym().toString() === 'version');
+ expect(versionEntry?.val().u32()).toBe(999);
+ });
+
+ it('generates malformed payload event', () => {
+ const event = EventFixtures.malformedPayload();
+ expect(event.ledger).toBe(-1);
+ });
+});
diff --git a/listener/src/test-utils/event-fixtures.ts b/listener/src/test-utils/event-fixtures.ts
new file mode 100644
index 00000000..f59e8614
--- /dev/null
+++ b/listener/src/test-utils/event-fixtures.ts
@@ -0,0 +1,54 @@
+import * as StellarSDK from '@stellar/stellar-sdk';
+import { xdr } from '@stellar/stellar-sdk';
+
+const defaultContract = 'CBIELTK6YBZJU5UP2WWQEUCYKLPU6AUNZ2BQ4WWFEIE3USCIU6KPNBAM';
+
+export const EventFixtures = {
+ valid: (
+ overrides: Partial = {},
+ ): StellarSDK.rpc.Api.EventResponse =>
+ ({
+ id: 'evt-valid-1',
+ type: 'contract',
+ ledger: 1000,
+ ledgerClosedAt: '2026-06-22T00:00:00Z',
+ transactionIndex: 0,
+ operationIndex: 0,
+ inSuccessfulContractCall: true,
+ txHash: 'tx-valid-abc',
+ topic: [xdr.scvSymbol('test_event')],
+ value: xdr.scvString('valid payload'),
+ contractId: { contractId: () => defaultContract } as any, // mock for contractId if needed
+ ...overrides,
+ }) as StellarSDK.rpc.Api.EventResponse,
+
+ duplicate: (): StellarSDK.rpc.Api.EventResponse[] => {
+ const base = EventFixtures.valid({ id: 'evt-dup-1', txHash: 'tx-dup-1' });
+ return [base, { ...base }];
+ },
+
+ missingFields: (): Partial => ({
+ // Missing id, topic, value, etc.
+ type: 'contract',
+ ledger: 1001,
+ inSuccessfulContractCall: true,
+ }),
+
+ unsupportedVersion: (): StellarSDK.rpc.Api.EventResponse =>
+ EventFixtures.valid({
+ id: 'evt-unsupported-1',
+ value: xdr.scvMap([
+ new xdr.ScMapEntry({
+ key: xdr.scvSymbol('version'),
+ val: xdr.scvU32(999), // Unsupported version
+ }),
+ ]),
+ }),
+
+ malformedPayload: (): StellarSDK.rpc.Api.EventResponse =>
+ EventFixtures.valid({
+ id: 'evt-malformed-1',
+ // validateEventPayload explicitly rejects ledger < 0 as malformed
+ ledger: -1,
+ }),
+};