diff --git a/.env.example b/.env.example index 1d7367c..2233ec6 100644 --- a/.env.example +++ b/.env.example @@ -48,3 +48,8 @@ EMAIL_FROM_ADDRESS= # How often the deposit ingest supervisor polls Horizon for all wallets, and the page size. INGEST_INTERVAL_SECS=5 INGEST_PAGE_LIMIT=50 + +# --- Graceful shutdown --- +# Max seconds to drain in-flight HTTP requests and the current ingest tick after SIGTERM/SIGINT +# before force-exiting. Keep below your orchestrator's kill deadline (k8s default: 30s). +SHUTDOWN_DRAIN_TIMEOUT_SECS=25 diff --git a/Cargo.lock b/Cargo.lock index e88c3c7..a80d0a7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1797,6 +1797,7 @@ dependencies = [ "sqlx", "thiserror 1.0.69", "tokio", + "tokio-util", "tracing", "uuid", "wiremock", @@ -1809,8 +1810,10 @@ dependencies = [ "anyhow", "base64", "dotenvy", + "hex", "octo-crypto", "octo-store", + "sha2", "tokio", "tracing", "tracing-subscriber", @@ -1840,6 +1843,7 @@ dependencies = [ "octo-wallet-core", "octo-webhooks", "tokio", + "tokio-util", "tracing", "tracing-subscriber", ] @@ -1853,6 +1857,7 @@ dependencies = [ "serde", "serde_json", "sqlx", + "subtle", "thiserror 1.0.69", "tokio", "uuid", diff --git a/Cargo.toml b/Cargo.toml index f8fb362..4bce4ea 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -39,11 +39,13 @@ sha2 = "0.10.9" hex = "0.4.3" base64 = "0.22" argon2 = "0.5" +subtle = "2.6" # Note: the MSRV-aware resolver (.cargo/config.toml: incompatible-rust-versions = "fallback") # keeps the whole tree on Rust-1.84-compatible versions, so no manual transitive pins are needed. # --- async runtime / web --- tokio = { version = "1", features = ["full"] } +tokio-util = "0.7" axum = "0.7" tower = "0.5" tower-http = { version = "0.6", features = ["trace", "cors", "limit"] } diff --git a/README.md b/README.md index d8f4684..44eb3ee 100644 --- a/README.md +++ b/README.md @@ -210,6 +210,32 @@ 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. +## Deployment + +`octo-server` runs the REST API and the deposit ingest worker in one process and is safe to +roll (Kubernetes, ECS, systemd) behind a load balancer. + +### Graceful shutdown + +On `SIGTERM` (or `SIGINT`) the server: + +1. logs `shutdown signal received` and stops accepting new connections; +2. logs `draining …` and lets in-flight HTTP requests complete, while the ingest supervisor + finishes its **current tick** (never aborting a page mid-processing) and then stops; +3. logs `drained` then `exiting` — or, if the drain outlasts the timeout, logs + `drain timeout elapsed; forcing exit` and exits anyway. A page cut short this way is safe: + the ingest cursor is only advanced after processing and deposit inserts are idempotent, so + the next instance re-runs it without double-crediting. + +| Variable | Default | Description | +|---|---|---| +| `SHUTDOWN_DRAIN_TIMEOUT_SECS` | `25` | Max seconds to drain before force-exiting | + +**Tuning:** keep `SHUTDOWN_DRAIN_TIMEOUT_SECS` a few seconds *below* the orchestrator's hard-kill +deadline (Kubernetes `terminationGracePeriodSeconds`, default 30; ECS `stopTimeout`, default +30), so the process exits on its own terms rather than being `SIGKILL`ed mid-drain. If ticks +routinely run long (many wallets, slow Horizon), raise both values together. + ## Roadmap - **Gas sponsorship** — *shipped.* App developers can sponsor their users' Stellar transactions diff --git a/bin/migrate-keys/Cargo.toml b/bin/migrate-keys/Cargo.toml index 60091ea..7578cd4 100644 --- a/bin/migrate-keys/Cargo.toml +++ b/bin/migrate-keys/Cargo.toml @@ -21,4 +21,6 @@ tracing.workspace = true tracing-subscriber.workspace = true base64.workspace = true uuid.workspace = true +sha2.workspace = true +hex.workspace = true dotenvy = "0.15" diff --git a/bin/migrate-keys/src/main.rs b/bin/migrate-keys/src/main.rs index 4631a3c..c69f2d5 100644 --- a/bin/migrate-keys/src/main.rs +++ b/bin/migrate-keys/src/main.rs @@ -48,6 +48,22 @@ //! MASTER_KEY= \ //! cargo run -p octo-migrate-keys -- --batch-size 100 //! ``` +//! +//! ## Checkpoint (resuming an interrupted run) +//! +//! After every fully-processed batch the tool writes the last wallet id it handled to a +//! checkpoint file, so a crash or Ctrl-C resumes from the last completed batch instead of +//! re-scanning the whole table. +//! +//! - **Location:** `migrate-keys.checkpoint` in the working directory, or the path in +//! `MIGRATE_KEYS_CHECKPOINT`. +//! - **Contents:** a fingerprint of the `(MASTER_KEY, MASTER_KEY_NEXT)` pair (a SHA-256 over +//! the keys — no key material) and the `after_id` cursor. A checkpoint written for a +//! different key pair is refused rather than silently skipping rows. +//! - **Lifecycle:** removed automatically on clean completion. +//! - **Forcing a full re-run:** stop the tool, then delete the file (`rm migrate-keys.checkpoint`). +//! This is always safe — the idempotency guard below makes re-scanned rows no-ops. Do not +//! hand-edit the cursor: moving it forward skips rows that were never migrated. #![forbid(unsafe_code)] @@ -55,11 +71,16 @@ use anyhow::{Context, Result}; use base64::Engine; use octo_crypto::{master_key_from_slice, reseal, MASTER_KEY_LEN, SCHEME_V1}; use octo_store::Store; +use sha2::{Digest, Sha256}; +use std::path::{Path, PathBuf}; use uuid::Uuid; /// Maximum rows per batch (hard cap, configurable via CLI). const DEFAULT_BATCH_SIZE: i64 = 100; +/// Checkpoint file used when `MIGRATE_KEYS_CHECKPOINT` is unset. +const DEFAULT_CHECKPOINT_PATH: &str = "migrate-keys.checkpoint"; + #[tokio::main] async fn main() -> Result<()> { let _ = dotenvy::dotenv(); @@ -78,7 +99,20 @@ async fn main() -> Result<()> { .context("connect to database")?; store.migrate().await.context("run migrations")?; - let mut after_id: Option = None; + let fingerprint = key_pair_fingerprint(&cfg.old_key, &cfg.new_key); + let mut after_id = read_checkpoint(&cfg.checkpoint_path, &fingerprint)?; + match after_id { + Some(id) => tracing::info!( + checkpoint = %cfg.checkpoint_path.display(), + after_id = %id, + "resuming from checkpoint" + ), + None => tracing::info!( + checkpoint = %cfg.checkpoint_path.display(), + "no checkpoint found; starting from the beginning" + ), + } + let mut batches_completed = 0usize; let mut total_migrated = 0usize; let mut total_skipped = 0usize; @@ -151,23 +185,106 @@ async fn main() -> Result<()> { } } - // Advance the cursor to the last wallet in this batch (ids are ordered ASC). + // Advance the cursor to the last wallet in this batch (ids are ordered ASC), and persist + // it only now that every row in the batch is done, so a resume never skips a row. after_id = batch.last().map(|w| w.id); + if let Some(id) = after_id { + write_checkpoint(&cfg.checkpoint_path, &fingerprint, id)?; + } + batches_completed += 1; + tracing::info!( + batches_completed, + total_migrated, + total_skipped, + after_id = ?after_id, + "batch complete" + ); } + remove_checkpoint(&cfg.checkpoint_path)?; tracing::info!( + batches_completed, total_migrated, total_skipped, - "migration complete — 0 wallets remaining on old scheme" + "migration complete — 0 wallets remaining on old scheme; checkpoint removed" ); Ok(()) } +/// Identify the key pair a checkpoint belongs to without writing any key material to disk. +fn key_pair_fingerprint(old_key: &[u8; MASTER_KEY_LEN], new_key: &[u8; MASTER_KEY_LEN]) -> String { + let mut h = Sha256::new(); + h.update(b"octo-migrate-keys/checkpoint/v1"); + h.update(old_key); + h.update(new_key); + hex::encode(h.finalize()) +} + +/// Read the resume cursor, refusing a malformed checkpoint or one from a different key pair. +fn read_checkpoint(path: &Path, fingerprint: &str) -> Result> { + let contents = match std::fs::read_to_string(path) { + Ok(c) => c, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None), + Err(e) => { + return Err(e).with_context(|| format!("read checkpoint {}", path.display())); + } + }; + let mut stored_fingerprint = None; + let mut stored_after_id = None; + for line in contents.lines() { + match line.split_once('=') { + Some(("fingerprint", v)) => stored_fingerprint = Some(v.trim()), + Some(("after_id", v)) => stored_after_id = Some(v.trim()), + _ => {} + } + } + let (Some(stored_fingerprint), Some(stored_after_id)) = (stored_fingerprint, stored_after_id) + else { + anyhow::bail!( + "checkpoint {} is malformed; delete it to restart from the beginning", + path.display() + ); + }; + if stored_fingerprint != fingerprint { + anyhow::bail!( + "checkpoint {} was written for a different MASTER_KEY/MASTER_KEY_NEXT pair; \ + delete it to restart from the beginning", + path.display() + ); + } + let id = Uuid::parse_str(stored_after_id).with_context(|| { + format!( + "checkpoint {} has an invalid after_id; delete it to restart from the beginning", + path.display() + ) + })?; + Ok(Some(id)) +} + +/// Persist the cursor via write-then-rename so a crash mid-write never leaves a torn file. +fn write_checkpoint(path: &Path, fingerprint: &str, after_id: Uuid) -> Result<()> { + let tmp = path.with_extension("checkpoint.tmp"); + std::fs::write(&tmp, format!("fingerprint={fingerprint}\nafter_id={after_id}\n")) + .with_context(|| format!("write checkpoint {}", tmp.display()))?; + std::fs::rename(&tmp, path).with_context(|| format!("replace checkpoint {}", path.display())) +} + +/// Delete the checkpoint after a clean run; a missing file is fine. +fn remove_checkpoint(path: &Path) -> Result<()> { + match std::fs::remove_file(path) { + Err(e) if e.kind() != std::io::ErrorKind::NotFound => { + Err(e).with_context(|| format!("remove checkpoint {}", path.display())) + } + _ => Ok(()), + } +} + struct Config { database_url: String, old_key: [u8; MASTER_KEY_LEN], new_key: [u8; MASTER_KEY_LEN], batch_size: i64, + checkpoint_path: PathBuf, } impl Config { @@ -192,12 +309,21 @@ impl Config { .nth(1) .and_then(|v| v.parse::().ok()) .unwrap_or(DEFAULT_BATCH_SIZE); + // LIMIT 0 returns an empty page, which the loop would misread as "migration complete". + if batch_size <= 0 { + anyhow::bail!("--batch-size must be a positive integer"); + } + + let checkpoint_path = std::env::var("MIGRATE_KEYS_CHECKPOINT") + .map(PathBuf::from) + .unwrap_or_else(|_| PathBuf::from(DEFAULT_CHECKPOINT_PATH)); Ok(Config { database_url, old_key, new_key, batch_size, + checkpoint_path, }) } } diff --git a/bin/server/Cargo.toml b/bin/server/Cargo.toml index 8b1d88c..a1c6de4 100644 --- a/bin/server/Cargo.toml +++ b/bin/server/Cargo.toml @@ -22,6 +22,7 @@ octo-email.workspace = true octo-wallet-core.workspace = true octo-resilience.workspace = true tokio.workspace = true +tokio-util.workspace = true anyhow.workspace = true axum.workspace = true tracing.workspace = true diff --git a/bin/server/src/main.rs b/bin/server/src/main.rs index ac5a8bc..3e93649 100644 --- a/bin/server/src/main.rs +++ b/bin/server/src/main.rs @@ -13,7 +13,9 @@ use octo_resilience::ResilienceConfig; use octo_store::Store; use octo_wallet_core::StellarNetwork; use octo_webhooks::WebhookSender; +use std::future::IntoFuture; use std::time::Duration; +use tokio_util::sync::CancellationToken; #[tokio::main] async fn main() -> Result<()> { @@ -76,14 +78,13 @@ async fn main() -> Result<()> { ingest_retry, ingest_circuit, ); - tokio::spawn(async move { - supervisor - .run( - Duration::from_secs(cfg.ingest_interval_secs), - cfg.ingest_page_limit, - ) - .await; - }); + // Cancelled on SIGTERM/SIGINT; both the ingest loop and the HTTP server drain on it. + let shutdown = CancellationToken::new(); + let ingest = tokio::spawn(supervisor.run_until_cancelled( + Duration::from_secs(cfg.ingest_interval_secs), + cfg.ingest_page_limit, + shutdown.clone(), + )); tracing::info!( interval_secs = cfg.ingest_interval_secs, "deposit ingest supervisor started" @@ -95,17 +96,94 @@ async fn main() -> Result<()> { .await .with_context(|| format!("bind {}", cfg.bind_addr))?; tracing::info!(addr = %cfg.bind_addr, "API listening"); + // Graceful shutdown stops accepting new connections and lets in-flight requests finish. // `into_make_service_with_connect_info` is what makes the peer address available to the // rate limiter's `ConnectInfo` extractor; without it every caller looks like one client. - axum::serve( - listener, - app.into_make_service_with_connect_info::(), - ) - .await - .context("serve API")?; + let mut server = tokio::spawn( + axum::serve( + listener, + app.into_make_service_with_connect_info::(), + ) + .with_graceful_shutdown(shutdown.clone().cancelled_owned()) + .into_future(), + ); + + // A server that dies on its own (not via a signal) is a hard failure, not a shutdown. + tokio::select! { + res = &mut server => { + shutdown.cancel(); + return res.context("API task panicked")?.context("serve API"); + } + signal = shutdown_signal() => { + tracing::info!(signal, "shutdown signal received"); + } + } + + shutdown.cancel(); + tracing::info!( + timeout_secs = cfg.shutdown_drain_timeout.as_secs(), + "draining in-flight HTTP requests and the current ingest tick" + ); + let ingest_abort = ingest.abort_handle(); + let server_abort = server.abort_handle(); + let drain = async { tokio::join!(server, ingest) }; + match tokio::time::timeout(cfg.shutdown_drain_timeout, drain).await { + Ok((http, ingest)) => { + let http = http + .context("API task panicked") + .and_then(|r| r.context("serve API")); + if let Err(e) = http { + tracing::error!(error = ?e, "API server errored while draining"); + } + if let Err(e) = ingest { + tracing::error!(error = ?e, "ingest supervisor task panicked while draining"); + } + tracing::info!("drained"); + } + Err(_) => { + // Deposit inserts are deduplicated, so a page cut short here re-runs safely. + tracing::warn!( + timeout_secs = cfg.shutdown_drain_timeout.as_secs(), + "drain timeout elapsed; forcing exit" + ); + ingest_abort.abort(); + server_abort.abort(); + } + } + tracing::info!("exiting"); Ok(()) } +/// Resolve on SIGTERM (what Kubernetes/ECS send on a rolling deploy) or SIGINT (Ctrl-C). +async fn shutdown_signal() -> &'static str { + let ctrl_c = async { + if let Err(e) = tokio::signal::ctrl_c().await { + tracing::error!(error = ?e, "failed to listen for SIGINT"); + std::future::pending::<()>().await; + } + }; + #[cfg(unix)] + let terminate = async { + use tokio::signal::unix::{signal, SignalKind}; + match signal(SignalKind::terminate()) { + Ok(mut sig) => { + sig.recv().await; + } + Err(e) => { + tracing::error!(error = ?e, "failed to listen for SIGTERM"); + std::future::pending::<()>().await; + } + } + }; + #[cfg(not(unix))] + let terminate = std::future::pending::<()>(); + + tokio::select! { + () = ctrl_c => "SIGINT", + () = terminate => "SIGTERM", + } +} + fn init_tracing() { let filter = std::env::var("RUST_LOG").unwrap_or_else(|_| "info,octo=debug".to_string()); tracing_subscriber::fmt() @@ -134,6 +212,10 @@ struct Config { bind_addr: String, ingest_interval_secs: u64, ingest_page_limit: u32, + /// Upper bound on the graceful-shutdown drain (in-flight HTTP requests + the current ingest + /// tick) before the process force-exits. `SHUTDOWN_DRAIN_TIMEOUT_SECS`, default 25 — keep it + /// below the orchestrator's kill deadline (Kubernetes `terminationGracePeriodSeconds`: 30). + shutdown_drain_timeout: Duration, /// Resilience settings for all Horizon clients (API + ingest). /// /// | Variable | Default | Description | @@ -201,6 +283,13 @@ impl Config { .and_then(|s| s.parse().ok()) .unwrap_or(50); + let shutdown_drain_timeout = Duration::from_secs( + std::env::var("SHUTDOWN_DRAIN_TIMEOUT_SECS") + .ok() + .and_then(|s| s.parse().ok()) + .unwrap_or(25), + ); + let resilience = ResilienceConfig::from_env(); Ok(Config { @@ -217,6 +306,7 @@ impl Config { bind_addr, ingest_interval_secs, ingest_page_limit, + shutdown_drain_timeout, resilience, }) } diff --git a/crates/email/src/templates.rs b/crates/email/src/templates.rs index 11dd603..5a3ba6e 100644 --- a/crates/email/src/templates.rs +++ b/crates/email/src/templates.rs @@ -5,6 +5,13 @@ //! data:image/svg+xml (and often data: images generally), so anything shown here has to be //! fetchable. The logo lives on Cloudinary; the small icon set is served from octohq.org/email/ //! (Octo-frontend's public/email/). +//! +//! ## Escaping convention +//! +//! Every value that originates outside this server — user input (email address, labels), a +//! user-signed transaction (asset code, destination), or a Horizon response (tx hash, error +//! detail) — MUST go through [`html_escape`] before it is interpolated into HTML. Only values +//! the server generates itself (the numeric OTP, the formatted amount, fixed copy) skip it. const BURGUNDY: &str = "#7b1733"; const BURGUNDY_BRIGHT: &str = "#b81f4d"; @@ -53,6 +60,23 @@ fn socials() -> String { .join("") } +/// Escape `&`, `<`, `>`, `"` and `'` so an untrusted value renders as text in HTML body and +/// attribute contexts, never as markup. +pub fn html_escape(s: &str) -> String { + let mut out = String::with_capacity(s.len()); + for c in s.chars() { + match c { + '&' => out.push_str("&"), + '<' => out.push_str("<"), + '>' => out.push_str(">"), + '"' => out.push_str("""), + '\'' => out.push_str("'"), + _ => out.push(c), + } + } + out +} + fn shell(body: &str, icon: Icon) -> String { let icon_url = format!("{ICON_BASE}/{}", icon_url(icon)); let socials = socials(); @@ -99,6 +123,8 @@ pub fn otp_email(code: &str, purpose: &str) -> String { /// Sent once, right after signup verification succeeds. pub fn welcome_email(email: &str) -> String { + // The address is user-supplied at signup; RFC 5322 local parts may legally contain `<"'&`. + let email = html_escape(email); shell( &format!( "

Welcome to Octo 🎉

\ @@ -116,6 +142,12 @@ pub fn withdrawal_success_email( destination: &str, tx_hash: &str, ) -> String { + // Asset and destination come from the user-signed XDR, the hash from Horizon's response. + let (asset, destination, tx_hash) = ( + html_escape(asset), + html_escape(destination), + html_escape(tx_hash), + ); shell( &format!( "

Withdrawal confirmed

\ @@ -137,6 +169,12 @@ pub fn withdrawal_failed_email( destination: &str, reason: &str, ) -> String { + // Asset and destination come from the user-signed XDR, the reason may echo Horizon's detail. + let (asset, destination, reason) = ( + html_escape(asset), + html_escape(destination), + html_escape(reason), + ); shell( &format!( "

Withdrawal attempt failed

\ diff --git a/crates/ingest/Cargo.toml b/crates/ingest/Cargo.toml index 612e5c4..73921df 100644 --- a/crates/ingest/Cargo.toml +++ b/crates/ingest/Cargo.toml @@ -15,6 +15,7 @@ octo-webhooks.workspace = true octo-resilience.workspace = true reqwest.workspace = true tokio.workspace = true +tokio-util.workspace = true serde.workspace = true serde_json.workspace = true thiserror.workspace = true diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index cf94637..60a796c 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -26,6 +26,7 @@ use octo_webhooks::{Event, WebhookSender}; use std::collections::HashMap; use std::sync::{Arc, Mutex}; use std::time::Duration; +use tokio_util::sync::CancellationToken; use uuid::Uuid; /// The widely-used testnet USDC issuer — must match `crates/api/src/routes/payment_links.rs`'s @@ -648,12 +649,31 @@ impl Supervisor { /// Run forever: every `interval`, poll all wallets on this network once. pub async fn run(self, interval: Duration, page_limit: u32) { - loop { + self.run_until_cancelled(interval, page_limit, CancellationToken::new()).await; + } + + /// Like [`Supervisor::run`], but returns once `shutdown` is cancelled. + /// + /// Cancellation is only observed *between* ticks: an in-progress tick always finishes its + /// current page (deposit rows + cursor), so a rolling deploy never aborts a page halfway. + /// Bounding how long that may take is the caller's job (see `bin/server`'s drain timeout). + pub async fn run_until_cancelled( + self, + interval: Duration, + page_limit: u32, + shutdown: CancellationToken, + ) { + while !shutdown.is_cancelled() { if let Err(e) = self.tick(page_limit).await { tracing::warn!(error = ?e, "ingest supervisor tick failed; will retry"); } - tokio::time::sleep(interval).await; + // Wake early on shutdown instead of sleeping out the full interval. + tokio::select! { + () = tokio::time::sleep(interval) => {} + () = shutdown.cancelled() => {} + } } + tracing::info!("ingest supervisor stopped after finishing its current tick"); } /// One supervision pass: poll every wallet on this network, draining each wallet's backlog diff --git a/crates/store/Cargo.toml b/crates/store/Cargo.toml index 8c8f48b..d86adae 100644 --- a/crates/store/Cargo.toml +++ b/crates/store/Cargo.toml @@ -15,6 +15,7 @@ serde_json.workspace = true uuid.workspace = true chrono.workspace = true thiserror.workspace = true +subtle.workspace = true [dev-dependencies] tokio.workspace = true diff --git a/crates/store/src/lib.rs b/crates/store/src/lib.rs index b2901b7..9c15b71 100644 --- a/crates/store/src/lib.rs +++ b/crates/store/src/lib.rs @@ -32,6 +32,20 @@ use uuid::Uuid; /// Embedded migrations, applied by [`Store::migrate`]. pub static MIGRATOR: sqlx::migrate::Migrator = sqlx::migrate!("./migrations"); +/// Compare a stored OTP hash with a candidate in constant time. +/// +/// Both sides are already hashes, but they are derived from a secret code drawn from a small +/// (6-digit) space, so a short-circuiting `!=` would leak how many leading bytes matched — a +/// signal an attacker could combine with precomputed code→hash tables. `subtle::ConstantTimeEq` +/// is the same primitive `hmac`'s `verify_slice` uses for the JWT and webhook signature checks. +/// Only the length check short-circuits, and hash length is public (always 64 hex chars). +/// The constant-time property is not unit-tested: timing tests are inherently flaky, so we rely +/// on using a well-reviewed primitive correctly instead. +fn otp_hash_matches(stored: &str, candidate: &str) -> bool { + use subtle::ConstantTimeEq; + stored.as_bytes().ct_eq(candidate.as_bytes()).into() +} + /// Hard ceiling on any list query's `LIMIT`, well above the API's max page (200 + 1 look-ahead). pub const MAX_LIST_LIMIT: i64 = 1000; @@ -39,6 +53,7 @@ pub const MAX_LIST_LIMIT: i64 = 1000; fn clamp_limit(limit: i64) -> i64 { limit.clamp(0, MAX_LIST_LIMIT) } +} /// A handle to the database (cloneable; wraps a connection pool). #[derive(Clone)] @@ -255,7 +270,7 @@ impl Store { { return Err(StoreError::InvalidOtp); } - if otp.code_hash != code_hash { + if !otp_hash_matches(&otp.code_hash, code_hash) { sqlx::query("UPDATE email_otps SET attempts = attempts + 1 WHERE id = $1") .bind(otp.id) .execute(&self.pool)