From 8131b08ea2570c6a8f491d863b7944b56b569049 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:59:34 +0100 Subject: [PATCH] feat: resolve issues #340, #343, #344, and #347 Resolves 4 issues simultaneously across test coverage, operational documentation, and observability: 1. Issue #340: Statistical smoke test for OTP code distribution - Verified that crates/email/src/lib.rs uses rand::rngs::OsRng (OS CSPRNG) to generate OTPs. - Added generate_otp_produces_a_roughly_uniform_distribution_of_digits_across_a_large_sample testing 100,000 generated OTP codes across all 6 digit positions to assert no systematic digit bias or modulo distortions. - Added generate_otp_duplicate_rate_across_a_large_sample_is_consistent_with_true_uniform_randomness testing duplicate collisions across 100,000 draws from the 1,000,000 code space against the birthday paradox expectation (~95,163 expected unique codes). - Documented that these tests serve as gross implementation bug smoke tests rather than full cryptographic certifications. - Closes #340 2. Issue #343: Proptest fuzz corpus for operation_index_from_toid - Added property-based fuzz tests using proptest! in crates/ingest/src/lib.rs. - Added operation_index_from_toid_never_panics_on_arbitrary_input fuzzing with arbitrary string inputs to guarantee crash freedom. - Added operation_index_from_toid_extracted_value_is_always_within_the_documented_valid_range_when_some across valid TOID patterns, whitespaced inputs, and boundary numbers asserting that any extracted operation index is non-negative and <= i32::MAX. - Closes #343 3. Issue #344: Operational runbooks for migrate-keys and backfill-operation-index - Created docs/runbook-migrate-keys.md providing end-to-end guidance for the dual-key rotation window, pre-flight checks, invocation examples (full rotation and cipher upgrade), healthy log indicators, store method references (Store::list_wallets_needing_reseal, Store::reseal_wallet), and interruption recovery/rollback procedures. - Created docs/runbook-backfill-operation-index.md providing guidance for historical deposit TOID operation index backfills, pre-flight candidate estimation, dry-run and live invocations, expected logs, store/ingest code cross-references (octo_ingest::operation_index_from_toid), and transactional idempotency guarantees. - Linked both runbooks in README.md under the Documentation section. - Closes #344 4. Issue #347: Named tracing spans around Horizon calls - Added distinct #[tracing::instrument] spans around public Horizon interaction methods in crates/api/src/horizon.rs (balances, account_sequence, account_info, submit_transaction, friendbot_fund) and crates/ingest/src/horizon.rs (payments_after). - Configured spans to record call metadata (call_type, account, cursor, limit) and dynamic outcome (success, circuit_open, not_found, rejected/tx_failed, failure) while skipping sensitive payloads (such as raw transaction XDR envelopes). - Added unit tests in crates/api/src/horizon.rs and crates/ingest/src/horizon.rs verifying named span emissions using custom test tracing subscriber layers. - Added tracing-subscriber to dev-dependencies for octo-api and octo-ingest. - Closes #347 --- README.md | 11 ++ crates/api/Cargo.toml | 1 + crates/api/src/horizon.rs | 129 ++++++++++++++++++++++- crates/email/src/lib.rs | 53 ++++++++++ crates/ingest/Cargo.toml | 1 + crates/ingest/src/horizon.rs | 56 ++++++++++ crates/ingest/src/lib.rs | 30 ++++++ docs/runbook-backfill-operation-index.md | 102 ++++++++++++++++++ docs/runbook-migrate-keys.md | 108 +++++++++++++++++++ 9 files changed, 489 insertions(+), 2 deletions(-) create mode 100644 docs/runbook-backfill-operation-index.md create mode 100644 docs/runbook-migrate-keys.md diff --git a/README.md b/README.md index d8f4684..6d154ff 100644 --- a/README.md +++ b/README.md @@ -210,6 +210,17 @@ Full mapping in **[docs/threat-model.md](docs/threat-model.md)**. Amounts are in end-to-end (never floats). Report vulnerabilities per **[SECURITY.md](SECURITY.md)** — **do not** open public issues for security reports. +## Documentation + +- [Architecture Overview](docs/architecture.md) +- [Threat Model & Security](docs/threat-model.md) +- [Deposit Model & Attribution](docs/deposit-model.md) +- [Non-Custodial Transaction Flow](docs/non-custodial-flow.md) +- [Operational Runbook: Key Migration](docs/runbook-migrate-keys.md) +- [Operational Runbook: Operation Index Backfill](docs/runbook-backfill-operation-index.md) +- [Backfill Constraint Safety Analysis](docs/backfill-constraint-analysis.md) +- [REST API Specification](docs/api.md) + ## Roadmap - **Gas sponsorship** — *shipped.* App developers can sponsor their users' Stellar transactions diff --git a/crates/api/Cargo.toml b/crates/api/Cargo.toml index 2f37fd0..0485e68 100644 --- a/crates/api/Cargo.toml +++ b/crates/api/Cargo.toml @@ -41,6 +41,7 @@ hex.workspace = true [dev-dependencies] tokio.workspace = true tower.workspace = true +tracing-subscriber.workspace = true sqlx.workspace = true stellar-base.workspace = true dotenvy = "0.15" diff --git a/crates/api/src/horizon.rs b/crates/api/src/horizon.rs index 0e43073..f5555f7 100644 --- a/crates/api/src/horizon.rs +++ b/crates/api/src/horizon.rs @@ -217,7 +217,13 @@ impl Horizon { /// Fetch an account's balances. Retried on transient failures (transport errors, 5xx). /// Returns `NotFound` if the account does not exist on-chain yet. + #[tracing::instrument( + name = "horizon_balances", + skip(self), + fields(call_type = "balances", account_g = %account_g, outcome = tracing::field::Empty) + )] pub async fn balances(&self, account_g: &str) -> Result, ApiError> { + let span = tracing::Span::current(); let url = format!( "{}/accounts/{}", self.base_url.trim_end_matches('/'), @@ -250,21 +256,48 @@ impl Horizon { }) .await; + match &result { + Ok(_) => span.record("outcome", "success"), + Err(ResilienceError::Circuit) => span.record("outcome", "circuit_open"), + Err(ResilienceError::Exhausted(FetchError::NotFound)) => { + span.record("outcome", "not_found") + } + Err(ResilienceError::Exhausted(_)) => span.record("outcome", "failure"), + }; + map_result(result) } /// Fetch an account's current sequence number. Retried on transient failures. /// Returns `NotFound` if the account doesn't exist. + #[tracing::instrument( + name = "horizon_account_sequence", + skip(self), + fields(call_type = "account_sequence", account_g = %account_g, outcome = tracing::field::Empty) + )] pub async fn account_sequence(&self, account_g: &str) -> Result { - self.account_info(account_g).await.map(|a| a.sequence) + let span = tracing::Span::current(); + let res = self.account_info(account_g).await.map(|a| a.sequence); + match &res { + Ok(_) => span.record("outcome", "success"), + Err(ApiError::NotFound) => span.record("outcome", "not_found"), + Err(_) => span.record("outcome", "failure"), + }; + res } /// Fetch balances, sequence, and reserve inputs for an account in a single Horizon call. /// `NotFound` if the account does not exist on-chain yet. + #[tracing::instrument( + name = "horizon_account_info", + skip(self), + fields(call_type = "account_info", account_g = %account_g, outcome = tracing::field::Empty) + )] pub async fn account_info(&self, account_g: &str) -> Result { // This is a read-only call, so it goes through the same retry + circuit-breaker path as // `balances`. (It previously issued a bare, unwrapped GET, so a transient 5xx from // Horizon failed immediately instead of being retried.) + let span = tracing::Span::current(); let url = format!( "{}/accounts/{}", self.base_url.trim_end_matches('/'), @@ -309,6 +342,15 @@ impl Horizon { }) .await; + match &result { + Ok(_) => span.record("outcome", "success"), + Err(ResilienceError::Circuit) => span.record("outcome", "circuit_open"), + Err(ResilienceError::Exhausted(FetchError::NotFound)) => { + span.record("outcome", "not_found") + } + Err(ResilienceError::Exhausted(_)) => span.record("outcome", "failure"), + }; + match result { Ok(info) => Ok(info), Err(ResilienceError::Circuit) => Err(ApiError::Internal), @@ -326,11 +368,17 @@ impl Horizon { /// /// Returns the result even when the transaction failed on-chain (`successful == false`) so the /// caller can record the failure; only transport/HTTP errors return `Err`. + #[tracing::instrument( + name = "horizon_submit_transaction", + skip(self, envelope_xdr), + fields(call_type = "submit", outcome = tracing::field::Empty) + )] pub async fn submit_transaction(&self, envelope_xdr: &str) -> Result { // NOTE: an eager, unwrapped POST used to sit here ahead of the resilience-wrapped call // below. It fired a *duplicate* submission on every call and referenced `http`/`xdr` // locals that were never bound (so this did not compile). Removed — the single submit // now happens inside `execute`, with SUBMIT_TIMEOUT applied to that request. + let span = tracing::Span::current(); let url = format!("{}/transactions", self.base_url.trim_end_matches('/')); let http = self.http.clone(); let xdr = envelope_xdr.to_string(); @@ -378,6 +426,21 @@ impl Horizon { }) .await; + match &result { + Ok(r) => { + if r.successful { + span.record("outcome", "success"); + } else { + span.record("outcome", "tx_failed"); + } + } + Err(ResilienceError::Circuit) => span.record("outcome", "circuit_open"), + Err(ResilienceError::Exhausted(FetchError::TxRejected)) => { + span.record("outcome", "rejected") + } + Err(ResilienceError::Exhausted(_)) => span.record("outcome", "failure"), + }; + match result { Ok(r) => Ok(r), Err(ResilienceError::Circuit) => Err(ApiError::Internal), @@ -443,8 +506,19 @@ fn map_result(r: Result>) -> Result Result<(), ApiError> { - friendbot_fund_with_timeout(friendbot_url, account_g, DEFAULT_TIMEOUT).await + let span = tracing::Span::current(); + let res = friendbot_fund_with_timeout(friendbot_url, account_g, DEFAULT_TIMEOUT).await; + match &res { + Ok(_) => span.record("outcome", "success"), + Err(_) => span.record("outcome", "failure"), + }; + res } async fn friendbot_fund_with_timeout( @@ -525,4 +599,55 @@ mod tests { assert!(matches!(result, Err(ApiError::Internal))); } + + // Asserts that every Horizon call type emits its distinct named tracing span. + #[tokio::test] + async fn horizon_calls_emit_named_tracing_spans() { + use std::sync::{Arc, Mutex}; + use tracing_subscriber::layer::SubscriberExt; + + let recorded_spans = Arc::new(Mutex::new(Vec::::new())); + let spans_clone = recorded_spans.clone(); + + struct SpanRecorder(Arc>>); + impl tracing_subscriber::Layer for SpanRecorder { + fn on_new_span( + &self, + attrs: &tracing::span::Attributes<'_>, + _id: &tracing::span::Id, + _ctx: tracing_subscriber::layer::Context<'_, S>, + ) { + self.0.lock().unwrap().push(attrs.metadata().name().to_string()); + } + } + + let subscriber = tracing_subscriber::registry().with(SpanRecorder(spans_clone)); + let _guard = tracing::subscriber::set_default(subscriber); + + let base_url = hanging_server().await; + let horizon = Horizon { + http: reqwest::Client::builder() + .timeout(Duration::from_millis(50)) + .build() + .unwrap(), + base_url: base_url.clone(), + retry: RetryPolicy { + max_attempts: 1, + ..Default::default() + }, + circuit: CircuitBreaker::new(u32::MAX, Duration::from_secs(60)), + }; + + let _ = horizon.balances("GABCDEFGHIJKLMNOPQRSTUVWXYZ").await; + let _ = horizon.account_sequence("GABCDEFGHIJKLMNOPQRSTUVWXYZ").await; + let _ = horizon.submit_transaction("AAAA").await; + let _ = friendbot_fund("http://127.0.0.1:1", "GABCDEFGHIJKLMNOPQRSTUVWXYZ").await; + + let spans = recorded_spans.lock().unwrap().clone(); + assert!(spans.contains(&"horizon_balances".to_string())); + assert!(spans.contains(&"horizon_account_sequence".to_string())); + assert!(spans.contains(&"horizon_account_info".to_string())); + assert!(spans.contains(&"horizon_submit_transaction".to_string())); + assert!(spans.contains(&"horizon_friendbot_fund".to_string())); + } } diff --git a/crates/email/src/lib.rs b/crates/email/src/lib.rs index a428519..954e939 100644 --- a/crates/email/src/lib.rs +++ b/crates/email/src/lib.rs @@ -139,3 +139,56 @@ pub fn hash_otp(code: &str) -> String { let digest = Sha256::digest(code.as_bytes()); hex::encode(digest) } + +#[cfg(test)] +mod tests { + use super::*; + use std::collections::HashSet; + + // Gross-bug smoke test: asserts digit distribution across 100k samples is roughly uniform (not a crypto audit). + #[test] + fn generate_otp_produces_a_roughly_uniform_distribution_of_digits_across_a_large_sample() { + const SAMPLES: usize = 100_000; + let mut digit_counts = [[0usize; 10]; 6]; + + for _ in 0..SAMPLES { + let otp = generate_otp(); + assert_eq!(otp.len(), 6, "OTP must always be exactly 6 characters"); + assert!(otp.chars().all(|c| c.is_ascii_digit()), "OTP must only contain digits"); + + for (pos, ch) in otp.chars().enumerate() { + let d = ch.to_digit(10).expect("valid digit") as usize; + digit_counts[pos][d] += 1; + } + } + + // Each digit at each position has expected frequency 10,000; assert within 8,000..=12,000. + for pos in 0..6 { + for digit in 0..10 { + let count = digit_counts[pos][digit]; + assert!( + (8_000..=12_000).contains(&count), + "digit {digit} at position {pos} occurred {count} times (expected ~10,000)" + ); + } + } + } + + // Gross-bug smoke test: duplicate rate across 100k draws in 1M space matches birthday-paradox expectation. + #[test] + fn generate_otp_duplicate_rate_across_a_large_sample_is_consistent_with_true_uniform_randomness() { + const SAMPLES: usize = 100_000; + let mut seen = HashSet::with_capacity(SAMPLES); + + for _ in 0..SAMPLES { + seen.insert(generate_otp()); + } + + // For 100,000 draws from 1,000,000 bins, expected unique is ~95,163 (assert within 92,000..=98,000). + let unique_count = seen.len(); + assert!( + (92_000..=98_000).contains(&unique_count), + "expected ~95,163 unique codes, got {unique_count}" + ); + } +} diff --git a/crates/ingest/Cargo.toml b/crates/ingest/Cargo.toml index 612e5c4..a05b788 100644 --- a/crates/ingest/Cargo.toml +++ b/crates/ingest/Cargo.toml @@ -23,6 +23,7 @@ uuid.workspace = true [dev-dependencies] tokio.workspace = true +tracing-subscriber.workspace = true dotenvy = "0.15" axum.workspace = true octo-webhooks.workspace = true diff --git a/crates/ingest/src/horizon.rs b/crates/ingest/src/horizon.rs index 4809d3c..7843e2b 100644 --- a/crates/ingest/src/horizon.rs +++ b/crates/ingest/src/horizon.rs @@ -134,12 +134,24 @@ impl HorizonPayments { /// Oldest-first (`order=asc`) so we process and advance the cursor monotonically. Transient /// failures are retried with exponential backoff; the circuit breaker opens after repeated /// failures so the ingest loop doesn't pile up independent timeouts. + #[tracing::instrument( + name = "horizon_payments_after", + skip(self), + fields( + call_type = "payments_after", + account_g = %account_g, + cursor = ?cursor, + limit = limit, + outcome = tracing::field::Empty + ) + )] pub async fn payments_after( &self, account_g: &str, cursor: Option<&str>, limit: u32, ) -> Result, HorizonError> { + let span = tracing::Span::current(); let mut url = format!( "{}/accounts/{}/payments?order=asc&limit={}&join=transactions", self.base_url.trim_end_matches('/'), @@ -178,6 +190,15 @@ impl HorizonPayments { }) .await; + match &result { + Ok(_) => span.record("outcome", "success"), + Err(ResilienceError::Circuit) => span.record("outcome", "circuit_open"), + Err(ResilienceError::Exhausted(IngestFetchError::Decode)) => { + span.record("outcome", "decode_error") + } + Err(ResilienceError::Exhausted(_)) => span.record("outcome", "failure"), + }; + match result { Ok(records) => Ok(records), Err(ResilienceError::Circuit) => Err(HorizonError::CircuitOpen), @@ -216,3 +237,38 @@ impl std::fmt::Display for IngestFetchError { } } } + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::{Arc, Mutex}; + use tracing_subscriber::layer::SubscriberExt; + + // Asserts that horizon_payments_after emits its distinct named tracing span. + #[tokio::test] + async fn horizon_payments_after_emits_tracing_span() { + let recorded_spans = Arc::new(Mutex::new(Vec::::new())); + let spans_clone = recorded_spans.clone(); + + struct SpanRecorder(Arc>>); + impl tracing_subscriber::Layer for SpanRecorder { + fn on_new_span( + &self, + attrs: &tracing::span::Attributes<'_>, + _id: &tracing::span::Id, + _ctx: tracing_subscriber::layer::Context<'_, S>, + ) { + self.0.lock().unwrap().push(attrs.metadata().name().to_string()); + } + } + + let subscriber = tracing_subscriber::registry().with(SpanRecorder(spans_clone)); + let _guard = tracing::subscriber::set_default(subscriber); + + let client = HorizonPayments::new("http://127.0.0.1:1"); + let _ = client.payments_after("GABCDEFGHIJKLMNOPQRSTUVWXYZ", None, 10).await; + + let spans = recorded_spans.lock().unwrap().clone(); + assert!(spans.contains(&"horizon_payments_after".to_string())); + } +} diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index 1ed809e..8d9671a 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -856,4 +856,34 @@ mod tests { // Numbers outside i32 range should fail assert_eq!(operation_index_from_toid("12345-1-2147483648"), None); } + + use proptest::prelude::*; + + proptest! { + #![proptest_config(ProptestConfig::with_cases(1000))] + + // Fuzz test asserting that operation_index_from_toid never panics on arbitrary string inputs. + #[test] + fn operation_index_from_toid_never_panics_on_arbitrary_input( + input in ".*" + ) { + let _ = operation_index_from_toid(&input); + } + + // Fuzz test asserting any successfully extracted operation index is non-negative and within valid bounds. + #[test] + fn operation_index_from_toid_extracted_value_is_always_within_the_documented_valid_range_when_some( + input in prop_oneof![ + ".*", + "[0-9]{1,19}-[0-9]{1,10}-[0-9]{1,10}", + "[ \t]*[0-9]+-[0-9]+-[0-9]+[ \t]*", + "-?[0-9]+--?[0-9]+--?[0-9]+", + "\\PC*", + ] + ) { + if let Some(idx) = operation_index_from_toid(&input) { + prop_assert!(idx >= 0 && idx <= i32::MAX); + } + } + } } diff --git a/docs/runbook-backfill-operation-index.md b/docs/runbook-backfill-operation-index.md new file mode 100644 index 0000000..8d82ae9 --- /dev/null +++ b/docs/runbook-backfill-operation-index.md @@ -0,0 +1,102 @@ +# Operational Runbook: Operation Index Backfill (`bin/backfill-operation-index`) + +## Overview + +The `octo-backfill-operation-index` binary is an operator tool designed to backfill historical deposit records in the `transactions` table. Early versions of deposit ingestion defaulted `operation_index` to `0` rather than extracting the true index from the Horizon Transaction Operation ID (`horizon_op_id` / TOID). + +Correcting `operation_index` ensures that multi-operation transactions satisfy the `uq_tx_onchain` partial unique index on `(stellar_tx_hash, operation_index)`. For a full theoretical and database constraint safety analysis, see [docs/backfill-constraint-analysis.md](file:///c:/Users/DELL/OneDrive/Desktop/drip/Octo-Protocol-6/docs/backfill-constraint-analysis.md). + +### When to Run +- Post-migration execution to update existing legacy deposits with accurate operation indices. +- Before enforcing strict schema constraints or auditing multi-operation deposit uniqueness. + +## Pre-flight Checks + +1. **Verify Database Connectivity & Snapshot**: + ```bash + pg_dump -Fc "$DATABASE_URL" -t transactions > "transactions_backup_$(date +%Y%m%d_%H%M%S).dump" + ``` +2. **Estimate Candidate Rows**: + Identify transactions requiring updates (where TOID indicates index > 0 but stored index is 0): + ```sql + SELECT count(*) + FROM transactions + WHERE direction = 'deposit' + AND horizon_op_id IS NOT NULL + AND operation_index = 0 + AND split_part(horizon_op_id, '-', 3) <> '0'; + ``` +3. **Execute a Dry Run**: + Run with `--dry-run` to preview planned updates without modifying any data. + +## Invocations + +### Dry Run (Non-destructive Preview) +```bash +DATABASE_URL="postgres://user:pass@localhost:5432/octo" \ + cargo run --release -p octo-backfill-operation-index -- --dry-run --batch-size 1000 +``` + +### Canary Run (Limited Sample) +Process a small batch (e.g. 50 records) to verify live database behavior: +```bash +DATABASE_URL="postgres://user:pass@localhost:5432/octo" \ + cargo run --release -p octo-backfill-operation-index -- --limit 50 --batch-size 50 +``` + +### Full Live Execution +```bash +DATABASE_URL="postgres://user:pass@localhost:5432/octo" \ + cargo run --release -p octo-backfill-operation-index -- --batch-size 1000 +``` + +## Expected Healthy Log Output + +```text +2026-09-26T22:05:00.000Z INFO backfill_operation_index: Starting operation_index backfill +2026-09-26T22:05:00.010Z INFO backfill_operation_index: Database URL: postgres://user:***... +2026-09-26T22:05:00.010Z INFO backfill_operation_index: Batch size: 1000 +2026-09-26T22:05:00.010Z INFO backfill_operation_index: Dry run: false +2026-09-26T22:05:00.010Z INFO backfill_operation_index: Limit: unlimited +2026-09-26T22:05:00.250Z INFO backfill_operation_index: Updated transaction 3fa85f64-...: 0 -> 1 (tx_hash: 7d2b4f...) +2026-09-26T22:05:00.255Z INFO backfill_operation_index: Updated transaction 8ce219a1-...: 0 -> 2 (tx_hash: 7d2b4f...) +... +2026-09-26T22:05:05.100Z INFO backfill_operation_index: Backfill Summary: +2026-09-26T22:05:05.100Z INFO backfill_operation_index: Total examined: 4500 +2026-09-26T22:05:05.100Z INFO backfill_operation_index: Needs update: 120 +2026-09-26T22:05:05.100Z INFO backfill_operation_index: Updated: 120 +2026-09-26T22:05:05.100Z INFO backfill_operation_index: Skipped (already correct): 4380 +2026-09-26T22:05:05.100Z INFO backfill_operation_index: Skipped (invalid TOID): 0 +2026-09-26T22:05:05.100Z INFO backfill_operation_index: Errors: 0 +``` + +## Implementation & Code Reference + +The backfill tool relies on: +- [`octo_ingest::operation_index_from_toid`](file:///c:/Users/DELL/OneDrive/Desktop/drip/Octo-Protocol-6/crates/ingest/src/lib.rs): + Extracts the 0-based operation index from TOID strings formatted as `{ledger}-{tx_index}-{op_index}`. +- Atomic SQL Transaction Updates: + Each candidate batch is updated inside an isolated SQL transaction with an optimistic guard: + ```sql + UPDATE transactions + SET operation_index = $1, updated_at = now() + WHERE id = $2 AND operation_index = $3 AND horizon_op_id = $4; + ``` +- Any rows modified concurrently will log a warning without aborting the batch. + +## Monitoring Progress + +Operators can monitor the remaining backfill volume during execution: +```sql +SELECT + count(*) FILTER (WHERE operation_index = 0 AND split_part(horizon_op_id, '-', 3) <> '0') AS pending_backfill, + count(*) FILTER (WHERE operation_index = split_part(horizon_op_id, '-', 3)::int) AS verified_correct +FROM transactions +WHERE direction = 'deposit' AND horizon_op_id IS NOT NULL; +``` + +## Interruption Recovery + +- **Idempotent**: Rows with matching `operation_index` and TOID component are identified as `already correct` and skipped on subsequent runs. +- **Transactional Batches**: Each batch commits atomically. If the process is terminated mid-execution, previously committed batches remain intact. +- **Recovery Action**: Simply re-execute the binary. It will query from offset or filter candidates and continue until all rows match their TOID index. diff --git a/docs/runbook-migrate-keys.md b/docs/runbook-migrate-keys.md new file mode 100644 index 0000000..42387a6 --- /dev/null +++ b/docs/runbook-migrate-keys.md @@ -0,0 +1,108 @@ +# Operational Runbook: Master Key Rotation (`bin/migrate-keys`) + +## Overview + +The `octo-migrate-keys` binary is an offline, resumable operator tool designed to re-seal HD master seeds under a new master key or cipher scheme without service downtime. + +### When to Run +- Routine cryptographic key rotation (e.g. quarterly or annual key roll). +- Secret incident mitigation (when the existing `MASTER_KEY` may have been exposed). +- Cipher upgrade (re-encrypting stored seeds under updated cipher parameters or schemes, such as `SCHEME_V1`). + +## Architecture & Dual-Key Window + +During rotation, `octo-server` supports a dual-key configuration: +- `MASTER_KEY`: The current/old 32-byte base64-encoded key used to decrypt existing seeds. +- `MASTER_KEY_NEXT`: The target/new 32-byte base64-encoded key. + +While `octo-migrate-keys` is executing, both keys must remain available to running API instances. Each row records its `sealed_scheme`, allowing decryption routines to identify the applicable key. Once `octo-migrate-keys` finishes with 0 remaining rows, `MASTER_KEY_NEXT` can be promoted to `MASTER_KEY` and the old key decommissioned. + +## Pre-flight Checks + +1. **Database Connectivity and Backup**: + Verify access to the production PostgreSQL cluster and take a snapshot: + ```bash + pg_dump -Fc "$DATABASE_URL" > "octo_backup_$(date +%Y%m%d_%H%M%S).dump" + ``` +2. **Key Material Validation**: + Ensure keys are valid 32-byte base64 strings: + ```bash + [ "$(echo -n "$MASTER_KEY" | base64 -d | wc -c)" -eq 32 ] || echo "Invalid MASTER_KEY length" + [ "$(echo -n "$MASTER_KEY_NEXT" | base64 -d | wc -c)" -eq 32 ] || echo "Invalid MASTER_KEY_NEXT length" + ``` +3. **Database Migration Status**: + Ensure the database schema is up-to-date (`Store::migrate` is also executed at binary startup). +4. **Current Unmigrated Row Count**: + Query rows currently requiring migration: + ```sql + SELECT sealed_scheme, count(*) + FROM wallets + WHERE sealed_ciphertext IS NOT NULL + GROUP BY sealed_scheme; + ``` + +## Invocations + +### Full Key Rotation +```bash +MASTER_KEY="" \ +MASTER_KEY_NEXT="" \ +DATABASE_URL="postgres://user:pass@localhost:5432/octo" \ +cargo run --release -p octo-migrate-keys -- --batch-size 100 +``` + +### Cipher Upgrade Only (Same Key) +If `MASTER_KEY_NEXT` is omitted, the tool defaults to re-sealing under `MASTER_KEY`: +```bash +MASTER_KEY="" \ +DATABASE_URL="postgres://user:pass@localhost:5432/octo" \ +cargo run --release -p octo-migrate-keys -- --batch-size 100 +``` + +## Expected Healthy Log Output + +```text +2026-09-26T22:00:00.000Z INFO octo_migrate_keys: batch_size=100 same_key=false octo-migrate-keys starting +2026-09-26T22:00:00.150Z INFO octo_migrate_keys: batch_len=100 after_id=None processing batch +2026-09-26T22:00:00.420Z DEBUG octo_migrate_keys: wallet_id=9d14... migrated +... +2026-09-26T22:00:01.200Z INFO octo_migrate_keys: batch_len=42 after_id=Some(a4f1...) processing batch +2026-09-26T22:00:01.350Z INFO octo_migrate_keys: total_migrated=142 total_skipped=0 migration complete — 0 wallets remaining on old scheme +``` + +## Store Method Implementation Reference + +The migration tool is backed by two primary methods in [`crates/store`](file:///c:/Users/DELL/OneDrive/Desktop/drip/Octo-Protocol-6/crates/store): +- [`Store::list_wallets_needing_reseal`](file:///c:/Users/DELL/OneDrive/Desktop/drip/Octo-Protocol-6/crates/store/src/wallets.rs): + Fetches batches of wallets where `sealed_scheme != target_scheme` ordered by `id ASC`, paginating via `after_id`. +- [`Store::reseal_wallet`](file:///c:/Users/DELL/OneDrive/Desktop/drip/Octo-Protocol-6/crates/store/src/wallets.rs): + Atomically updates `sealed_ciphertext`, `sealed_nonce`, `sealed_salt`, and `sealed_scheme` with an optimistic concurrency guard (`WHERE id = $1 AND sealed_scheme = $expected_old_scheme`). +- Cryptographic re-encryption is executed via `octo_crypto::reseal` in [`crates/crypto`](file:///c:/Users/DELL/OneDrive/Desktop/drip/Octo-Protocol-6/crates/crypto). + +## Monitoring Progress + +Operators can monitor ongoing execution in another terminal: +```sql +SELECT + count(*) FILTER (WHERE sealed_scheme = 1) AS migrated_v1, + count(*) FILTER (WHERE sealed_scheme != 1) AS remaining_old +FROM wallets +WHERE sealed_ciphertext IS NOT NULL; +``` + +## Interruption Recovery & Rollback + +### Resuming an Interrupted Run +- **Safe Interruption**: The tool operates using atomic per-wallet updates and paginated batches. +- If stopped (SIGINT, network timeout, process termination), simply re-run the same command. +- Wallets already migrated will have `sealed_scheme == target_scheme` and will not be re-processed by `Store::list_wallets_needing_reseal`. + +### Rollback Procedure +- If rotation needs to be reverted before decommissioning the old key: + Swap the values of `MASTER_KEY` and `MASTER_KEY_NEXT`: + ```bash + MASTER_KEY="" \ + MASTER_KEY_NEXT="" \ + DATABASE_URL="postgres://user:pass@localhost:5432/octo" \ + cargo run --release -p octo-migrate-keys -- --batch-size 100 + ```