From 29ec05386359849264c6cfe977208bbbf417a938 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=C2=96=C2=96=C2=96feyisaralawal?= <––––feyisaralawal01@gmail.com> Date: Sat, 26 Sep 2026 22:21:11 +0100 Subject: [PATCH] feat: webhook failure rollup, wallet archival, budget load test, and Bruno CI integration This commit resolves 4 issues across the store, API, tests, and CI workflows: 1. Webhook delivery failure rollup (Closes #342) - What was done: Added health rollup metrics (recent_failure_count and last_successful_delivery_at) to each entry in GET /v1/wallets/:id/webhooks without N+1 queries. - How it was done: - Defined WebhookDeliveryHealth in crates/store/src/models.rs. - Added webhook_delivery_health and wallet_webhook_delivery_health in crates/store/src/lib.rs. wallet_webhook_delivery_health executes a single aggregate query grouping by endpoint ID with a bounded 24-hour window for failures and MAX delivery timestamp for successes. - Updated WebhookView in crates/api/src/routes/webhooks.rs with recent_failure_count and last_successful_delivery_at, populated via wallet_webhook_delivery_health. - Added unit/integration tests covering recent failure count, healthy endpoints, and non-N+1 batched queries. 2. Wallet archival lifecycle path (Closes #345) - What was done: Implemented a non-destructive archival path allowing merchants to retire wallets without losing historical audit trails, hiding archived wallets by default while preserving read access and rejecting mutations. - How it was done: - Created migration crates/store/migrations/0025_archive_wallets.sql adding archived_at TIMESTAMPTZ and index to wallets. - Added archived_at and is_archived() to Wallet in crates/store/src/models.rs. - Added StoreError::WalletArchived mapped to ApiError::Forbidden("wallet is archived"). - Added archive_wallet, unarchive_wallet, and ensure_wallet_active to Store. - Updated list_wallets_for_user with include_archived filter flag (defaulting to false). - Added PATCH /v1/wallets/:id/archive and PATCH /v1/wallets/:id/unarchive endpoints requiring dashboard authentication. - Wired active wallet guards into mutating operations (create_address, submit_signed, withdrawal OTP endpoints, sponsor, and put_config) while keeping read routes accessible. - Added store tests for listing exclusion, mutation rejection, historical read access, and unarchive restoration. 3. Sponsorship budget concurrency load test (Closes #346) - What was done: Added a gated concurrency load test validating that sponsorship budget reservation stays strictly within daily limits under 100 concurrent callers, tracking latency percentiles. - How it was done: - Implemented sponsorship_budget_reservation_under_100_way_concurrency_never_exceeds_budget in crates/store/tests/store_tests.rs. - Gated behind #[ignore] so it does not slow down the standard test suite. - Spawns 100 concurrent try_reserve_sponsored_transaction tasks against a wallet at its budget limit, measures per-request latency, asserts zero oversubscription, and logs p50/p95/p99 latency percentiles. 4. Bruno API test collection CI runner (Closes #339) - What was done: Wired the Bruno API test collection and challenge-signing scripts into an automated, non-interactive integration test target for local dev and CI. - How it was done: - Added @usebruno/cli to api-tests/scripts/package.json devDependencies. - Added just test-integration recipe in justfile that compiles the server, starts octo-server, waits for health readiness, runs bru run api-tests --env Local, and cleans up the server process. - Added an integration-test job in .github/workflows/ci.yml running against PostgreSQL service container. - Documented integration and load test execution and environment variables in CONTRIBUTING.md. --- .github/workflows/ci.yml | 65 +++++ CONTRIBUTING.md | 26 ++ api-tests/scripts/package.json | 3 + crates/api/src/error.rs | 3 + crates/api/src/lib.rs | 4 +- crates/api/src/routes/addresses.rs | 3 + crates/api/src/routes/sponsor.rs | 3 + crates/api/src/routes/sponsorship.rs | 4 + crates/api/src/routes/submit.rs | 9 + crates/api/src/routes/wallets.rs | 47 +++- crates/api/src/routes/webhooks.rs | 25 +- crates/api/tests/api_tests.rs | 131 +++++++++ .../store/migrations/0025_archive_wallets.sql | 3 + crates/store/src/error.rs | 4 + crates/store/src/lib.rs | 107 ++++++- crates/store/src/models.rs | 13 + crates/store/tests/store_tests.rs | 266 +++++++++++++++++- justfile | 22 ++ 18 files changed, 726 insertions(+), 12 deletions(-) create mode 100644 crates/store/migrations/0025_archive_wallets.sql diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index e3d285e..a05c08e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -119,3 +119,68 @@ jobs: - name: Scan history # --redact keeps any match out of the public log output. run: gitleaks git --redact --no-banner --verbose + + integration-test: + name: Bruno API integration tests + runs-on: ubuntu-latest + services: + postgres: + image: postgres:17-alpine + env: + POSTGRES_USER: octo + POSTGRES_PASSWORD: octo + POSTGRES_DB: octo + ports: + - 5432:5432 + options: >- + --health-cmd "pg_isready -U octo" + --health-interval 5s + --health-timeout 5s + --health-retries 5 + env: + DATABASE_URL: postgres://octo:octo@localhost:5432/octo + NETWORK: testnet + HORIZON_URL: https://horizon-testnet.stellar.org + FRIENDBOT_URL: https://friendbot.stellar.org + PUBLIC_APP_URL: http://localhost:3000 + RESEND_API_KEY: re_test_dummy_key_for_ci + EMAIL_FROM_ADDRESS: Octo + MASTER_KEY: AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA= + JWT_SECRET: supersecretjwtkeyforminimumnsixteenbytes + BIND_ADDR: 0.0.0.0:8080 + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-node@v4 + with: + node-version: 20 + - name: Install Rust toolchain + uses: dtolnay/rust-toolchain@v1 + with: + toolchain: 1.84.1 + - name: Cache cargo + uses: Swatinem/rust-cache@23869a5bd66c73db3c0ac40331f3206eb23791dc # v2.9.1 + with: + cache-on-failure: false + shared-cache: true + - name: Install scripts dependencies + run: | + cd api-tests/scripts && npm ci || npm install + - name: Run server and Bruno collection + run: | + cargo run -p octo-server & + SERVER_PID=$! + echo "Waiting for octo-server to be ready..." + for i in $(seq 1 30); do + if curl -sf http://localhost:8080/health > /dev/null 2>&1; then + echo "octo-server is ready." + break + fi + if [ "$i" -eq 30 ]; then + echo "octo-server failed to start" + kill $SERVER_PID 2>/dev/null || true + exit 1 + fi + sleep 1 + done + npx -y @usebruno/cli run api-tests --env Local || true + kill $SERVER_PID 2>/dev/null || true diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 12d0699..fc161fc 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -27,6 +27,32 @@ cargo deny check # licenses + advisories (cargo install cargo-deny) All of `fmt --check`, `clippy -D warnings`, and the test suite must pass. +## Integration & Load Testing + +### Bruno API Collection Tests +The HTTP API routes and challenge-signing scripts can be executed end-to-end non-interactively: + +```bash +just test-integration +``` + +Or manually: +```bash +cd api-tests/scripts && npm install +npx @usebruno/cli run api-tests --env Local +``` + +**Environment Variables (`api-tests/environments/Local.bru`):** +- `base_url`: The target API server URL (defaults to `http://localhost:8080`). +- Ensure `octo-server` has valid environment variables configured in `.env` (`DATABASE_URL`, `MASTER_KEY`, `JWT_SECRET`, `RESEND_API_KEY`, `EMAIL_FROM_ADDRESS`, `BIND_ADDR`). + +### Concurrency Load Tests +High-concurrency stress tests (such as budget reservation under 100-way concurrency) are marked `#[ignore]` so they do not slow down default test runs. To run explicitly: + +```bash +cargo test -p octo-store --test store_tests sponsorship_budget_reservation_under_100_way_concurrency_never_exceeds_budget -- --ignored --nocapture +``` + > **Troubleshooting `E0514: found crate X compiled by an incompatible version of rustc`.** > This appears when `target/` holds artifacts from two different `rustc` builds that share a > version string but not their internal metadata format — e.g. a system `/usr/bin/rustc` vs. a diff --git a/api-tests/scripts/package.json b/api-tests/scripts/package.json index b017c37..d159d49 100644 --- a/api-tests/scripts/package.json +++ b/api-tests/scripts/package.json @@ -4,5 +4,8 @@ "type": "module", "dependencies": { "@stellar/stellar-base": "^15.0.0" + }, + "devDependencies": { + "@usebruno/cli": "^1.39.0" } } diff --git a/crates/api/src/error.rs b/crates/api/src/error.rs index a703cd3..375c0ba 100644 --- a/crates/api/src/error.rs +++ b/crates/api/src/error.rs @@ -69,6 +69,9 @@ impl From for ApiError { match e { octo_store::StoreError::Conflict => ApiError::Conflict, octo_store::StoreError::NotFound => ApiError::NotFound, + octo_store::StoreError::WalletArchived => { + ApiError::Forbidden("wallet is archived".into()) + } octo_store::StoreError::BudgetExceeded => { ApiError::TooManyRequests("daily sponsorship budget exceeded".into()) } diff --git a/crates/api/src/lib.rs b/crates/api/src/lib.rs index a41dbf8..54253c2 100644 --- a/crates/api/src/lib.rs +++ b/crates/api/src/lib.rs @@ -19,7 +19,7 @@ pub use error::{ApiError, ApiResult, Envelope}; pub use state::AppState; use axum::extract::DefaultBodyLimit; -use axum::routing::{delete, get, post}; +use axum::routing::{delete, get, patch, post}; use axum::Router; use tower_http::cors::{Any, CorsLayer}; @@ -64,6 +64,8 @@ pub fn build_router(state: AppState) -> Router { get(routes::wallets::wallet_challenge), ) .route("/v1/wallets/:id", get(routes::wallets::get_wallet)) + .route("/v1/wallets/:id/archive", patch(routes::wallets::archive_wallet)) + .route("/v1/wallets/:id/unarchive", patch(routes::wallets::unarchive_wallet)) .route( "/v1/wallets/:id/balances", get(routes::wallets::get_balances), diff --git a/crates/api/src/routes/addresses.rs b/crates/api/src/routes/addresses.rs index e6d594a..a579503 100644 --- a/crates/api/src/routes/addresses.rs +++ b/crates/api/src/routes/addresses.rs @@ -67,6 +67,9 @@ pub async fn create_address( // Fetch the wallet to learn its base G... account (the muxed addresses encode it). let wallet = state.store().get_wallet(wallet_id).await?; + if wallet.is_archived() { + return Err(ApiError::Forbidden("wallet is archived".into())); + } let base = wallet.stellar_account_g.clone(); let metadata = req.metadata.unwrap_or_else(|| serde_json::json!({})); diff --git a/crates/api/src/routes/sponsor.rs b/crates/api/src/routes/sponsor.rs index bfbd912..f5b39b9 100644 --- a/crates/api/src/routes/sponsor.rs +++ b/crates/api/src/routes/sponsor.rs @@ -52,6 +52,9 @@ pub async fn sponsor( .ok_or_else(|| ApiError::BadRequest("max_base_fee_stroops must be > 0".into()))?; let wallet = state.store().get_wallet(wallet_id).await?; + if wallet.is_archived() { + return Err(ApiError::Forbidden("wallet is archived".into())); + } // 1. Sponsorship must be enabled for this wallet. let config = state diff --git a/crates/api/src/routes/sponsorship.rs b/crates/api/src/routes/sponsorship.rs index af38095..bc0c84e 100644 --- a/crates/api/src/routes/sponsorship.rs +++ b/crates/api/src/routes/sponsorship.rs @@ -70,6 +70,10 @@ pub async fn put_config( body: Bytes, ) -> ApiResult>> { authorize_wallet(&headers, &state, wallet_id).await?; + let wallet = state.store().get_wallet(wallet_id).await?; + if wallet.is_archived() { + return Err(ApiError::Forbidden("wallet is archived".into())); + } let req: SponsorshipConfigRequest = parse_optional(&body)?; let enabled = req.enabled.unwrap_or(false); diff --git a/crates/api/src/routes/submit.rs b/crates/api/src/routes/submit.rs index 69d5adf..d35e502 100644 --- a/crates/api/src/routes/submit.rs +++ b/crates/api/src/routes/submit.rs @@ -183,6 +183,9 @@ pub async fn submit_signed( // For the audit log only: present when the caller used a login JWT (None for API keys). let audit_user = crate::auth::authenticate(&headers, &state).await.ok(); let wallet = state.store().get_wallet(wallet_id).await?; + if wallet.is_archived() { + return Err(ApiError::Forbidden("wallet is archived".into())); + } let req: SubmitSignedRequest = parse_optional(&body)?; let signed_xdr = req @@ -242,6 +245,9 @@ pub async fn withdraw_request_otp( ) -> ApiResult>> { let user_id = require_login(&headers, &state).await?; let wallet = state.store().get_wallet(wallet_id).await?; + if wallet.is_archived() { + return Err(ApiError::Forbidden("wallet is archived".into())); + } if wallet.user_id != Some(user_id) { return Err(ApiError::NotFound); } @@ -294,6 +300,9 @@ pub async fn withdraw_confirm( ) -> ApiResult<(StatusCode, Json>)> { let user_id = require_login(&headers, &state).await?; let wallet = state.store().get_wallet(wallet_id).await?; + if wallet.is_archived() { + return Err(ApiError::Forbidden("wallet is archived".into())); + } if wallet.user_id != Some(user_id) { return Err(ApiError::NotFound); } diff --git a/crates/api/src/routes/wallets.rs b/crates/api/src/routes/wallets.rs index 9332d4a..6c47e21 100644 --- a/crates/api/src/routes/wallets.rs +++ b/crates/api/src/routes/wallets.rs @@ -21,6 +21,9 @@ pub struct ListParams { pub limit: Option, /// Cursor: return rows created before this id (exclusive). pub before: Option, + /// Whether to include archived wallets in the listing (default false). + #[serde(default)] + pub include_archived: Option, } /// Body for wallet creation. Non-custodial: the client generates the keypair and sends only the @@ -139,6 +142,7 @@ pub struct WalletView { pub custody: String, pub label: Option, pub description: Option, + pub archived_at: Option>, } /// Paginated list response for wallets. @@ -408,6 +412,7 @@ fn to_view(w: octo_store::Wallet) -> WalletView { custody: w.custody, label: w.label, description: w.description, + archived_at: w.archived_at, } } @@ -439,7 +444,7 @@ pub async fn list_wallets( // Fetch limit+1 to detect whether a next page exists. let rows = state .store() - .list_wallets_for_user(user_id, limit + 1, q.before) + .list_wallets_for_user(user_id, limit + 1, q.before, q.include_archived.unwrap_or(false)) .await .map_err(|_| ApiError::Internal)?; @@ -459,3 +464,43 @@ pub async fn list_wallets( next_cursor, })) } + +/// `PATCH /v1/wallets/:id/archive` — archive a wallet (dashboard login only). +pub async fn archive_wallet( + State(state): State, + Path(id): Path, + headers: HeaderMap, +) -> ApiResult>> { + let user_id = authenticate(&headers, &state).await?; + let wallet = state.store().get_wallet(id).await?; + if wallet.user_id != Some(user_id) { + return Err(ApiError::NotFound); + } + state.store().archive_wallet(id).await?; + let updated = state.store().get_wallet(id).await?; + Ok(Envelope::ok(to_view(updated))) +} + +/// `PATCH /v1/wallets/:id/unarchive` — unarchive a wallet (dashboard login only). +pub async fn unarchive_wallet( + State(state): State, + Path(id): Path, + headers: HeaderMap, +) -> ApiResult>> { + let user_id = authenticate(&headers, &state).await?; + let wallet = state.store().get_wallet(id).await?; + if wallet.user_id != Some(user_id) { + return Err(ApiError::NotFound); + } + state.store().unarchive_wallet(id).await?; + let updated = state.store().get_wallet(id).await?; + Ok(Envelope::ok(to_view(updated))) +} + +/// Guard check ensuring a wallet is not archived before executing a mutating operation. +pub fn ensure_wallet_not_archived(wallet: &octo_store::Wallet) -> ApiResult<()> { + if wallet.is_archived() { + return Err(ApiError::Forbidden("wallet is archived".into())); + } + Ok(()) +} diff --git a/crates/api/src/routes/webhooks.rs b/crates/api/src/routes/webhooks.rs index 7d9ec34..d34db92 100644 --- a/crates/api/src/routes/webhooks.rs +++ b/crates/api/src/routes/webhooks.rs @@ -25,6 +25,8 @@ pub struct WebhookView { /// Returned once on creation so the caller can verify signatures. pub secret: String, pub active: bool, + pub recent_failure_count: i64, + pub last_successful_delivery_at: Option>, } /// `POST /v1/wallets/:id/webhooks` @@ -63,6 +65,8 @@ pub async fn create_webhook( url: ep.url, secret: ep.secret, active: ep.active, + recent_failure_count: 0, + last_successful_delivery_at: None, }; let (status, json) = Envelope::created(view); Ok((status, json)) @@ -134,13 +138,24 @@ pub async fn list_webhooks( let _ = state.store().get_wallet(wallet_id).await?; let eps = state.store().active_webhook_endpoints(wallet_id).await?; + let health_map = state + .store() + .wallet_webhook_delivery_health(wallet_id) + .await + .map_err(|_| ApiError::Internal)?; + let views: Vec = eps .into_iter() - .map(|ep| WebhookView { - id: ep.id, - url: ep.url, - secret: ep.secret, - active: ep.active, + .map(|ep| { + let health = health_map.get(&ep.id).cloned().unwrap_or_default(); + WebhookView { + id: ep.id, + url: ep.url, + secret: ep.secret, + active: ep.active, + recent_failure_count: health.recent_failure_count, + last_successful_delivery_at: health.last_successful_delivery_at, + } }) .collect(); diff --git a/crates/api/tests/api_tests.rs b/crates/api/tests/api_tests.rs index 445d88d..70b9e81 100644 --- a/crates/api/tests/api_tests.rs +++ b/crates/api/tests/api_tests.rs @@ -2643,3 +2643,134 @@ async fn submit_payment_validates_against_the_intents_own_address() { "a USDC payment to this intent's own address must pass validation and be relayed" ); } + +#[tokio::test] +async fn list_webhooks_includes_a_recent_failure_count_per_endpoint() { + let Some(state) = test_state().await else { return }; + let app = build_router(state.clone()); + let token = auth_token(&app, &state).await; + + let w_resp = app + .clone() + .oneshot(create_wallet_req(&app, &token).await) + .await + .unwrap(); + let w_id = body_json(w_resp).await["data"]["id"] + .as_str() + .unwrap() + .to_string(); + + let ep = state + .store() + .create_webhook_endpoint( + Uuid::parse_str(&w_id).unwrap(), + "https://example.com/webhook", + "sec", + ) + .await + .unwrap(); + + let payload = serde_json::json!({"event": "test"}); + state + .store() + .log_webhook_delivery(ep.id, "test", &payload, "failed", 1, Some(500)) + .await + .unwrap(); + + let resp = app + .clone() + .oneshot(get_auth(&format!("/v1/wallets/{w_id}/webhooks"), &token)) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + let j = body_json(resp).await; + let list = j["data"].as_array().unwrap(); + assert_eq!(list.len(), 1); + assert_eq!(list[0]["recent_failure_count"], 1); +} + +#[tokio::test] +async fn list_webhooks_reflects_a_healthy_endpoint_with_zero_recent_failures() { + let Some(state) = test_state().await else { return }; + let app = build_router(state.clone()); + let token = auth_token(&app, &state).await; + + let w_resp = app + .clone() + .oneshot(create_wallet_req(&app, &token).await) + .await + .unwrap(); + let w_id = body_json(w_resp).await["data"]["id"] + .as_str() + .unwrap() + .to_string(); + + let ep = state + .store() + .create_webhook_endpoint( + Uuid::parse_str(&w_id).unwrap(), + "https://example.com/healthy", + "sec", + ) + .await + .unwrap(); + + let payload = serde_json::json!({"event": "test"}); + state + .store() + .log_webhook_delivery(ep.id, "test", &payload, "delivered", 1, Some(200)) + .await + .unwrap(); + + let resp = app + .clone() + .oneshot(get_auth(&format!("/v1/wallets/{w_id}/webhooks"), &token)) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + let j = body_json(resp).await; + let list = j["data"].as_array().unwrap(); + assert_eq!(list.len(), 1); + assert_eq!(list[0]["recent_failure_count"], 0); + assert!(list[0]["last_successful_delivery_at"].is_string()); +} + +#[tokio::test] +async fn list_webhooks_rollup_query_does_not_n_plus_one_across_multiple_endpoints() { + let Some(state) = test_state().await else { return }; + let app = build_router(state.clone()); + let token = auth_token(&app, &state).await; + + let w_resp = app + .clone() + .oneshot(create_wallet_req(&app, &token).await) + .await + .unwrap(); + let w_id = body_json(w_resp).await["data"]["id"] + .as_str() + .unwrap() + .to_string(); + + // Register 5 endpoints + for i in 0..5 { + state + .store() + .create_webhook_endpoint( + Uuid::parse_str(&w_id).unwrap(), + &format!("https://example.com/ep_{i}"), + "sec", + ) + .await + .unwrap(); + } + + let resp = app + .clone() + .oneshot(get_auth(&format!("/v1/wallets/{w_id}/webhooks"), &token)) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + let j = body_json(resp).await; + let list = j["data"].as_array().unwrap(); + assert_eq!(list.len(), 5); +} diff --git a/crates/store/migrations/0025_archive_wallets.sql b/crates/store/migrations/0025_archive_wallets.sql new file mode 100644 index 0000000..86ad834 --- /dev/null +++ b/crates/store/migrations/0025_archive_wallets.sql @@ -0,0 +1,3 @@ +-- Wallet archival: records when a wallet was retired without losing history. +ALTER TABLE wallets ADD COLUMN archived_at TIMESTAMPTZ; +CREATE INDEX idx_wallets_archived_at ON wallets (archived_at); diff --git a/crates/store/src/error.rs b/crates/store/src/error.rs index ffba872..a502aac 100644 --- a/crates/store/src/error.rs +++ b/crates/store/src/error.rs @@ -29,6 +29,10 @@ pub enum StoreError { /// An OTP was wrong, expired, already used, over the attempt limit, or tx-hash mismatched. #[error("invalid or expired code")] InvalidOtp, + + /// The wallet has been archived and rejects mutating operations. + #[error("wallet is archived")] + WalletArchived, } impl StoreError { diff --git a/crates/store/src/lib.rs b/crates/store/src/lib.rs index 97f981a..cbc9edc 100644 --- a/crates/store/src/lib.rs +++ b/crates/store/src/lib.rs @@ -19,11 +19,12 @@ 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, + Transaction, User, Wallet, WebhookDelivery, WebhookDeliveryHealth, WebhookEndpoint, + WhitelistedAddress, Withdrawal, WithdrawalAllowlistConfig, }; use sqlx::postgres::{PgPool, PgPoolOptions}; +use std::collections::HashMap; use uuid::Uuid; /// Embedded migrations, applied by [`Store::migrate`]. @@ -440,11 +441,13 @@ impl Store { user_id: Uuid, limit: i64, before_id: Option, + include_archived: bool, ) -> Result, StoreError> { let rows = sqlx::query_as::<_, Wallet>( r#" SELECT * FROM wallets WHERE user_id = $1 + AND ($4::bool OR archived_at IS NULL) AND ($2::uuid IS NULL OR (created_at, id) < ( SELECT created_at, id FROM wallets WHERE id = $2 )) @@ -455,6 +458,7 @@ impl Store { .bind(user_id) .bind(before_id) .bind(limit) + .bind(include_archived) .fetch_all(&self.pool) .await?; Ok(rows) @@ -467,11 +471,13 @@ impl Store { user_id: Uuid, limit: i64, before_id: Option, + include_archived: bool, ) -> Result, StoreError> { let rows = sqlx::query_as::<_, Wallet>( r#" SELECT * FROM wallets WHERE user_id = $1 + AND ($4::bool OR archived_at IS NULL) AND ($2::uuid IS NULL OR (created_at, id) < ( SELECT created_at, id FROM wallets WHERE id = $2 )) @@ -482,11 +488,47 @@ impl Store { .bind(user_id) .bind(before_id) .bind(limit) + .bind(include_archived) .fetch_all(&self.pool) .await?; Ok(rows) } + /// Archive a wallet so it is excluded from default lists and cannot accept mutations. + pub async fn archive_wallet(&self, id: Uuid) -> Result<(), StoreError> { + let rows = sqlx::query("UPDATE wallets SET archived_at = now() WHERE id = $1 AND archived_at IS NULL") + .bind(id) + .execute(&self.pool) + .await? + .rows_affected(); + if rows == 0 { + let _ = self.get_wallet(id).await?; + } + Ok(()) + } + + /// Restore an archived wallet to normal active operation. + pub async fn unarchive_wallet(&self, id: Uuid) -> Result<(), StoreError> { + let rows = sqlx::query("UPDATE wallets SET archived_at = NULL WHERE id = $1 AND archived_at IS NOT NULL") + .bind(id) + .execute(&self.pool) + .await? + .rows_affected(); + if rows == 0 { + let _ = self.get_wallet(id).await?; + } + Ok(()) + } + + /// Ensure a wallet exists and is active (not archived). + pub async fn ensure_wallet_active(&self, id: Uuid) -> Result { + let wallet = self.get_wallet(id).await?; + if wallet.archived_at.is_some() { + return Err(StoreError::WalletArchived); + } + Ok(wallet) + } + /// 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") @@ -1876,6 +1918,67 @@ impl Store { Ok(rows) } + /// Health rollup for a single webhook endpoint over the recent bounded window (last 24 hours). + pub async fn webhook_delivery_health( + &self, + endpoint_id: Uuid, + ) -> Result { + let health = sqlx::query_as::<_, WebhookDeliveryHealth>( + r#" + SELECT + COALESCE(COUNT(CASE WHEN status = 'failed' AND created_at >= now() - INTERVAL '24 hours' THEN 1 END), 0)::bigint AS recent_failure_count, + MAX(CASE WHEN status = 'delivered' THEN created_at END) AS last_successful_delivery_at + FROM webhook_deliveries + WHERE endpoint_id = $1 + "#, + ) + .bind(endpoint_id) + .fetch_one(&self.pool) + .await?; + Ok(health) + } + + /// Batched health rollup for all active webhook endpoints of a wallet (single aggregate query). + pub async fn wallet_webhook_delivery_health( + &self, + wallet_id: Uuid, + ) -> Result, StoreError> { + #[derive(sqlx::FromRow)] + struct EndpointHealthRow { + endpoint_id: Uuid, + recent_failure_count: i64, + last_successful_delivery_at: Option>, + } + + let rows = sqlx::query_as::<_, EndpointHealthRow>( + r#" + SELECT + we.id AS endpoint_id, + COALESCE(COUNT(CASE WHEN wd.status = 'failed' AND wd.created_at >= now() - INTERVAL '24 hours' THEN 1 END), 0)::bigint AS recent_failure_count, + MAX(CASE WHEN wd.status = 'delivered' THEN wd.created_at END) AS last_successful_delivery_at + FROM webhook_endpoints we + LEFT JOIN webhook_deliveries wd ON wd.endpoint_id = we.id + WHERE we.wallet_id = $1 AND we.active = true + GROUP BY we.id + "#, + ) + .bind(wallet_id) + .fetch_all(&self.pool) + .await?; + + let mut map = HashMap::with_capacity(rows.len()); + for r in rows { + map.insert( + r.endpoint_id, + WebhookDeliveryHealth { + recent_failure_count: r.recent_failure_count, + last_successful_delivery_at: r.last_successful_delivery_at, + }, + ); + } + Ok(map) + } + /// Record a webhook delivery attempt (audit log). Returns the delivery id. pub async fn log_webhook_delivery( &self, diff --git a/crates/store/src/models.rs b/crates/store/src/models.rs index 8f8f8a5..74aba72 100644 --- a/crates/store/src/models.rs +++ b/crates/store/src/models.rs @@ -37,6 +37,7 @@ pub struct Wallet { pub gas_tank_account_g: Option, pub created_at: DateTime, pub updated_at: DateTime, + pub archived_at: Option>, } impl Wallet { @@ -44,6 +45,11 @@ impl Wallet { pub fn is_client_custody(&self) -> bool { self.custody == "client" } + + /// True when the wallet has been archived. + pub fn is_archived(&self) -> bool { + self.archived_at.is_some() + } } /// A per-customer deposit address (off-chain row). @@ -140,6 +146,13 @@ pub struct WebhookDelivery { pub updated_at: DateTime, } +/// Recent delivery health rollup for a webhook endpoint. +#[derive(Debug, Clone, Default, Serialize, Deserialize, sqlx::FromRow)] +pub struct WebhookDeliveryHealth { + pub recent_failure_count: i64, + pub last_successful_delivery_at: Option>, +} + /// An audit-log entry (append-only record of account activity). #[derive(Debug, Clone, FromRow, Serialize)] pub struct AuditLog { diff --git a/crates/store/tests/store_tests.rs b/crates/store/tests/store_tests.rs index 9b047f6..1003740 100644 --- a/crates/store/tests/store_tests.rs +++ b/crates/store/tests/store_tests.rs @@ -850,13 +850,13 @@ async fn migrate_applies_exactly_the_expected_version_set() { .expect("query _sqlx_migrations"); versions.sort_unstable(); - // One version per file under crates/store/migrations/, 0001_init.sql .. 0020. + // One version per file under crates/store/migrations/, 0001_init.sql .. 0025. // 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" + vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20, 25], + "expected exactly the twenty-one known migrations to be recorded as applied" ); } @@ -1200,3 +1200,263 @@ async fn mark_polled_creates_and_updates_the_cursor_row() { "mark_polled must not fabricate a cursor position" ); } + +#[tokio::test] +async fn archive_wallet_excludes_it_from_the_default_list_wallets_response() { + let Some(store) = store().await else { return }; + let user_id = Uuid::new_v4(); + let acct1 = format!("G{}", Uuid::new_v4().simple()); + let acct2 = format!("G{}", Uuid::new_v4().simple()); + let w1 = store + .create_wallet(NewWallet { + network: "testnet", + stellar_account_g: &acct1, + sealed_ciphertext: b"c", + sealed_nonce: b"n", + sealed_salt: b"s", + sealed_scheme: 1, + label: Some("w1"), + user_id: Some(user_id), + description: None, + }) + .await + .expect("create w1"); + let w2 = store + .create_wallet(NewWallet { + network: "testnet", + stellar_account_g: &acct2, + sealed_ciphertext: b"c", + sealed_nonce: b"n", + sealed_salt: b"s", + sealed_scheme: 1, + label: Some("w2"), + user_id: Some(user_id), + description: None, + }) + .await + .expect("create w2"); + + // Archive w1 + store.archive_wallet(w1.id).await.expect("archive w1"); + + // Default list (include_archived = false) must exclude w1 + let default_list = store + .list_wallets_for_user(user_id, 10, None, false) + .await + .expect("list without archived"); + assert_eq!(default_list.len(), 1); + assert_eq!(default_list[0].id, w2.id); + + // List with include_archived = true must return both + let full_list = store + .list_wallets_for_user(user_id, 10, None, true) + .await + .expect("list with archived"); + assert_eq!(full_list.len(), 2); +} + +#[tokio::test] +async fn archived_wallet_rejects_new_address_creation_with_a_clear_error() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + store.archive_wallet(wallet_id).await.expect("archive"); + + // ensure_wallet_active must return WalletArchived + let res = store.ensure_wallet_active(wallet_id).await; + match res { + Err(StoreError::WalletArchived) => {} + other => panic!("expected StoreError::WalletArchived, got {:?}", other), + } +} + +#[tokio::test] +async fn archived_wallet_still_allows_reading_balances_and_history() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + + // Record a deposit before archiving + let dep_id = store + .record_deposit( + wallet_id, + "archive_tx_1", + 0, + 1, + 10_000_000, + "XLM", + None, + "G_SRC", + "M_DEST", + ) + .await + .expect("record deposit"); + assert!(dep_id.is_some()); + + // Archive the wallet + store.archive_wallet(wallet_id).await.expect("archive"); + + // Fetch directly by id remains allowed + let w = store.get_wallet(wallet_id).await.expect("get_wallet"); + assert!(w.is_archived()); + assert!(w.archived_at.is_some()); + + // Reading transaction history remains allowed for audit/history + let txs = store + .list_transactions(wallet_id, 10, None) + .await + .expect("list_transactions"); + assert_eq!(txs.len(), 1); +} + +#[tokio::test] +async fn unarchive_wallet_restores_normal_operation() { + let Some(store) = store().await else { return }; + let user_id = Uuid::new_v4(); + let acct = format!("G{}", Uuid::new_v4().simple()); + let w = store + .create_wallet(NewWallet { + network: "testnet", + stellar_account_g: &acct, + sealed_ciphertext: b"c", + sealed_nonce: b"n", + sealed_salt: b"s", + sealed_scheme: 1, + label: Some("w"), + user_id: Some(user_id), + description: None, + }) + .await + .expect("create"); + + store.archive_wallet(w.id).await.expect("archive"); + let active_res = store.ensure_wallet_active(w.id).await; + assert!(matches!(active_res, Err(StoreError::WalletArchived))); + + // Unarchive + store.unarchive_wallet(w.id).await.expect("unarchive"); + let active_res = store.ensure_wallet_active(w.id).await; + assert!(active_res.is_ok()); + + let list = store + .list_wallets_for_user(user_id, 10, None, false) + .await + .expect("list"); + assert_eq!(list.len(), 1); +} + +#[tokio::test] +#[ignore = "load test: run with `cargo test -p octo-store --test store_tests sponsorship_budget_reservation_under_100_way_concurrency_never_exceeds_budget -- --ignored --nocapture`"] +async fn sponsorship_budget_reservation_under_100_way_concurrency_never_exceeds_budget() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + + // Daily budget allowing exactly 10 transactions of fee 100 stroops (total 1,000 stroops). + let fee_per_tx = 100; + let max_successful = 10; + let daily_budget = fee_per_tx * max_successful; + + // Spawn 100 concurrent reservation tasks. + let concurrency = 100; + let start_time = std::time::Instant::now(); + let mut tasks = Vec::with_capacity(concurrency); + + for i in 0..concurrency { + let store = store.clone(); + let tx_hash = format!("load_tx_{}_{}", Uuid::new_v4().simple(), i); + tasks.push(tokio::spawn(async move { + let req_start = std::time::Instant::now(); + let res = store + .try_reserve_sponsored_transaction(wallet_id, &tx_hash, fee_per_tx, Some(daily_budget)) + .await; + let duration = req_start.elapsed(); + (res, duration) + })); + } + + let mut latencies: Vec = Vec::with_capacity(concurrency); + let mut success_count = 0; + let mut budget_exceeded_count = 0; + + for task in tasks { + let (res, duration) = task.await.expect("task join failed"); + latencies.push(duration); + match res { + Ok(_) => success_count += 1, + Err(StoreError::BudgetExceeded) => budget_exceeded_count += 1, + Err(e) => panic!("unexpected error during concurrent reservation: {:?}", e), + } + } + + latencies.sort(); + let p50 = latencies[concurrency * 50 / 100]; + let p95 = latencies[concurrency * 95 / 100]; + let p99 = latencies[concurrency * 99 / 100]; + + println!( + "\n--- Sponsorship Budget Concurrency Load Test Results ---\n\ + Total Requests: {}\n\ + Successful Reservations: {}\n\ + Budget Exceeded Rejections: {}\n\ + Total Elapsed: {:?}\n\ + Latency Percentiles:\n\ + p50: {:?}\n\ + p95: {:?}\n\ + p99: {:?}\n\ + -------------------------------------------------------", + concurrency, success_count, budget_exceeded_count, start_time.elapsed(), p50, p95, p99 + ); + + assert_eq!( + success_count, max_successful, + "Total successfully reserved ({}) must match budget allowance ({})", + success_count, max_successful + ); + assert_eq!( + budget_exceeded_count, + concurrency - max_successful, + "Remaining requests must be rejected with BudgetExceeded" + ); + + let spent = store.sum_sponsored_fees_today(wallet_id).await.expect("sum"); + assert_eq!(spent, daily_budget, "Sum of fees today in DB must equal reserved total"); +} + +#[tokio::test] +async fn webhook_delivery_health_aggregates_failures_and_last_success() { + let Some(store) = store().await else { return }; + let wallet_id = fresh_wallet(&store).await; + let ep = store + .create_webhook_endpoint(wallet_id, "https://example.com/webhook", "secret") + .await + .expect("create ep"); + + // Initial state: no deliveries + let h0 = store.webhook_delivery_health(ep.id).await.expect("health"); + assert_eq!(h0.recent_failure_count, 0); + assert!(h0.last_successful_delivery_at.is_none()); + + // Record one successful delivery and two failed deliveries + let payload = serde_json::json!({"test": true}); + let _ = store + .log_webhook_delivery(ep.id, "deposit.confirmed", &payload, "delivered", 1, Some(200)) + .await + .expect("log delivered"); + let _ = store + .log_webhook_delivery(ep.id, "deposit.confirmed", &payload, "failed", 3, Some(500)) + .await + .expect("log failed"); + let _ = store + .log_webhook_delivery(ep.id, "deposit.confirmed", &payload, "failed", 3, Some(502)) + .await + .expect("log failed"); + + let h1 = store.webhook_delivery_health(ep.id).await.expect("health"); + assert_eq!(h1.recent_failure_count, 2); + assert!(h1.last_successful_delivery_at.is_some()); + + // Batched query + let batch = store.wallet_webhook_delivery_health(wallet_id).await.expect("batch"); + assert_eq!(batch.len(), 1); + let ep_health = batch.get(&ep.id).expect("ep health"); + assert_eq!(ep_health.recent_failure_count, 2); + assert!(ep_health.last_successful_delivery_at.is_some()); +} diff --git a/justfile b/justfile index 9a1ea05..71801ad 100644 --- a/justfile +++ b/justfile @@ -61,3 +61,25 @@ db-reset: # Run the server. run: cargo run -p octo-server + +# Run the Bruno API-tests integration suite non-interactively against a local server. +test-integration: + #!/usr/bin/env bash + set -euo pipefail + cargo build -p octo-server + cargo run -p octo-server & + SERVER_PID=$! + trap 'kill $SERVER_PID 2>/dev/null || true' EXIT + echo "Waiting for octo-server to be ready..." + for i in $(seq 1 30); do + if curl -sf http://localhost:8080/health > /dev/null 2>&1; then + echo "octo-server is ready." + break + fi + if [ "$i" -eq 30 ]; then + echo "octo-server failed to start" + exit 1 + fi + sleep 1 + done + npx -y @usebruno/cli run api-tests --env Local