Skip to content
Draft
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
32 changes: 32 additions & 0 deletions apps/api/db/migrations/029_add_atomic_swap_dispute_bridge.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
-- Migration: 029_add_atomic_swap_dispute_bridge.sql
-- Description: Cross-Ledger Settlement Time-Lock Atomic Swap Dispute Bridge (#446)

DO $$ BEGIN
CREATE TYPE swap_dispute_state AS ENUM (
'ACTIVE',
'SECRET_EXTRACTED',
'REFUND_CLAIMABLE',
'RESOLVED'
);
EXCEPTION
WHEN duplicate_object THEN NULL;
END $$;

CREATE TABLE IF NOT EXISTS atomic_swap_dispute_bridges (
swap_id VARCHAR(64) PRIMARY KEY,
initiator_address VARCHAR(56) NOT NULL,
counterparty_address VARCHAR(56) NOT NULL,
secret_hash VARCHAR(64) NOT NULL,
secret_preimage VARCHAR(64),
expiration_ledger BIGINT NOT NULL,
state swap_dispute_state NOT NULL DEFAULT 'ACTIVE',
execution_proof TEXT,
resolved_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX IF NOT EXISTS idx_atomic_swap_dispute_bridges_state ON atomic_swap_dispute_bridges(state);
CREATE INDEX IF NOT EXISTS idx_atomic_swap_dispute_bridges_expiration_ledger ON atomic_swap_dispute_bridges(expiration_ledger);
CREATE INDEX IF NOT EXISTS idx_atomic_swap_dispute_bridges_initiator ON atomic_swap_dispute_bridges(initiator_address);
CREATE INDEX IF NOT EXISTS idx_atomic_swap_dispute_bridges_counterparty ON atomic_swap_dispute_bridges(counterparty_address);
29 changes: 29 additions & 0 deletions apps/api/db/migrations/030_add_distributed_webhook_pipeline.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
-- 030_add_distributed_webhook_pipeline.sql
-- Distributed Multi-Node Webhook Event Delivery Engine & Dead-Letter Queue (DLQ) Recovery System

CREATE TYPE webhook_delivery_status AS ENUM ('QUEUED', 'DELIVERED', 'FAILED', 'DEAD_LETTER');

CREATE TABLE webhook_endpoints (
endpoint_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
user_id VARCHAR(64) NOT NULL,
target_url TEXT NOT NULL,
secret_key VARCHAR(64) NOT NULL,
is_active BOOLEAN NOT NULL DEFAULT TRUE,
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);

CREATE TABLE webhook_delivery_logs (
delivery_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
endpoint_id UUID NOT NULL REFERENCES webhook_endpoints(endpoint_id) ON DELETE CASCADE,
event_type VARCHAR(64) NOT NULL,
payload JSONB NOT NULL,
signature_header VARCHAR(64) NOT NULL,
attempt_count INT NOT NULL DEFAULT 0,
status webhook_delivery_status NOT NULL DEFAULT 'QUEUED',
last_response_code INT NULL,
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);

CREATE INDEX IF NOT EXISTS idx_webhook_endpoints_user_id ON webhook_endpoints(user_id);
CREATE INDEX IF NOT EXISTS idx_webhook_delivery_logs_endpoint_id ON webhook_delivery_logs(endpoint_id);
CREATE INDEX IF NOT EXISTS idx_webhook_delivery_logs_status ON webhook_delivery_logs(status);
22 changes: 21 additions & 1 deletion apps/api/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -45,8 +45,13 @@ import { collateralRoutes } from "./routes/collateral.js";
import { CollateralGuardStore } from "./lib/collateralGuard.js";
import { multisigEscrowRoutes } from "./routes/multisig-escrow.js";
import { MultisigEscrowStore } from "./lib/multisigEscrowStore.js";
import { getChatInfrastructure } from "./lib/chat-infrastructure.js";
import { juryArbitrationRoutes } from "./routes/jury-arbitration.js";
import { webhooksRoutes } from "./routes/webhooks.js";
import { WebhookStore } from "./lib/webhook-store.js";
import { swapDisputeRoutes } from "./routes/swap-dispute.js";
import { SwapDisputeStore } from "./lib/workers/swapDisputeWorker.js";



const MAX_PAYMENTS_CACHE = 10000;
const usedPayments = new Map<string, number>();
Expand Down Expand Up @@ -443,3 +448,18 @@ app.register(multisigEscrowRoutes, {
// (#404) Decentralized Jury Dispute Arbitration: commit-reveal voting,
// VRF juror selection, and automated escrow resolution with stake slashing.
app.register(juryArbitrationRoutes, { prefix: "/api/v1" });
// (#445) Distributed Multi-Node Webhook Event Delivery Engine & DLQ Recovery System.
export const webhookStore = new WebhookStore(pgPool ?? undefined);
app.register(webhooksRoutes, {
prefix: "/api/v1",
store: webhookStore,
});
// (#446) Cross-Ledger Settlement Time-Lock Atomic Swap Dispute Bridge
export const swapDisputeStore = new SwapDisputeStore(pgPool ?? undefined);
app.register(swapDisputeRoutes, {
prefix: "/api/v1",
store: swapDisputeStore,
});



32 changes: 32 additions & 0 deletions apps/api/src/db/migrations/029_add_atomic_swap_dispute_bridge.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
-- Migration: 029_add_atomic_swap_dispute_bridge.sql
-- Description: Cross-Ledger Settlement Time-Lock Atomic Swap Dispute Bridge (#446)

DO $$ BEGIN
CREATE TYPE swap_dispute_state AS ENUM (
'ACTIVE',
'SECRET_EXTRACTED',
'REFUND_CLAIMABLE',
'RESOLVED'
);
EXCEPTION
WHEN duplicate_object THEN NULL;
END $$;

CREATE TABLE IF NOT EXISTS atomic_swap_dispute_bridges (
swap_id VARCHAR(64) PRIMARY KEY,
initiator_address VARCHAR(56) NOT NULL,
counterparty_address VARCHAR(56) NOT NULL,
secret_hash VARCHAR(64) NOT NULL,
secret_preimage VARCHAR(64),
expiration_ledger BIGINT NOT NULL,
state swap_dispute_state NOT NULL DEFAULT 'ACTIVE',
execution_proof TEXT,
resolved_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX IF NOT EXISTS idx_atomic_swap_dispute_bridges_state ON atomic_swap_dispute_bridges(state);
CREATE INDEX IF NOT EXISTS idx_atomic_swap_dispute_bridges_expiration_ledger ON atomic_swap_dispute_bridges(expiration_ledger);
CREATE INDEX IF NOT EXISTS idx_atomic_swap_dispute_bridges_initiator ON atomic_swap_dispute_bridges(initiator_address);
CREATE INDEX IF NOT EXISTS idx_atomic_swap_dispute_bridges_counterparty ON atomic_swap_dispute_bridges(counterparty_address);
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
-- 030_add_distributed_webhook_pipeline.sql
-- Distributed Multi-Node Webhook Event Delivery Engine & Dead-Letter Queue (DLQ) Recovery System

CREATE TYPE webhook_delivery_status AS ENUM ('QUEUED', 'DELIVERED', 'FAILED', 'DEAD_LETTER');

CREATE TABLE webhook_endpoints (
endpoint_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
user_id VARCHAR(64) NOT NULL,
target_url TEXT NOT NULL,
secret_key VARCHAR(64) NOT NULL,
is_active BOOLEAN NOT NULL DEFAULT TRUE,
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);

CREATE TABLE webhook_delivery_logs (
delivery_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
endpoint_id UUID NOT NULL REFERENCES webhook_endpoints(endpoint_id) ON DELETE CASCADE,
event_type VARCHAR(64) NOT NULL,
payload JSONB NOT NULL,
signature_header VARCHAR(64) NOT NULL,
attempt_count INT NOT NULL DEFAULT 0,
status webhook_delivery_status NOT NULL DEFAULT 'QUEUED',
last_response_code INT NULL,
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);

CREATE INDEX IF NOT EXISTS idx_webhook_endpoints_user_id ON webhook_endpoints(user_id);
CREATE INDEX IF NOT EXISTS idx_webhook_delivery_logs_endpoint_id ON webhook_delivery_logs(endpoint_id);
CREATE INDEX IF NOT EXISTS idx_webhook_delivery_logs_status ON webhook_delivery_logs(status);
14 changes: 5 additions & 9 deletions apps/api/src/lib/kms/aws-kms-driver.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { createHash } from "node:crypto";
import type { KmsDriver, KmsSignRequest, KmsSignResult } from "./kms-driver.interface.js";

/**
Expand All @@ -22,17 +23,12 @@ export class AwsKmsDriver implements KmsDriver {
// Deterministic mock: hash(payload+keyId) expanded to 64 bytes.
// Real implementation: const cmd = new SignCommand({ KeyId, Message, MessageType:"DIGEST", SigningAlgorithm:"ECDSA_SHA_512" })
const seed = `${request.keyId}:${request.payloadHex}`;
const sig = await mockEd25519(seed);
const sig = mockEd25519(seed);
return { signatureHex: sig, keyId: request.keyId };
}
}

async function mockEd25519(seed: string): Promise<string> {
const enc = new TextEncoder().encode(seed);
const hash = await crypto.subtle.digest("SHA-512", enc);
const bytes = new Uint8Array(hash);
// SHA-512 is 64 bytes already — hex-encode it as mock signature
return Array.from(bytes)
.map((b) => b.toString(16).padStart(2, "0"))
.join("");
function mockEd25519(seed: string): string {
return createHash("sha512").update(seed).digest("hex");
}

13 changes: 5 additions & 8 deletions apps/api/src/lib/kms/gcp-kms-driver.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { createHash } from "node:crypto";
import type { KmsDriver, KmsSignRequest, KmsSignResult } from "./kms-driver.interface.js";

/**
Expand All @@ -19,16 +20,12 @@ export class GcpKmsDriver implements KmsDriver {
throw new Error("GcpKmsDriver: payloadHex must be hex");
}
const seed = `gcp:${request.keyId}:${request.payloadHex}`;
const sig = await mockEd25519(seed);
const sig = mockEd25519(seed);
return { signatureHex: sig, keyId: request.keyId };
}
}

async function mockEd25519(seed: string): Promise<string> {
const enc = new TextEncoder().encode(seed);
const hash = await crypto.subtle.digest("SHA-512", enc);
const bytes = new Uint8Array(hash);
return Array.from(bytes)
.map((b) => b.toString(16).padStart(2, "0"))
.join("");
function mockEd25519(seed: string): string {
return createHash("sha512").update(seed).digest("hex");
}

13 changes: 5 additions & 8 deletions apps/api/src/lib/kms/vault-kms-driver.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { createHash } from "node:crypto";
import type { KmsDriver, KmsSignRequest, KmsSignResult } from "./kms-driver.interface.js";

/**
Expand All @@ -19,16 +20,12 @@ export class VaultKmsDriver implements KmsDriver {
throw new Error("VaultKmsDriver: payloadHex must be hex");
}
const seed = `vault:${request.keyId}:${request.payloadHex}`;
const sig = await mockEd25519(seed);
const sig = mockEd25519(seed);
return { signatureHex: sig, keyId: request.keyId };
}
}

async function mockEd25519(seed: string): Promise<string> {
const enc = new TextEncoder().encode(seed);
const hash = await crypto.subtle.digest("SHA-512", enc);
const bytes = new Uint8Array(hash);
return Array.from(bytes)
.map((b) => b.toString(16).padStart(2, "0"))
.join("");
function mockEd25519(seed: string): string {
return createHash("sha512").update(seed).digest("hex");
}

5 changes: 5 additions & 0 deletions apps/api/src/lib/stellar-timeout.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,10 @@ import { describe, expect, it, vi, beforeEach } from "vitest";
// Hoist RPC server mocks so they are available before the module import.
// ---------------------------------------------------------------------------
const h = vi.hoisted(() => {
const nodeCrypto = require("node:crypto");
if (!globalThis.crypto) {
Object.defineProperty(globalThis, "crypto", { value: nodeCrypto.webcrypto });
}
const preparedTx = {
sign: () => {},
hash: () => Buffer.from("00".repeat(32), "hex"),
Expand All @@ -20,6 +24,7 @@ const h = vi.hoisted(() => {
};
});


vi.mock("@stellar/stellar-sdk/rpc", async (importOriginal) => {
const actual = await importOriginal<typeof import("@stellar/stellar-sdk/rpc")>();
return {
Expand Down
69 changes: 68 additions & 1 deletion apps/api/src/lib/stellar.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
import nodeCrypto from "node:crypto";
if (!globalThis.crypto) {
Object.defineProperty(globalThis, "crypto", { value: nodeCrypto.webcrypto });
}
import {
Address,
BASE_FEE,
Expand All @@ -14,11 +18,12 @@ import {
Account,
} from "@stellar/stellar-sdk";
import { Server, Api, assembleTransaction } from "@stellar/stellar-sdk/rpc";
export { RpcTimeoutError } from "./rpc-errors.js";
import { RpcTimeoutError } from "./rpc-errors.js";
export { RpcTimeoutError };
import { createHash } from "node:crypto";
import nacl from "tweetnacl";


// Re-export commonly used SDK types and constants
export { BASE_FEE, Keypair, Operation, TransactionBuilder, xdr, Account, nativeToScVal, scValToNative };
export { Server, Api, assembleTransaction };
Expand Down Expand Up @@ -1807,3 +1812,65 @@ export async function getRotationProposal(
);
return proposal ?? null;
}

export interface ClaimAtomicSwapRefundParams {
contractId: string;
swapId: string;
}

export interface ReleaseAtomicSwapParams {
contractId: string;
swapId: string;
secretHex: string;
}

/**
* Claim refund for an atomic swap whose expiration ledger has passed.
*/
export async function claimAtomicSwapRefund(
params: ClaimAtomicSwapRefundParams,
): Promise<{ hash: string }> {
const signer = loadSignerKeypair();
const txHash = await invokeContract(
params.contractId,
"refund",
[nativeToScVal(Buffer.from(params.swapId, "hex"), { type: "bytes" })],
signer,
);
return { hash: typeof txHash === "string" ? txHash : "tx_success" };
}

/**
* Release atomic swap funds by providing the revealed secret preimage.
*/
export async function releaseAtomicSwap(
params: ReleaseAtomicSwapParams,
): Promise<{ hash: string }> {
const signer = loadSignerKeypair();
const txHash = await invokeContract(
params.contractId,
"release",
[
nativeToScVal(Buffer.from(params.swapId, "hex"), { type: "bytes" }),
nativeToScVal(Buffer.from(params.secretHex, "hex"), { type: "bytes" }),
],
signer,
);
return { hash: typeof txHash === "string" ? txHash : "tx_success" };
}

/**
* Read-only accessor for an atomic swap trade's on-chain state.
*/
export async function getAtomicSwapTrade(
contractId: string,
swapId: string,
): Promise<any | null> {
const trade = await simulateContractRead<any | null>(
contractId,
"get_trade",
[nativeToScVal(Buffer.from(swapId, "hex"), { type: "bytes" })],
);
return trade ?? null;
}

Loading
Loading