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
5 changes: 5 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
5 changes: 5 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"] }
Expand Down
26 changes: 26 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions bin/migrate-keys/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
132 changes: 129 additions & 3 deletions bin/migrate-keys/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,18 +48,39 @@
//! MASTER_KEY=<b64> \
//! 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)]

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();
Expand All @@ -78,7 +99,20 @@ async fn main() -> Result<()> {
.context("connect to database")?;
store.migrate().await.context("run migrations")?;

let mut after_id: Option<Uuid> = 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;

Expand Down Expand Up @@ -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<Option<Uuid>> {
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 {
Expand All @@ -192,12 +309,21 @@ impl Config {
.nth(1)
.and_then(|v| v.parse::<i64>().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,
})
}
}
Expand Down
1 change: 1 addition & 0 deletions bin/server/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading