From 041a00c151c606b6212d39fe4fec0f40331a277a Mon Sep 17 00:00:00 2001 From: Justice Date: Sat, 26 Sep 2026 07:55:11 +0100 Subject: [PATCH] feat(transactions): add blockchain recovery and confirmation tracking --- __tests__/api/tx-recovery.test.ts | 40 +++ __tests__/tx-recovery/service.test.ts | 75 ++++++ app/api/tx-recovery/[hash]/route.ts | 48 ++++ app/api/tx-recovery/poll/route.ts | 26 ++ app/api/tx-recovery/submit/route.ts | 52 ++++ env.example | 12 + .../010_transaction_recovery_system.sql | 55 ++++ lib/tx-recovery/index.ts | 2 + lib/tx-recovery/service.ts | 252 ++++++++++++++++++ lib/tx-recovery/types.ts | 73 +++++ package.json | 1 + scripts/tx-recovery-worker.ts | 57 ++++ 12 files changed, 693 insertions(+) create mode 100644 __tests__/api/tx-recovery.test.ts create mode 100644 __tests__/tx-recovery/service.test.ts create mode 100644 app/api/tx-recovery/[hash]/route.ts create mode 100644 app/api/tx-recovery/poll/route.ts create mode 100644 app/api/tx-recovery/submit/route.ts create mode 100644 lib/db/migrations/010_transaction_recovery_system.sql create mode 100644 lib/tx-recovery/index.ts create mode 100644 lib/tx-recovery/service.ts create mode 100644 lib/tx-recovery/types.ts create mode 100644 scripts/tx-recovery-worker.ts diff --git a/__tests__/api/tx-recovery.test.ts b/__tests__/api/tx-recovery.test.ts new file mode 100644 index 0000000..5c5884e --- /dev/null +++ b/__tests__/api/tx-recovery.test.ts @@ -0,0 +1,40 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest' +import { NextRequest } from 'next/server' + +const { query } = vi.hoisted(() => ({ query: vi.fn() })) +vi.mock('@/lib/db', () => ({ sql: { query } })) +vi.mock('@/lib/auth/middleware', () => ({ + withAuth: (handler: Function) => (request: NextRequest) => handler(request, { walletAddress: 'GTEST' }), +})) + +import { POST as submit } from '@/app/api/tx-recovery/submit/route' +import { POST as poll } from '@/app/api/tx-recovery/poll/route' + +describe('transaction recovery routes', () => { + beforeEach(() => { + query.mockReset() + process.env.TX_RECOVERY_WORKER_SECRET = 'test-secret' + }) + + it('returns 422 for invalid submission input', async () => { + const response = await submit(new NextRequest('http://local/api/tx-recovery/submit', { + method: 'POST', body: JSON.stringify({ txHash: 'bad', actionType: 'escrow_fund' }), + })) + expect(response.status).toBe(422) + expect((await response.json()).code).toBe('INVALID_TRANSACTION') + }) + + it('protects the polling endpoint', async () => { + const response = await poll(new NextRequest('http://local/api/tx-recovery/poll', { method: 'POST' })) + expect(response.status).toBe(401) + }) + + it('accepts an authorized empty polling cycle', async () => { + query.mockResolvedValueOnce([]) + const response = await poll(new NextRequest('http://local/api/tx-recovery/poll', { + method: 'POST', headers: { authorization: 'Bearer test-secret' }, body: '{}', + })) + expect(response.status).toBe(200) + expect((await response.json()).data.processed).toBe(0) + }) +}) diff --git a/__tests__/tx-recovery/service.test.ts b/__tests__/tx-recovery/service.test.ts new file mode 100644 index 0000000..e691991 --- /dev/null +++ b/__tests__/tx-recovery/service.test.ts @@ -0,0 +1,75 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest' + +const { query } = vi.hoisted(() => ({ query: vi.fn() })) +vi.mock('@/lib/db', () => ({ sql: { query } })) + +import { TxRecoveryService } from '@/lib/tx-recovery/service' + +const hash = 'a'.repeat(64) +const contractId = '11111111-1111-4111-8111-111111111111' +const userId = '22222222-2222-4222-8222-222222222222' +const row = (overrides: Record = {}) => ({ + id: '33333333-3333-4333-8333-333333333333', tx_hash: hash, network: 'stellar', + contract_id: contractId, milestone_id: null, user_id: userId, action_type: 'escrow_fund', + status: 'pending', retry_count: 0, max_retries: 3, last_polled_at: null, + next_poll_at: '2026-09-26T07:35:00Z', confirmed_at: null, error_code: null, + error_message: null, metadata: {}, created_at: '2026-09-26T07:30:00Z', + updated_at: '2026-09-26T07:30:00Z', ...overrides, +}) + +describe('TxRecoveryService', () => { + beforeEach(() => query.mockReset()) + + it('rejects malformed transaction hashes', async () => { + const service = new TxRecoveryService() + await expect(service.submitTransaction({ txHash: 'bad', network: 'stellar', actionType: 'escrow_fund', contractId, userId })) + .rejects.toMatchObject({ code: 'INVALID_TX_HASH', httpStatus: 422 }) + }) + + it('registers a transaction and pending audit event in one statement', async () => { + query.mockResolvedValueOnce([row()]) + const service = new TxRecoveryService() + const result = await service.submitTransaction({ txHash: hash, network: 'stellar', actionType: 'escrow_fund', contractId, userId }) + expect(result.isDuplicate).toBe(false) + expect(query.mock.calls[0][0]).toContain('transaction_lifecycle_events') + expect(query.mock.calls[0][0]).toContain('ON CONFLICT (tx_hash) DO NOTHING') + }) + + it('rejects replay of a hash with different context', async () => { + query.mockResolvedValueOnce([]).mockResolvedValueOnce([row({ contract_id: contractId })]) + const service = new TxRecoveryService() + await expect(service.submitTransaction({ + txHash: hash, network: 'stellar', actionType: 'contract_deploy', contractId, userId, + })).rejects.toMatchObject({ code: 'TX_HASH_CONTEXT_CONFLICT', httpStatus: 409 }) + }) + + it('keeps an unconfirmed transaction pending and schedules a retry', async () => { + query.mockResolvedValueOnce([]) + const service = new TxRecoveryService({ verifier: async () => ({ status: 'pending' }) }) + await expect(service.verifyAndSyncTransaction(service['mapRow'](row()))).resolves.toBe('pending') + expect(query.mock.calls[0][0]).toContain("status = 'pending'") + }) + + it('expires a transaction when the retry budget is exhausted', async () => { + query.mockResolvedValueOnce([{ id: row().id }]) + const service = new TxRecoveryService({ verifier: async () => ({ status: 'pending' }) }) + await expect(service.verifyAndSyncTransaction(service['mapRow'](row({ retry_count: 2 })))).resolves.toBe('expired') + expect(query.mock.calls[0][1][1]).toBe('expired') + }) + + it('atomically gates domain sync and audit on the pending-to-success transition', async () => { + query.mockResolvedValueOnce([{ id: row().id }]) + const service = new TxRecoveryService({ verifier: async () => ({ status: 'success', ledger: 42 }) }) + await expect(service.verifyAndSyncTransaction(service['mapRow'](row()))).resolves.toBe('success') + const statement = query.mock.calls[0][0] as string + expect(statement).toContain("status='pending'") + expect(statement).toContain('domain_update AS') + expect(statement).toContain('transaction_lifecycle_events') + }) + + it('returns a structured not-found error', async () => { + query.mockResolvedValueOnce([]) + await expect(new TxRecoveryService().getTransactionStatus(hash)) + .rejects.toMatchObject({ code: 'TX_NOT_FOUND', httpStatus: 404 }) + }) +}) diff --git a/app/api/tx-recovery/[hash]/route.ts b/app/api/tx-recovery/[hash]/route.ts new file mode 100644 index 0000000..20adb60 --- /dev/null +++ b/app/api/tx-recovery/[hash]/route.ts @@ -0,0 +1,48 @@ +export const dynamic = 'force-dynamic' + +import { NextRequest, NextResponse } from 'next/server' +import { txRecoveryService, TxRecoveryError } from '@/lib/tx-recovery' + +export async function GET( + _request: NextRequest, + { params }: { params: Promise<{ hash: string }> } +) { + try { + const { hash } = await params + const cleanHash = (hash || '').trim().toLowerCase() + + if (!/^[0-9a-f]{64}$/.test(cleanHash)) { + return NextResponse.json( + { error: 'Invalid transaction hash format. Must be a 64-character hex string.', code: 'INVALID_HASH' }, + { status: 422 } + ) + } + + const txStatus = await txRecoveryService.getTransactionStatus(cleanHash) + + if (txStatus.status === 'expired') { + return NextResponse.json( + { + error: 'Transaction processing timed out or expired without blockchain confirmation.', + code: 'TX_EXPIRED_TIMEOUT', + data: txStatus, + }, + { status: 408 } + ) + } + + return NextResponse.json({ + data: txStatus, + }) + } catch (error) { + if (error instanceof TxRecoveryError) return NextResponse.json( + { error: error.message, code: error.code }, { status: error.httpStatus } + ) + + console.error('[api/tx-recovery/[hash]] Error:', error) + return NextResponse.json( + { error: 'Internal server error', code: 'TX_QUERY_FAILED' }, + { status: 500 } + ) + } +} diff --git a/app/api/tx-recovery/poll/route.ts b/app/api/tx-recovery/poll/route.ts new file mode 100644 index 0000000..565919e --- /dev/null +++ b/app/api/tx-recovery/poll/route.ts @@ -0,0 +1,26 @@ +export const dynamic = 'force-dynamic' + +import { timingSafeEqual } from 'crypto' +import { NextRequest, NextResponse } from 'next/server' +import { txRecoveryService } from '@/lib/tx-recovery' + +function authorized(request: NextRequest): boolean { + const expected = process.env.TX_RECOVERY_WORKER_SECRET + const provided = request.headers.get('authorization')?.replace(/^Bearer\s+/i, '') + if (!expected || !provided) return false + const a = Buffer.from(expected) + const b = Buffer.from(provided) + return a.length === b.length && timingSafeEqual(a, b) +} + +export async function POST(request: NextRequest) { + if (!authorized(request)) return NextResponse.json({ error: 'Unauthorized', code: 'AUTH_REQUIRED' }, { status: 401 }) + try { + const body = await request.json().catch(() => ({})) + const batchSize = Number.isInteger(body.batchSize) ? body.batchSize : 50 + return NextResponse.json({ data: await txRecoveryService.pollPendingQueue(batchSize) }) + } catch (error) { + console.error('[api/tx-recovery/poll] Error:', error) + return NextResponse.json({ error: 'Polling cycle failed', code: 'TX_POLL_FAILED' }, { status: 500 }) + } +} diff --git a/app/api/tx-recovery/submit/route.ts b/app/api/tx-recovery/submit/route.ts new file mode 100644 index 0000000..df076ea --- /dev/null +++ b/app/api/tx-recovery/submit/route.ts @@ -0,0 +1,52 @@ +export const dynamic = 'force-dynamic' + +import { NextRequest, NextResponse } from 'next/server' +import { z } from 'zod' +import { withAuth } from '@/lib/auth/middleware' +import { sql } from '@/lib/db' +import { txActionTypes, txRecoveryService, TxRecoveryError } from '@/lib/tx-recovery' + +const schema = z.object({ + txHash: z.string().trim().regex(/^[0-9a-fA-F]{64}$/), + network: z.enum(['stellar', 'soroban']).default('stellar'), + actionType: z.enum(txActionTypes), + contractId: z.string().uuid().optional(), + milestoneId: z.string().uuid().optional(), + metadata: z.record(z.unknown()).default({}), + maxRetries: z.number().int().min(1).max(50).default(10), +}) + +export const POST = withAuth(async (request: NextRequest, auth) => { + let input: z.infer + try { + input = schema.parse(await request.json()) + } catch (error) { + return NextResponse.json( + { error: error instanceof z.ZodError ? error.issues : 'Request body must be valid JSON', code: 'INVALID_TRANSACTION' }, + { status: 422 } + ) + } + const users = await sql.query('SELECT id FROM users WHERE wallet_address=$1 LIMIT 1', [auth.walletAddress]) as Array<{ id: string }> + if (!users.length) return NextResponse.json({ error: 'Authenticated wallet has no platform account', code: 'USER_NOT_FOUND' }, { status: 401 }) + const userId = users[0].id + if (input.contractId || input.milestoneId) { + const allowed = await sql.query( + `SELECT 1 FROM contracts c LEFT JOIN milestones m ON m.contract_id=c.id + WHERE (c.id=$1::uuid OR m.id=$2::uuid) AND (c.client_id=$3::uuid OR c.freelancer_id=$3::uuid) LIMIT 1`, + [input.contractId ?? null, input.milestoneId ?? null, userId] + ) as Array> + if (!allowed.length) return NextResponse.json({ error: 'Transaction context is not accessible', code: 'FORBIDDEN_CONTEXT' }, { status: 403 }) + } + try { + const result = await txRecoveryService.submitTransaction({ ...input, userId, txHash: input.txHash.toLowerCase() }) + const data = await txRecoveryService.getTransactionStatus(input.txHash) + return NextResponse.json( + { message: result.isDuplicate ? 'Transaction is already tracked.' : 'Transaction submitted for tracking.', isDuplicate: result.isDuplicate, data }, + { status: result.isDuplicate ? 409 : 201 } + ) + } catch (error) { + if (error instanceof TxRecoveryError) return NextResponse.json({ error: error.message, code: error.code }, { status: error.httpStatus }) + console.error('[api/tx-recovery/submit] Error:', error) + return NextResponse.json({ error: 'Failed to submit transaction', code: 'TX_SUBMIT_FAILED' }, { status: 500 }) + } +}) diff --git a/env.example b/env.example index b4fa902..de25651 100644 --- a/env.example +++ b/env.example @@ -14,6 +14,18 @@ JWT_SECRET=replace_with_a_long_random_secret_minimum_32_characters # Horizon RPC endpoint. Defaults to testnet if not set. STELLAR_HORIZON_URL=https://horizon-testnet.stellar.org +# Soroban JSON-RPC endpoint used for contract transaction confirmation. +SOROBAN_RPC_URL=https://soroban-testnet.stellar.org + +# Explorer transaction prefix; change to /public/tx for mainnet. +STELLAR_EXPLORER_TX_URL=https://stellar.expert/explorer/testnet/tx + +# Shared bearer secret for POST /api/tx-recovery/poll. +TX_RECOVERY_WORKER_SECRET=replace_with_a_long_random_secret + +# Optional worker loop interval in milliseconds (default: 10000). +# TX_POLL_INTERVAL_MS=10000 + # Network passphrase — must match the Horizon endpoint above. # Testnet : "Test SDF Network ; September 2015" # Mainnet : "Public Global Stellar Network ; September 2015" diff --git a/lib/db/migrations/010_transaction_recovery_system.sql b/lib/db/migrations/010_transaction_recovery_system.sql new file mode 100644 index 0000000..c530ee2 --- /dev/null +++ b/lib/db/migrations/010_transaction_recovery_system.sql @@ -0,0 +1,55 @@ +-- Durable state and append-only audit history for Stellar/Soroban transactions. +DO $$ BEGIN + CREATE TYPE tx_recovery_status AS ENUM ('pending', 'success', 'failed', 'expired'); +EXCEPTION WHEN duplicate_object THEN NULL; +END $$; + +DO $$ BEGIN + CREATE TYPE tx_network AS ENUM ('stellar', 'soroban'); +EXCEPTION WHEN duplicate_object THEN NULL; +END $$; + +CREATE TABLE IF NOT EXISTS tracked_transactions ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + tx_hash TEXT NOT NULL UNIQUE CHECK (tx_hash ~ '^[0-9a-f]{64}$'), + network tx_network NOT NULL DEFAULT 'stellar', + contract_id UUID REFERENCES contracts (id) ON DELETE SET NULL, + milestone_id UUID REFERENCES milestones (id) ON DELETE SET NULL, + user_id UUID NOT NULL REFERENCES users (id) ON DELETE RESTRICT, + action_type TEXT NOT NULL CHECK (action_type IN ( + 'escrow_fund', 'milestone_submit', 'milestone_approve', 'payment_release', + 'refund', 'dispute_raise', 'dispute_resolve', 'contract_deploy' + )), + status tx_recovery_status NOT NULL DEFAULT 'pending', + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), + max_retries INTEGER NOT NULL DEFAULT 10 CHECK (max_retries BETWEEN 1 AND 50), + last_polled_at TIMESTAMPTZ, + next_poll_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + confirmed_at TIMESTAMPTZ, + error_code TEXT, + error_message TEXT, + metadata JSONB NOT NULL DEFAULT '{}', + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +CREATE INDEX IF NOT EXISTS idx_tracked_transactions_queue + ON tracked_transactions (next_poll_at) WHERE status = 'pending'; +CREATE INDEX IF NOT EXISTS idx_tracked_transactions_contract + ON tracked_transactions (contract_id, created_at DESC) WHERE contract_id IS NOT NULL; +CREATE INDEX IF NOT EXISTS idx_tracked_transactions_milestone + ON tracked_transactions (milestone_id, created_at DESC) WHERE milestone_id IS NOT NULL; + +CREATE TABLE IF NOT EXISTS transaction_lifecycle_events ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + tracked_transaction_id UUID NOT NULL REFERENCES tracked_transactions (id) ON DELETE CASCADE, + status tx_recovery_status NOT NULL, + code TEXT, + message TEXT NOT NULL, + metadata JSONB NOT NULL DEFAULT '{}', + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + UNIQUE (tracked_transaction_id, status) +); + +CREATE INDEX IF NOT EXISTS idx_transaction_lifecycle_events_transaction + ON transaction_lifecycle_events (tracked_transaction_id, created_at DESC); diff --git a/lib/tx-recovery/index.ts b/lib/tx-recovery/index.ts new file mode 100644 index 0000000..256d3e0 --- /dev/null +++ b/lib/tx-recovery/index.ts @@ -0,0 +1,2 @@ +export * from './types' +export * from './service' diff --git a/lib/tx-recovery/service.ts b/lib/tx-recovery/service.ts new file mode 100644 index 0000000..24cf88c --- /dev/null +++ b/lib/tx-recovery/service.ts @@ -0,0 +1,252 @@ +import { Horizon, rpc } from '@stellar/stellar-sdk' +import { sql } from '@/lib/db' +import type { + PollBatchResult, SubmitTransactionParams, TrackedTransaction, + TransactionStatusResponse, TxActionType, TxNetwork, TxRecoveryStatus, +} from './types' +import { TxRecoveryError } from './types' + +type ChainResult = + | { status: 'pending' } + | { status: 'success'; ledger?: number; createdAt?: string | number; fee?: string } + | { status: 'failed'; code: string; message: string } + +type Verifier = (tx: TrackedTransaction) => Promise +const HASH = /^[0-9a-f]{64}$/ +const SELECT_COLUMNS = `id, tx_hash, network, contract_id, milestone_id, user_id, + action_type, status, retry_count, max_retries, last_polled_at, next_poll_at, + confirmed_at, error_code, error_message, metadata, created_at, updated_at` + +export class TxRecoveryService { + private readonly verifyOverride?: Verifier + private readonly horizonUrl: string + private readonly sorobanRpcUrl: string + private readonly explorerBaseUrl: string + + constructor(options: { + horizonUrl?: string + sorobanRpcUrl?: string + explorerBaseUrl?: string + verifier?: Verifier + } = {}) { + this.horizonUrl = options.horizonUrl ?? process.env.STELLAR_HORIZON_URL ?? 'https://horizon-testnet.stellar.org' + this.sorobanRpcUrl = options.sorobanRpcUrl ?? process.env.SOROBAN_RPC_URL ?? 'https://soroban-testnet.stellar.org' + this.explorerBaseUrl = options.explorerBaseUrl ?? process.env.STELLAR_EXPLORER_TX_URL ?? 'https://stellar.expert/explorer/testnet/tx' + this.verifyOverride = options.verifier + } + + async submitTransaction(params: SubmitTransactionParams): Promise<{ transaction: TrackedTransaction; isDuplicate: boolean }> { + const txHash = params.txHash.trim().toLowerCase() + if (!HASH.test(txHash)) throw new TxRecoveryError('INVALID_TX_HASH', 'Transaction hash must be 64 hexadecimal characters.', 422) + this.validateContext(params.actionType, params.contractId, params.milestoneId) + + const inserted = await sql.query( + `WITH inserted AS ( + INSERT INTO tracked_transactions + (tx_hash, network, contract_id, milestone_id, user_id, action_type, max_retries, metadata) + VALUES ($1, $2::tx_network, $3::uuid, $4::uuid, $5::uuid, $6, $7, $8::jsonb) + ON CONFLICT (tx_hash) DO NOTHING + RETURNING ${SELECT_COLUMNS} + ), audited AS ( + INSERT INTO transaction_lifecycle_events (tracked_transaction_id, status, message) + SELECT id, 'pending', 'Transaction submitted for confirmation tracking' FROM inserted + ) SELECT * FROM inserted`, + [txHash, params.network, params.contractId ?? null, params.milestoneId ?? null, + params.userId, params.actionType, params.maxRetries ?? 10, JSON.stringify(params.metadata ?? {})] + ) as Array> + + if (inserted.length) return { transaction: this.mapRow(inserted[0]), isDuplicate: false } + + const existing = await sql.query( + `SELECT ${SELECT_COLUMNS} FROM tracked_transactions WHERE tx_hash = $1 LIMIT 1`, [txHash] + ) as Array> + if (!existing.length) throw new TxRecoveryError('TX_SUBMIT_FAILED', 'Transaction registration failed.', 500) + const transaction = this.mapRow(existing[0]) + const sameContext = transaction.network === params.network && + transaction.actionType === params.actionType && + transaction.contractId === (params.contractId ?? null) && + transaction.milestoneId === (params.milestoneId ?? null) && + transaction.userId === params.userId + if (!sameContext) throw new TxRecoveryError('TX_HASH_CONTEXT_CONFLICT', 'Transaction hash is already bound to a different operation.', 409) + return { transaction, isDuplicate: true } + } + + async getTransactionStatus(txHashInput: string): Promise { + const txHash = txHashInput.trim().toLowerCase() + if (!HASH.test(txHash)) throw new TxRecoveryError('INVALID_TX_HASH', 'Transaction hash must be 64 hexadecimal characters.', 422) + const rows = await sql.query( + `SELECT ${SELECT_COLUMNS} FROM tracked_transactions WHERE tx_hash = $1 LIMIT 1`, [txHash] + ) as Array> + if (!rows.length) throw new TxRecoveryError('TX_NOT_FOUND', 'Transaction is not tracked.', 404) + const tx = this.mapRow(rows[0]) + return { + hash: tx.txHash, network: tx.network, actionType: tx.actionType, + contractId: tx.contractId, milestoneId: tx.milestoneId, status: tx.status, + retryCount: tx.retryCount, maxRetries: tx.maxRetries, + explorerUrl: `${this.explorerBaseUrl}/${tx.txHash}`, + errorCode: tx.errorCode, errorMessage: tx.errorMessage, + confirmedAt: tx.confirmedAt, createdAt: tx.createdAt, updatedAt: tx.updatedAt, + } + } + + async pollPendingQueue(batchSize = 50): Promise { + const limit = Number.isInteger(batchSize) ? Math.min(100, Math.max(1, batchSize)) : 50 + // Claim rows for five minutes so overlapping workers do not poll the same hash. + const rows = await sql.query( + `WITH due AS ( + SELECT id FROM tracked_transactions + WHERE status = 'pending' AND next_poll_at <= NOW() + ORDER BY next_poll_at LIMIT $1 FOR UPDATE SKIP LOCKED + ) + UPDATE tracked_transactions t SET next_poll_at = NOW() + INTERVAL '5 minutes', updated_at = NOW() + FROM due WHERE t.id = due.id RETURNING ${SELECT_COLUMNS.split(', ').map(c => `t.${c.trim()}`).join(', ')}`, + [limit] + ) as Array> + const result: PollBatchResult = { processed: rows.length, succeeded: 0, failed: 0, expired: 0, stillPending: 0, errors: [] } + for (const row of rows) { + const tx = this.mapRow(row) + try { + const status = await this.verifyAndSyncTransaction(tx) + if (status === 'success') result.succeeded++ + else if (status === 'failed') result.failed++ + else if (status === 'expired') result.expired++ + else result.stillPending++ + } catch (error) { + result.errors.push(`Tx ${tx.txHash}: ${error instanceof Error ? error.message : String(error)}`) + } + } + return result + } + + async verifyAndSyncTransaction(txOrHash: TrackedTransaction | string): Promise { + const tx = typeof txOrHash === 'string' ? await this.getTracked(txOrHash) : txOrHash + if (tx.status !== 'pending') return tx.status + try { + const result = this.verifyOverride ? await this.verifyOverride(tx) : await this.verifyOnChain(tx) + if (result.status === 'success') return this.syncSuccess(tx, result) + if (result.status === 'failed') return this.transitionTerminal(tx, 'failed', result.code, result.message) + return this.scheduleRetry(tx, 'TX_NOT_CONFIRMED', 'Transaction is not yet present in a confirmed ledger') + } catch (error) { + return this.scheduleRetry(tx, 'RPC_UNAVAILABLE', error instanceof Error ? error.message : 'Blockchain RPC request failed') + } + } + + private async verifyOnChain(tx: TrackedTransaction): Promise { + if (tx.network === 'soroban') { + const response = await new rpc.Server(this.sorobanRpcUrl).getTransaction(tx.txHash) + if (response.txHash.toLowerCase() !== tx.txHash) return { status: 'failed', code: 'RPC_HASH_MISMATCH', message: 'Soroban RPC returned a different transaction hash' } + if (response.status === rpc.Api.GetTransactionStatus.SUCCESS) return { status: 'success', ledger: response.ledger, createdAt: response.createdAt } + if (response.status === rpc.Api.GetTransactionStatus.FAILED) return { status: 'failed', code: 'TX_FAILED_ON_CHAIN', message: 'Soroban transaction execution failed' } + return { status: 'pending' } + } + try { + const response = await new Horizon.Server(this.horizonUrl).transactions().transaction(tx.txHash).call() + const record = response as unknown as { hash: string; successful: boolean; ledger?: number; created_at?: string; fee_charged?: string } + if (record.hash.toLowerCase() !== tx.txHash) return { status: 'failed', code: 'RPC_HASH_MISMATCH', message: 'Horizon returned a different transaction hash' } + return record.successful + ? { status: 'success', ledger: record.ledger, createdAt: record.created_at, fee: record.fee_charged } + : { status: 'failed', code: 'TX_FAILED_ON_CHAIN', message: 'Stellar transaction execution failed' } + } catch (error) { + if (this.httpStatus(error) === 404) return { status: 'pending' } + throw error + } + } + + private async scheduleRetry(tx: TrackedTransaction, code: string, message: string): Promise { + const retryCount = tx.retryCount + 1 + if (retryCount >= tx.maxRetries) return this.transitionTerminal(tx, 'expired', 'TX_EXPIRED_TIMEOUT', `${message}; retry limit reached`) + const base = Math.min(300, 3 * 2 ** tx.retryCount) + const delay = Math.round(base * (0.75 + Math.random() * 0.5)) + await sql.query( + `UPDATE tracked_transactions SET retry_count = $2, last_polled_at = NOW(), + next_poll_at = NOW() + ($3 * INTERVAL '1 second'), error_code = $4, + error_message = $5, updated_at = NOW() WHERE id = $1::uuid AND status = 'pending'`, + [tx.id, retryCount, delay, code, message] + ) + return 'pending' + } + + private async syncSuccess(tx: TrackedTransaction, chain: Extract): Promise { + const metadata = JSON.stringify({ ...tx.metadata, ledger: chain.ledger, onChainCreatedAt: chain.createdAt, fee: chain.fee }) + let domainCte = '' + if (tx.actionType === 'escrow_fund') domainCte = `, domain_update AS (UPDATE contracts SET escrow_status='funded', status='active', funded_at=NOW(), funding_tx_hash=$2, updated_at=NOW() WHERE id=$3::uuid AND EXISTS (SELECT 1 FROM transitioned))` + else if (tx.actionType === 'contract_deploy') domainCte = `, domain_update AS (UPDATE contracts SET contract_tx_hash=$2, updated_at=NOW() WHERE id=$3::uuid AND EXISTS (SELECT 1 FROM transitioned))` + else if (tx.actionType === 'dispute_raise') domainCte = `, domain_update AS (UPDATE contracts SET status='disputed', updated_at=NOW() WHERE id=$3::uuid AND EXISTS (SELECT 1 FROM transitioned))` + else { + const state: Partial> = { milestone_submit: 'submitted', milestone_approve: 'approved', payment_release: 'paid', refund: 'rejected' } + const next = state[tx.actionType] + if (tx.actionType === 'refund') domainCte = `, + domain_update AS (UPDATE milestones SET status='rejected', rejection_reason='Refund confirmed on-chain', updated_at=NOW() WHERE id=$4::uuid AND EXISTS (SELECT 1 FROM transitioned) RETURNING contract_id), + contract_refund AS (UPDATE contracts SET escrow_status='refunded', updated_at=NOW() WHERE id IN (SELECT contract_id FROM domain_update))` + else if (next) domainCte = `, domain_update AS (UPDATE milestones SET status='${next}', updated_at=NOW() WHERE id=$4::uuid AND EXISTS (SELECT 1 FROM transitioned))` + } + const rows = await sql.query( + `WITH transitioned AS ( + UPDATE tracked_transactions SET status='success', confirmed_at=NOW(), last_polled_at=NOW(), + error_code=NULL, error_message=NULL, metadata=$5::jsonb, updated_at=NOW() + WHERE id=$1::uuid AND status='pending' RETURNING id + )${domainCte}, audited AS ( + INSERT INTO transaction_lifecycle_events (tracked_transaction_id,status,message,metadata) + SELECT id,'success','Blockchain confirmation received',$5::jsonb FROM transitioned + ON CONFLICT (tracked_transaction_id,status) DO NOTHING + ) SELECT id FROM transitioned`, + [tx.id, tx.txHash, tx.contractId, tx.milestoneId, metadata] + ) as Array<{ id: string }> + return rows.length ? 'success' : (await this.getTracked(tx.txHash)).status + } + + private async transitionTerminal(tx: TrackedTransaction, status: 'failed' | 'expired', code: string, message: string): Promise { + const rows = await sql.query( + `WITH transitioned AS ( + UPDATE tracked_transactions SET status=$2::tx_recovery_status, retry_count=CASE WHEN $2='expired' THEN retry_count+1 ELSE retry_count END, + last_polled_at=NOW(), error_code=$3, error_message=$4, updated_at=NOW() + WHERE id=$1::uuid AND status='pending' RETURNING id + ), audited AS ( + INSERT INTO transaction_lifecycle_events (tracked_transaction_id,status,code,message) + SELECT id,$2::tx_recovery_status,$3,$4 FROM transitioned + ON CONFLICT (tracked_transaction_id,status) DO NOTHING + ) SELECT id FROM transitioned`, [tx.id, status, code, message] + ) as Array<{ id: string }> + return rows.length ? status : (await this.getTracked(tx.txHash)).status + } + + private async getTracked(hash: string): Promise { + const clean = hash.trim().toLowerCase() + if (!HASH.test(clean)) throw new TxRecoveryError('INVALID_TX_HASH', 'Transaction hash must be 64 hexadecimal characters.', 422) + const rows = await sql.query(`SELECT ${SELECT_COLUMNS} FROM tracked_transactions WHERE tx_hash=$1 LIMIT 1`, [clean]) as Array> + if (!rows.length) throw new TxRecoveryError('TX_NOT_FOUND', 'Transaction is not tracked.', 404) + return this.mapRow(rows[0]) + } + + private validateContext(action: TxActionType, contractId?: string, milestoneId?: string): void { + const milestoneActions: TxActionType[] = ['milestone_submit', 'milestone_approve', 'payment_release', 'refund'] + if (milestoneActions.includes(action) && !milestoneId) throw new TxRecoveryError('MILESTONE_REQUIRED', 'This action requires a milestoneId.', 422) + if (['escrow_fund', 'contract_deploy'].includes(action) && !contractId) throw new TxRecoveryError('CONTRACT_REQUIRED', 'This action requires a contractId.', 422) + if (action.startsWith('dispute_') && !contractId && !milestoneId) throw new TxRecoveryError('ENTITY_REQUIRED', 'Dispute actions require a contractId or milestoneId.', 422) + } + + private httpStatus(error: unknown): number | undefined { + return typeof error === 'object' && error !== null && 'response' in error + ? (error as { response?: { status?: number } }).response?.status : undefined + } + + private mapRow(row: Record): TrackedTransaction { + return { + id: String(row.id), txHash: String(row.tx_hash), network: String(row.network) as TxNetwork, + contractId: row.contract_id ? String(row.contract_id) : null, + milestoneId: row.milestone_id ? String(row.milestone_id) : null, + userId: String(row.user_id), actionType: String(row.action_type) as TxActionType, + status: String(row.status) as TxRecoveryStatus, retryCount: Number(row.retry_count), + maxRetries: Number(row.max_retries), lastPolledAt: this.iso(row.last_polled_at), + nextPollAt: this.iso(row.next_poll_at)!, confirmedAt: this.iso(row.confirmed_at), + errorCode: row.error_code ? String(row.error_code) : null, + errorMessage: row.error_message ? String(row.error_message) : null, + metadata: typeof row.metadata === 'object' && row.metadata ? row.metadata as Record : {}, + createdAt: this.iso(row.created_at)!, updatedAt: this.iso(row.updated_at)!, + } + } + + private iso(value: unknown): string | null { return value ? new Date(value as string | number | Date).toISOString() : null } +} + +export const txRecoveryService = new TxRecoveryService() diff --git a/lib/tx-recovery/types.ts b/lib/tx-recovery/types.ts new file mode 100644 index 0000000..5b2ae77 --- /dev/null +++ b/lib/tx-recovery/types.ts @@ -0,0 +1,73 @@ +export const txActionTypes = [ + 'escrow_fund', 'milestone_submit', 'milestone_approve', 'payment_release', + 'refund', 'dispute_raise', 'dispute_resolve', 'contract_deploy', +] as const + +export type TxActionType = (typeof txActionTypes)[number] +export type TxRecoveryStatus = 'pending' | 'success' | 'failed' | 'expired' +export type TxNetwork = 'stellar' | 'soroban' + +export interface SubmitTransactionParams { + txHash: string + actionType: TxActionType + network: TxNetwork + contractId?: string + milestoneId?: string + userId: string + metadata?: Record + maxRetries?: number +} + +export interface TrackedTransaction { + id: string + txHash: string + network: TxNetwork + contractId: string | null + milestoneId: string | null + userId: string + actionType: TxActionType + status: TxRecoveryStatus + retryCount: number + maxRetries: number + lastPolledAt: string | null + nextPollAt: string + confirmedAt: string | null + errorCode: string | null + errorMessage: string | null + metadata: Record + createdAt: string + updatedAt: string +} + +export interface TransactionStatusResponse { + hash: string + network: TxNetwork + actionType: TxActionType + contractId: string | null + milestoneId: string | null + status: TxRecoveryStatus + retryCount: number + maxRetries: number + explorerUrl: string + errorCode: string | null + errorMessage: string | null + confirmedAt: string | null + createdAt: string + updatedAt: string +} + +export interface PollBatchResult { + processed: number + succeeded: number + failed: number + expired: number + stillPending: number + errors: string[] +} + +export class TxRecoveryError extends Error { + constructor(public readonly code: string, message: string, public readonly httpStatus: number) { + super(message) + this.name = 'TxRecoveryError' + } +} diff --git a/package.json b/package.json index be669c5..ba1e576 100644 --- a/package.json +++ b/package.json @@ -13,6 +13,7 @@ "worker": "tsx scripts/worker.ts", "sync-worker": "tsx scripts/sync-worker.ts", "deadline-monitor": "tsx scripts/deadline-monitor.ts", + "tx-worker": "tsx scripts/tx-recovery-worker.ts", "migrate": "tsx scripts/migrate.ts", "build:production": "npm run migrate && next build", "start:production": "next start" diff --git a/scripts/tx-recovery-worker.ts b/scripts/tx-recovery-worker.ts new file mode 100644 index 0000000..a476b5b --- /dev/null +++ b/scripts/tx-recovery-worker.ts @@ -0,0 +1,57 @@ +import { txRecoveryService } from '@/lib/tx-recovery' +import * as dotenv from 'dotenv' + +dotenv.config() + +if (!process.env.DATABASE_URL) { + console.error('FATAL: DATABASE_URL is not set. Transaction recovery worker cannot connect to database.') + process.exit(1) +} + +const POLL_INTERVAL_MS = Number(process.env.TX_POLL_INTERVAL_MS) || 10_000 + +async function startTxRecoveryWorker() { + console.log('[TxRecoveryWorker] Starting Web3 Transaction Recovery & Confirmation Service...') + console.log(`[TxRecoveryWorker] Poll interval set to ${POLL_INTERVAL_MS}ms`) + + const runPollCycle = async () => { + try { + const result = await txRecoveryService.pollPendingQueue(50) + if (result.processed > 0) { + console.log( + `[TxRecoveryWorker] Polled ${result.processed} tx(s) — Succeeded: ${result.succeeded}, Failed: ${result.failed}, Expired: ${result.expired}, Still Pending: ${result.stillPending}` + ) + } + if (result.errors.length > 0) { + console.error('[TxRecoveryWorker] Errors during poll cycle:', result.errors) + } + } catch (err) { + console.error('[TxRecoveryWorker] Unhandled poll cycle error:', err) + } + } + + // Initial execution + await runPollCycle() + + // Recurring loop + const intervalHandle = setInterval(runPollCycle, POLL_INTERVAL_MS) + + // Heartbeat logger every 60s + setInterval(() => { + console.log(`[TxRecoveryWorker HEARTBEAT] ${new Date().toISOString()} — Daemon active`) + }, 60_000) + + const shutdown = () => { + console.log('[TxRecoveryWorker] Gracefully shutting down...') + clearInterval(intervalHandle) + process.exit(0) + } + + process.on('SIGINT', shutdown) + process.on('SIGTERM', shutdown) +} + +startTxRecoveryWorker().catch((err) => { + console.error('[FATAL TxRecoveryWorker ERROR]', err) + process.exit(1) +})