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
40 changes: 40 additions & 0 deletions __tests__/api/tx-recovery.test.ts
Original file line number Diff line number Diff line change
@@ -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)
})
})
75 changes: 75 additions & 0 deletions __tests__/tx-recovery/service.test.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown> = {}) => ({
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 })
})
})
48 changes: 48 additions & 0 deletions app/api/tx-recovery/[hash]/route.ts
Original file line number Diff line number Diff line change
@@ -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 }
)
}
}
26 changes: 26 additions & 0 deletions app/api/tx-recovery/poll/route.ts
Original file line number Diff line number Diff line change
@@ -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 })
}
}
52 changes: 52 additions & 0 deletions app/api/tx-recovery/submit/route.ts
Original file line number Diff line number Diff line change
@@ -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<typeof schema>
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<Record<string, unknown>>
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 })
}
})
12 changes: 12 additions & 0 deletions env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
55 changes: 55 additions & 0 deletions lib/db/migrations/010_transaction_recovery_system.sql
Original file line number Diff line number Diff line change
@@ -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);
2 changes: 2 additions & 0 deletions lib/tx-recovery/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
export * from './types'
export * from './service'
Loading
Loading