diff --git a/crates/api/src/routes/sponsorship.rs b/crates/api/src/routes/sponsorship.rs index 87d9c41..bf56d54 100644 --- a/crates/api/src/routes/sponsorship.rs +++ b/crates/api/src/routes/sponsorship.rs @@ -62,7 +62,7 @@ pub async fn get_config( Ok(Envelope::ok(view)) } -/// `PUT /v1/wallets/:id/sponsorship` +/// `PUT /v1/wallets/:id/sponsorship`; requires at least one field (see `docs/api.md`). pub async fn put_config( State(state): State, Path(wallet_id): Path, @@ -71,6 +71,14 @@ pub async fn put_config( ) -> ApiResult>> { authorize_wallet(&headers, &state, wallet_id).await?; let req: SponsorshipConfigRequest = parse_optional(&body)?; + if req.enabled.is_none() + && req.per_tx_fee_cap_stroops.is_none() + && req.daily_budget_stroops.is_none() + { + return Err(ApiError::BadRequest( + "at least one field must be provided".into(), + )); + } let enabled = req.enabled.unwrap_or(false); if let Some(cap) = req.per_tx_fee_cap_stroops { diff --git a/crates/api/src/routes/whitelist.rs b/crates/api/src/routes/whitelist.rs index 38ba18b..c174d9e 100644 --- a/crates/api/src/routes/whitelist.rs +++ b/crates/api/src/routes/whitelist.rs @@ -57,7 +57,7 @@ pub async fn get_config( Ok(Envelope::ok(AllowlistConfigView { enabled })) } -/// `PUT /v1/wallets/:id/whitelist/config` +/// `PUT /v1/wallets/:id/whitelist/config`; requires at least one field (see `docs/api.md`). pub async fn put_config( State(state): State, Path(wallet_id): Path, @@ -66,6 +66,11 @@ pub async fn put_config( ) -> ApiResult>> { owned_wallet(&state, &headers, wallet_id).await?; let req: AllowlistConfigRequest = parse_optional(&body)?; + if req.enabled.is_none() { + return Err(ApiError::BadRequest( + "at least one field must be provided".into(), + )); + } let enabled = req.enabled.unwrap_or(false); // Enabling with an empty list would lock the wallet out of every destination — including diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index 60a796c..9daab69 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -693,8 +693,10 @@ impl Supervisor { // of the shared concurrency budget. Paged by id so memory doesn't scale with wallet count. let semaphore = Arc::new(tokio::sync::Semaphore::new(Self::MAX_CONCURRENT_POLLS)); let mut tasks = tokio::task::JoinSet::new(); + let mut task_wallets = HashMap::new(); for w in wallets { + let wallet_id = w.id; let store = self.store.clone(); let store_for_mark = self.store.clone(); let horizon_url = self.horizon_url.clone(); @@ -703,7 +705,7 @@ impl Supervisor { let retry = self.retry.clone(); let circuit = self.circuit.clone(); let semaphore = semaphore.clone(); - tasks.spawn(async move { + let task_id = tasks.spawn(async move { // Held for the duration of this wallet's poll; bounds how many Horizon requests // are in flight at once without limiting how many wallets we *queue*. let _permit = semaphore.acquire_owned().await; @@ -741,6 +743,7 @@ impl Supervisor { } (w.id, result) }); + task_wallets.insert(task_id.id(), wallet_id); } let mut total = 0; @@ -794,6 +797,23 @@ impl Supervisor { } } + /// Queue one wallet's poll on `tasks`, gated by the shared concurrency `semaphore`. + fn spawn_poll( + &self, + tasks: &mut tokio::task::JoinSet<(Uuid, Result)>, + semaphore: &Arc, + } + } + + while let Some(joined) = tasks.join_next().await { + total += Self::tally(joined); + } + match fetch_error { + Some(e) => Err(e.into()), + None => Ok(total), + } + } + /// Queue one wallet's poll on `tasks`, gated by the shared concurrency `semaphore`. fn spawn_poll( &self, diff --git a/crates/store/src/lib.rs b/crates/store/src/lib.rs index e69de29..cf3173f 100644 --- a/crates/store/src/lib.rs +++ b/crates/store/src/lib.rs @@ -0,0 +1,1966 @@ +//! Postgres persistence for octo (sqlx). +//! +//! Tables: `wallets`, `addresses`, `transactions`, `withdrawals`, `webhook_endpoints`, +//! `webhook_deliveries`, `ingest_cursor` — see `migrations/0001_init.sql`. +//! +//! Security-relevant guarantees implemented here (see `docs/threat-model.md`): +//! - All queries are parameterized (no string-built SQL) → no SQL injection. +//! - [`Store::allocate_address`] increments the per-wallet muxed-id counter **atomically** inside a +//! transaction, so concurrent address creation can't collide or reuse an id. +//! - [`Store::record_deposit`] is **idempotent** on the immutable `(tx_hash, operation_index)` +//! unique index, so a replayed/reorged Horizon event cannot double-credit. +//! - [`Store::create_withdrawal`] is idempotent on `(wallet_id, idempotency_key)`. +#![forbid(unsafe_code)] + +mod error; +mod models; + +pub use error::StoreError; +pub use models::{ + Address, ApiKey, AuditLog, DenylistedToken, EmailOtp, GasSponsorshipConfig, NewDeposit, + NewPaymentLink, NewSponsoredTx, PaymentLink, PaymentLinkPayment, SponsoredTransaction, + Transaction, User, Wallet, WebhookDelivery, WebhookEndpoint, WhitelistedAddress, Withdrawal, + WithdrawalAllowlistConfig, +}; + +use sqlx::postgres::{PgPool, PgPoolOptions}; +use uuid::Uuid; + +/// Embedded migrations, applied by [`Store::migrate`]. +pub static MIGRATOR: sqlx::migrate::Migrator = sqlx::migrate!("./migrations"); + +/// A handle to the database (cloneable; wraps a connection pool). +#[derive(Clone)] +pub struct Store { + pool: PgPool, +} + +/// Parameters for creating a server-custody wallet (legacy wallets and gas-tank fee accounts — +/// the only rows that carry a server-held sealed seed). +pub struct NewWallet<'a> { + pub network: &'a str, + pub stellar_account_g: &'a str, + pub sealed_ciphertext: &'a [u8], + pub sealed_nonce: &'a [u8], + pub sealed_salt: &'a [u8], + /// Scheme version tag for the sealed seed. Use `octo_crypto::SCHEME_V1`. + pub sealed_scheme: i16, + pub label: Option<&'a str>, + pub user_id: Option, + pub description: Option<&'a str>, +} + +/// Parameters for creating a non-custodial (client-custody) wallet: the client generated the +/// keypair and sends only the public account plus an opaque password-encrypted backup blob the +/// server cannot decrypt. +pub struct NewClientWallet<'a> { + pub network: &'a str, + pub stellar_account_g: &'a str, + pub encrypted_backup: Option<&'a str>, + pub label: Option<&'a str>, + pub user_id: Option, + pub description: Option<&'a str>, +} + +/// Parameters for creating a withdrawal intent. +pub struct NewWithdrawal<'a> { + pub wallet_id: Uuid, + pub idempotency_key: &'a str, + pub destination_account: &'a str, + pub asset_code: &'a str, + pub asset_issuer: Option<&'a str>, + pub amount_stroops: i64, + pub memo_id: Option, +} + +impl Store { + /// Connect to Postgres at `database_url` and return a pooled handle. + pub async fn connect(database_url: &str) -> Result { + let pool = PgPoolOptions::new() + .max_connections(10) + .connect(database_url) + .await?; + Ok(Self { pool }) + } + + /// Build a store from an existing pool (useful in tests). + pub fn from_pool(pool: PgPool) -> Self { + Self { pool } + } + + /// Apply all pending migrations. + pub async fn migrate(&self) -> Result<(), StoreError> { + MIGRATOR.run(&self.pool).await?; + Ok(()) + } + + /// Borrow the underlying pool. + pub fn pool(&self) -> &PgPool { + &self.pool + } + + // --- users ------------------------------------------------------------ + + /// Create a user. `email` should already be lowercased by the caller. Returns + /// [`StoreError::Conflict`] if the email is already registered. + pub async fn create_user(&self, email: &str, password_hash: &str) -> Result { + sqlx::query_as::<_, User>( + "INSERT INTO users (email, password_hash) VALUES ($1, $2) RETURNING *", + ) + .bind(email) + .bind(password_hash) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from_sqlx_conflict) + } + + /// Set a user's display username. Returns [`StoreError::Conflict`] if another user already + /// has it (compared case-insensitively, per the `users_username_unique_idx` index). + pub async fn update_username(&self, user_id: Uuid, username: &str) -> Result { + sqlx::query_as::<_, User>( + "UPDATE users SET username = $2, updated_at = now() WHERE id = $1 RETURNING *", + ) + .bind(user_id) + .bind(username) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from_sqlx_conflict) + } + + /// Delete a user outright. Only safe pre-verification — used to roll back a signup whose + /// OTP email never went out, so the email isn't stuck as "already registered" forever. + pub async fn delete_unverified_user(&self, user_id: Uuid) -> Result<(), StoreError> { + sqlx::query("DELETE FROM users WHERE id = $1 AND email_verified_at IS NULL") + .bind(user_id) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Look up a user by email (caller lowercases). + pub async fn find_user_by_email(&self, email: &str) -> Result, StoreError> { + let row = sqlx::query_as::<_, User>("SELECT * FROM users WHERE email = $1") + .bind(email) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + /// Fetch a user by id. + pub async fn get_user(&self, id: Uuid) -> Result, StoreError> { + let row = sqlx::query_as::<_, User>("SELECT * FROM users WHERE id = $1") + .bind(id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + /// Mark a user's email as verified. + pub async fn mark_email_verified(&self, user_id: Uuid) -> Result<(), StoreError> { + sqlx::query("UPDATE users SET email_verified_at = now() WHERE id = $1") + .bind(user_id) + .execute(&self.pool) + .await?; + Ok(()) + } + + // --- email OTP ---------------------------------------------------------- + + /// Issue a fresh OTP row. Callers hash the code themselves before calling this. + pub async fn create_otp( + &self, + user_id: Uuid, + purpose: &str, + code_hash: &str, + tx_hash_bound: Option<&str>, + ttl: chrono::Duration, + ) -> Result { + let id: Uuid = sqlx::query_scalar( + "INSERT INTO email_otps (user_id, purpose, code_hash, tx_hash_bound, expires_at) + VALUES ($1, $2, $3, $4, now() + $5) RETURNING id", + ) + .bind(user_id) + .bind(purpose) + .bind(code_hash) + .bind(tx_hash_bound) + .bind(ttl) + .fetch_one(&self.pool) + .await?; + Ok(id) + } + + /// Verify an already-hashed code against the most recent unconsumed OTP for + /// `(user_id, purpose)`. On a wrong code, increments `attempts` and returns `InvalidOtp` + /// rather than panicking — callers should surface a generic "invalid or expired code" either + /// way, so guessing can't distinguish "wrong code" from "no such code exists". + pub async fn verify_and_consume_otp( + &self, + user_id: Uuid, + purpose: &str, + code_hash: &str, + tx_hash_bound: Option<&str>, + ) -> Result<(), StoreError> { + const MAX_ATTEMPTS: i16 = 5; + + let otp = sqlx::query_as::<_, EmailOtp>( + "SELECT * FROM email_otps WHERE user_id = $1 AND purpose = $2 + ORDER BY created_at DESC LIMIT 1", + ) + .bind(user_id) + .bind(purpose) + .fetch_optional(&self.pool) + .await? + .ok_or(StoreError::InvalidOtp)?; + + if otp.consumed_at.is_some() + || otp.attempts >= MAX_ATTEMPTS + || otp.expires_at < chrono::Utc::now() + || otp.tx_hash_bound.as_deref() != tx_hash_bound + { + return Err(StoreError::InvalidOtp); + } + if otp.code_hash != code_hash { + sqlx::query("UPDATE email_otps SET attempts = attempts + 1 WHERE id = $1") + .bind(otp.id) + .execute(&self.pool) + .await?; + return Err(StoreError::InvalidOtp); + } + + sqlx::query("UPDATE email_otps SET consumed_at = now() WHERE id = $1") + .bind(otp.id) + .execute(&self.pool) + .await?; + Ok(()) + } + + // --- audit logs ------------------------------------------------------- + + /// Append an audit-log entry. Best-effort: failures are surfaced to the caller, which logs and + /// continues (auditing must never block the primary operation). + pub async fn record_audit( + &self, + user_id: Uuid, + action: &str, + category: &str, + target: Option<&str>, + ip_address: Option<&str>, + ) -> Result<(), StoreError> { + sqlx::query( + "INSERT INTO audit_logs (user_id, action, category, target, ip_address) + VALUES ($1, $2, $3, $4, $5)", + ) + .bind(user_id) + .bind(action) + .bind(category) + .bind(target) + .bind(ip_address) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// List a user's audit logs (most recent first), optionally filtered by `category` and a + /// case-insensitive `search` over the action/target. Capped at `limit` rows. + pub async fn list_audit_logs( + &self, + user_id: Uuid, + category: Option<&str>, + search: Option<&str>, + limit: i64, + ) -> Result, StoreError> { + // Build with optional filters; `$2`/`$3` are NULL when not provided. + let rows = sqlx::query_as::<_, AuditLog>( + r#" + SELECT * FROM audit_logs + WHERE user_id = $1 + AND ($2::text IS NULL OR category = $2) + AND ($3::text IS NULL OR action ILIKE '%' || $3 || '%' + OR coalesce(target, '') ILIKE '%' || $3 || '%') + ORDER BY created_at DESC + LIMIT $4 + "#, + ) + .bind(user_id) + .bind(category) + .bind(search) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + // --- api keys --------------------------------------------------------- + + /// Create or replace the wallet's API key (regenerate). Stores only the hash + display prefix. + pub async fn upsert_api_key( + &self, + wallet_id: Uuid, + prefix: &str, + key_hash: &str, + ) -> Result { + sqlx::query_as::<_, ApiKey>( + r#" + INSERT INTO api_keys (wallet_id, prefix, key_hash) + VALUES ($1, $2, $3) + ON CONFLICT (wallet_id) + DO UPDATE SET prefix = EXCLUDED.prefix, key_hash = EXCLUDED.key_hash, + created_at = now() + RETURNING * + "#, + ) + .bind(wallet_id) + .bind(prefix) + .bind(key_hash) + .fetch_one(&self.pool) + .await + .map_err(StoreError::Database) + } + + /// Get the wallet's API key metadata (prefix only — never the secret), if one exists. + pub async fn get_api_key(&self, wallet_id: Uuid) -> Result, StoreError> { + let row = sqlx::query_as::<_, ApiKey>("SELECT * FROM api_keys WHERE wallet_id = $1") + .bind(wallet_id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + /// Look up the wallet that owns a key by its hash (for API-key authentication later). + pub async fn wallet_id_for_key_hash(&self, key_hash: &str) -> Result, StoreError> { + let row: Option<(Uuid,)> = + sqlx::query_as("SELECT wallet_id FROM api_keys WHERE key_hash = $1") + .bind(key_hash) + .fetch_optional(&self.pool) + .await?; + Ok(row.map(|r| r.0)) + } + + /// Delete (revoke) the API key for a wallet. Returns `Ok(())` even if no key existed. + pub async fn delete_api_key(&self, wallet_id: Uuid) -> Result<(), StoreError> { + sqlx::query("DELETE FROM api_keys WHERE wallet_id = $1") + .bind(wallet_id) + .execute(&self.pool) + .await?; + Ok(()) + } + + // --- wallets ---------------------------------------------------------- + + /// Create a master wallet. Fails with [`StoreError::Conflict`] if the account already exists. + pub async fn create_wallet(&self, new: NewWallet<'_>) -> Result { + sqlx::query_as::<_, Wallet>( + r#" + INSERT INTO wallets + (network, stellar_account_g, sealed_ciphertext, sealed_nonce, sealed_salt, + sealed_scheme, label, user_id, description, custody) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, 'server') + RETURNING * + "#, + ) + .bind(new.network) + .bind(new.stellar_account_g) + .bind(new.sealed_ciphertext) + .bind(new.sealed_nonce) + .bind(new.sealed_salt) + .bind(new.sealed_scheme) + .bind(new.label) + .bind(new.user_id) + .bind(new.description) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from_sqlx_conflict) + } + + /// Attach a gas-tank fee account to a client-custody wallet: stores the tank's sealed seed + /// and public account. The tank only ever holds fee float — never customer funds. + /// + /// `sealed_scheme` must be written alongside the seed: the `wallets_gas_tank_has_seed` CHECK + /// requires it, and key rotation (`bin/migrate-keys`) needs the tag to know how to open it. + pub async fn set_gas_tank( + &self, + wallet_id: Uuid, + gas_tank_account_g: &str, + sealed_ciphertext: &[u8], + sealed_nonce: &[u8], + sealed_salt: &[u8], + sealed_scheme: i16, + ) -> Result { + sqlx::query_as::<_, Wallet>( + r#" + UPDATE wallets + SET gas_tank_account_g = $2, sealed_ciphertext = $3, sealed_nonce = $4, + sealed_salt = $5, sealed_scheme = $6, updated_at = now() + WHERE id = $1 AND custody = 'client' AND gas_tank_account_g IS NULL + RETURNING * + "#, + ) + .bind(wallet_id) + .bind(gas_tank_account_g) + .bind(sealed_ciphertext) + .bind(sealed_nonce) + .bind(sealed_salt) + .bind(sealed_scheme) + .fetch_optional(&self.pool) + .await? + .ok_or(StoreError::Conflict) // already has a tank, or not a client wallet + } + + /// Create a non-custodial wallet: no seed is stored; the server can never sign for it. + pub async fn create_client_wallet( + &self, + new: NewClientWallet<'_>, + ) -> Result { + sqlx::query_as::<_, Wallet>( + r#" + INSERT INTO wallets + (network, stellar_account_g, label, user_id, description, custody, + encrypted_backup) + VALUES ($1, $2, $3, $4, $5, 'client', $6) + RETURNING * + "#, + ) + .bind(new.network) + .bind(new.stellar_account_g) + .bind(new.label) + .bind(new.user_id) + .bind(new.description) + .bind(new.encrypted_backup) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from_sqlx_conflict) + } + + /// List a user's wallets (most recent first), with optional cursor-based pagination. + /// + /// Fetching `limit + 1` rows lets the caller detect whether a next page exists without a + /// separate COUNT query — the same pattern used by `list_sponsored_transactions`. + pub async fn list_wallets_for_user( + &self, + user_id: Uuid, + limit: i64, + before_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, Wallet>( + r#" + SELECT * FROM wallets + WHERE user_id = $1 + AND ($2::uuid IS NULL OR (created_at, id) < ( + SELECT created_at, id FROM wallets WHERE id = $2 + )) + ORDER BY created_at DESC, id DESC + LIMIT $3 + "#, + ) + .bind(user_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Paginated version of [`list_wallets_for_user`]: returns at most `limit` rows, newest first. + /// Pass the last page's final wallet id as `before_id` to fetch the next page. + pub async fn list_wallets_for_user_page( + &self, + user_id: Uuid, + limit: i64, + before_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, Wallet>( + r#" + SELECT * FROM wallets + WHERE user_id = $1 + AND ($2::uuid IS NULL OR (created_at, id) < ( + SELECT created_at, id FROM wallets WHERE id = $2 + )) + ORDER BY created_at DESC, id DESC + LIMIT $3 + "#, + ) + .bind(user_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// List all wallets (used by the ingest supervisor to fan out poll loops). + pub async fn list_wallets(&self) -> Result, StoreError> { + let rows = sqlx::query_as::<_, Wallet>("SELECT * FROM wallets ORDER BY created_at") + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Wallets on `network` that are due for an ingest poll, given activity-based backoff. + /// + /// A dev/production database accumulates wallets that never see another deposit. Polling all + /// of them on the same short cycle spends the concurrency budget on dead accounts and delays + /// the ones that are actually transacting. Idleness is measured by `ingest_cursor.updated_at`, + /// which is only bumped when a record is actually processed: + /// + /// - active (last activity < `active_after_secs`): every tick + /// - idle: at most once per `idle_interval_secs` + /// - dormant (last activity older than `dormant_after_secs`): at most once per + /// `dormant_interval_secs` + /// + /// A wallet with no cursor row has never been polled, so it is always due. + pub async fn wallets_due_for_poll( + &self, + network: &str, + active_after_secs: i64, + idle_interval_secs: i64, + dormant_after_secs: i64, + dormant_interval_secs: i64, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, Wallet>( + r#" + SELECT w.* FROM wallets w + LEFT JOIN ingest_cursor c ON c.wallet_id = w.id + WHERE w.network = $1 + -- Never polled, or never saw activity => always due. + AND ( + c.last_polled_at IS NULL + OR c.updated_at IS NULL + OR c.last_polled_at < now() - make_interval(secs => + CASE + -- Active: no extra wait, poll every tick. + WHEN c.updated_at > now() - make_interval(secs => $2) THEN 0 + -- Dormant: longest wait between polls. + WHEN c.updated_at <= now() - make_interval(secs => $4) THEN $5 + -- Idle: in between. + ELSE $3 + END) + ) + ORDER BY w.created_at + "#, + ) + .bind(network) + .bind(active_after_secs as f64) + .bind(idle_interval_secs as f64) + .bind(dormant_after_secs as f64) + .bind(dormant_interval_secs as f64) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Record that a wallet was polled (whether or not anything new arrived). + /// + /// Distinct from [`Store::set_cursor`], which only advances on real activity — the backoff + /// tiers need both "when did we last see money" and "when did we last look". + pub async fn mark_polled(&self, wallet_id: Uuid) -> Result<(), StoreError> { + // `updated_at` is deliberately backdated to the epoch on INSERT: it means "last time this + // wallet saw activity", and merely looking at a wallet is not activity. Letting it take + // its `DEFAULT now()` would mark every never-used wallet as freshly active and the + // backoff tiers would never engage. `set_cursor` is the only writer that advances it. + sqlx::query( + r#" + INSERT INTO ingest_cursor (wallet_id, last_polled_at, updated_at) + VALUES ($1, now(), 'epoch') + ON CONFLICT (wallet_id) DO UPDATE SET last_polled_at = now() + "#, + ) + .bind(wallet_id) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Fetch a wallet by id. + pub async fn get_wallet(&self, id: Uuid) -> Result { + sqlx::query_as::<_, Wallet>("SELECT * FROM wallets WHERE id = $1") + .bind(id) + .fetch_optional(&self.pool) + .await? + .ok_or(StoreError::NotFound) + } + + /// Atomically swap the sealed seed material for a single wallet after a reseal/key-rotation. + /// + /// The caller (typically `bin/migrate-keys`) opens the old seed with the old master key, + /// re-seals it with the new master key via `octo_crypto::reseal`, and then calls this method + /// to persist the result. The `expected_scheme` guard ensures idempotency: if the row was + /// already migrated (e.g. by a concurrent runner) the update is silently skipped rather than + /// overwriting a newer record. + /// + /// Returns `true` if the row was updated, `false` if it was already on the target scheme. + pub async fn reseal_wallet( + &self, + wallet_id: Uuid, + new_ciphertext: &[u8], + new_nonce: &[u8], + new_salt: &[u8], + new_scheme: i16, + expected_old_scheme: i16, + ) -> Result { + // Only update the row if it still carries the old scheme — this is the idempotency guard. + // A concurrent runner that already migrated this wallet will have set sealed_scheme to + // `new_scheme`, so the WHERE clause won't match and no double-reseal can occur. + let result = sqlx::query( + r#" + UPDATE wallets + SET sealed_ciphertext = $2, + sealed_nonce = $3, + sealed_salt = $4, + sealed_scheme = $5, + updated_at = now() + WHERE id = $1 + AND sealed_scheme = $6 + "#, + ) + .bind(wallet_id) + .bind(new_ciphertext) + .bind(new_nonce) + .bind(new_salt) + .bind(new_scheme) + .bind(expected_old_scheme) + .execute(&self.pool) + .await?; + + Ok(result.rows_affected() > 0) + } + + /// Fetch a page of wallets whose `sealed_scheme` does not equal `target_scheme`, for the + /// migration backfill job. Returns at most `batch_size` rows ordered by `id` (stable for + /// resumable cursored iteration). Pass the last returned wallet's `id` as `after_id` on + /// subsequent calls to page through the full table without re-scanning already-migrated rows. + pub async fn list_wallets_needing_reseal( + &self, + target_scheme: i16, + batch_size: i64, + after_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, Wallet>( + r#" + SELECT * FROM wallets + WHERE sealed_scheme <> $1 + AND ($2::uuid IS NULL OR id > $2) + ORDER BY id + LIMIT $3 + "#, + ) + .bind(target_scheme) + .bind(after_id) + .bind(batch_size) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + // --- addresses -------------------------------------------------------- + + /// Atomically allocate the next muxed id for `wallet_id` and insert the address row. + /// + /// The counter bump and the insert happen in one transaction with a row lock, so two + /// concurrent callers always get distinct, gap-free-enough ids and never collide. + pub async fn allocate_address( + &self, + wallet_id: Uuid, + muxed_address_for: impl FnOnce(i64) -> Result, + customer_ref: Option<&str>, + metadata: serde_json::Value, + ) -> Result { + let mut tx = self.pool.begin().await?; + + // Lock the wallet row and read+bump the counter. + let next_id: i64 = + sqlx::query_scalar("SELECT next_muxed_id FROM wallets WHERE id = $1 FOR UPDATE") + .bind(wallet_id) + .fetch_optional(&mut *tx) + .await? + .ok_or(StoreError::NotFound)?; + + sqlx::query("UPDATE wallets SET next_muxed_id = next_muxed_id + 1, updated_at = now() WHERE id = $1") + .bind(wallet_id) + .execute(&mut *tx) + .await?; + + // Derive the muxed address for this id via the caller-provided closure (wallet-core). + let muxed_address = muxed_address_for(next_id).map_err(|_| StoreError::NotFound)?; + + let address = sqlx::query_as::<_, Address>( + r#" + INSERT INTO addresses (wallet_id, muxed_id, muxed_address, customer_ref, metadata) + VALUES ($1, $2, $3, $4, $5) + RETURNING * + "#, + ) + .bind(wallet_id) + .bind(next_id) + .bind(&muxed_address) + .bind(customer_ref) + .bind(metadata) + .fetch_one(&mut *tx) + .await + .map_err(StoreError::from_sqlx_conflict)?; + + tx.commit().await?; + Ok(address) + } + + /// List addresses for a wallet (most recent first), with optional cursor-based pagination. + pub async fn list_addresses( + &self, + wallet_id: Uuid, + limit: i64, + before_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, Address>( + r#" + SELECT * FROM addresses + WHERE wallet_id = $1 + AND ($2::uuid IS NULL OR (created_at, id) < ( + SELECT created_at, id FROM addresses WHERE id = $2 + )) + ORDER BY created_at DESC, id DESC + LIMIT $3 + "#, + ) + .bind(wallet_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Paginated version of [`list_addresses`]: returns at most `limit` rows, newest first. + /// Pass the last page's final address id as `before_id` to fetch the next page. + pub async fn list_addresses_page( + &self, + wallet_id: Uuid, + limit: i64, + before_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, Address>( + r#" + SELECT * FROM addresses + WHERE wallet_id = $1 + AND ($2::uuid IS NULL OR (created_at, id) < ( + SELECT created_at, id FROM addresses WHERE id = $2 + )) + ORDER BY created_at DESC, id DESC + LIMIT $3 + "#, + ) + .bind(wallet_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Fetch an address by id. + pub async fn get_address(&self, id: Uuid) -> Result, StoreError> { + let row = sqlx::query_as::<_, Address>("SELECT * FROM addresses WHERE id = $1") + .bind(id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + /// Find the address for a given `(wallet_id, muxed_id)`, if any. + pub async fn address_by_muxed_id( + &self, + wallet_id: Uuid, + muxed_id: i64, + ) -> Result, StoreError> { + let row = sqlx::query_as::<_, Address>( + "SELECT * FROM addresses WHERE wallet_id = $1 AND muxed_id = $2", + ) + .bind(wallet_id) + .bind(muxed_id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + // --- transactions (deposits) ------------------------------------------ + + /// Idempotently record a confirmed deposit. + /// + /// Returns `Ok(Some(tx))` on first insert and `Ok(None)` if this exact on-chain operation was + /// already recorded (the `(tx_hash, operation_index)` unique index fired) — so replays and + /// reorged re-deliveries never double-credit. + pub async fn record_deposit(&self, d: &NewDeposit) -> Result, StoreError> { + let result = sqlx::query_as::<_, Transaction>( + r#" + INSERT INTO transactions + (wallet_id, address_id, direction, asset_code, asset_issuer, amount_stroops, + source_account, destination_account, stellar_tx_hash, operation_index, + horizon_op_id, ledger, memo_id, status) + VALUES ($1, $2, 'deposit', $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, 'confirmed') + RETURNING * + "#, + ) + .bind(d.wallet_id) + .bind(d.address_id) + .bind(&d.asset_code) + .bind(&d.asset_issuer) + .bind(d.amount_stroops) + .bind(&d.source_account) + .bind(&d.destination_account) + .bind(&d.stellar_tx_hash) + .bind(d.operation_index) + .bind(&d.horizon_op_id) + .bind(d.ledger) + .bind(d.memo_id) + .fetch_one(&self.pool) + .await; + + match result { + Ok(tx) => Ok(Some(tx)), + Err(e) => match StoreError::from_sqlx_conflict(e) { + StoreError::Conflict => Ok(None), // already recorded — benign + other => Err(other), + }, + } + } + + /// List transactions for a wallet (most recent first), with optional cursor-based pagination. + pub async fn list_transactions( + &self, + wallet_id: Uuid, + limit: i64, + before_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, Transaction>( + r#" + SELECT * FROM transactions + WHERE wallet_id = $1 + AND ($2::uuid IS NULL OR (created_at, id) < ( + SELECT created_at, id FROM transactions WHERE id = $2 + )) + ORDER BY created_at DESC, id DESC + LIMIT $3 + "#, + ) + .bind(wallet_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Paginated version of [`list_transactions`]: returns at most `limit` rows, newest first. + /// Pass the last page's final transaction id as `before_id` to fetch the next page. + pub async fn list_transactions_page( + &self, + wallet_id: Uuid, + limit: i64, + before_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, Transaction>( + r#" + SELECT * FROM transactions + WHERE wallet_id = $1 + AND ($2::uuid IS NULL OR (created_at, id) < ( + SELECT created_at, id FROM transactions WHERE id = $2 + )) + ORDER BY created_at DESC, id DESC + LIMIT $3 + "#, + ) + .bind(wallet_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Fetch a single transaction by id. + pub async fn get_transaction(&self, id: Uuid) -> Result, StoreError> { + let row = sqlx::query_as::<_, Transaction>("SELECT * FROM transactions WHERE id = $1") + .bind(id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + // --- withdrawals ------------------------------------------------------ + + /// Cheap existence check on `(wallet_id, idempotency_key)`, used to short-circuit a retried + /// request with a 409 **before** running any pre-flight Horizon checks — a key that has + /// already been consumed doesn't need its request re-validated against the chain. + pub async fn withdrawal_exists( + &self, + wallet_id: Uuid, + idempotency_key: &str, + ) -> Result { + let found: Option = sqlx::query_scalar( + "SELECT id FROM withdrawals WHERE wallet_id = $1 AND idempotency_key = $2", + ) + .bind(wallet_id) + .bind(idempotency_key) + .fetch_optional(&self.pool) + .await?; + Ok(found.is_some()) + } + + /// Create a withdrawal intent. Idempotent on `(wallet_id, idempotency_key)`: a retried request + /// with the same key returns [`StoreError::Conflict`] instead of creating a second payout. + /// Record a confirmed/failed outbound transfer in the `transactions` history (the table the + /// dashboard lists). Withdrawals previously lived only in `withdrawals`, which is why they + /// never showed up in "recent transactions". + #[allow(clippy::too_many_arguments)] + pub async fn record_withdrawal_transaction( + &self, + wallet_id: Uuid, + asset_code: &str, + asset_issuer: Option<&str>, + amount_stroops: i64, + source_account: &str, + destination_account: &str, + stellar_tx_hash: Option<&str>, + status: &str, + ) -> Result { + let row = sqlx::query_as::<_, Transaction>( + r#" + INSERT INTO transactions + (wallet_id, direction, asset_code, asset_issuer, amount_stroops, + source_account, destination_account, stellar_tx_hash, status) + VALUES ($1, 'withdrawal', $2, $3, $4, $5, $6, $7, $8) + RETURNING * + "#, + ) + .bind(wallet_id) + .bind(asset_code) + .bind(asset_issuer) + .bind(amount_stroops) + .bind(source_account) + .bind(destination_account) + .bind(stellar_tx_hash) + .bind(status) + .fetch_one(&self.pool) + .await?; + Ok(row) + } + + pub async fn create_withdrawal( + &self, + new: NewWithdrawal<'_>, + ) -> Result { + sqlx::query_as::<_, Withdrawal>( + r#" + INSERT INTO withdrawals + (wallet_id, idempotency_key, destination_account, asset_code, asset_issuer, + amount_stroops, memo_id) + VALUES ($1, $2, $3, $4, $5, $6, $7) + RETURNING * + "#, + ) + .bind(new.wallet_id) + .bind(new.idempotency_key) + .bind(new.destination_account) + .bind(new.asset_code) + .bind(new.asset_issuer) + .bind(new.amount_stroops) + .bind(new.memo_id) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from_sqlx_conflict) + } + + /// Update a withdrawal's status (and optional tx hash) after submission. + pub async fn update_withdrawal_status( + &self, + id: Uuid, + status: &str, + stellar_tx_hash: Option<&str>, + ) -> Result<(), StoreError> { + sqlx::query( + "UPDATE withdrawals SET status = $2, stellar_tx_hash = $3, updated_at = now() WHERE id = $1", + ) + .bind(id) + .bind(status) + .bind(stellar_tx_hash) + .execute(&self.pool) + .await?; + Ok(()) + } + + // --- sponsored transactions ------------------------------------------- + + /// List sponsored transactions for a wallet (most recent first), with + /// optional status filter and cursor-based pagination. + pub async fn list_sponsored_transactions( + &self, + wallet_id: Uuid, + limit: i64, + status_filter: Option<&str>, + before_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, SponsoredTransaction>( + r#" + SELECT * FROM sponsored_transactions + WHERE wallet_id = $1 + AND ($2::text IS NULL OR status = $2) + AND ($3::uuid IS NULL OR (created_at, id) < (SELECT created_at, id FROM sponsored_transactions WHERE id = $3)) + ORDER BY created_at DESC, id DESC + LIMIT $4 + "#, + ) + .bind(wallet_id) + .bind(status_filter) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + // --- gas sponsorship config ------------------------------------------- + + /// Fetch a wallet's sponsorship config, or `None` if none has been saved. + pub async fn get_gas_sponsorship_config( + &self, + wallet_id: Uuid, + ) -> Result, StoreError> { + let row = sqlx::query_as::<_, GasSponsorshipConfig>( + "SELECT * FROM gas_sponsorship_configs WHERE wallet_id = $1", + ) + .bind(wallet_id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + /// Create or replace a wallet's sponsorship config. + pub async fn upsert_gas_sponsorship_config( + &self, + wallet_id: Uuid, + enabled: bool, + per_tx_fee_cap_stroops: Option, + daily_budget_stroops: Option, + ) -> Result { + sqlx::query_as::<_, GasSponsorshipConfig>( + r#" + INSERT INTO gas_sponsorship_configs + (wallet_id, enabled, per_tx_fee_cap_stroops, daily_budget_stroops) + VALUES ($1, $2, $3, $4) + ON CONFLICT (wallet_id) DO UPDATE SET + enabled = EXCLUDED.enabled, + per_tx_fee_cap_stroops = EXCLUDED.per_tx_fee_cap_stroops, + daily_budget_stroops = EXCLUDED.daily_budget_stroops, + updated_at = now() + RETURNING * + "#, + ) + .bind(wallet_id) + .bind(enabled) + .bind(per_tx_fee_cap_stroops) + .bind(daily_budget_stroops) + .fetch_one(&self.pool) + .await + .map_err(StoreError::Database) + } + + /// Sum of sponsored fees reserved (pending + confirmed) for a wallet so far today (UTC). + /// Used to enforce the rolling daily budget and to report `spent_today`. + pub async fn sum_sponsored_fees_reserved_today( + &self, + wallet_id: Uuid, + ) -> Result { + let total: Option = sqlx::query_scalar( + r#" + SELECT COALESCE(SUM(fee_stroops), 0)::bigint + FROM sponsored_transactions + WHERE wallet_id = $1 + AND status IN ('pending', 'confirmed') + AND created_at >= date_trunc('day', now() AT TIME ZONE 'UTC') + "#, + ) + .bind(wallet_id) + .fetch_one(&self.pool) + .await?; + Ok(total.unwrap_or(0)) + } + + // --- withdrawal allowlist ---------------------------------------------- + + /// Fetch a wallet's withdrawal-allowlist config, if one has ever been set. `None` means the + /// wallet has never touched this feature — treat that the same as `enabled = false`. + pub async fn get_withdrawal_allowlist_config( + &self, + wallet_id: Uuid, + ) -> Result, StoreError> { + let row = sqlx::query_as::<_, WithdrawalAllowlistConfig>( + "SELECT * FROM withdrawal_allowlist_configs WHERE wallet_id = $1", + ) + .bind(wallet_id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + /// Create or replace a wallet's withdrawal-allowlist toggle. + pub async fn upsert_withdrawal_allowlist_config( + &self, + wallet_id: Uuid, + enabled: bool, + ) -> Result { + sqlx::query_as::<_, WithdrawalAllowlistConfig>( + r#" + INSERT INTO withdrawal_allowlist_configs (wallet_id, enabled) + VALUES ($1, $2) + ON CONFLICT (wallet_id) DO UPDATE SET + enabled = EXCLUDED.enabled, + updated_at = now() + RETURNING * + "#, + ) + .bind(wallet_id) + .bind(enabled) + .fetch_one(&self.pool) + .await + .map_err(StoreError::Database) + } + + /// Add an address to a wallet's withdrawal allowlist. `Conflict` if already present. + pub async fn add_whitelisted_address( + &self, + wallet_id: Uuid, + address: &str, + label: Option<&str>, + ) -> Result { + sqlx::query_as::<_, WhitelistedAddress>( + r#" + INSERT INTO whitelisted_addresses (wallet_id, address, label) + VALUES ($1, $2, $3) + RETURNING * + "#, + ) + .bind(wallet_id) + .bind(address) + .bind(label) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from_sqlx_conflict) + } + + /// List a wallet's whitelisted addresses, newest first. + pub async fn list_whitelisted_addresses( + &self, + wallet_id: Uuid, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, WhitelistedAddress>( + "SELECT * FROM whitelisted_addresses WHERE wallet_id = $1 ORDER BY created_at DESC", + ) + .bind(wallet_id) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Remove a whitelisted address. `NotFound` if it doesn't belong to `wallet_id`. + pub async fn remove_whitelisted_address( + &self, + wallet_id: Uuid, + entry_id: Uuid, + ) -> Result<(), StoreError> { + let result = + sqlx::query("DELETE FROM whitelisted_addresses WHERE id = $1 AND wallet_id = $2") + .bind(entry_id) + .bind(wallet_id) + .execute(&self.pool) + .await?; + if result.rows_affected() == 0 { + return Err(StoreError::NotFound); + } + Ok(()) + } + + /// `true` if `address` (already normalized to its base `G...` form by the caller) is on + /// `wallet_id`'s allowlist. Pure existence check — callers first check whether the allowlist + /// is even `enabled` via [`Store::get_withdrawal_allowlist_config`]. + pub async fn is_address_whitelisted( + &self, + wallet_id: Uuid, + address: &str, + ) -> Result { + let exists: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM whitelisted_addresses WHERE wallet_id = $1 AND address = $2)", + ) + .bind(wallet_id) + .bind(address) + .fetch_one(&self.pool) + .await?; + Ok(exists) + } + + // --- per-address received totals --------------------------------------- + + /// Lifetime total (in stroops) of confirmed deposits credited to one generated address. + /// This is historical bookkeeping, not a live on-chain balance — deposits to any address + /// land in the wallet's single master account (that's the point of muxed addresses; there is + /// nothing to sweep), so this number will not match a per-address Horizon balance query. + pub async fn sum_deposits_for_address(&self, address_id: Uuid) -> Result { + let total: Option = sqlx::query_scalar( + r#" + SELECT COALESCE(SUM(amount_stroops), 0)::bigint + FROM transactions + WHERE address_id = $1 AND direction = 'deposit' AND status = 'confirmed' + "#, + ) + .bind(address_id) + .fetch_one(&self.pool) + .await?; + Ok(total.unwrap_or(0)) + } + + /// Batched version of [`Store::sum_deposits_for_address`] for an address list page: returns + /// `(address_id, total_stroops)` pairs in one round trip instead of N. + pub async fn sum_deposits_for_addresses( + &self, + address_ids: &[Uuid], + ) -> Result, StoreError> { + if address_ids.is_empty() { + return Ok(Vec::new()); + } + let rows: Vec<(Uuid, i64)> = sqlx::query_as( + r#" + SELECT address_id, COALESCE(SUM(amount_stroops), 0)::bigint AS total + FROM transactions + WHERE address_id = ANY($1) AND direction = 'deposit' AND status = 'confirmed' + GROUP BY address_id + "#, + ) + .bind(address_ids) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + // --- payment links ------------------------------------------------------- + + /// Create a payment link backed by an already-allocated deposit address. + pub async fn create_payment_link( + &self, + link: NewPaymentLink<'_>, + ) -> Result { + let row = sqlx::query_as::<_, PaymentLink>( + r#" + INSERT INTO payment_links + (wallet_id, address_id, slug, name, description, image_url, redirect_url, amount_usdc_stroops) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8) + RETURNING * + "#, + ) + .bind(link.wallet_id) + .bind(link.address_id) + .bind(link.slug) + .bind(link.name) + .bind(link.description) + .bind(link.image_url) + .bind(link.redirect_url) + .bind(link.amount_usdc_stroops) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from_sqlx_conflict)?; + Ok(row) + } + + /// Fetch a payment link owned by `wallet_id` (scoped so one merchant can't read another's). + pub async fn get_payment_link( + &self, + wallet_id: Uuid, + id: Uuid, + ) -> Result { + sqlx::query_as::<_, PaymentLink>( + "SELECT * FROM payment_links WHERE id = $1 AND wallet_id = $2", + ) + .bind(id) + .bind(wallet_id) + .fetch_optional(&self.pool) + .await? + .ok_or(StoreError::NotFound) + } + + /// Public lookup by slug — the UNIQUE constraint supplies the index for this equality lookup. + /// No wallet scoping; this is the pay-page entry point. + pub async fn get_payment_link_by_slug(&self, slug: &str) -> Result { + sqlx::query_as::<_, PaymentLink>("SELECT * FROM payment_links WHERE slug = $1") + .bind(slug) + .fetch_optional(&self.pool) + .await? + .ok_or(StoreError::NotFound) + } + + /// Unscoped lookup by id — for internal (non-owner-facing) callers that already know which + /// row they want, e.g. the expiry sweep resolving a payment's link to build its webhook. + pub async fn get_payment_link_by_id( + &self, + id: Uuid, + ) -> Result, StoreError> { + let row = sqlx::query_as::<_, PaymentLink>("SELECT * FROM payment_links WHERE id = $1") + .bind(id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + /// The payment link whose dedicated deposit address is `address_id`, if any. + pub async fn get_payment_link_by_address( + &self, + address_id: Uuid, + ) -> Result, StoreError> { + let row = + sqlx::query_as::<_, PaymentLink>("SELECT * FROM payment_links WHERE address_id = $1") + .bind(address_id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + pub async fn list_payment_links( + &self, + wallet_id: Uuid, + limit: i64, + before_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, PaymentLink>( + r#" + SELECT * FROM payment_links + WHERE wallet_id = $1 + AND ($2::uuid IS NULL OR (created_at, id) < ( + SELECT created_at, id FROM payment_links WHERE id = $2 + )) + ORDER BY created_at DESC, id DESC + LIMIT $3 + "#, + ) + .bind(wallet_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + pub async fn set_payment_link_active( + &self, + wallet_id: Uuid, + id: Uuid, + active: bool, + ) -> Result { + sqlx::query_as::<_, PaymentLink>( + r#" + UPDATE payment_links SET active = $1, updated_at = now() + WHERE id = $2 AND wallet_id = $3 + RETURNING * + "#, + ) + .bind(active) + .bind(id) + .bind(wallet_id) + .fetch_optional(&self.pool) + .await? + .ok_or(StoreError::NotFound) + } + + /// Record a payer's intent to pay (the "Continue" step, before any on-chain payment lands). + pub async fn record_payment_link_intent( + &self, + payment_link_id: Uuid, + payer_name: Option<&str>, + payer_email: Option<&str>, + amount_usdc_stroops: i64, + address_id: Option, + ) -> Result { + let row = sqlx::query_as::<_, PaymentLinkPayment>( + r#" + INSERT INTO payment_link_payments + (payment_link_id, payer_name, payer_email, amount_usdc_stroops, address_id) + VALUES ($1, $2, $3, $4, $5) + RETURNING * + "#, + ) + .bind(payment_link_id) + .bind(payer_name) + .bind(payer_email) + .bind(amount_usdc_stroops) + .bind(address_id) + .fetch_one(&self.pool) + .await?; + Ok(row) + } + + /// The pending intent owning `address_id`, if any — ingest's exact deposit match. + pub async fn pending_payment_by_address( + &self, + address_id: Uuid, + ) -> Result, StoreError> { + let row = sqlx::query_as::<_, PaymentLinkPayment>( + r#" + SELECT * FROM payment_link_payments + WHERE address_id = $1 AND status = 'pending' + ORDER BY created_at ASC + LIMIT 1 + "#, + ) + .bind(address_id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + pub async fn get_payment_link_payment( + &self, + payment_link_id: Uuid, + id: Uuid, + ) -> Result { + sqlx::query_as::<_, PaymentLinkPayment>( + "SELECT * FROM payment_link_payments WHERE id = $1 AND payment_link_id = $2", + ) + .bind(id) + .bind(payment_link_id) + .fetch_optional(&self.pool) + .await? + .ok_or(StoreError::NotFound) + } + + /// The oldest still-pending payment on a link — ingest matches deposits against this one. + pub async fn oldest_pending_payment_link_payment( + &self, + payment_link_id: Uuid, + ) -> Result, StoreError> { + let row = sqlx::query_as::<_, PaymentLinkPayment>( + r#" + SELECT * FROM payment_link_payments + WHERE payment_link_id = $1 AND status = 'pending' + ORDER BY created_at ASC + LIMIT 1 + "#, + ) + .bind(payment_link_id) + .fetch_optional(&self.pool) + .await?; + Ok(row) + } + + pub async fn confirm_payment_link_payment( + &self, + id: Uuid, + transaction_id: Uuid, + ) -> Result<(), StoreError> { + sqlx::query( + r#" + UPDATE payment_link_payments + SET status = 'confirmed', transaction_id = $1 + WHERE id = $2 + "#, + ) + .bind(transaction_id) + .bind(id) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Record a deposit that landed on this payment's address but for the wrong amount. + /// `status` must be `"underpaid"` or `"overpaid"` — the transaction is still linked (so the + /// merchant/payer can see what actually arrived) but the payment is deliberately NOT marked + /// `confirmed`. + pub async fn mark_payment_link_payment_mismatched( + &self, + id: Uuid, + transaction_id: Uuid, + status: &str, + ) -> Result<(), StoreError> { + sqlx::query( + r#" + UPDATE payment_link_payments + SET status = $1, transaction_id = $2 + WHERE id = $3 + "#, + ) + .bind(status) + .bind(transaction_id) + .bind(id) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Mark payments still `pending` past a 1-hour deadline as `expired`, returning the rows that + /// were flipped so the caller can fire one webhook per expiry without a second query. + pub async fn expire_stale_payment_link_payments( + &self, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, PaymentLinkPayment>( + r#" + UPDATE payment_link_payments + SET status = 'expired' + WHERE status = 'pending' AND created_at < now() - interval '1 hour' + RETURNING * + "#, + ) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Payments recorded against a link (newest first), with cursor pagination. + /// + /// Includes pending intents, not just confirmed ones — a merchant wants to see that someone + /// started paying, and pending rows are how an abandoned checkout shows up. + pub async fn list_payment_link_payments( + &self, + payment_link_id: Uuid, + limit: i64, + before_id: Option, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, PaymentLinkPayment>( + r#" + SELECT * FROM payment_link_payments + WHERE payment_link_id = $1 + AND ($2::uuid IS NULL OR (created_at, id) < ( + SELECT created_at, id FROM payment_link_payments WHERE id = $2 + )) + ORDER BY created_at DESC, id DESC + LIMIT $3 + "#, + ) + .bind(payment_link_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Lifetime total (in USDC stroops) confirmed on a payment link. + pub async fn sum_payment_link_collected( + &self, + payment_link_id: Uuid, + ) -> Result { + let total: Option = sqlx::query_scalar( + r#" + SELECT COALESCE(SUM(amount_usdc_stroops), 0)::bigint + FROM payment_link_payments + WHERE payment_link_id = $1 AND status = 'confirmed' + "#, + ) + .bind(payment_link_id) + .fetch_one(&self.pool) + .await?; + Ok(total.unwrap_or(0)) + } + + /// Batched version of [`Store::sum_payment_link_collected`] for a link list page. + pub async fn sum_payment_link_collected_batch( + &self, + payment_link_ids: &[Uuid], + ) -> Result, StoreError> { + if payment_link_ids.is_empty() { + return Ok(Vec::new()); + } + let rows: Vec<(Uuid, i64)> = sqlx::query_as( + r#" + SELECT payment_link_id, COALESCE(SUM(amount_usdc_stroops), 0)::bigint AS total + FROM payment_link_payments + WHERE payment_link_id = ANY($1) AND status = 'confirmed' + GROUP BY payment_link_id + "#, + ) + .bind(payment_link_ids) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Atomically reserve budget and record a sponsored transaction. + /// + /// Inserts a `pending` row **only if** doing so keeps today's reserved fees within + /// `daily_budget_stroops` (a `NULL` budget means unlimited). The check and insert happen in one + /// statement (a conditional CTE), so concurrent sponsorships can't oversubscribe the budget. + /// Returns `StoreError::BudgetExceeded` if the budget would be exceeded, or + /// `StoreError::Conflict` if this `inner_tx_hash` was already sponsored (double-submit). + pub async fn try_reserve_sponsored_transaction( + &self, + wallet_id: Uuid, + inner_tx_hash: &str, + fee_stroops: i64, + daily_budget_stroops: Option, + ) -> Result { + // The read-then-insert below must be serialized per wallet. A bare conditional CTE is NOT + // enough: under READ COMMITTED every concurrent transaction computes `spent` from a + // snapshot taken before the others' inserts are visible, so N requests can each see the + // same total and all pass the budget guard (observed: 11 reservations against a 10-slot + // budget under 20 concurrent requests). + // + // A transaction-scoped advisory lock keyed on the wallet id makes the check-and-insert + // mutually exclusive for that wallet, while leaving other wallets fully parallel. The + // lock is released automatically when the transaction commits or rolls back. + let mut tx = self.pool.begin().await?; + + // Fold the wallet UUID into a stable i64 lock key. + let lock_key = { + let b = wallet_id.as_bytes(); + i64::from_be_bytes([b[0], b[1], b[2], b[3], b[4], b[5], b[6], b[7]]) + ^ i64::from_be_bytes([b[8], b[9], b[10], b[11], b[12], b[13], b[14], b[15]]) + }; + sqlx::query("SELECT pg_advisory_xact_lock($1)") + .bind(lock_key) + .execute(&mut *tx) + .await?; + + let result = sqlx::query_as::<_, SponsoredTransaction>( + r#" + WITH spent AS ( + SELECT COALESCE(SUM(fee_stroops), 0)::bigint AS total + FROM sponsored_transactions + WHERE wallet_id = $1 + AND status IN ('pending', 'confirmed') + AND created_at >= date_trunc('day', now() AT TIME ZONE 'UTC') + ) + INSERT INTO sponsored_transactions (wallet_id, inner_tx_hash, fee_stroops, status) + SELECT $1, $2, $3, 'pending' + FROM spent + WHERE $4::bigint IS NULL OR spent.total + $3 <= $4 + RETURNING * + "#, + ) + .bind(wallet_id) + .bind(inner_tx_hash) + .bind(fee_stroops) + .bind(daily_budget_stroops) + .fetch_optional(&mut *tx) + .await; + + // Commit before returning so the reservation (and the lock release) are durable. + if result.is_ok() { + tx.commit().await?; + } + + match result { + // A row means the insert (and budget check) succeeded. + Ok(Some(row)) => Ok(row), + // No row means the WHERE budget guard rejected the insert. + Ok(None) => Err(StoreError::BudgetExceeded), + // Unique violation on inner_tx_hash => already sponsored. + Err(e) => Err(StoreError::from_sqlx_conflict(e)), + } + } + + /// Update a sponsored transaction's outcome after submission. + pub async fn finalize_sponsored_transaction( + &self, + id: Uuid, + status: &str, + fee_bump_tx_hash: Option<&str>, + error: Option<&str>, + ) -> Result<(), StoreError> { + self.update_sponsored_tx_status(id, status, fee_bump_tx_hash, error) + .await + } + + /// Insert a sponsored transaction as `pending` (no budget check — see + /// [`Store::try_reserve_sponsored_transaction`] for the atomic budget-aware insert). + /// Fails with [`StoreError::Conflict`] if this `inner_tx_hash` was already recorded. + pub async fn record_sponsored_tx( + &self, + new: NewSponsoredTx<'_>, + ) -> Result { + sqlx::query_as::<_, SponsoredTransaction>( + r#" + INSERT INTO sponsored_transactions (wallet_id, inner_tx_hash, fee_stroops, status) + VALUES ($1, $2, $3, 'pending') + RETURNING * + "#, + ) + .bind(new.wallet_id) + .bind(new.inner_tx_hash) + .bind(new.fee_stroops) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from_sqlx_conflict) + } + + /// Update a sponsored transaction's status, fee-bump hash, and error. + pub async fn update_sponsored_tx_status( + &self, + id: Uuid, + status: &str, + fee_bump_tx_hash: Option<&str>, + error: Option<&str>, + ) -> Result<(), StoreError> { + sqlx::query( + "UPDATE sponsored_transactions SET status = $2, fee_bump_tx_hash = $3, error = $4 WHERE id = $1", + ) + .bind(id) + .bind(status) + .bind(fee_bump_tx_hash) + .bind(error) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Sum of **confirmed** sponsored fees for a wallet so far today (UTC) — i.e. actually spent. + /// (Pending rows are excluded; for budget *reservation* use + /// [`Store::sum_sponsored_fees_reserved_today`].) + pub async fn sum_sponsored_fees_today(&self, wallet_id: Uuid) -> Result { + let total: Option = sqlx::query_scalar( + r#" + SELECT COALESCE(SUM(fee_stroops), 0)::bigint + FROM sponsored_transactions + WHERE wallet_id = $1 + AND status = 'confirmed' + AND created_at >= date_trunc('day', now() AT TIME ZONE 'UTC') + "#, + ) + .bind(wallet_id) + .fetch_one(&self.pool) + .await?; + Ok(total.unwrap_or(0)) + } + + // --- token deny-list ------------------------------------------------- + + /// Add a token to the deny-list so it cannot be replayed after logout. + /// + /// `token_hash` must be the **SHA-256 hex** of the raw JWT (never the token itself). + /// `expires_at` should mirror the token's own `exp` claim so that rows can be pruned once + /// they are past their natural expiry and cannot match any valid token anyway. + /// + /// Inserting the same hash twice is harmless (ON CONFLICT DO NOTHING). + pub async fn denylist_token( + &self, + token_hash: &str, + user_id: Uuid, + expires_at: chrono::DateTime, + ) -> Result<(), StoreError> { + sqlx::query( + r#" + INSERT INTO token_denylist (token_hash, user_id, expires_at) + VALUES ($1, $2, $3) + ON CONFLICT (token_hash) DO NOTHING + "#, + ) + .bind(token_hash) + .bind(user_id) + .bind(expires_at) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Returns `true` if the token hash is present in the deny-list **and** has not yet expired. + /// + /// Expired rows are logically irrelevant (the token itself would fail `verify_token`'s expiry + /// check), but this query skips them so a slow pruning job doesn't affect correctness. + pub async fn is_token_denylisted(&self, token_hash: &str) -> Result { + let found: Option = sqlx::query_scalar( + "SELECT true FROM token_denylist WHERE token_hash = $1 AND expires_at > now() LIMIT 1", + ) + .bind(token_hash) + .fetch_optional(&self.pool) + .await?; + Ok(found.is_some()) + } + + // --- ingest cursor ---------------------------------------------------- + + /// Read the saved Horizon paging token for a wallet, if any. + pub async fn get_cursor(&self, wallet_id: Uuid) -> Result, StoreError> { + let token: Option = + sqlx::query_scalar("SELECT paging_token FROM ingest_cursor WHERE wallet_id = $1") + .bind(wallet_id) + .fetch_optional(&self.pool) + .await? + .flatten(); + Ok(token) + } + + /// Upsert the Horizon paging token for a wallet (durable resume point). + pub async fn set_cursor(&self, wallet_id: Uuid, paging_token: &str) -> Result<(), StoreError> { + sqlx::query( + r#" + INSERT INTO ingest_cursor (wallet_id, paging_token, updated_at) + VALUES ($1, $2, now()) + ON CONFLICT (wallet_id) + DO UPDATE SET paging_token = EXCLUDED.paging_token, updated_at = now() + "#, + ) + .bind(wallet_id) + .bind(paging_token) + .execute(&self.pool) + .await?; + Ok(()) + } + + // --- webhooks --------------------------------------------------------- + + /// Register a webhook endpoint for a wallet. + pub async fn create_webhook_endpoint( + &self, + wallet_id: Uuid, + url: &str, + secret: &str, + ) -> Result { + sqlx::query_as::<_, WebhookEndpoint>( + r#" + INSERT INTO webhook_endpoints (wallet_id, url, secret) + VALUES ($1, $2, $3) + RETURNING * + "#, + ) + .bind(wallet_id) + .bind(url) + .bind(secret) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from_sqlx_conflict) + } + + /// List the active webhook endpoints for a wallet. + pub async fn active_webhook_endpoints( + &self, + wallet_id: Uuid, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, WebhookEndpoint>( + "SELECT * FROM webhook_endpoints WHERE wallet_id = $1 AND active = true", + ) + .bind(wallet_id) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Deactivate a webhook endpoint by setting its active status to false. + pub async fn deactivate_webhook_endpoint(&self, id: Uuid) -> Result<(), StoreError> { + sqlx::query("UPDATE webhook_endpoints SET active = false WHERE id = $1") + .bind(id) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Fetch a single webhook endpoint by id. `NotFound` if it does not exist. + /// + /// Callers must still check `wallet_id` before returning data, so that an endpoint belonging + /// to another wallet is reported as 404 rather than 403 (no existence leak). + pub async fn get_webhook_endpoint(&self, id: Uuid) -> Result { + sqlx::query_as::<_, WebhookEndpoint>("SELECT * FROM webhook_endpoints WHERE id = $1") + .bind(id) + .fetch_optional(&self.pool) + .await? + .ok_or(StoreError::NotFound) + } + + /// Check whether a webhook endpoint is still active using its indexed id. + pub async fn is_webhook_endpoint_active(&self, id: Uuid) -> Result { + sqlx::query_scalar::<_, bool>( + "SELECT EXISTS (SELECT 1 FROM webhook_endpoints WHERE id = $1 AND active = true)", + ) + .bind(id) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from) + } + + /// An endpoint's delivery history, newest first, capped at `limit` rows. + pub async fn list_webhook_deliveries( + &self, + endpoint_id: Uuid, + limit: i64, + ) -> Result, StoreError> { + let rows = sqlx::query_as::<_, WebhookDelivery>( + r#" + SELECT * FROM webhook_deliveries + WHERE endpoint_id = $1 + ORDER BY created_at DESC, id DESC + LIMIT $2 + "#, + ) + .bind(endpoint_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; + Ok(rows) + } + + /// Record a webhook delivery attempt (audit log). Returns the delivery id. + pub async fn log_webhook_delivery( + &self, + endpoint_id: Uuid, + event_type: &str, + payload: &serde_json::Value, + status: &str, + attempts: i32, + response_code: Option, + ) -> Result { + let id: Uuid = sqlx::query_scalar( + r#" + INSERT INTO webhook_deliveries + (endpoint_id, event_type, payload, status, attempts, response_code) + VALUES ($1, $2, $3, $4, $5, $6) + RETURNING id + "#, + ) + .bind(endpoint_id) + .bind(event_type) + .bind(payload) + .bind(status) + .bind(attempts) + .bind(response_code) + .fetch_one(&self.pool) + .await?; + Ok(id) + } + + // --- token deny-list -------------------------------------------------- + + /// Revoke a JWT by inserting it into the deny-list. + /// + /// `expires_at` should match the token's `exp` claim (converted from Unix seconds). Duplicate + /// revocations (same token) are silently ignored via `ON CONFLICT DO NOTHING`. + pub async fn revoke_token( + &self, + token: &str, + expires_at: chrono::DateTime, + ) -> Result<(), StoreError> { + sqlx::query( + r#" + INSERT INTO token_denylist (token, expires_at) + VALUES ($1, $2) + ON CONFLICT (token) DO NOTHING + "#, + ) + .bind(token) + .bind(expires_at) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Return `true` if the token has been revoked (is in the deny-list). + pub async fn is_token_revoked(&self, token: &str) -> Result { + let exists: bool = + sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM token_denylist WHERE token = $1)") + .bind(token) + .fetch_one(&self.pool) + .await?; + Ok(exists) + } + + /// Delete expired deny-list entries (those whose `expires_at` is in the past). + /// + /// Intended to be called periodically (e.g. once per hour in a background task) to prevent + /// unbounded table growth. Safe to skip — expired tokens are rejected by `verify_token()` + /// regardless of the deny-list. + pub async fn purge_expired_tokens(&self) -> Result { + let result = sqlx::query("DELETE FROM token_denylist WHERE expires_at < now()") + .execute(&self.pool) + .await?; + Ok(result.rows_affected()) + } +} diff --git a/crates/webhooks/src/lib.rs b/crates/webhooks/src/lib.rs index 85917db..2059a6f 100644 --- a/crates/webhooks/src/lib.rs +++ b/crates/webhooks/src/lib.rs @@ -179,6 +179,12 @@ impl WebhookSender { } }), ) + .await; + + } + } + }), + ) .await; let (status, outcome) = match result { diff --git a/docs/api.md b/docs/api.md index 139612a..e69de29 100644 --- a/docs/api.md +++ b/docs/api.md @@ -1,223 +0,0 @@ -# API reference - -The machine-readable contract is **[openapi.yaml](openapi.yaml)**, and it is enforced: the -`drift_tests` integration test validates live responses against that spec, so the two cannot -silently diverge. This page is the human-readable tour. - -All responses use a consistent envelope — **including errors**, where `data` is `null`: - -```json -{ "statusCode": 200, "message": "OK", "data": { } } -``` - -## Authentication - -Two credential types, with deliberately different power: - -| Credential | Header | Can do | -|---|---|---| -| **Dashboard JWT** | `Authorization: Bearer ` | Everything: create wallets, provision a gas tank, read the key backup | -| **Wallet API key** | `Authorization: Bearer ` | Per-wallet operations only. **Cannot** provision a gas tank or read a backup | - -Neither can move funds — see below. Tokens carry a unique `jti`; `logout` and `refresh` both -deny-list the presented token, and every authenticated request checks that deny-list. - -- `POST /v1/auth/signup` — create an account, returns a JWT. -- `POST /v1/auth/login` — returns a JWT. -- `POST /v1/auth/refresh` — issue a new token **and revoke the presented one**. -- `POST /v1/auth/logout` — revoke the presented token (a second logout is `401`, not `200`). -- `POST /v1/auth/change-password` — `{current_password, new_password}`; re-verifies the current - password, revokes **every** session issued before the change (per-user `session_epoch`), and - returns a fresh token. Login-JWT only (not API keys); rate-limited per IP and per user. -- `GET /v1/auth/me` — the current user. -- `POST /v1/auth/request-password-reset` — `{email}`; emails a 10-minute OTP if a verified account - exists. Always returns the same `200` either way (no account enumeration). -- `POST /v1/auth/confirm-password-reset` — `{email, code, new_password}`; sets the password and - **revokes every existing session**. Any failure is `400 invalid or expired code`. -- `POST /v1/auth/change-email` — `{new_email, password}` (login required): verifies the current - password and emails an OTP to the **new** address. Nothing changes yet. -- `POST /v1/auth/change-email/confirm` — `{new_email, code}` (login required): applies the change - once the new address's OTP is confirmed, and notifies the old address. - -## Custody model — read this before the wallet endpoints - -octo is **non-custodial**. The wallet's private key is generated and held **client-side**; the -server stores only the public account and an opaque, client-encrypted backup blob it cannot -decrypt. Consequently: - -- There is **no endpoint that signs a payment for you.** You build and sign locally, then relay. -- `POST /v1/wallets/:id/withdraw` is a **`410 Gone` tombstone** pointing integrators at - `submit-signed`. -- `POST /v1/wallets/:id/trustlines` takes `{asset_code, asset_issuer, limit_stroops?}`, validates - them, and returns ChangeTrust signing info (`account`, `sequence`, `network_passphrase`, - `base_fee_stroops`, `limit_stroops`, `submit_url`). The server never signs it — the client builds - and signs the ChangeTrust locally and relays it via `submit-signed`. - -## Wallets - -- `POST /v1/wallets` — register a wallet from a **client-generated** keypair. - Body: `{ "public_key": "G...", "encrypted_backup"?: string, "label"?: string, - "description"?: string }`. `public_key` is required; a body without it is `400`. - Returns `201` with `{ id, network, address, custody, funded }`. - **Never returns a mnemonic** — the client generated it and the server never saw it. -- `GET /v1/wallets` — list your wallets (paginated). -- `GET /v1/wallets/{id}` — wallet details. -- `GET /v1/wallets/{id}/balances` — live on-chain balances. Fetched synchronously from - Horizon under a **10 s** route timeout (independent of per-attempt retries); if Horizon is - slower than that, the request ends with `504` in the standard envelope — safe to retry. -- `GET /v1/wallets/{id}/transactions` — deposits + outbound transfers (paginated, optional `?direction=deposit|withdrawal`). -- `GET /v1/wallets/{id}/backup` — the opaque client-encrypted backup blob, for new-device - recovery. **Dashboard JWT only.** Useless without the user's password. - -## Moving funds (the non-custodial path) - -1. `GET /v1/wallets/{id}/signing-info` — returns the account `sequence`, the network - passphrase, and the base fee, so you can build a transaction without talking to Horizon. -2. Build and **sign locally**. -3. `POST /v1/wallets/{id}/submit-signed` — body `{ "transaction_xdr": "" }`. - The server validates (v1 envelope, at least one signature, source account == this wallet, - operation-type allowlist) and relays it **unmodified**. On failure it returns Horizon's - result codes (`tx_bad_seq`, `op_no_trust`, …) so you can correct and re-sign. - -## Addresses - -- `POST /v1/wallets/{id}/addresses` — generate a dedicated customer address. - Returns `muxed_address` (`M...`) **and** the `{ base_address, memo_id }` fallback. -- `GET /v1/wallets/{id}/addresses` — list addresses (paginated). - -## Gas sponsorship - -Lets you pay your users' Stellar fees. The **gas tank** is a separate, server-held account that -carries fee float only — the one server-held key in the system, bounded by your gas budget. - -- `POST /v1/wallets/{id}/gas-tank` — provision the gas tank. **Dashboard JWT only** (an API key - gets `401`). Idempotent: a second call returns the existing tank. -- `GET /v1/wallets/{id}/gas-tank` — the tank's public account (`gas_tank_address`), whether it is - `provisioned`, and today's `spent_today_stroops` against `daily_budget_stroops`. A wallet with - no tank returns `200` with `provisioned: false`. Never includes the sealed seed. -- `GET /v1/wallets/{id}/sponsorship` / `PUT` — read/update `enabled`, the per-transaction fee - cap, and the daily budget. - - `daily_budget_stroops`: `null`/omitted = **unlimited**; `0` = sponsorship **fully disabled** - for the day (every sponsor request gets `429`); negative → `400`. - - `per_tx_fee_cap_stroops`: `null`/omitted = no per-transaction cap; negative → `400`. - - When both are set, `per_tx_fee_cap_stroops` must be `<=` `daily_budget_stroops`, otherwise - `400` naming both values. Either field may be left unset independently. -- `POST /v1/wallets/{id}/sponsor` — fee-bump a user's **already-signed** inner transaction. - The gas tank signs only the outer fee-bump envelope; the inner transaction is passed through - untouched. Over budget → `429`; duplicate inner tx → `409`. -- `GET /v1/wallets/{id}/sponsored-transactions` — sponsorship history (paginated, filterable - by status). - -## Webhooks - -- `POST /v1/wallets/{id}/webhooks` — register an endpoint (URL + generated secret). -- `GET /v1/wallets/{id}/webhooks` — list active endpoints. -- `DELETE /v1/wallets/{id}/webhooks/{endpoint_id}` — deactivate (soft delete, so the delivery - history survives as an audit trail). -- `GET /v1/wallets/{id}/webhooks/{endpoint_id}/deliveries` — delivery history (`?limit=`, - default 50, max 200). Each row carries `response_code` (HTTP status of the last attempt, `null` on - a connection error/timeout) and `response_body_snippet` (first ≤ 1 KiB of the response body, with - the signature and secret redacted). Transport errors and 5xx are retried with backoff (3 attempts, - 20 s ceiling); other non-2xx responses are not retried. - -Deliveries are signed `HMAC-SHA256` over the raw body. Endpoint URLs are SSRF-screened: -loopback, private and link-local targets are rejected, IPv4 and bracketed IPv6 alike — -including IPv4-mapped IPv6 (`[::ffff:127.0.0.1]`) and the unspecified address (`0.0.0.0`, `[::]`). - -## API keys - -All three require a **dashboard JWT** and wallet ownership — an API key can never manage keys, -so it cannot escalate or revoke itself. - -- `POST /v1/wallets/{id}/api-key` — generate (the plaintext key is shown **once**; only a - SHA-256 hash is stored). The first key needs no body. Once a key exists, rotating it requires - `{"confirm": true}` — otherwise `409` ("an API key already exists; pass confirm=true to rotate - it"). Rotation immediately invalidates the previous key. -- `GET /v1/wallets/{id}/api-key` — metadata (prefix, created_at) — never the key itself. -- `DELETE /v1/wallets/{id}/api-key` — revoke. - -## Payment link checkout flow - -A payment link is a shareable, USDC-only checkout page. The merchant creates it once (authenticated); -every payer step after that is **public — no credential** — and keyed by the link's `slug`. Public -routes are rate-limited per client IP (per minute): read 60, **intent 5**, signing-info 60, -submit 20, status 60. Over the limit → `429`. - -Merchant, once (dashboard JWT or wallet API key; `$TOKEN` as in the README): - -```bash -curl -s -X POST localhost:8080/v1/wallets//payment-links \ - -H "authorization: Bearer $TOKEN" -H 'content-type: application/json' \ - -d '{"name":"Order #1042","amount_usdc_stroops":50000000}' | jq # omit the amount for a flexible link -# -> data.slug (e.g. "3f9c1a7d2e"), data.url (hosted checkout page) -``` - -Payer, from there (`$SLUG` is `data.slug`; a fixed-amount link ignores any amount the payer sends, -a flexible link requires `amount_usdc_stroops > 0`): - -```bash -# 1. Fetch the link: what is being paid and where. 404 if the slug is unknown or the link is inactive. -curl -s localhost:8080/v1/pay/$SLUG | jq -# -> { name, description, image_url, redirect_url, amount_usdc_stroops, deposit_address, asset_code: "USDC" } - -# 2. Create a payment intent. Each intent gets its OWN muxed deposit address, so a deposit maps to -# exactly one payment. payer_name / payer_email are optional. -curl -s -X POST localhost:8080/v1/pay/$SLUG/intent \ - -H 'content-type: application/json' -d '{"payer_name":"Ada","payer_email":"ada@example.com"}' | jq -# -> 201 { payment_id, deposit_address, amount_usdc_stroops } (keep payment_id) - -# 3. Get what you need to build the transaction. `account` is the PAYER's own G... account, so the -# returned sequence is the payer's. Omit it and the merchant wallet's account is used. -# A payer account that does not exist on the network yet (unfunded) returns 404. -curl -s "localhost:8080/v1/pay/$SLUG/signing-info?account=" | jq -# -> { account, sequence, network_passphrase, base_fee_stroops } - -# 4. Build and SIGN LOCALLY (e.g. in Freighter): exactly one USDC Payment to the intent's -# deposit_address, nothing else. Then relay it, passing the payment_id from step 2. -curl -s -X POST localhost:8080/v1/pay/$SLUG/submit-signed \ - -H 'content-type: application/json' \ - -d '{"transaction_xdr":"","payment_id":""}' | jq -# -> 201 { status: "confirmed" | "failed", stellar_tx_hash, detail } - -# 5. Poll until the deposit is matched (the pay page polls about every 3s). -curl -s localhost:8080/v1/pay/$SLUG/payments/ | jq -# -> { status, transaction_id, expected_usdc_stroops, received_usdc_stroops } -``` - -Things an integrator should know: - -- **`submit-signed` returns `201` even when `status` is `"failed"`** — check `status` and `detail` - (a Horizon result code such as `op_underfunded`), not just the HTTP code. `400` means the - transaction was rejected before relay: not a v1 envelope, unsigned, or not exactly one USDC - `Payment` to this intent's `deposit_address`. -- **The relay is deliberately narrow**: it never signs and cannot spend anything except that one - payment. Without `payment_id` it falls back to the link's own address (legacy clients). -- **Status** is one of `pending`, `confirmed`, `expired`, `underpaid`, `overpaid`. `received_usdc_stroops` - is `null` until a deposit is matched; `expected_usdc_stroops` is always present so a client can - show "you sent X, expected Y". A `confirmed` status comes from the ingest worker seeing the - deposit, not from the submit response. -- **Payer PII is write-only.** `payer_name` and `payer_email` are stored for the merchant but no - public route returns them, and none of the public responses above carries merchant-internal - fields (wallet id, link id, collected totals). If a public response ever gains or loses a field, - update the shapes shown here. -- Amounts are integer stroops: `50000000` = 5 USDC. - -Merchant-side management (list/get/deactivate links and list a link's payments) lives under -`/v1/wallets/{id}/payment-links` — see [openapi.yaml](openapi.yaml). - -## Audit logs - -- `GET /v1/audit-logs` — your account's activity, filterable by `category` and a free-text - `search`. Valid categories: `authentication`, `wallet`, `address`, `credentials`, `configuration`, `sponsorship`. The category set and every event that emits one is catalogued in - [audit-log.md](audit-log.md). - -## Conventions - -- **Pagination:** list endpoints take `?limit=` (default 50, max 200) and `?before=` for - keyset pagination. They return `{ "data": [...], "next_cursor": }` — note this - sits *inside* the response envelope, so the full shape is - `{ statusCode, message, data: { data: [...], next_cursor } }`. -- **Amounts** are integer **stroops** (1 XLM = 10,000,000) end-to-end — never floats. -- **Errors** map to `400` (validation), `401`, `403`, `404`, `409` (conflict), `410` (removed - custodial endpoints), `413` (body over 64 KiB), `429` (budget exceeded), `504` (upstream - Horizon exceeded a route timeout). There is no `422`.