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
106 changes: 106 additions & 0 deletions src/jobs/process-webhook-retries.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,106 @@
import { logger } from "../utils/logger.js";
import { redis } from "../config/redis.js";

const QUEUE_KEY = "chainlearn:retry:webhooks";
const MAX_RETRIES = 5;
const BASE_DELAY_MS = 1000;
const MAX_DELAY_MS = 3600000;

export interface WebhookAttempt {
id: string;
webhookId: string;
event: string;
payload: unknown;
statusCode: number | null;
errorMessage: string | null;
retryCount: number;
nextRetryAt: string | null;
succeededAt: string | null;
failedAt: string | null;
}

function isTransientFailure(statusCode: number | null, error: string | null): boolean {
if (!statusCode) return true;
return statusCode >= 500 || statusCode === 429;
}

function calculateNextRetry(retryCount: number): Date {
const delay = Math.min(BASE_DELAY_MS * Math.pow(2, retryCount), MAX_DELAY_MS);
const jitter = delay * (0.5 + Math.random() * 0.5);
return new Date(Date.now() + jitter);
}

export async function enqueueWebhookRetry(attempt: WebhookAttempt): Promise<void> {
await redis.lpush(QUEUE_KEY, JSON.stringify(attempt));
logger.info({ webhookId: attempt.webhookId, event: attempt.event }, "Webhook retry enqueued");
}

export async function dequeueWebhookRetry(): Promise<WebhookAttempt | null> {
const raw = await redis.rpop(QUEUE_KEY);
if (!raw) return null;
return JSON.parse(raw) as WebhookAttempt;
}

export function shouldRetry(attempt: WebhookAttempt): boolean {
if (attempt.retryCount >= MAX_RETRIES) return false;
if (!isTransientFailure(attempt.statusCode, attempt.errorMessage)) return false;
return true;
}

export function markForRetry(attempt: WebhookAttempt): WebhookAttempt {
const nextRetry = calculateNextRetry(attempt.retryCount);
return {
...attempt,
retryCount: attempt.retryCount + 1,
nextRetryAt: nextRetry.toISOString(),
};
}

export function markAsFailed(attempt: WebhookAttempt): WebhookAttempt {
return {
...attempt,
failedAt: new Date().toISOString(),
};
}

let processorRunning = false;
let processorTimer: ReturnType<typeof setTimeout> | null = null;

export async function startWebhookRetryProcessor(
processFn: (attempt: WebhookAttempt) => Promise<boolean>
): Promise<void> {
if (processorRunning) return;
processorRunning = true;

const tick = async () => {
if (!processorRunning) return;
try {
const attempt = await dequeueWebhookRetry();
if (attempt) {
if (!shouldRetry(attempt)) {
const failed = markAsFailed(attempt);
logger.warn(
{ webhookId: failed.webhookId, retryCount: failed.retryCount },
"Webhook permanently failed"
);
return;
}

const success = await processFn(attempt);
if (!success) {
const retried = markForRetry(attempt);
await enqueueWebhookRetry(retried);
}
}
} catch (err) {
logger.error({ err }, "Webhook retry processor tick failed");
}
if (processorRunning) {
processorTimer = setTimeout(tick, 30000);
}
};

tick();
import { processWebhookRetries } from "../services/webhook-dispatcher.js";

let retryProcessorRunning = false;
Expand Down Expand Up @@ -29,6 +131,10 @@ export async function startWebhookRetryProcessor(): Promise<void> {
}

export function stopWebhookRetryProcessor(): void {
processorRunning = false;
if (processorTimer) {
clearTimeout(processorTimer);
processorTimer = null;
retryProcessorRunning = false;
retryProcessorGeneration++;
if (retryProcessorTimer) {
Expand Down
36 changes: 36 additions & 0 deletions tests/unit/services/reconcile-pending-rewards.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,41 @@
import { describe, it, expect, vi, beforeEach } from "vitest";

const { mockUpdate, mockSelect } = vi.hoisted(() => ({
mockUpdate: vi.fn(),
mockSelect: vi.fn(),
}));

vi.mock("../../../src/config/database.js", () => ({
db: {
update: mockUpdate,
select: mockSelect,
},
}));

vi.mock("../../../src/utils/logger.js", () => ({
logger: {
info: vi.fn(),
warn: vi.fn(),
error: vi.fn(),
},
}));

vi.mock("../../../src/config/redis.js", () => ({
redis: {
get: vi.fn(),
set: vi.fn(),
del: vi.fn(),
},
}));

describe("reconcile-pending-rewards", () => {
beforeEach(() => {
vi.clearAllMocks();
});

it("should have proper mock structure", () => {
expect(mockUpdate).toBeDefined();
expect(mockSelect).toBeDefined();
// ─── Mocks ────────────────────────────────────────────────────────────────────

const mockUpdate = vi.fn();
Expand Down
121 changes: 121 additions & 0 deletions tests/unit/stellar/transactions.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,16 @@
import { describe, it, expect, vi, beforeEach } from "vitest";
import { StellarError } from "../../../src/utils/errors.js";

const VALID_PUBLIC_KEY = "GAXK5L7G7U7YGH4WGONDXK3ZO2RMGET6ANTTSBCN7CNNRWBFUO5NGXAM";
const VALID_CONTRACT_ID = "CAZ7YFMK5FGL3NF2E2NJ3Q7XXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXX";

const { mockSubmitTransaction, mockGetAccount, mockSimulateTransaction, mockInvalidate, mockGetNextSequence, mockContractCall } = vi.hoisted(() => ({
mockSubmitTransaction: vi.fn(),
mockGetAccount: vi.fn(),
mockSimulateTransaction: vi.fn(),
mockInvalidate: vi.fn(),
mockGetNextSequence: vi.fn(),
mockContractCall: vi.fn(),
import { describe, it, expect, vi, beforeEach, type Mock } from "vitest";

vi.mock("../../../src/config/stellar.js", () => ({
Expand All @@ -12,6 +25,15 @@ vi.mock("../../../src/config/index.js", () => ({

vi.mock("../../../src/stellar/client.js", () => ({
stellarClient: {
submitTransaction: mockSubmitTransaction,
getAccount: mockGetAccount,
},
}));

vi.mock("../../../src/stellar/sequence-cache.js", () => ({
sequenceCache: {
getNextSequence: mockGetNextSequence,
invalidate: mockInvalidate,
submitTransaction: vi.fn(),
},
}));
Expand All @@ -28,6 +50,105 @@ vi.mock("../../../src/stellar/sequence-cache.js", () => ({
}));

vi.mock("../../../src/utils/account-lock.js", () => ({
withAccountLock: vi.fn((_id: string, fn: () => Promise<any>) => fn()),
}));

vi.mock("../../../src/config/stellar.js", () => ({
getPlatformKeypair: vi.fn(() => ({
publicKey: () => VALID_PUBLIC_KEY,
sign: vi.fn(),
})),
getNetworkPassphrase: vi.fn(() => "Test SDF Network ; September 2015"),
getSorobanServer: vi.fn(() => ({
simulateTransaction: mockSimulateTransaction,
})),
}));

vi.mock("../../../src/config/index.js", () => ({
config: {
STELLAR_QUIZ_CONTRACT_ID: VALID_CONTRACT_ID,
},
}));

vi.mock("../../../src/utils/logger.js", () => ({
logger: {
info: vi.fn(),
warn: vi.fn(),
error: vi.fn(),
},
}));

vi.mock("@stellar/stellar-sdk", async () => {
const actual = await vi.importActual<typeof import("@stellar/stellar-sdk")>("@stellar/stellar-sdk");
return {
...actual,
Contract: vi.fn().mockImplementation(() => ({
call: mockContractCall,
})),
TransactionBuilder: vi.fn().mockImplementation(() => ({
addOperation: vi.fn().mockReturnThis(),
setTimeout: vi.fn().mockReturnThis(),
build: vi.fn().mockReturnValue({
sign: vi.fn(),
}),
})),
Account: vi.fn(),
BASE_FEE: "100",
Operation: {
payment: vi.fn(),
},
Asset: {
native: vi.fn(),
},
rpc: {
Api: {
isSimulationError: vi.fn().mockReturnValue(false),
},
assembleTransaction: vi.fn().mockReturnValue({
build: vi.fn().mockReturnValue({
sign: vi.fn(),
}),
}),
},
};
});

import { invokeContract } from "../../../src/stellar/transactions.js";

describe("invokeContract", () => {
beforeEach(() => {
vi.clearAllMocks();
mockGetNextSequence.mockResolvedValue("100");
mockSimulateTransaction.mockResolvedValue({
transactionData: {},
results: [],
});
mockContractCall.mockReturnValue([]);
});

it("should submit a transaction successfully", async () => {
mockSubmitTransaction.mockResolvedValue({ hash: "abc123" });
mockSimulateTransaction.mockResolvedValue({
transactionData: { toXDR: () => "AAAA" },
results: [{}],
});

const result = await invokeContract(VALID_CONTRACT_ID, "method", []);
expect(result).toBe("abc123");
expect(mockSubmitTransaction).toHaveBeenCalledTimes(1);
});

it("should retry on sequence conflict and invalidate cache", async () => {
mockSubmitTransaction
.mockRejectedValueOnce(new StellarError('tx failed: ["tx_bad_seq"]'))
.mockResolvedValueOnce({ hash: "def456" });

const result = await invokeContract(VALID_CONTRACT_ID, "method", []);
expect(result).toBe("def456");
expect(mockInvalidate).toHaveBeenCalledWith(VALID_PUBLIC_KEY);
expect(mockSubmitTransaction).toHaveBeenCalledTimes(2);
});
});
withAccountLock: vi.fn((publicKey, fn) => fn()),
}));

Expand Down
Loading