diff --git a/crates/api/src/auth.rs b/crates/api/src/auth.rs index cc8f59f..4691b9e 100644 --- a/crates/api/src/auth.rs +++ b/crates/api/src/auth.rs @@ -172,7 +172,7 @@ fn check_auth_rate_limit( let ip = crate::rate_limit::client_ip(headers, peer); if state .rate_limiter() - .check(&ip, "auth", 10, std::time::Duration::from_secs(60)) + .check(&ip, "auth", crate::rate_limit::AUTH_RATE_LIMIT, crate::rate_limit::AUTH_RATE_WINDOW) { Ok(()) } else { @@ -333,8 +333,8 @@ pub async fn resend_otp( if !state.rate_limiter().check( &format!("otp:{user_id}"), "otp_resend", - 3, - std::time::Duration::from_secs(60 * 60), + crate::rate_limit::OTP_RESEND_USER_LIMIT, + crate::rate_limit::OTP_RESEND_USER_WINDOW, ) { return Err(ApiError::TooManyRequests( "too many resend attempts — wait a while and try again".into(), @@ -344,8 +344,8 @@ pub async fn resend_otp( if !state.rate_limiter().check( &ip, "otp_resend_ip", - 10, - std::time::Duration::from_secs(60 * 60), + crate::rate_limit::OTP_RESEND_IP_LIMIT, + crate::rate_limit::OTP_RESEND_IP_WINDOW, ) { return Err(ApiError::TooManyRequests( "too many resend attempts — wait a while and try again".into(), diff --git a/crates/api/src/horizon.rs b/crates/api/src/horizon.rs index 0e43073..5d046ce 100644 --- a/crates/api/src/horizon.rs +++ b/crates/api/src/horizon.rs @@ -387,6 +387,23 @@ impl Horizon { Err(ResilienceError::Exhausted(_)) => Err(ApiError::Internal), } } + + // Probe Horizon root endpoint to verify node reachability. + pub async fn check_reachability(&self) -> Result<(), String> { + let url = self.base_url.trim_end_matches('/'); + let resp = self + .http + .get(url) + .timeout(Duration::from_secs(3)) + .send() + .await + .map_err(|e| format!("horizon unreachable: {e}"))?; + if resp.status().is_success() { + Ok(()) + } else { + Err(format!("horizon returned HTTP {}", resp.status())) + } + } } // --------------------------------------------------------------------------- diff --git a/crates/api/src/lib.rs b/crates/api/src/lib.rs index 4da04bd..a61f771 100644 --- a/crates/api/src/lib.rs +++ b/crates/api/src/lib.rs @@ -18,11 +18,12 @@ pub mod submit_validation; pub use error::{ApiError, ApiResult, Envelope}; pub use state::AppState; -use axum::extract::{DefaultBodyLimit, Request}; +use axum::extract::{DefaultBodyLimit, Request, State}; +use axum::http::StatusCode; use axum::middleware::{self, Next}; use axum::response::{IntoResponse, Response}; use axum::routing::{delete, get, post}; -use axum::Router; +use axum::{Json, Router}; use std::time::Duration; use tower_http::cors::{Any, CorsLayer}; @@ -52,6 +53,7 @@ pub fn build_router(state: AppState) -> Router { // together with the error handler that turns an oversized body into a 413 envelope. Router::new() .route("/health", get(health)) + .route("/health/ready", get(health_ready)) .route("/v1/auth/signup", post(auth::signup)) .route("/v1/auth/verify-email", post(auth::verify_email)) .route("/v1/auth/resend-otp", post(auth::resend_otp)) @@ -234,6 +236,53 @@ async fn health() -> &'static str { "ok" } +// Readiness probe checking database and Horizon reachability. +async fn health_ready(State(state): State) -> impl IntoResponse { + let mut db_ok = false; + let mut horizon_ok = false; + let mut db_err = None; + let mut horizon_err = None; + + match state.store().ping().await { + Ok(_) => db_ok = true, + Err(e) => db_err = Some(e.to_string()), + } + + match state.horizon().check_reachability().await { + Ok(_) => horizon_ok = true, + Err(e) => horizon_err = Some(e), + } + + if db_ok && horizon_ok { + ( + StatusCode::OK, + Json(serde_json::json!({ + "status": "ready", + "database": "ok", + "horizon": "ok" + })), + ) + } else { + let mut failed = Vec::new(); + if !db_ok { + failed.push("database"); + } + if !horizon_ok { + failed.push("horizon"); + } + ( + StatusCode::SERVICE_UNAVAILABLE, + Json(serde_json::json!({ + "status": "not_ready", + "database": if db_ok { "ok".to_string() } else { db_err.unwrap_or_else(|| "unreachable".into()) }, + "horizon": if horizon_ok { "ok".to_string() } else { horizon_err.unwrap_or_else(|| "unreachable".into()) }, + "failed": failed, + "error": format!("unreachable dependencies: {}", failed.join(", ")) + })), + ) + } +} + // NOTE: a `handle_errors` HandleErrorLayer helper lived here to convert oversized-body errors // into a 413 envelope. It is unnecessary with `DefaultBodyLimit` (axum renders that rejection as // 413 itself) and did not satisfy `Router::layer`'s Service bounds, so it was removed. diff --git a/crates/api/src/rate_limit.rs b/crates/api/src/rate_limit.rs index 28f4812..476e8b5 100644 --- a/crates/api/src/rate_limit.rs +++ b/crates/api/src/rate_limit.rs @@ -12,6 +12,24 @@ use std::time::{Duration, Instant}; /// Cap on tracked (ip, class) buckets before expired entries are swept. const SWEEP_THRESHOLD: usize = 10_000; +// Rate limit thresholds and fixed windows for API endpoints. +pub const AUTH_RATE_LIMIT: u32 = 10; +pub const AUTH_RATE_WINDOW: Duration = Duration::from_secs(60); +pub const OTP_RESEND_USER_LIMIT: u32 = 3; +pub const OTP_RESEND_USER_WINDOW: Duration = Duration::from_secs(3600); +pub const OTP_RESEND_IP_LIMIT: u32 = 10; +pub const OTP_RESEND_IP_WINDOW: Duration = Duration::from_secs(3600); +pub const PAY_READ_LIMIT: u32 = 60; +pub const PAY_READ_WINDOW: Duration = Duration::from_secs(60); +pub const PAY_INTENT_LIMIT: u32 = 5; +pub const PAY_INTENT_WINDOW: Duration = Duration::from_secs(60); +pub const PAY_STATUS_LIMIT: u32 = 60; +pub const PAY_STATUS_WINDOW: Duration = Duration::from_secs(60); +pub const PAY_SIGNING_INFO_LIMIT: u32 = 60; +pub const PAY_SIGNING_INFO_WINDOW: Duration = Duration::from_secs(60); +pub const PAY_SUBMIT_LIMIT: u32 = 20; +pub const PAY_SUBMIT_WINDOW: Duration = Duration::from_secs(60); + /// Bucket key: the client IP plus the endpoint class it is being limited against. type BucketKey = (String, &'static str); /// Bucket value: when the current fixed window started, and hits so far within it. diff --git a/crates/api/src/routes/payment_links.rs b/crates/api/src/routes/payment_links.rs index 487f5b8..61f33c1 100644 --- a/crates/api/src/routes/payment_links.rs +++ b/crates/api/src/routes/payment_links.rs @@ -326,11 +326,12 @@ fn check_public_rate_limit( peer: Option>, class: &'static str, limit: u32, + window: std::time::Duration, ) -> Result<(), ApiError> { let ip = crate::rate_limit::client_ip(headers, peer.map(|c| c.0)); if state .rate_limiter() - .check(&ip, class, limit, std::time::Duration::from_secs(60)) + .check(&ip, class, limit, window) { Ok(()) } else { @@ -347,7 +348,14 @@ pub async fn get_public_payment_link( peer: Option>, headers: HeaderMap, ) -> ApiResult>> { - check_public_rate_limit(&state, &headers, peer, "pay_read", 60)?; + check_public_rate_limit( + &state, + &headers, + peer, + "pay_read", + crate::rate_limit::PAY_READ_LIMIT, + crate::rate_limit::PAY_READ_WINDOW, + )?; let link = state.store().get_payment_link_by_slug(&slug).await?; if !link.active { return Err(ApiError::NotFound); @@ -390,7 +398,14 @@ pub async fn create_payment_intent( headers: HeaderMap, body: Bytes, ) -> ApiResult<(StatusCode, Json>)> { - check_public_rate_limit(&state, &headers, peer, "pay_intent", 5)?; + check_public_rate_limit( + &state, + &headers, + peer, + "pay_intent", + crate::rate_limit::PAY_INTENT_LIMIT, + crate::rate_limit::PAY_INTENT_WINDOW, + )?; let link = state.store().get_payment_link_by_slug(&slug).await?; if !link.active { return Err(ApiError::NotFound); @@ -463,7 +478,14 @@ pub async fn get_payment_status( headers: HeaderMap, ) -> ApiResult>> { // The pay page polls this every ~3s while waiting, so the ceiling is generous. - check_public_rate_limit(&state, &headers, peer, "pay_status", 60)?; + check_public_rate_limit( + &state, + &headers, + peer, + "pay_status", + crate::rate_limit::PAY_STATUS_LIMIT, + crate::rate_limit::PAY_STATUS_WINDOW, + )?; let link = state.store().get_payment_link_by_slug(&slug).await?; let payment = state .store() @@ -501,7 +523,14 @@ pub async fn public_signing_info( peer: Option>, headers: HeaderMap, ) -> ApiResult>> { - check_public_rate_limit(&state, &headers, peer, "pay_signing_info", 60)?; + check_public_rate_limit( + &state, + &headers, + peer, + "pay_signing_info", + crate::rate_limit::PAY_SIGNING_INFO_LIMIT, + crate::rate_limit::PAY_SIGNING_INFO_WINDOW, + )?; // Confirms the link exists/is active before doing any Horizon work on the caller's behalf. let link = state.store().get_payment_link_by_slug(&slug).await?; if !link.active { @@ -552,7 +581,14 @@ pub async fn submit_payment( headers: HeaderMap, body: Bytes, ) -> ApiResult<(StatusCode, Json>)> { - check_public_rate_limit(&state, &headers, peer, "pay_submit", 20)?; + check_public_rate_limit( + &state, + &headers, + peer, + "pay_submit", + crate::rate_limit::PAY_SUBMIT_LIMIT, + crate::rate_limit::PAY_SUBMIT_WINDOW, + )?; let link = state.store().get_payment_link_by_slug(&slug).await?; if !link.active { return Err(ApiError::NotFound); diff --git a/crates/api/tests/api_tests.rs b/crates/api/tests/api_tests.rs index 439c086..7027681 100644 --- a/crates/api/tests/api_tests.rs +++ b/crates/api/tests/api_tests.rs @@ -435,6 +435,97 @@ async fn health_is_public_and_ok() { assert_eq!(resp.status(), StatusCode::OK); } +// Local mock Horizon server for readiness testing. +async fn start_mock_horizon_ok() -> String { + let app = Router::new().route("/", axum::routing::get(|| async { "horizon ok" })); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind mock horizon"); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + format!("http://{addr}") +} + +#[tokio::test] +async fn health_ready_returns_200_when_db_and_horizon_are_both_reachable() { + let Some(_) = test_state().await else { + return; + }; + let mock_horizon = start_mock_horizon_ok().await; + let url = database_url().unwrap(); + let store = Store::connect(&url).await.expect("connect"); + let state = AppState::new( + store, + [42u8; 32], + StellarNetwork::Testnet, + mock_horizon, + None, + octo_email::EmailSender::new_captured(), + ); + let app = build_router(state); + let resp = app.oneshot(get("/health/ready")).await.unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + let json = body_json(resp).await; + assert_eq!(json["status"], "ready"); + assert_eq!(json["database"], "ok"); + assert_eq!(json["horizon"], "ok"); +} + +#[tokio::test] +async fn health_ready_returns_a_clear_503_naming_the_db_when_the_database_is_unreachable() { + let mock_horizon = start_mock_horizon_ok().await; + let dead_pool = sqlx::postgres::PgPoolOptions::new() + .acquire_timeout(std::time::Duration::from_millis(100)) + .connect_lazy("postgres://postgres:wrong@127.0.0.1:1/nonexistent") + .unwrap(); + let store = Store::from_pool(dead_pool); + let state = AppState::new( + store, + [42u8; 32], + StellarNetwork::Testnet, + mock_horizon, + None, + octo_email::EmailSender::new_captured(), + ); + let app = build_router(state); + let resp = app.oneshot(get("/health/ready")).await.unwrap(); + assert_eq!(resp.status(), StatusCode::SERVICE_UNAVAILABLE); + let json = body_json(resp).await; + assert_eq!(json["status"], "not_ready"); + assert_eq!(json["horizon"], "ok"); + assert!(json["database"] != "ok"); + let error_str = json["error"].as_str().unwrap(); + assert!(error_str.contains("database")); +} + +#[tokio::test] +async fn health_ready_returns_a_clear_503_naming_horizon_when_horizon_is_unreachable() { + let Some(_) = test_state().await else { + return; + }; + let url = database_url().unwrap(); + let store = Store::connect(&url).await.expect("connect"); + let state = AppState::new( + store, + [42u8; 32], + StellarNetwork::Testnet, + "http://127.0.0.1:1".into(), + None, + octo_email::EmailSender::new_captured(), + ); + let app = build_router(state); + let resp = app.oneshot(get("/health/ready")).await.unwrap(); + assert_eq!(resp.status(), StatusCode::SERVICE_UNAVAILABLE); + let json = body_json(resp).await; + assert_eq!(json["status"], "not_ready"); + assert_eq!(json["database"], "ok"); + assert!(json["horizon"] != "ok"); + let error_str = json["error"].as_str().unwrap(); + assert!(error_str.contains("horizon")); +} + #[tokio::test] async fn backup_round_trips_the_opaque_blob_verbatim() { let Some(state) = test_state().await else { diff --git a/crates/api/tests/sponsor_e2e_tests.rs b/crates/api/tests/sponsor_e2e_tests.rs index e266606..d20551b 100644 --- a/crates/api/tests/sponsor_e2e_tests.rs +++ b/crates/api/tests/sponsor_e2e_tests.rs @@ -134,6 +134,29 @@ async fn create_wallet(app: &Router, token: &str) -> (String, String) { (id, address) } +// Create a client-custody wallet without provisioning a gas tank. +async fn create_wallet_without_gas_tank(app: &Router, token: &str) -> (String, String) { + let kp = DalekKeyPair::random().unwrap(); + let reg_body = common::wallet_body(app, token, &kp).await; + let resp = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/v1/wallets") + .header("content-type", "application/json") + .header("authorization", format!("Bearer {token}")) + .body(Body::from(reg_body)) + .unwrap(), + ) + .await + .unwrap(); + let w = body_json(resp).await; + let id = w["data"]["id"].as_str().unwrap().to_string(); + let address = w["data"]["address"].as_str().unwrap().to_string(); + (id, address) +} + /// Local mock Horizon that accepts POST /transactions and returns a successful submission. async fn start_mock_horizon() -> String { async fn submit() -> axum::Json { @@ -813,3 +836,213 @@ async fn e2e_concurrent_sponsor_requests_respect_budget() { "total reserved fees {total:?} must not exceed budget {daily_budget}" ); } + +#[tokio::test] +async fn sponsor_rejects_when_sponsorship_is_disabled_for_the_wallet() { + let horizon = start_mock_horizon().await; + let Some(state) = test_state(horizon).await else { + return; + }; + let app = build_router(state.clone()); + let token = auth_token(&app, &state).await; + let (wallet_id, master_g) = create_wallet(&app, &token).await; + + let uri = format!("/v1/wallets/{wallet_id}/sponsorship"); + let resp = app + .clone() + .oneshot(put_json_auth( + &uri, + r#"{"enabled":false,"daily_budget_stroops":1000000}"#, + &token, + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let xdr = random_payment_xdr(&master_g); + let body = format!(r#"{{"transaction_xdr":"{xdr}","max_base_fee_stroops":200}}"#); + let resp = app + .oneshot(post_json_auth( + &format!("/v1/wallets/{wallet_id}/sponsor"), + &body, + &token, + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::FORBIDDEN); + let json = body_json(resp).await; + assert_eq!( + json["message"], + "gas sponsorship is not enabled for this wallet" + ); +} + +#[tokio::test] +async fn sponsor_rejects_a_max_fee_above_the_per_tx_cap() { + let horizon = start_mock_horizon().await; + let Some(state) = test_state(horizon).await else { + return; + }; + let app = build_router(state.clone()); + let token = auth_token(&app, &state).await; + let (wallet_id, master_g) = create_wallet(&app, &token).await; + + let cap = 250_000_i64; + let uri = format!("/v1/wallets/{wallet_id}/sponsorship"); + let resp = app + .clone() + .oneshot(put_json_auth( + &uri, + &format!(r#"{{"enabled":true,"per_tx_fee_cap_stroops":{cap}}}"#), + &token, + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let xdr = random_payment_xdr(&master_g); + let body = format!(r#"{{"transaction_xdr":"{xdr}","max_base_fee_stroops":{}}}"#, cap + 1); + let resp = app + .oneshot(post_json_auth( + &format!("/v1/wallets/{wallet_id}/sponsor"), + &body, + &token, + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); + let json = body_json(resp).await; + assert!( + json["message"] + .as_str() + .unwrap() + .contains("exceeds the per-transaction cap"), + "expected cap error, got: {}", + json["message"] + ); +} + +#[tokio::test] +async fn sponsor_rejects_a_self_sponsoring_inner_transaction() { + let horizon = start_mock_horizon().await; + let Some(state) = test_state(horizon).await else { + return; + }; + let app = build_router(state.clone()); + let token = auth_token(&app, &state).await; + let (wallet_id, master_g) = create_wallet(&app, &token).await; + enable_sponsorship_via_api(&app, &token, &wallet_id).await; + + let dest_g = "GBAW5XGWORWVFE2XTJYDTLDHXTY2Q2MO73HYCGB3XMFMQ562Q2W2GJQX"; + let master_g_key = stellar_base::crypto::PublicKey::from_account_id(&master_g).unwrap(); + let dest = stellar_base::crypto::PublicKey::from_account_id(dest_g).unwrap(); + let op = Operation::new_payment() + .with_destination(dest) + .with_amount(stellar_base::amount::Stroops::new(100)) + .unwrap() + .with_asset(stellar_base::asset::Asset::new_native()) + .build() + .unwrap(); + let tx = Transaction::builder(master_g_key, 1, MIN_BASE_FEE) + .add_operation(op) + .into_transaction() + .unwrap(); + let xdr = tx.into_envelope().xdr_base64().unwrap(); + + let body = format!(r#"{{"transaction_xdr":"{xdr}","max_base_fee_stroops":200}}"#); + let resp = app + .oneshot(post_json_auth( + &format!("/v1/wallets/{wallet_id}/sponsor"), + &body, + &token, + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); + let json = body_json(resp).await; + assert_eq!( + json["message"], + "inner transaction source must not be the master account" + ); +} + +#[tokio::test] +async fn sponsor_rejects_once_the_daily_budget_is_exhausted() { + let horizon = start_mock_horizon().await; + let Some(state) = test_state(horizon).await else { + return; + }; + let app = build_router(state.clone()); + let token = auth_token(&app, &state).await; + let (wallet_id, master_g) = create_wallet(&app, &token).await; + + let uri = format!("/v1/wallets/{wallet_id}/sponsorship"); + let resp = app + .clone() + .oneshot(put_json_auth( + &uri, + r#"{"enabled":true,"daily_budget_stroops":10000000}"#, + &token, + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let fee = 6_000_000; + let body1 = format!( + r#"{{"transaction_xdr":"{}","max_base_fee_stroops":{fee}}}"#, + random_payment_xdr(&master_g) + ); + let sponsor_uri = format!("/v1/wallets/{wallet_id}/sponsor"); + let resp = app + .clone() + .oneshot(post_json_auth(&sponsor_uri, &body1, &token)) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::CREATED); + + let body2 = format!( + r#"{{"transaction_xdr":"{}","max_base_fee_stroops":{fee}}}"#, + random_payment_xdr(&master_g) + ); + let resp = app + .oneshot(post_json_auth(&sponsor_uri, &body2, &token)) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS); + let json = body_json(resp).await; + assert_eq!(json["message"], "daily sponsorship budget exceeded"); +} + +#[tokio::test] +async fn sponsor_rejects_for_a_client_custody_wallet_with_no_gas_tank() { + let horizon = start_mock_horizon().await; + let Some(state) = test_state(horizon).await else { + return; + }; + let app = build_router(state.clone()); + let token = auth_token(&app, &state).await; + let (wallet_id, master_g) = create_wallet_without_gas_tank(&app, &token).await; + enable_sponsorship_via_api(&app, &token, &wallet_id).await; + + let xdr = random_payment_xdr(&master_g); + let body = format!(r#"{{"transaction_xdr":"{xdr}","max_base_fee_stroops":200}}"#); + let resp = app + .oneshot(post_json_auth( + &format!("/v1/wallets/{wallet_id}/sponsor"), + &body, + &token, + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::FORBIDDEN); + let json = body_json(resp).await; + assert!( + json["message"] + .as_str() + .unwrap() + .contains("no gas-tank account"), + "expected gas-tank error, got: {}", + json["message"] + ); +} diff --git a/crates/store/src/lib.rs b/crates/store/src/lib.rs index cf3173f..8e58090 100644 --- a/crates/store/src/lib.rs +++ b/crates/store/src/lib.rs @@ -99,6 +99,12 @@ impl Store { &self.pool } + // Ping the database to verify pool reachability. + pub async fn ping(&self) -> Result<(), StoreError> { + sqlx::query("SELECT 1").execute(&self.pool).await?; + Ok(()) + } + // --- users ------------------------------------------------------------ /// Create a user. `email` should already be lowercased by the caller. Returns @@ -468,22 +474,13 @@ impl Store { 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?; + let query = cursor_pagination_query("wallets", "user_id"); + let rows = sqlx::query_as::<_, Wallet>(&query) + .bind(user_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; Ok(rows) } @@ -736,22 +733,13 @@ impl Store { 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?; + let query = cursor_pagination_query("addresses", "wallet_id"); + let rows = sqlx::query_as::<_, Address>(&query) + .bind(wallet_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; Ok(rows) } @@ -856,22 +844,13 @@ impl Store { 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?; + let query = cursor_pagination_query("transactions", "wallet_id"); + let rows = sqlx::query_as::<_, Transaction>(&query) + .bind(wallet_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; Ok(rows) } @@ -1325,22 +1304,13 @@ impl Store { 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?; + let query = cursor_pagination_query("payment_links", "wallet_id"); + let rows = sqlx::query_as::<_, PaymentLink>(&query) + .bind(wallet_id) + .bind(before_id) + .bind(limit) + .fetch_all(&self.pool) + .await?; Ok(rows) } @@ -1964,3 +1934,18 @@ impl Store { Ok(result.rows_affected()) } } + +// Builds keyset cursor pagination query using (created_at, id) tuple comparison for deterministic descending order. +pub fn cursor_pagination_query(table: &str, filter_column: &str) -> String { + format!( + r#" + SELECT * FROM {table} + WHERE {filter_column} = $1 + AND ($2::uuid IS NULL OR (created_at, id) < ( + SELECT created_at, id FROM {table} WHERE id = $2 + )) + ORDER BY created_at DESC, id DESC + LIMIT $3 + "# + ) +} diff --git a/crates/store/tests/store_tests.rs b/crates/store/tests/store_tests.rs index e69de29..5d275f5 100644 --- a/crates/store/tests/store_tests.rs +++ b/crates/store/tests/store_tests.rs @@ -0,0 +1,1213 @@ +//! Integration tests for octo-store. Require a running Postgres. +//! +//! Run with: `docker compose up -d db` then `cargo test -p octo-store`. +//! +//! `DATABASE_URL` is read from the workspace `.env` automatically (via dotenvy), so the plain +//! `cargo test -p octo-store` works without exporting anything. If no URL can be found, the tests +//! print a clear SKIPPED message and pass (so a DB-less `cargo test` of the whole workspace is +//! green). If a URL is found but the DB is unreachable, the test fails loudly with the reason. + +use octo_store::{ + NewDeposit, NewPaymentLink, NewSponsoredTx, NewWallet, NewWithdrawal, Store, StoreError, +}; +use std::sync::Once; +use uuid::Uuid; + +static LOAD_ENV: Once = Once::new(); + +/// Resolve `DATABASE_URL`, loading the workspace `.env` first. Returns `None` only if no URL is +/// configured anywhere (in which case tests skip with a message). +fn database_url() -> Option { + LOAD_ENV.call_once(|| { + // Search upward from the crate dir for a .env (workspace root holds it). + let _ = dotenvy::dotenv(); + }); + std::env::var("DATABASE_URL").ok() +} + +async fn store() -> Option { + let Some(url) = database_url() else { + eprintln!( + "SKIPPED: DATABASE_URL is not set (no .env found). \ + Run `docker compose up -d db` and ensure .env exists to run store tests." + ); + return None; + }; + let store = Store::connect(&url) + .await + .unwrap_or_else(|e| panic!("could not connect to {url}: {e}")); + store.migrate().await.expect("migrate"); + Some(store) +} + +/// Create a throwaway wallet with a unique account id (so tests don't collide). +async fn fresh_wallet(store: &Store) -> Uuid { + let acct = format!("G{}", Uuid::new_v4().simple()); // unique, not a real strkey (fine for store tests) + let w = store + .create_wallet(NewWallet { + network: "testnet", + stellar_account_g: &acct, + sealed_ciphertext: b"ciphertext", + sealed_nonce: b"nonce12bytes", + sealed_salt: b"saltsaltsaltsalt", + sealed_scheme: 1, // octo_crypto::SCHEME_V1 + label: Some("test"), + user_id: None, + description: None, + }) + .await + .expect("create wallet"); + w.id +} + +#[tokio::test] +async fn create_and_get_wallet() { + let Some(store) = store().await else { return }; + let id = fresh_wallet(&store).await; + let w = store.get_wallet(id).await.expect("get"); + assert_eq!(w.network, "testnet"); + assert_eq!(w.next_muxed_id, 1); +} + +#[tokio::test] +async fn allocate_address_increments_atomically() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + + // muxed_address is globally unique in the schema (real ones encode the base account), so make + // the test value unique per wallet too. + let wid = wallet_id.simple(); + let a = store + .allocate_address( + wallet_id, + |id| Ok(format!("M{wid}-{id}")), + Some("user-a"), + serde_json::json!({}), + ) + .await + .expect("alloc a"); + let b = store + .allocate_address( + wallet_id, + |id| Ok(format!("M{wid}-{id}")), + Some("user-b"), + serde_json::json!({}), + ) + .await + .expect("alloc b"); + + assert_eq!(a.muxed_id, 1); + assert_eq!(b.muxed_id, 2); + assert_ne!(a.muxed_address, b.muxed_address); + + let list = store + .list_addresses(wallet_id, 100, None) + .await + .expect("list"); + assert_eq!(list.len(), 2); +} + +#[tokio::test] +async fn record_deposit_is_idempotent() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + let tx_hash = Uuid::new_v4().to_string(); + + let dep = NewDeposit { + wallet_id, + address_id: None, + asset_code: "native".into(), + asset_issuer: None, + amount_stroops: 10_000_000, + source_account: Some("Gsender".into()), + destination_account: Some("Gmaster".into()), + stellar_tx_hash: tx_hash.clone(), + operation_index: 0, + horizon_op_id: format!("{tx_hash}-0"), + ledger: Some(123), + memo_id: None, + }; + + // First insert credits. + let first = store.record_deposit(&dep).await.expect("first"); + assert!(first.is_some(), "first deposit must be recorded"); + + // Replaying the SAME horizon_op_id must NOT double-credit. + let second = store.record_deposit(&dep).await.expect("second"); + assert!( + second.is_none(), + "duplicate deposit must be a no-op (anti double-credit)" + ); + + let txs = store + .list_transactions(wallet_id, 100, None) + .await + .expect("list"); + assert_eq!(txs.len(), 1, "exactly one ledger entry for one on-chain op"); +} + +#[tokio::test] +async fn different_op_index_same_tx_is_distinct() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + let tx_hash = Uuid::new_v4().to_string(); + + let base = NewDeposit { + wallet_id, + address_id: None, + asset_code: "native".into(), + asset_issuer: None, + amount_stroops: 5, + source_account: None, + destination_account: None, + stellar_tx_hash: tx_hash.clone(), + operation_index: 0, + horizon_op_id: format!("{tx_hash}-0"), + ledger: None, + memo_id: None, + }; + let op1 = NewDeposit { + operation_index: 1, + horizon_op_id: format!("{tx_hash}-1"), + ..base.clone() + }; + + assert!(store.record_deposit(&base).await.expect("op0").is_some()); + assert!(store.record_deposit(&op1).await.expect("op1").is_some()); + assert_eq!( + store + .list_transactions(wallet_id, 100, None) + .await + .unwrap() + .len(), + 2 + ); +} + +#[tokio::test] +async fn sum_deposits_for_address_totals_only_that_addresss_confirmed_deposits() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + let wid = wallet_id.simple(); + + let addr_a = store + .allocate_address( + wallet_id, + |id| Ok(format!("M{wid}-a-{id}")), + Some("a"), + serde_json::json!({}), + ) + .await + .expect("alloc a"); + let addr_b = store + .allocate_address( + wallet_id, + |id| Ok(format!("M{wid}-b-{id}")), + Some("b"), + serde_json::json!({}), + ) + .await + .expect("alloc b"); + + // Two deposits to A, one to B — A's total must be the sum of only its own two, not B's. + for (i, amount) in [(0, 10_000_000i64), (1, 2_500_000)] { + let tx_hash = Uuid::new_v4().to_string(); + store + .record_deposit(&NewDeposit { + wallet_id, + address_id: Some(addr_a.id), + asset_code: "native".into(), + asset_issuer: None, + amount_stroops: amount, + source_account: Some("Gsender".into()), + destination_account: Some("Gmaster".into()), + stellar_tx_hash: tx_hash.clone(), + operation_index: i, + horizon_op_id: format!("{tx_hash}-{i}"), + ledger: Some(1), + memo_id: None, + }) + .await + .expect("record deposit to a"); + } + let tx_hash_b = Uuid::new_v4().to_string(); + store + .record_deposit(&NewDeposit { + wallet_id, + address_id: Some(addr_b.id), + asset_code: "native".into(), + asset_issuer: None, + amount_stroops: 999_000_000, + source_account: Some("Gsender".into()), + destination_account: Some("Gmaster".into()), + stellar_tx_hash: tx_hash_b.clone(), + operation_index: 0, + horizon_op_id: format!("{tx_hash_b}-0"), + ledger: Some(1), + memo_id: None, + }) + .await + .expect("record deposit to b"); + + assert_eq!( + store + .sum_deposits_for_address(addr_a.id) + .await + .expect("sum a"), + 12_500_000, + "A's total must be the sum of its own two deposits, unaffected by B's" + ); + assert_eq!( + store + .sum_deposits_for_address(addr_b.id) + .await + .expect("sum b"), + 999_000_000 + ); + + // A brand-new address with no deposits sums to 0, not an error. + let addr_c = store + .allocate_address( + wallet_id, + |id| Ok(format!("M{wid}-c-{id}")), + Some("c"), + serde_json::json!({}), + ) + .await + .expect("alloc c"); + assert_eq!( + store + .sum_deposits_for_address(addr_c.id) + .await + .expect("sum c"), + 0 + ); + + // The batched form must agree with the per-address form, and only return entries that + // actually have deposits (address C has none, so it's absent rather than a zero row). + let batched = store + .sum_deposits_for_addresses(&[addr_a.id, addr_b.id, addr_c.id]) + .await + .expect("batched sum"); + let totals: std::collections::HashMap = batched.into_iter().collect(); + assert_eq!(totals.get(&addr_a.id), Some(&12_500_000)); + assert_eq!(totals.get(&addr_b.id), Some(&999_000_000)); + assert_eq!( + totals.get(&addr_c.id), + None, + "an address with zero deposits has no row in the batched result (GROUP BY yields nothing)" + ); + + // Empty id list must short-circuit to an empty result, not error or scan the whole table. + assert_eq!( + store + .sum_deposits_for_addresses(&[]) + .await + .expect("empty batch"), + Vec::new() + ); +} + +#[tokio::test] +async fn payment_link_lifecycle_intent_confirm_and_sum() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + let wid = wallet_id.simple(); + + let addr = store + .allocate_address( + wallet_id, + |id| Ok(format!("M{wid}-{id}")), + None, + serde_json::json!({}), + ) + .await + .expect("alloc address"); + + let slug = format!("link-{wid}"); + let link = store + .create_payment_link(NewPaymentLink { + wallet_id, + address_id: addr.id, + slug: &slug, + name: "Support octo", + description: Some("donations"), + image_url: None, + redirect_url: None, + amount_usdc_stroops: None, + }) + .await + .expect("create link"); + assert_eq!(link.slug, slug); + assert!(link.active); + + // Public lookup by slug must work with no wallet_id in hand. + let by_slug = store + .get_payment_link_by_slug(&slug) + .await + .expect("by slug"); + assert_eq!(by_slug.id, link.id); + + // A fresh link has nothing collected yet. + assert_eq!( + store + .sum_payment_link_collected(link.id) + .await + .expect("sum"), + 0 + ); + + let intent = store + .record_payment_link_intent( + link.id, + Some("Ada"), + Some("ada@example.com"), + 10_000_000, + Some(addr.id), + ) + .await + .expect("record intent"); + assert_eq!(intent.status, "pending"); + + let oldest = store + .oldest_pending_payment_link_payment(link.id) + .await + .expect("oldest pending") + .expect("one pending row"); + assert_eq!(oldest.id, intent.id); + + // Exact-address lookup is how ingest matches a deposit to one specific intent. + let by_address = store + .pending_payment_by_address(addr.id) + .await + .expect("by address") + .expect("pending intent on this address"); + assert_eq!(by_address.id, intent.id); + assert_eq!(by_address.address_id, Some(addr.id)); + + let tx_hash = Uuid::new_v4().to_string(); + let dep = store + .record_deposit(&NewDeposit { + wallet_id, + address_id: Some(addr.id), + asset_code: "USDC".into(), + asset_issuer: Some("GISSUER".into()), + amount_stroops: 10_000_000, + source_account: Some("Gpayer".into()), + destination_account: Some("Gmaster".into()), + stellar_tx_hash: tx_hash.clone(), + operation_index: 0, + horizon_op_id: format!("{tx_hash}-0"), + ledger: Some(1), + memo_id: None, + }) + .await + .expect("record deposit") + .expect("first insert"); + + store + .confirm_payment_link_payment(intent.id, dep.id) + .await + .expect("confirm payment"); + + let confirmed = store + .get_payment_link_payment(link.id, intent.id) + .await + .expect("get payment"); + assert_eq!(confirmed.status, "confirmed"); + assert_eq!(confirmed.transaction_id, Some(dep.id)); + + // Once confirmed, it's no longer the oldest pending (there is none left). + assert!(store + .oldest_pending_payment_link_payment(link.id) + .await + .expect("oldest pending after confirm") + .is_none()); + + assert_eq!( + store + .sum_payment_link_collected(link.id) + .await + .expect("sum after confirm"), + 10_000_000 + ); + + let batch = store + .sum_payment_link_collected_batch(&[link.id]) + .await + .expect("batch sum"); + assert_eq!(batch, vec![(link.id, 10_000_000)]); + + // Deactivating is scoped to the owning wallet. + let deactivated = store + .set_payment_link_active(wallet_id, link.id, false) + .await + .expect("deactivate"); + assert!(!deactivated.active); +} + +#[tokio::test] +async fn payment_link_mismatched_deposit_records_the_transaction_but_does_not_confirm() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + let wid = wallet_id.simple(); + + let addr = store + .allocate_address( + wallet_id, + |id| Ok(format!("M{wid}-{id}")), + None, + serde_json::json!({}), + ) + .await + .expect("alloc address"); + + let link = store + .create_payment_link(NewPaymentLink { + wallet_id, + address_id: addr.id, + slug: &format!("link-mismatch-{wid}"), + name: "Underpaid test", + description: None, + image_url: None, + redirect_url: None, + amount_usdc_stroops: Some(10_000_000), + }) + .await + .expect("create link"); + + let intent = store + .record_payment_link_intent(link.id, None, None, 10_000_000, Some(addr.id)) + .await + .expect("record intent"); + + let tx_hash = Uuid::new_v4().to_string(); + let dep = store + .record_deposit(&NewDeposit { + wallet_id, + address_id: Some(addr.id), + asset_code: "USDC".into(), + asset_issuer: Some("GISSUER".into()), + amount_stroops: 5_000_000, // half of what was expected + source_account: Some("Gpayer".into()), + destination_account: Some("Gmaster".into()), + stellar_tx_hash: tx_hash.clone(), + operation_index: 0, + horizon_op_id: format!("{tx_hash}-0"), + ledger: Some(1), + memo_id: None, + }) + .await + .expect("record deposit") + .expect("first insert"); + + store + .mark_payment_link_payment_mismatched(intent.id, dep.id, "underpaid") + .await + .expect("mark mismatched"); + + let mismatched = store + .get_payment_link_payment(link.id, intent.id) + .await + .expect("get payment"); + assert_eq!(mismatched.status, "underpaid"); + assert_eq!( + mismatched.transaction_id, + Some(dep.id), + "the short deposit must still be linked, so the merchant can see what actually arrived" + ); + + // A mismatched payment is not "pending" any more, so it must not still be matchable — ingest + // must not later confuse a second, correct deposit with this already-resolved intent. + assert!(store + .pending_payment_by_address(addr.id) + .await + .expect("by address") + .is_none()); +} + +#[tokio::test] +async fn expire_stale_payment_link_payments_only_sweeps_old_pending_rows() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + let wid = wallet_id.simple(); + + let addr = store + .allocate_address( + wallet_id, + |id| Ok(format!("M{wid}-{id}")), + None, + serde_json::json!({}), + ) + .await + .expect("alloc address"); + + let link = store + .create_payment_link(NewPaymentLink { + wallet_id, + address_id: addr.id, + slug: &format!("link-expiry-{wid}"), + name: "Expiry test", + description: None, + image_url: None, + redirect_url: None, + amount_usdc_stroops: Some(10_000_000), + }) + .await + .expect("create link"); + + let stale = store + .record_payment_link_intent(link.id, None, None, 10_000_000, Some(addr.id)) + .await + .expect("record stale intent"); + // Backdate it past the 1-hour deadline directly — this test can't wait an hour. + sqlx::query( + "UPDATE payment_link_payments SET created_at = now() - interval '2 hours' WHERE id = $1", + ) + .bind(stale.id) + .execute(store.pool()) + .await + .expect("backdate"); + + let fresh = store + .record_payment_link_intent(link.id, None, None, 10_000_000, Some(addr.id)) + .await + .expect("record fresh intent"); + + let expired = store + .expire_stale_payment_link_payments() + .await + .expect("sweep"); + let expired_ids: Vec = expired.iter().map(|p| p.id).collect(); + assert!( + expired_ids.contains(&stale.id), + "the >1hr-old pending row must be swept" + ); + assert!( + !expired_ids.contains(&fresh.id), + "a freshly-created pending row must not be swept" + ); + + let stale_after = store + .get_payment_link_payment(link.id, stale.id) + .await + .expect("get stale"); + assert_eq!(stale_after.status, "expired"); + + let fresh_after = store + .get_payment_link_payment(link.id, fresh.id) + .await + .expect("get fresh"); + assert_eq!(fresh_after.status, "pending"); + + // Running the sweep again must be a no-op for already-expired rows (idempotent). + let expired_again = store + .expire_stale_payment_link_payments() + .await + .expect("sweep again"); + assert!(!expired_again.iter().any(|p| p.id == stale.id)); +} + +#[tokio::test] +async fn withdrawal_idempotency_key_blocks_double_spend() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + + let mk = |key: &'static str| NewWithdrawal { + wallet_id, + idempotency_key: key, + destination_account: "Gdest", + asset_code: "native", + asset_issuer: None, + amount_stroops: 1_000, + memo_id: None, + }; + + let first = store.create_withdrawal(mk("key-1")).await; + assert!(first.is_ok(), "first withdrawal accepted"); + + // Same idempotency key => conflict, not a second payout. + let second = store.create_withdrawal(mk("key-1")).await; + assert!( + matches!(second, Err(StoreError::Conflict)), + "retry must conflict" + ); + + // A different key is a different withdrawal. + let third = store.create_withdrawal(mk("key-2")).await; + assert!(third.is_ok()); +} + +/// Insert a minimal gas_sponsorship_configs row (no limits) for `wallet_id`. +async fn insert_sponsorship_config(store: &Store, wallet_id: Uuid) { + sqlx::query("INSERT INTO gas_sponsorship_configs (wallet_id, enabled) VALUES ($1, true)") + .bind(wallet_id) + .execute(store.pool()) + .await + .expect("insert gas_sponsorship_configs"); +} + +#[tokio::test] +async fn record_and_update_sponsored_tx() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + insert_sponsorship_config(&store, wallet_id).await; + + let hash = format!("inner-{}", Uuid::new_v4().simple()); + let row = store + .record_sponsored_tx(NewSponsoredTx { + wallet_id, + inner_tx_hash: &hash, + fee_stroops: 500, + }) + .await + .expect("record"); + + assert_eq!(row.wallet_id, wallet_id); + assert_eq!(row.inner_tx_hash, hash); + assert_eq!(row.fee_stroops, 500); + assert_eq!(row.status, "pending"); + assert!(row.fee_bump_tx_hash.is_none()); + + // Update to confirmed. + let bump_hash = format!("bump-{}", Uuid::new_v4().simple()); + store + .update_sponsored_tx_status(row.id, "confirmed", Some(&bump_hash), None) + .await + .expect("update"); + + // Verify via pool (the store has no get_sponsored_tx yet; query directly). + let updated: (String, Option) = + sqlx::query_as("SELECT status, fee_bump_tx_hash FROM sponsored_transactions WHERE id = $1") + .bind(row.id) + .fetch_one(store.pool()) + .await + .expect("fetch updated"); + + assert_eq!(updated.0, "confirmed"); + assert_eq!(updated.1.as_deref(), Some(bump_hash.as_str())); +} + +#[tokio::test] +async fn sum_fees_today_counts_only_confirmed() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + insert_sponsorship_config(&store, wallet_id).await; + + // No rows → 0. + let initial = store + .sum_sponsored_fees_today(wallet_id) + .await + .expect("sum"); + assert_eq!(initial, 0); + + // Insert a pending tx (fee 200): should not count. + let pending = store + .record_sponsored_tx(NewSponsoredTx { + wallet_id, + inner_tx_hash: &format!("pending-{}", Uuid::new_v4().simple()), + fee_stroops: 200, + }) + .await + .expect("pending record"); + // Still 0 — pending doesn't count. + assert_eq!(store.sum_sponsored_fees_today(wallet_id).await.unwrap(), 0); + + // Confirm the tx → now it counts. + store + .update_sponsored_tx_status(pending.id, "confirmed", None, None) + .await + .expect("update to confirmed"); + assert_eq!( + store.sum_sponsored_fees_today(wallet_id).await.unwrap(), + 200 + ); + + // A second confirmed tx adds to the total. + let second = store + .record_sponsored_tx(NewSponsoredTx { + wallet_id, + inner_tx_hash: &format!("second-{}", Uuid::new_v4().simple()), + fee_stroops: 300, + }) + .await + .expect("second record"); + store + .update_sponsored_tx_status(second.id, "confirmed", None, None) + .await + .unwrap(); + assert_eq!( + store.sum_sponsored_fees_today(wallet_id).await.unwrap(), + 500 + ); +} + +#[tokio::test] +async fn sum_fees_today_can_use_wallet_status_created_at_index() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + + let mut tx = store.pool().begin().await.expect("begin transaction"); + sqlx::query("SET LOCAL enable_seqscan = off") + .execute(&mut *tx) + .await + .expect("disable sequential scans for index eligibility check"); + let plan: Vec = sqlx::query_scalar( + r#"EXPLAIN (COSTS OFF) + 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_all(&mut *tx) + .await + .expect("explain sum_sponsored_fees_today"); + let plan = plan.join("\n"); + + assert!( + plan.contains("idx_sponsored_wallet_status_"), + "expected the wallet/status/created_at index, got:\n{plan}" + ); + assert!( + !plan.contains("Seq Scan"), + "sum query must not require a full table scan:\n{plan}" + ); +} + +#[tokio::test] +async fn duplicate_inner_tx_hash_is_conflict() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + insert_sponsorship_config(&store, wallet_id).await; + + let hash = format!("dup-{}", Uuid::new_v4().simple()); + + let first = store + .record_sponsored_tx(NewSponsoredTx { + wallet_id, + inner_tx_hash: &hash, + fee_stroops: 100, + }) + .await; + assert!(first.is_ok(), "first record must succeed"); + + // Same inner_tx_hash → UNIQUE violation → Conflict. + let second = store + .record_sponsored_tx(NewSponsoredTx { + wallet_id, + inner_tx_hash: &hash, + fee_stroops: 100, + }) + .await; + assert!( + matches!(second, Err(StoreError::Conflict)), + "duplicate inner_tx_hash must conflict, got: {second:?}" + ); +} + +#[tokio::test] +async fn cursor_roundtrip() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + + assert_eq!(store.get_cursor(wallet_id).await.unwrap(), None); + store.set_cursor(wallet_id, "token-1").await.unwrap(); + assert_eq!( + store.get_cursor(wallet_id).await.unwrap().as_deref(), + Some("token-1") + ); + // Upsert overwrites. + store.set_cursor(wallet_id, "token-2").await.unwrap(); + assert_eq!( + store.get_cursor(wallet_id).await.unwrap().as_deref(), + Some("token-2") + ); +} + +#[tokio::test] +async fn migrate_is_idempotent_when_run_twice() { + let Some(store) = store().await else { return }; + // `store()` already ran migrate() once during setup; running it again against the same + // already-migrated database mirrors a server restart (bin/server/src/main.rs calls + // store.migrate().await on every boot) and must be a safe no-op, not an error. + store + .migrate() + .await + .expect("second migrate() call must succeed with no error"); +} + +#[tokio::test] +async fn migrate_applies_exactly_the_expected_version_set() { + let Some(store) = store().await else { return }; + + let mut versions: Vec = sqlx::query_scalar( + "SELECT version FROM _sqlx_migrations WHERE success = true ORDER BY version", + ) + .fetch_all(store.pool()) + .await + .expect("query _sqlx_migrations"); + versions.sort_unstable(); + + // One version per file under crates/store/migrations/, 0001_init.sql .. 0020. + // Guards against silent version collisions — sqlx keys migrations by version, so a repeated + // number means only one of the colliding pair actually ran. + assert_eq!( + versions, + vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20], + "expected exactly the twenty known migrations to be recorded as applied" + ); +} + +#[tokio::test] +async fn upsert_gas_sponsorship_config_works() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + let cfg = store + .upsert_gas_sponsorship_config(wallet_id, true, Some(500_000), Some(10_000_000)) + .await + .expect("upsert"); + assert!(cfg.enabled); + let spent = store + .sum_sponsored_fees_reserved_today(wallet_id) + .await + .expect("sum"); + assert_eq!(spent, 0); +} + +/// Create a throwaway user with a unique email (so tests don't collide). +async fn fresh_user(store: &Store) -> Uuid { + let email = format!("test-{}@example.invalid", Uuid::new_v4().simple()); + store + .create_user(&email, "not-a-real-hash") + .await + .expect("create user") + .id +} + +// --- indexing-overhaul correctness regressions (hard/store/indexing-overhaul-with-load-test) --- +// +// These assert result *correctness* (ordering, filtering) for the query shapes the new indices in +// migrations/0008_sponsored_and_audit_indexing.sql target. An index change must never change which +// rows come back or in what order — if either of these starts failing, the index migration altered +// query semantics, not just performance, and that's a bug in the migration. + +#[tokio::test] +async fn list_sponsored_transactions_orders_filters_and_paginates_correctly() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + insert_sponsorship_config(&store, wallet_id).await; + + // Three rows, two different statuses, with `created_at` pinned to strictly increasing values + // (rather than relying on wall-clock ordering, which is too coarse to guarantee distinct + // timestamps for back-to-back inserts and would make the ORDER BY assertions flaky). + let mut ids = Vec::new(); + for (i, (label, status)) in [("a", "pending"), ("b", "confirmed"), ("c", "confirmed")] + .into_iter() + .enumerate() + { + let row = store + .record_sponsored_tx(NewSponsoredTx { + wallet_id, + inner_tx_hash: &format!("order-{label}-{}", Uuid::new_v4().simple()), + fee_stroops: 100, + }) + .await + .expect("record"); + if status == "confirmed" { + store + .update_sponsored_tx_status(row.id, "confirmed", None, None) + .await + .expect("confirm"); + } + sqlx::query("UPDATE sponsored_transactions SET created_at = now() - make_interval(secs => $2) WHERE id = $1") + .bind(row.id) + .bind((10 - i) as f64) + .execute(store.pool()) + .await + .expect("pin created_at"); + ids.push(row.id); + } + + // Unfiltered: most-recent-first (created_at DESC, id DESC — insertion order reversed). + let all = store + .list_sponsored_transactions(wallet_id, 10, None, None) + .await + .expect("list all"); + let all_ids: Vec = all.iter().map(|r| r.id).collect(); + assert_eq!(all_ids, vec![ids[2], ids[1], ids[0]]); + + // Status filter: only the two confirmed rows, same relative order. + let confirmed = store + .list_sponsored_transactions(wallet_id, 10, Some("confirmed"), None) + .await + .expect("list confirmed"); + let confirmed_ids: Vec = confirmed.iter().map(|r| r.id).collect(); + assert_eq!(confirmed_ids, vec![ids[2], ids[1]]); + + // Cursor pagination: page of 1 starting after the newest row returns the next one down. + let page = store + .list_sponsored_transactions(wallet_id, 1, None, Some(ids[2])) + .await + .expect("list after cursor"); + assert_eq!(page.len(), 1); + assert_eq!(page[0].id, ids[1]); +} + +#[tokio::test] +async fn list_audit_logs_filters_by_category_and_search_correctly() { + let Some(store) = store().await else { return }; + let user_id = fresh_user(&store).await; + + store + .record_audit( + user_id, + "signed in", + "authentication", + None, + Some("203.0.113.1"), + ) + .await + .expect("record 1"); + store + .record_audit( + user_id, + "created wallet octo master wallet", + "wallet", + Some("octo master wallet"), + None, + ) + .await + .expect("record 2"); + store + .record_audit(user_id, "rotated api key", "credentials", None, None) + .await + .expect("record 3"); + + // Pin `created_at` to strictly increasing values in insertion order (see the sponsored-tx test + // above for why wall-clock ordering alone isn't reliable enough for the ORDER BY assertions). + for (offset_secs, action) in [ + (10.0, "signed in"), + (9.0, "created wallet octo master wallet"), + (8.0, "rotated api key"), + ] { + sqlx::query( + "UPDATE audit_logs SET created_at = now() - make_interval(secs => $2) \ + WHERE user_id = $1 AND action = $3", + ) + .bind(user_id) + .bind(offset_secs) + .bind(action) + .execute(store.pool()) + .await + .expect("pin created_at"); + } + + // Category filter: only the "wallet" row. + let by_category = store + .list_audit_logs(user_id, Some("wallet"), None, 10) + .await + .expect("list by category"); + assert_eq!(by_category.len(), 1); + assert_eq!(by_category[0].category, "wallet"); + + // Search filter (the ILIKE / trigram-index case): matches action OR target, case-insensitive. + let by_search = store + .list_audit_logs(user_id, None, Some("MASTER"), 10) + .await + .expect("list by search"); + assert_eq!(by_search.len(), 1); + assert_eq!(by_search[0].action, "created wallet octo master wallet"); + + // No match. + let no_match = store + .list_audit_logs(user_id, None, Some("nonexistent-term"), 10) + .await + .expect("list no match"); + assert!(no_match.is_empty()); + + // Unfiltered: all three, most-recent-first. + let all = store + .list_audit_logs(user_id, None, None, 10) + .await + .expect("list all"); + assert_eq!(all.len(), 3); + assert_eq!(all[0].action, "rotated api key"); +} + +#[tokio::test] +async fn wallets_due_for_poll_applies_activity_backoff() { + let Some(store) = store().await else { return }; + + // `network` is CHECK-constrained to mainnet/testnet, so this test can't invent its own. It + // uses mainnet (a handful of inert rows) and filters results down to the ids it created. + let network = "mainnet"; + let mut ids = Vec::new(); + for label in ["never-polled", "active", "idle", "dormant"] { + let acct = format!("G{}", Uuid::new_v4().simple()); + let w = store + .create_wallet(NewWallet { + network, + stellar_account_g: &acct, + sealed_ciphertext: b"ct", + sealed_nonce: b"nonce", + sealed_salt: b"salt", + sealed_scheme: 1, + label: Some(label), + user_id: None, + description: None, + }) + .await + .expect("create wallet"); + ids.push(w.id); + } + let (never, active, idle, dormant) = (ids[0], ids[1], ids[2], ids[3]); + + // Tiers for this test: active < 60s, idle polled at most every 100s, dormant (> 300s since + // activity) polled at most every 100_000s. + let mine = ids.clone(); + let due = |store: &Store| { + let store = store.clone(); + let mine = mine.clone(); + async move { + store + .wallets_due_for_poll(network, 60, 100, 300, 100_000) + .await + .expect("due query") + .into_iter() + .map(|w| w.id) + // Other mainnet rows may exist in a shared dev DB; only assert on our own. + .filter(|id| mine.contains(id)) + .collect::>() + } + }; + + // Nothing has a cursor row yet: every wallet is due. + let ids_due = due(&store).await; + assert_eq!( + ids_due.len(), + 4, + "wallets with no cursor row are always due" + ); + + // Give each wallet a cursor row with a distinct activity/poll profile. All were *just* + // polled, so only the active one should come back as due again immediately. + for (id, activity_secs) in [(active, 10i64), (idle, 200), (dormant, 100_000)] { + sqlx::query( + "INSERT INTO ingest_cursor (wallet_id, paging_token, updated_at, last_polled_at) + VALUES ($1, 'tok', now() - make_interval(secs => $2), now())", + ) + .bind(id) + .bind(activity_secs as f64) + .execute(store.pool()) + .await + .expect("seed cursor"); + } + + let ids_due = due(&store).await; + assert!( + ids_due.contains(&active), + "an actively-transacting wallet must be polled every tick" + ); + assert!( + !ids_due.contains(&idle), + "an idle wallet polled just now must wait for its interval" + ); + assert!( + !ids_due.contains(&dormant), + "a dormant wallet polled just now must wait for its (longer) interval" + ); + assert!( + ids_due.contains(&never), + "a wallet that has never been polled is still due" + ); + + // Move the idle wallet's last poll past its 100s interval — it becomes due, while the + // dormant one (100_000s interval) is still not. + sqlx::query("UPDATE ingest_cursor SET last_polled_at = now() - make_interval(secs => 150) WHERE wallet_id = $1") + .bind(idle) + .execute(store.pool()) + .await + .expect("age idle poll"); + sqlx::query("UPDATE ingest_cursor SET last_polled_at = now() - make_interval(secs => 150) WHERE wallet_id = $1") + .bind(dormant) + .execute(store.pool()) + .await + .expect("age dormant poll"); + + let ids_due = due(&store).await; + assert!( + ids_due.contains(&idle), + "idle wallet is due once its interval elapses" + ); + assert!( + !ids_due.contains(&dormant), + "dormant wallet needs much longer than the idle interval before it is due" + ); +} + +#[tokio::test] +async fn mark_polled_creates_and_updates_the_cursor_row() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + + // No cursor row yet — mark_polled must create one rather than silently no-op. + store.mark_polled(wallet_id).await.expect("first mark"); + let first: Option> = + sqlx::query_scalar("SELECT last_polled_at FROM ingest_cursor WHERE wallet_id = $1") + .bind(wallet_id) + .fetch_one(store.pool()) + .await + .expect("read cursor"); + let first = first.expect("last_polled_at set"); + + tokio::time::sleep(std::time::Duration::from_millis(20)).await; + store.mark_polled(wallet_id).await.expect("second mark"); + let second: Option> = + sqlx::query_scalar("SELECT last_polled_at FROM ingest_cursor WHERE wallet_id = $1") + .bind(wallet_id) + .fetch_one(store.pool()) + .await + .expect("read cursor again"); + assert!( + second.expect("still set") > first, + "repeat polls advance the timestamp" + ); + + // Marking a poll must NOT look like activity. If it did, every never-used wallet would count + // as freshly active and the backoff tiers would never engage at all. + let activity: chrono::DateTime = + sqlx::query_scalar("SELECT updated_at FROM ingest_cursor WHERE wallet_id = $1") + .bind(wallet_id) + .fetch_one(store.pool()) + .await + .expect("read updated_at"); + assert!( + activity < chrono::Utc::now() - chrono::Duration::days(365), + "mark_polled must not advance updated_at (last-activity); got {activity}" + ); + + // Marking a poll must not invent a paging token — that only advances on real activity. + let token: Option = + sqlx::query_scalar("SELECT paging_token FROM ingest_cursor WHERE wallet_id = $1") + .bind(wallet_id) + .fetch_one(store.pool()) + .await + .expect("read token"); + assert!( + token.is_none(), + "mark_polled must not fabricate a cursor position" + ); +} + +#[test] +fn shared_cursor_pagination_helper_encodes_invariant() { + let query = octo_store::cursor_pagination_query("wallets", "user_id"); + assert!(query.contains("SELECT * FROM wallets")); + assert!(query.contains("WHERE user_id = $1")); + assert!(query.contains("($2::uuid IS NULL OR (created_at, id) < (")); + assert!(query.contains("SELECT created_at, id FROM wallets WHERE id = $2")); + assert!(query.contains("ORDER BY created_at DESC, id DESC")); + assert!(query.contains("LIMIT $3")); +} diff --git a/docs/api.md b/docs/api.md index e69de29..a37adbb 100644 --- a/docs/api.md +++ b/docs/api.md @@ -0,0 +1,143 @@ +# 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`). +- `GET /v1/auth/me` — the current user. + +## 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` and `POST /v1/wallets/:id/trustlines` are **`410 Gone` + tombstones**. They exist only to give integrators a clear error pointing at `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. +- `GET /v1/wallets/{id}/transactions` — deposits + outbound transfers (paginated). +- `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}/sponsorship` / `PUT` — read/update `enabled`, the per-transaction fee + cap, and the daily budget. +- `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). + +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. + +## 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/regenerate (the plaintext key is shown **once**; + only a SHA-256 hash is stored). +- `GET /v1/wallets/{id}/api-key` — metadata (prefix, created_at) — never the key itself. +- `DELETE /v1/wallets/{id}/api-key` — revoke. + +## Audit logs + +- `GET /v1/audit-logs` — your account's activity, filterable by `category` and a free-text + `search`. + +## 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). There is no `422`. + +## Rate Limits + +The API enforces fixed-window rate limiting on unauthenticated and authentication endpoints to protect against brute-force and resource-exhaustion attacks. When a rate limit is exceeded, the server responds with HTTP `429 Too Many Requests`. + +| Endpoint | Method | Key Scope | Limit | Window | Code Constant | +|---|---|---|---|---|---| +| `/v1/auth/signup` | POST | Per-IP | 10 req | 60s (1m) | `AUTH_RATE_LIMIT` / `AUTH_RATE_WINDOW` | +| `/v1/auth/verify-email` | POST | Per-IP | 10 req | 60s (1m) | `AUTH_RATE_LIMIT` / `AUTH_RATE_WINDOW` | +| `/v1/auth/login` | POST | Per-IP | 10 req | 60s (1m) | `AUTH_RATE_LIMIT` / `AUTH_RATE_WINDOW` | +| `/v1/auth/refresh` | POST | Per-IP | 10 req | 60s (1m) | `AUTH_RATE_LIMIT` / `AUTH_RATE_WINDOW` | +| `/v1/auth/resend-otp` | POST | Per-IP | 10 req | 60s (1m) | `AUTH_RATE_LIMIT` / `AUTH_RATE_WINDOW` | +| `/v1/auth/resend-otp` | POST | Per-User (`otp:{user_id}`) | 3 req | 3600s (1h) | `OTP_RESEND_USER_LIMIT` / `OTP_RESEND_USER_WINDOW` | +| `/v1/auth/resend-otp` | POST | Per-IP | 10 req | 3600s (1h) | `OTP_RESEND_IP_LIMIT` / `OTP_RESEND_IP_WINDOW` | +| `/v1/pay/:slug` | GET | Per-IP | 60 req | 60s (1m) | `PAY_READ_LIMIT` / `PAY_READ_WINDOW` | +| `/v1/pay/:slug/intent` | POST | Per-IP | 5 req | 60s (1m) | `PAY_INTENT_LIMIT` / `PAY_INTENT_WINDOW` | +| `/v1/pay/:slug/payments/:payment_id` | GET | Per-IP | 60 req | 60s (1m) | `PAY_STATUS_LIMIT` / `PAY_STATUS_WINDOW` | +| `/v1/pay/:slug/signing-info` | GET | Per-IP | 60 req | 60s (1m) | `PAY_SIGNING_INFO_LIMIT` / `PAY_SIGNING_INFO_WINDOW` | +| `/v1/pay/:slug/submit-signed` | POST | Per-IP | 20 req | 60s (1m) | `PAY_SUBMIT_LIMIT` / `PAY_SUBMIT_WINDOW` | + +> [!NOTE] +> **Process Note:** Any new rate limit added to the API must update this table and reference named constants in `crates/api/src/rate_limit.rs` within the same pull request. diff --git a/docs/architecture.md b/docs/architecture.md index 5cc291d..e69de29 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -1,107 +0,0 @@ -# Architecture - -octo is a Cargo workspace. The guiding rule: **secret material is confined to one crate** -(`wallet-core`), decrypted only in-memory at signing time, and zeroized immediately after. - -## Crates - -``` -crates/ - crypto/ AES-256-GCM seal/open of a gas-tank seed (random nonce + salt). No Stellar knowledge. - wallet-core/ The only code that touches secret keys (server-side: gas tank only): - - SEP-0005 (SLIP-0010 ed25519) derivation: m/44'/148'/' - - muxed address (M...) encode/decode - - build + sign fee-bump envelopes, then zeroize - resilience/ Retry with backoff + circuit breaker for outbound Horizon calls. - store/ Postgres models + migrations (sqlx). - webhooks/ HMAC-SHA256 signed outbound webhooks with retry + delivery log. - ingest/ Horizon payment streaming + durable cursor → deposit detection & attribution. - api/ axum REST API (wallets, addresses, submit-signed, sponsorship, webhooks). -bin/ - server/ Composes api + ingest into one process (splittable later to scale). - migrate-keys/ Offline backfill that re-seals gas-tank seeds under a new master key - (zero-downtime rotation; skips client-custody rows, which hold no seed). -``` - -## Request flows - -### Create master wallet (non-custodial) -The **client** generates the BIP39 mnemonic and derives the base keypair (`m/44'/148'/0'`) in the -browser/SDK. It sends `api` only the public account (`G...`), plus an optional `encrypted_backup` -blob it encrypted under the user's password. `store` persists the public key, the opaque blob and -`custody = 'client'` — **no seed, no mnemonic, ever.** On testnet, friendbot funds the account so -it exists on-chain. - -### Generate a customer address -`api` atomically increments the wallet's id counter → `wallet-core` encodes a muxed `M...` from -the base `G...` + id → `store` saves the row. **No on-chain operation.** The response also returns -the `G...` + numeric-memo fallback for senders that don't support muxed. - -### Detect a deposit -`ingest` streams the master account's payments from Horizon (with a persisted cursor). Each -payment is attributed to a customer by its **muxed id** or **memo id**, recorded as a `deposit` -transaction, and a signed webhook fires. - -For the full contract around cursor resume, dedup, reorg handling, and the quarantine path -see [`docs/ingest-integration.md`](ingest-integration.md). - -### Move funds out (client-signed) -The client fetches `GET /signing-info` (sequence, network passphrase, base fee), builds and -**signs the transaction locally**, then relays it via `POST /submit-signed`. `api` validates the -envelope and submits it to Horizon **unmodified** → record + webhook on confirmation. Horizon's -result codes are passed back so the client can correct and re-sign. - -The custodial `POST /withdraw` endpoint is a `410 Gone` tombstone. `POST /trustlines` validates the -asset and returns ChangeTrust signing info (sequence, passphrase, fee, limit); the client signs -locally and relays via `submit-signed`. - -## Signing safety - -The user's key is never on the server, so there is no server-side signing path for user funds — -and therefore no signing oracle to abuse. What `api` does on the submit path is *validate*: - -1. Envelope is a v1 `Tx` (not a fee-bump wrapper smuggled in). -2. At least one signature is present. -3. The source account **is this wallet**. -4. Every operation is on the allowlist (payment / path-payment / change-trust). -5. Submit verbatim — the server never re-signs or alters the transaction. - -### The one server-held key: the gas tank -Fee sponsorship still needs a server signature, so a wallet may provision a **gas tank**: a -separate account holding fee float only. Its seed is the only plaintext key material on the -server, and it is confined to one crate: - -1. Retrieve the encrypted gas-tank seed from `store`. -2. `crypto::open` decrypts in-memory (AES-256-GCM; tag verifies integrity, network bound as AAD). -3. `wallet-core` derives the private key via SEP-0005. -4. Sign **only the outer fee-bump envelope** — the user's inner transaction is untouched. -5. `zeroize` the seed and key buffers. - -Keys are never written to disk or logs and are never persisted in derived form. Worst-case -exposure of this key is the gas budget — never customer balances. - -### Entropy source (load-bearing) -Every server-generated mnemonic comes from `WalletSeed::generate` (`wallet-core/src/derive.rs`), -which fills 128 bits of entropy from `rand::rngs::OsRng` (the OS CSPRNG, `getrandom(2)`) and -calls `Mnemonic::from_entropy`. It deliberately bypasses tiny-bip39's `Mnemonic::new`, whose -`thread_rng()` source depends on a default crate feature and a `rand` implementation detail. -`crypto::seal` uses the same `OsRng` for nonces and salts. Any bump of `tiny-bip39` or `rand` -must re-confirm this path stays OS-backed. - -## Wallet Foreign Key Constraints - -Wallets are intended to be permanent master records. To prevent accidental cascading deletions or silent orphaned records, all tables referencing `wallets(id)` enforce `ON DELETE RESTRICT`: - -| Table | Column | Initial Migration Constraint | Intended & Enforced Constraint | -| --- | --- | --- | --- | -| `addresses` | `wallet_id` | `ON DELETE CASCADE` (0001) | `ON DELETE RESTRICT` (0021) | -| `transactions` | `wallet_id` | `ON DELETE CASCADE` (0001) | `ON DELETE RESTRICT` (0021) | -| `withdrawals` | `wallet_id` | `ON DELETE CASCADE` (0001) | `ON DELETE RESTRICT` (0021) | -| `webhook_endpoints` | `wallet_id` | `ON DELETE CASCADE` (0001) | `ON DELETE RESTRICT` (0021) | -| `ingest_cursor` | `wallet_id` | `ON DELETE CASCADE` (0001) | `ON DELETE RESTRICT` (0021) | -| `api_keys` | `wallet_id` | `ON DELETE CASCADE` (0005) | `ON DELETE RESTRICT` (0021) | -| `gas_sponsorship_configs` | `wallet_id` | `ON DELETE CASCADE` (0007) | `ON DELETE RESTRICT` (0021) | -| `sponsored_transactions` | `wallet_id` | `ON DELETE CASCADE` (0007) | `ON DELETE RESTRICT` (0021) | -| `withdrawal_allowlist_configs` | `wallet_id` | `ON DELETE CASCADE` (0013) | `ON DELETE RESTRICT` (0021) | -| `whitelisted_addresses` | `wallet_id` | `ON DELETE CASCADE` (0013) | `ON DELETE RESTRICT` (0021) | -| `payment_links` | `wallet_id` | `ON DELETE CASCADE` (0014) | `ON DELETE RESTRICT` (0021) |