Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -212,6 +212,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)

## Deployment

`octo-server` runs the REST API and the deposit ingest worker in one process and is safe to
Expand Down
129 changes: 127 additions & 2 deletions crates/api/src/horizon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Vec<Balance>, ApiError> {
let span = tracing::Span::current();
let url = format!(
"{}/accounts/{}",
self.base_url.trim_end_matches('/'),
Expand Down Expand Up @@ -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<i64, ApiError> {
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<AccountInfo, ApiError> {
// 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('/'),
Expand Down Expand Up @@ -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),
Expand All @@ -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<SubmitResult, ApiError> {
// 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();
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -460,8 +523,19 @@ fn map_result<T>(r: Result<T, ResilienceError<FetchError>>) -> Result<T, ApiErro

/// Fund a testnet account via friendbot. Best-effort; a single retry is safe because friendbot
/// is idempotent (re-funding an already-funded account is a no-op on testnet).
#[tracing::instrument(
name = "horizon_friendbot_fund",
skip(friendbot_url),
fields(call_type = "friendbot", account_g = %account_g, outcome = tracing::field::Empty)
)]
pub async fn friendbot_fund(friendbot_url: &str, account_g: &str) -> 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(
Expand Down Expand Up @@ -542,4 +616,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::<String>::new()));
let spans_clone = recorded_spans.clone();

struct SpanRecorder(Arc<Mutex<Vec<String>>>);
impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> 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()));
}
}
53 changes: 53 additions & 0 deletions crates/email/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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}"
);
}
}
1 change: 1 addition & 0 deletions crates/ingest/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ uuid.workspace = true

[dev-dependencies]
tokio.workspace = true
tracing-subscriber.workspace = true
dotenvy = "0.15"
axum.workspace = true
octo-webhooks.workspace = true
Expand Down
56 changes: 56 additions & 0 deletions crates/ingest/src/horizon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Vec<PaymentRecord>, HorizonError> {
let span = tracing::Span::current();
let mut url = format!(
"{}/accounts/{}/payments?order=asc&limit={}&join=transactions",
self.base_url.trim_end_matches('/'),
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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::<String>::new()));
let spans_clone = recorded_spans.clone();

struct SpanRecorder(Arc<Mutex<Vec<String>>>);
impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> 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()));
}
}
30 changes: 30 additions & 0 deletions crates/ingest/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1095,4 +1095,34 @@ mod tests {
// Ledger i32::MAX is representable but implausible.
assert_eq!(operation_index_from_toid("9223372032559812609"), 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);
}
}
}
}
Loading
Loading