diff --git a/.containerignore b/.containerignore new file mode 100644 index 0000000..3a1b4f6 --- /dev/null +++ b/.containerignore @@ -0,0 +1,4 @@ +target/ +.git/ +.jj/ +*.tar diff --git a/Cargo.lock b/Cargo.lock index 7577dcd..54ff9d1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -51,6 +51,56 @@ dependencies = [ "libc", ] +[[package]] +name = "anstream" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "824a212faf96e9acacdbd09febd34438f8f711fb84e09a8916013cd7815ca28d" +dependencies = [ + "anstyle", + "anstyle-parse", + "anstyle-query", + "anstyle-wincon", + "colorchoice", + "is_terminal_polyfill", + "utf8parse", +] + +[[package]] +name = "anstyle" +version = "1.0.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "940b3a0ca603d1eade50a4846a2afffd5ef57a9feac2c0e2ec2e14f9ead76000" + +[[package]] +name = "anstyle-parse" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52ce7f38b242319f7cabaa6813055467063ecdc9d355bbb4ce0c68908cd8130e" +dependencies = [ + "utf8parse", +] + +[[package]] +name = "anstyle-query" +version = "1.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc" +dependencies = [ + "windows-sys 0.61.2", +] + +[[package]] +name = "anstyle-wincon" +version = "3.0.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d" +dependencies = [ + "anstyle", + "once_cell_polyfill", + "windows-sys 0.61.2", +] + [[package]] name = "anyhow" version = "1.0.102" @@ -242,12 +292,15 @@ dependencies = [ "bobbin-record-lru", "bobbin-search", "bobbin-slingshot-client", - "bobbin-types", "bobbin-xrpc", - "bytes", - "jacquard-common", - "serde_json", + "clap", + "confique", + "futures", + "socket2", + "thiserror 2.0.18", "tokio", + "tokio-util", + "toml", "tracing", "tracing-subscriber", "url", @@ -272,17 +325,22 @@ version = "0.0.1" dependencies = [ "bobbin-edge-index", "bobbin-record-lru", + "bobbin-slingshot-client", "bobbin-types", + "bytes", "futures", "jacquard-common", + "scc", "serde", "serde_json", "thiserror 2.0.18", "tokio", "tokio-tungstenite 0.29.0", + "tokio-util", "tracing", "tracing-subscriber", "url", + "wiremock", ] [[package]] @@ -377,6 +435,8 @@ dependencies = [ "thiserror 2.0.18", "tokio", "tower", + "tower-http", + "tracing", "url", "wiremock", ] @@ -542,6 +602,46 @@ dependencies = [ "unsigned-varint", ] +[[package]] +name = "clap" +version = "4.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ddb117e43bbf7dacf0a4190fef4d345b9bad68dfc649cb349e7d17d28428e51" +dependencies = [ + "clap_builder", + "clap_derive", +] + +[[package]] +name = "clap_builder" +version = "4.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "714a53001bf66416adb0e2ef5ac857140e7dc3a0c48fb28b2f10762fc4b5069f" +dependencies = [ + "anstream", + "anstyle", + "clap_lex", + "strsim", +] + +[[package]] +name = "clap_derive" +version = "4.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2ce8604710f6733aa641a2b3731eaa1e8b3d9973d5e3565da11800813f997a9" +dependencies = [ + "heck 0.5.0", + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "clap_lex" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9" + [[package]] name = "cobs" version = "0.3.0" @@ -551,6 +651,12 @@ dependencies = [ "thiserror 2.0.18", ] +[[package]] +name = "colorchoice" +version = "1.0.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570" + [[package]] name = "compression-codecs" version = "0.4.38" @@ -568,6 +674,29 @@ version = "0.4.32" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cc14f565cf027a105f7a44ccf9e5b424348421a1d8952a8fc9d499d313107789" +[[package]] +name = "confique" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "06b4f5ec222421e22bb0a8cbaa36b1d2b50fd45cdd30c915ded34108da78b29f" +dependencies = [ + "confique-macro", + "serde", + "toml", +] + +[[package]] +name = "confique-macro" +version = "0.0.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e4d1754680cd218e7bcb4c960cc9bae3444b5197d64563dccccfdf83cab9e1a7" +dependencies = [ + "heck 0.5.0", + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "const-oid" version = "0.9.6" @@ -1667,6 +1796,12 @@ dependencies = [ "serde", ] +[[package]] +name = "is_terminal_polyfill" +version = "1.70.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" + [[package]] name = "itertools" version = "0.14.0" @@ -2128,6 +2263,12 @@ version = "1.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" +[[package]] +name = "once_cell_polyfill" +version = "1.70.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "384b8ab6d37215f3c5301a95a4accb5d64aa607f1fcb26a11b5303878451b4fe" + [[package]] name = "oneshot" version = "0.1.13" @@ -3024,6 +3165,15 @@ dependencies = [ "syn", ] +[[package]] +name = "serde_spanned" +version = "1.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6662b5879511e06e8999a8a235d848113e942c9124f211511b16466ee2995f26" +dependencies = [ + "serde_core", +] + [[package]] name = "serde_urlencoded" version = "0.7.1" @@ -3429,7 +3579,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.2", "once_cell", "rustix", "windows-sys 0.61.2", @@ -3642,6 +3792,45 @@ dependencies = [ "tokio", ] +[[package]] +name = "toml" +version = "0.9.12+spec-1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cf92845e79fc2e2def6a5d828f0801e29a2f8acc037becc5ab08595c7d5e9863" +dependencies = [ + "indexmap", + "serde_core", + "serde_spanned", + "toml_datetime", + "toml_parser", + "toml_writer", + "winnow 0.7.15", +] + +[[package]] +name = "toml_datetime" +version = "0.7.5+spec-1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "92e1cfed4a3038bc5a127e35a2d360f145e1f4b971b551a2ba5fd7aedf7e1347" +dependencies = [ + "serde_core", +] + +[[package]] +name = "toml_parser" +version = "1.1.2+spec-1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a2abe9b86193656635d2411dc43050282ca48aa31c2451210f4202550afb7526" +dependencies = [ + "winnow 1.0.2", +] + +[[package]] +name = "toml_writer" +version = "1.1.1+spec-1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "756daf9b1013ebe47a8776667b466417e2d4c5679d441c26230efd9ef78692db" + [[package]] name = "tower" version = "0.5.3" @@ -3679,6 +3868,7 @@ dependencies = [ "tower", "tower-layer", "tower-service", + "tracing", ] [[package]] @@ -3737,6 +3927,16 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-serde" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "704b1aeb7be0d0a84fc9828cae51dab5970fee5088f83d1dd7ee6f6246fc6ff1" +dependencies = [ + "serde", + "tracing-core", +] + [[package]] name = "tracing-subscriber" version = "0.3.23" @@ -3747,12 +3947,15 @@ dependencies = [ "nu-ansi-term", "once_cell", "regex-automata", + "serde", + "serde_json", "sharded-slab", "smallvec", "thread_local", "tracing", "tracing-core", "tracing-log", + "tracing-serde", ] [[package]] @@ -3892,6 +4095,7 @@ dependencies = [ "idna", "percent-encoding", "serde", + "serde_derive", ] [[package]] @@ -3912,6 +4116,12 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b6c140620e7ffbb22c2dee59cafe6084a59b5ffc27a8859a5f0d494b5d52b6be" +[[package]] +name = "utf8parse" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" + [[package]] name = "uuid" version = "1.23.1" @@ -4385,6 +4595,18 @@ version = "0.53.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d6bbff5f0aada427a1e5a6da5f1f98158182f26556f345ac9e04d36d0ebed650" +[[package]] +name = "winnow" +version = "0.7.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "df79d97927682d2fd8adb29682d1140b343be4ac0f08fd68b7765d9c059d3945" + +[[package]] +name = "winnow" +version = "1.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2ee1708bef14716a11bae175f579062d4554d95be2c6829f518df847b7b3fdd0" + [[package]] name = "wiremock" version = "0.6.5" diff --git a/Cargo.toml b/Cargo.toml index 4e5bd58..18339b1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -37,6 +37,7 @@ miette = "7" walkdir = "2" tokio = { version = "1.52", features = ["macros", "rt-multi-thread", "time", "signal", "io-util", "sync"] } +tokio-util = "0.7" tokio-tungstenite = { version = "0.29", features = ["rustls-tls-webpki-roots"] } futures = "0.3" @@ -46,7 +47,7 @@ serde_ipld_dagcbor = "0.6" chrono = { version = "0.4", features = ["serde"] } cid = "0.11" -url = "2" +url = { version = "2", features = ["serde"] } bytes = "1" scc = "3" @@ -57,6 +58,7 @@ quick_cache = "0.6" reqwest = { version = "0.12", default-features = false, features = ["rustls-tls-webpki-roots", "http2", "json", "gzip", "stream"] } axum = "0.8" tower = { version = "0.5", features = ["util"] } +tower-http = { version = "0.6", features = ["trace"] } http = "1" wiremock = "0.6" @@ -64,9 +66,14 @@ wiremock = "0.6" tantivy = "0.26" tracing = "0.1" +tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt", "json"] } thiserror = "2" +confique = { version = "0.4", default-features = false, features = ["toml"] } +clap = { version = "4", features = ["derive", "env"] } +toml = { version = "0.9", default-features = false, features = ["parse"] } + [profile.release] lto = "fat" strip = true diff --git a/containerfiles/bobbin.Containerfile b/containerfiles/bobbin.Containerfile new file mode 100644 index 0000000..b59eb6c --- /dev/null +++ b/containerfiles/bobbin.Containerfile @@ -0,0 +1,14 @@ +FROM docker.io/library/rust:1-alpine3.23 AS builder +RUN apk add --no-cache build-base musl-dev cmake perl pkgconfig +WORKDIR /src +COPY . ./ +RUN rm -f .cargo/config.toml +RUN cargo build --release --bin bobbin --package bobbin +RUN strip target/release/bobbin + +FROM docker.io/library/alpine:3.23 +RUN apk add --no-cache ca-certificates +COPY --from=builder /src/target/release/bobbin /usr/local/bin/bobbin +ENV BOBBIN_BIND=0.0.0.0:8090 +EXPOSE 8090 +ENTRYPOINT ["/usr/local/bin/bobbin"] diff --git a/containerfiles/hydrant.Containerfile b/containerfiles/hydrant.Containerfile index 8eb64ac..ab7fb33 100644 --- a/containerfiles/hydrant.Containerfile +++ b/containerfiles/hydrant.Containerfile @@ -1,7 +1,7 @@ FROM docker.io/library/rust:1-alpine3.23 AS builder RUN apk add --no-cache build-base musl-dev cmake perl pkgconfig WORKDIR /src -COPY src/ ./ +COPY . ./ RUN rm -f .cargo/config.toml RUN cargo build --release --bin hydrant RUN strip target/release/hydrant diff --git a/containerfiles/slingshot.Containerfile b/containerfiles/slingshot.Containerfile index 130ee66..96ccafc 100644 --- a/containerfiles/slingshot.Containerfile +++ b/containerfiles/slingshot.Containerfile @@ -1,7 +1,7 @@ FROM docker.io/library/rust:1-alpine3.23 AS builder RUN apk add --no-cache build-base musl-dev cmake perl pkgconfig WORKDIR /src -COPY src/ ./ +COPY . ./ RUN rm -f .cargo/config.toml RUN cargo build --release --bin slingshot --package slingshot RUN strip target/release/slingshot @@ -12,6 +12,6 @@ WORKDIR /app COPY --from=builder /src/target/release/slingshot /usr/local/bin/slingshot COPY --from=builder /src/slingshot/static /app/static ENV SLINGSHOT_CACHE_DIR=/var/lib/slingshot -ENV SLINGSHOT_BIND=0.0.0.0:8080 +ENV SLINGSHOT_BIND=[::]:8080 EXPOSE 8080 ENTRYPOINT ["/usr/local/bin/slingshot"] diff --git a/crates/bobbin/Cargo.toml b/crates/bobbin/Cargo.toml index 652000f..158ab60 100644 --- a/crates/bobbin/Cargo.toml +++ b/crates/bobbin/Cargo.toml @@ -16,18 +16,18 @@ bobbin-knot-proxy = { workspace = true } bobbin-record-lru = { workspace = true } bobbin-search = { workspace = true } bobbin-slingshot-client = { workspace = true } -bobbin-types = { workspace = true } bobbin-xrpc = { workspace = true } -jacquard-common = { workspace = true } - axum = { workspace = true } -serde_json = { workspace = true } -tokio = { workspace = true, features = ["macros", "rt-multi-thread", "signal"] } +futures = { workspace = true } +socket2 = "0.6" +tokio = { workspace = true, features = ["macros", "rt-multi-thread", "net", "signal"] } +tokio-util = { workspace = true } url = { workspace = true } tracing = { workspace = true } -tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt"] } +tracing-subscriber = { workspace = true } anyhow = { workspace = true } - -[dev-dependencies] -bytes = { workspace = true } +clap = { workspace = true } +confique = { workspace = true } +thiserror = { workspace = true } +toml = { workspace = true } diff --git a/crates/bobbin/src/config.rs b/crates/bobbin/src/config.rs new file mode 100644 index 0000000..ce58e9d --- /dev/null +++ b/crates/bobbin/src/config.rs @@ -0,0 +1,397 @@ +use std::collections::HashSet; +use std::net::SocketAddr; +use std::path::{Path, PathBuf}; +use std::str::FromStr; + +use anyhow::{Context, anyhow}; +use confique::Config; +use url::Url; + +const SYSTEM_CONFIG_PATH: &str = "/etc/bobbin/config.toml"; +const ENV_PREFIX: &str = "BOBBIN_"; + +const KNOWN_KEYS: &[&str] = &[ + "server.binds", + "server.shutdown_grace_secs", + "hydrant.url", + "hydrant.start_cursor", + "slingshot.url", + "record_cache.lru_bytes", + "search.heap_bytes", + "knot.allow_private", + "knot.require_https", + "log.format", + "log.filter", +]; + +const KNOWN_ENVS: &[&str] = &[ + "BOBBIN_CONFIG", + "BOBBIN_BIND", + "BOBBIN_SHUTDOWN_GRACE_SECS", + "BOBBIN_HYDRANT_URL", + "BOBBIN_START_CURSOR", + "BOBBIN_SLINGSHOT_URL", + "BOBBIN_RECORD_LRU_BYTES", + "BOBBIN_SEARCH_HEAP_BYTES", + "BOBBIN_KNOT_ALLOW_PRIVATE", + "BOBBIN_KNOT_REQUIRE_HTTPS", + "BOBBIN_LOG_FORMAT", + "BOBBIN_LOG", +]; + +#[derive(Debug, Config)] +pub struct BobbinConfig { + #[config(nested)] + pub server: ServerConfig, + + #[config(nested)] + pub hydrant: HydrantConfig, + + #[config(nested)] + pub slingshot: SlingshotConfig, + + #[config(nested)] + pub record_cache: RecordCacheConfig, + + #[config(nested)] + pub search: SearchConfig, + + #[config(nested)] + pub knot: KnotConfig, + + #[config(nested)] + pub log: LogConfig, +} + +#[derive(Debug, Config)] +pub struct ServerConfig { + /// Addresses the XRPC server listens on. When using as an env var, comma-separated. + #[config( + env = "BOBBIN_BIND", + parse_env = parse_binds, + default = ["127.0.0.1:8090", "[::1]:8090"] + )] + pub binds: Vec, + + /// The amount of time in seconds to allow in-flight requests to drain after sigterm + /// before forcing the listener closed. + #[config(env = "BOBBIN_SHUTDOWN_GRACE_SECS", default = 30)] + pub shutdown_grace_secs: u64, +} + +#[derive(Debug, thiserror::Error)] +pub enum BindParseError { + #[error("BOBBIN_BIND must list at least one address")] + Empty, + #[error("invalid bind entry `{0}`: {1}")] + Invalid(String, std::net::AddrParseError), +} + +fn parse_binds(raw: &str) -> Result, BindParseError> { + let addrs: Vec = raw + .split(',') + .map(str::trim) + .filter(|s| !s.is_empty()) + .map(|s| SocketAddr::from_str(s).map_err(|e| BindParseError::Invalid(s.to_owned(), e))) + .collect::>()?; + if addrs.is_empty() { + return Err(BindParseError::Empty); + } + Ok(addrs) +} + +#[derive(Debug, Config)] +pub struct HydrantConfig { + /// Base URL of the hydrant instance - the cursor-replayable /stream lives + /// under this. Use `ws://` or `wss://` - `http://` and `https://` + /// are rewritten to the corresponding ws scheme at connection time! + #[config(env = "BOBBIN_HYDRANT_URL", default = "http://127.0.0.1:13010")] + pub url: Url, + + /// The cursor to request on first connect. Reconnects ignore this and resume + /// strictly after the last cursor seen internally. + #[config(env = "BOBBIN_START_CURSOR", default = 0)] + pub start_cursor: u64, +} + +#[derive(Debug, Config)] +pub struct SlingshotConfig { + /// Base URL of a slingshot instance. Used for record bodies and identity. + #[config(env = "BOBBIN_SLINGSHOT_URL", default = "http://127.0.0.1:13011")] + pub url: Url, +} + +#[derive(Debug, Config)] +pub struct RecordCacheConfig { + /// Bound of bytes on the in-process record LRU. Records evict on a weighted + /// LRU policy keyed on URI plus payload length. + #[config(env = "BOBBIN_RECORD_LRU_BYTES", default = 67_108_864)] + pub lru_bytes: u64, +} + +#[derive(Debug, Config)] +pub struct SearchConfig { + /// The heap size in bytes for the in-mem tantivy writer. Larger values + /// trade RAM for fewer segment merges - the index itself lives in + /// `RamDirectory` and is rebuilt from hydrant replay on every restart. + #[config(env = "BOBBIN_SEARCH_HEAP_BYTES", default = 50_000_000)] + pub heap_bytes: u64, +} + +#[derive(Debug, Config)] +pub struct KnotConfig { + /// Whether to allow the knot proxy to dial private/loopback addresses. Off in + /// production - on for local testing against a knotserver on localhost. + #[config(env = "BOBBIN_KNOT_ALLOW_PRIVATE", default = false)] + pub allow_private: bool, + + /// Require https on knot hosts. Disable only when proxying to a local + /// knot for development. + #[config(env = "BOBBIN_KNOT_REQUIRE_HTTPS", default = true)] + pub require_https: bool, +} + +#[derive(Debug, Config)] +pub struct LogConfig { + /// Log emitter format. `text` produces human-readable output for local + /// development. `json` emits one structured object per line for log + /// shippers in production. + #[config(env = "BOBBIN_LOG_FORMAT", default = "text")] + pub format: String, + + /// `tracing-subscriber` env-filter directive. Defaults to `info` across + /// every span - override with `BOBBIN_LOG=bobbin_xrpc=debug,info` etc. + #[config(env = "BOBBIN_LOG", default = "info")] + pub filter: String, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum LogFormat { + Text, + Json, +} + +impl std::str::FromStr for LogFormat { + type Err = String; + + fn from_str(s: &str) -> Result { + match s { + "text" => Ok(Self::Text), + "json" => Ok(Self::Json), + other => Err(format!( + "log format must be `text` or `json`, got `{other}`" + )), + } + } +} + +pub fn load(path: Option<&PathBuf>) -> anyhow::Result { + check_envs(std::env::vars().map(|(k, _)| k))?; + if let Some(p) = path { + check_keys(p)?; + } + check_keys(Path::new(SYSTEM_CONFIG_PATH))?; + let mut builder = BobbinConfig::builder().env(); + if let Some(p) = path { + builder = builder.file(p); + } + builder + .file(SYSTEM_CONFIG_PATH) + .load() + .context("load configuration") +} + +pub fn template() -> String { + confique::toml::template::(confique::toml::FormatOptions::default()) +} + +fn check_keys(path: &Path) -> anyhow::Result<()> { + let bytes = match std::fs::read_to_string(path) { + Ok(b) => b, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()), + Err(e) => return Err(anyhow!("read {}: {e}", path.display())), + }; + let value: toml::Value = + toml::from_str(&bytes).map_err(|e| anyhow!("parse {}: {e}", path.display()))?; + let known: HashSet<&str> = KNOWN_KEYS.iter().copied().collect(); + let unknown: Vec = collect_paths(&value, "") + .into_iter() + .filter(|p| !known.contains(p.as_str())) + .collect(); + if unknown.is_empty() { + Ok(()) + } else { + Err(anyhow!( + "unknown keys in {}: {}\nrun `bobbin config-template` for the canonical schema", + path.display(), + unknown.join(", "), + )) + } +} + +fn check_envs(iter: I) -> anyhow::Result<()> +where + I: IntoIterator, + S: AsRef, +{ + let known: HashSet<&str> = KNOWN_ENVS.iter().copied().collect(); + let unknown: Vec = iter + .into_iter() + .filter_map(|var| { + let var = var.as_ref(); + (var.starts_with(ENV_PREFIX) && !known.contains(var)).then(|| var.to_owned()) + }) + .collect(); + if unknown.is_empty() { + Ok(()) + } else { + Err(anyhow!( + "unknown {ENV_PREFIX}* environment variables: {}\nrun `bobbin config-template` for the canonical schema", + unknown.join(", "), + )) + } +} + +fn collect_paths(value: &toml::Value, prefix: &str) -> Vec { + match value { + toml::Value::Table(table) => table + .iter() + .flat_map(|(key, child)| { + let path = if prefix.is_empty() { + key.clone() + } else { + format!("{prefix}.{key}") + }; + match child { + toml::Value::Table(_) => collect_paths(child, &path), + _ => vec![path], + } + }) + .collect(), + _ => Vec::new(), + } +} + +#[cfg(test)] +fn template_paths(template: &str) -> Vec { + template + .lines() + .scan(String::new(), |section, line| { + let trimmed = line.trim_start(); + if let Some(name) = trimmed + .strip_prefix('[') + .and_then(|rest| rest.trim_end().strip_suffix(']')) + { + *section = name.to_owned(); + return Some(None); + } + let path = trimmed + .strip_prefix('#') + .and_then(|rest| rest.split_once('=')) + .map(|(key, _)| key.trim()) + .filter(|key| { + !key.is_empty() && key.chars().all(|c| c.is_ascii_alphanumeric() || c == '_') + }) + .map(|key| format!("{section}.{key}")); + Some(path) + }) + .flatten() + .collect() +} + +#[cfg(test)] +mod tests { + use super::*; + + fn write(name: &str, body: &str) -> PathBuf { + let dir = std::env::temp_dir().join(format!( + "bobbin-config-test-{}-{}", + std::process::id(), + name + )); + std::fs::create_dir_all(&dir).unwrap(); + let path = dir.join("config.toml"); + std::fs::write(&path, body).unwrap(); + path + } + + #[test] + fn known_keys_exactly_match_template_paths() { + let template = template(); + let from_template: HashSet = template_paths(&template).into_iter().collect(); + let known: HashSet = KNOWN_KEYS.iter().map(|s| (*s).to_owned()).collect(); + assert_eq!( + from_template, known, + "KNOWN_KEYS must equal the set of paths the confique template generates", + ); + } + + #[test] + fn template_paths_extracts_commented_leaf_assignments_under_each_section() { + let raw = "[server]\n# Doc comment with = inside.\n#binds = [\"127.0.0.1:8090\"]\n#shutdown_grace_secs = 30\n\n[log]\n#format = \"text\"\n"; + let mut paths = template_paths(raw); + paths.sort(); + assert_eq!( + paths, + vec![ + "log.format".to_owned(), + "server.binds".to_owned(), + "server.shutdown_grace_secs".to_owned(), + ], + ); + } + + #[test] + fn unknown_section_rejected() { + let path = write( + "unknown_section", + "[server]\nbinds = [\"127.0.0.1:9000\"]\n\n[mystery]\nfoo = 1\n", + ); + let err = check_keys(&path).expect_err("must reject"); + assert!(err.to_string().contains("mystery.foo"), "got {err}"); + } + + #[test] + fn unknown_field_rejected() { + let path = write("typo", "[server]\nbidns = [\"127.0.0.1:9000\"]\n"); + let err = check_keys(&path).expect_err("must reject"); + assert!(err.to_string().contains("server.bidns"), "got {err}"); + } + + #[test] + fn missing_file_passes_silently() { + check_keys(Path::new("/definitely/does/not/exist.toml")).expect("no file is fine"); + } + + #[test] + fn full_known_config_passes() { + let path = write("ok", &template()); + check_keys(&path).expect("template must validate against itself"); + } + + #[test] + fn known_envs_carry_the_env_prefix() { + KNOWN_ENVS.iter().for_each(|name| { + assert!( + name.starts_with(ENV_PREFIX), + "KNOWN_ENVS entry {name:?} must start with {ENV_PREFIX:?}" + ); + }); + } + + #[test] + fn unknown_bobbin_env_rejected() { + let err = check_envs(["BOBBIN_BIDNS"]).expect_err("typo must surface"); + assert!(err.to_string().contains("BOBBIN_BIDNS"), "got {err}"); + } + + #[test] + fn known_bobbin_env_passes() { + check_envs(["BOBBIN_BIND", "BOBBIN_LOG"]).expect("known names must pass"); + } + + #[test] + fn non_bobbin_env_ignored() { + check_envs(["PATH", "HOME", "RUST_LOG"]).expect("only BOBBIN_* names are validated"); + } +} diff --git a/crates/bobbin/src/main.rs b/crates/bobbin/src/main.rs index d14898f..0253274 100644 --- a/crates/bobbin/src/main.rs +++ b/crates/bobbin/src/main.rs @@ -1,90 +1,300 @@ -use std::env; use std::net::SocketAddr; +use std::path::PathBuf; +use std::process::ExitCode; use std::sync::Arc; +use std::time::Duration; use anyhow::{Context, anyhow}; use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor}; -use bobbin_ingest::{IngestConfig, run as run_ingest}; +use bobbin_ingest::{IngestConfig, IngestRuntime, RepoIdResolver, run as run_ingest}; use bobbin_knot_proxy::{KnotProxy, KnotProxyConfig}; use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; -use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex}; +use bobbin_search::SearchIndex; use bobbin_slingshot_client::SlingshotClient; use bobbin_xrpc::{AppState, router}; +use clap::{Parser, Subcommand}; +use tokio::signal::unix::{SignalKind, signal}; +use tokio::task::JoinHandle; +use tokio_util::sync::CancellationToken; +use tracing::level_filters::LevelFilter; use tracing_subscriber::EnvFilter; -use url::Url; + +mod config; + +use config::{BobbinConfig, LogFormat}; + +#[derive(Parser)] +#[command(name = "bobbin", about = "Read-only AppView for Tangled records")] +struct Cli { + /// Path to a TOML config file. Environment variables override file values + /// for any `BOBBIN_*` setting; `/etc/bobbin/config.toml` is consulted as a + /// final fallback so distro packaging can drop a default in place. + #[arg(short, long, value_name = "FILE", env = "BOBBIN_CONFIG")] + config: Option, + + #[command(subcommand)] + command: Option, +} + +#[derive(Subcommand)] +enum Command { + /// Print a fully-commented TOML template to stdout. Use this to seed + /// `config.toml` for a fresh deploy. + ConfigTemplate, + /// Load and validate the configuration without starting the server. + Validate, +} #[tokio::main] -async fn main() -> anyhow::Result<()> { - tracing_subscriber::fmt() - .with_env_filter( - EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")), - ) - .init(); - - let hydrant_url = - env::var("BOBBIN_HYDRANT_URL").unwrap_or_else(|_| "http://127.0.0.1:13010".into()); - let slingshot_url = - env::var("BOBBIN_SLINGSHOT_URL").unwrap_or_else(|_| "http://127.0.0.1:13011".into()); - let bind: SocketAddr = env::var("BOBBIN_BIND") - .unwrap_or_else(|_| "127.0.0.1:8090".into()) - .parse()?; - let cache_bytes: u64 = env::var("BOBBIN_RECORD_LRU_BYTES") - .ok() - .and_then(|v| v.parse().ok()) - .unwrap_or(64 * 1024 * 1024); - let search_heap_bytes: usize = env::var("BOBBIN_SEARCH_HEAP_BYTES") - .ok() - .and_then(|v| v.parse().ok()) - .unwrap_or(DEFAULT_WRITER_HEAP_BYTES); - let start_cursor: u64 = env::var("BOBBIN_START_CURSOR") - .ok() - .and_then(|v| v.parse().ok()) - .unwrap_or(0); - - let records: Arc = - Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(cache_bytes))); - let slingshot = SlingshotClient::new(Url::parse(&slingshot_url)?)?; +async fn main() -> ExitCode { + let cli = Cli::parse(); + + if let Some(Command::ConfigTemplate) = cli.command { + print!("{}", config::template()); + return ExitCode::SUCCESS; + } + + let cfg = match config::load(cli.config.as_ref()) { + Ok(c) => c, + Err(e) => { + eprintln!("failed to load configuration: {e:#}"); + return ExitCode::FAILURE; + } + }; + + if let Err(e) = init_tracing(&cfg) { + eprintln!("failed to install tracing subscriber: {e}"); + return ExitCode::FAILURE; + } + + if matches!(cli.command, Some(Command::Validate)) { + println!("configuration is valid"); + return ExitCode::SUCCESS; + } + + match run(cfg).await { + Ok(()) => ExitCode::SUCCESS, + Err(e) => { + tracing::error!(error = ?e, "fatal"); + ExitCode::FAILURE + } + } +} + +fn init_tracing(cfg: &BobbinConfig) -> Result<(), String> { + let combined = format!("{},{}", LevelFilter::INFO, cfg.log.filter); + let filter = EnvFilter::try_new(&combined) + .map_err(|e| format!("invalid log filter `{}`: {e}", cfg.log.filter))?; + let format: LogFormat = cfg.log.format.parse()?; + let builder = tracing_subscriber::fmt().with_env_filter(filter); + match format { + LogFormat::Text => builder.try_init().map_err(|e| e.to_string()), + LogFormat::Json => builder.json().try_init().map_err(|e| e.to_string()), + } +} + +async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { + let records: Arc = Arc::new(LruRecordStore::new(CacheCapacity::from_bytes( + cfg.record_cache.lru_bytes, + ))); + let slingshot = SlingshotClient::new(cfg.slingshot.url.clone())?; + let resolver = Arc::new(RepoIdResolver::with_slingshot(slingshot.clone())); let edges = Arc::new(EdgeStore::new()); let coverage = Arc::new(CoverageWatch::new()); - let knots = Arc::new(KnotProxy::new(KnotProxyConfig::default())?); - let search = Arc::new(SearchIndex::new(search_heap_bytes)?); + let knots = Arc::new(KnotProxy::new(KnotProxyConfig { + allow_private_hosts: cfg.knot.allow_private, + require_https: cfg.knot.require_https, + ..KnotProxyConfig::default() + })?); + let search_heap = usize::try_from(cfg.search.heap_bytes) + .with_context(|| format!("search.heap_bytes {} exceeds usize", cfg.search.heap_bytes))?; + let search = Arc::new(SearchIndex::new(search_heap)?); let ingest_cfg = IngestConfig { - hydrant_base: Url::parse(&hydrant_url)?, - start_cursor: HydrantCursor::new(start_cursor), + hydrant_base: cfg.hydrant.url.clone(), + start_cursor: HydrantCursor::new(cfg.hydrant.start_cursor), }; - let ingest_edges = edges.clone(); + let cancel = CancellationToken::new(); let ingest_coverage = coverage.clone(); - let ingest_search = search.clone(); - let ingest_records = records.clone(); - let ingest_handle = tokio::spawn(async move { - run_ingest( - ingest_cfg, - ingest_edges, - ingest_coverage, - ingest_search, - ingest_records, - ) - .await - }); + let ingest_runtime = IngestRuntime { + store: edges.clone(), + coverage: coverage.clone(), + search: search.clone(), + records: records.clone(), + resolver: resolver.clone(), + cancel: cancel.clone(), + }; + let mut ingest_handle = tokio::spawn(run_ingest(ingest_cfg, ingest_runtime)); let state = AppState::new(records, slingshot, edges, coverage, knots, search); let app = router(state); - tracing::info!(%bind, %hydrant_url, %slingshot_url, "bobbin listening"); - let listener = tokio::net::TcpListener::bind(&bind).await?; - let server_handle = tokio::spawn(async move { axum::serve(listener, app).await }); + let binds = cfg.server.binds.clone(); + let hydrant_url = cfg.hydrant.url.as_str().to_owned(); + let slingshot_url = cfg.slingshot.url.as_str().to_owned(); + let grace = Duration::from_secs(cfg.server.shutdown_grace_secs); + let bind_display = binds + .iter() + .map(SocketAddr::to_string) + .collect::>() + .join(","); + tracing::info!(binds = %bind_display, %hydrant_url, %slingshot_url, "bobbin listening"); + + let signal_cancel = cancel.clone(); + let mut server_handle = tokio::spawn(serve_all(binds, app, signal_cancel)); + + tokio::select! { + res = &mut server_handle => { + cancel.cancel(); + let cursor = ingest_coverage.snapshot().last_cursor().raw(); + tracing::info!(grace_secs = grace.as_secs(), cursor, "draining ingest"); + drain_with_grace("ingest", grace, &mut ingest_handle).await; + match res { + Ok(Ok(())) => Ok(()), + Ok(Err(e)) => Err(anyhow::Error::from(e)).context("axum server failed"), + Err(join) => Err(anyhow!("server task panicked: {join}")), + } + } + res = &mut ingest_handle => { + cancel.cancel(); + let cursor = ingest_coverage.snapshot().last_cursor().raw(); + tracing::info!(grace_secs = grace.as_secs(), cursor, "draining server"); + drain_with_grace("server", grace, &mut server_handle).await; + match res { + Ok(Ok(())) => Err(anyhow!("ingest run loop exited; loop is supposed to be infinite")), + Ok(Err(e)) => Err(anyhow::Error::from(e)).context("ingest exited"), + Err(join) => Err(anyhow!("ingest task panicked: {join}")), + } + } + } +} + +async fn drain_with_grace(label: &'static str, grace: Duration, handle: &mut JoinHandle) { + if tokio::time::timeout(grace, &mut *handle).await.is_err() { + tracing::warn!( + grace_secs = grace.as_secs(), + label, + "task did not stop within grace, aborting" + ); + handle.abort(); + } +} + +async fn serve_all( + binds: Vec, + app: axum::Router, + cancel: CancellationToken, +) -> std::io::Result<()> { + let listeners = futures::future::try_join_all(binds.into_iter().map(|addr| async move { + let listener = bind_listener(addr)?; + tracing::info!(%addr, "bobbin listener bound"); + Ok::<_, std::io::Error>(listener) + })) + .await?; + + let trigger = cancel.clone(); + let signal_task = tokio::spawn(async move { + wait_for_shutdown().await; + tracing::info!("shutdown signal received, draining server"); + trigger.cancel(); + }); + + let services = listeners.into_iter().map(|listener| { + let app = app.clone(); + let cancel = cancel.clone(); + async move { + axum::serve(listener, app) + .with_graceful_shutdown(async move { cancel.cancelled().await }) + .await + } + }); + let result = futures::future::try_join_all(services).await.map(|_| ()); + signal_task.abort(); + result +} + +fn bind_listener(addr: SocketAddr) -> std::io::Result { + let domain = match addr { + SocketAddr::V4(_) => socket2::Domain::IPV4, + SocketAddr::V6(_) => socket2::Domain::IPV6, + }; + let socket = socket2::Socket::new(domain, socket2::Type::STREAM, Some(socket2::Protocol::TCP))?; + if matches!(addr, SocketAddr::V6(_)) { + socket.set_only_v6(true)?; + } + socket.set_reuse_address(true)?; + socket.set_nonblocking(true)?; + socket.bind(&addr.into())?; + socket.listen(1024)?; + let std_listener: std::net::TcpListener = socket.into(); + tokio::net::TcpListener::from_std(std_listener) +} + +async fn wait_for_shutdown() { + let ctrl_c = tokio::signal::ctrl_c(); + let mut sigterm = match signal(SignalKind::terminate()) { + Ok(s) => s, + Err(e) => { + tracing::warn!( + ?e, + "could not install SIGTERM handler, shutdown will only honor ctrl-c" + ); + ctrl_c.await.ok(); + return; + } + }; tokio::select! { - res = server_handle => match res { - Ok(Ok(())) => Ok(()), - Ok(Err(e)) => Err(e).context("axum server failed"), - Err(join) => Err(anyhow!("server task panicked: {join}")), - }, - res = ingest_handle => match res { - Ok(Ok(())) => Err(anyhow!("ingest run loop exited; loop is supposed to be infinite")), - Ok(Err(e)) => Err(e).context("ingest exited"), - Err(join) => Err(anyhow!("ingest task panicked: {join}")), - }, + _ = ctrl_c => {} + _ = sigterm.recv() => {} + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test(start_paused = true)] + async fn drain_returns_immediately_when_task_already_done() { + let mut handle = tokio::spawn(async { 7u32 }); + tokio::time::advance(Duration::from_millis(1)).await; + let start = tokio::time::Instant::now(); + drain_with_grace("test", Duration::from_secs(60), &mut handle).await; + assert!(start.elapsed() < Duration::from_millis(10)); + } + + #[tokio::test(start_paused = true)] + async fn drain_aborts_runaway_task_after_grace() { + let mut handle = tokio::spawn(async { + std::future::pending::<()>().await; + }); + let grace = Duration::from_secs(5); + let start = tokio::time::Instant::now(); + drain_with_grace("test", grace, &mut handle).await; + assert!(start.elapsed() >= grace); + let outcome = handle.await; + assert!(outcome.is_err() && outcome.unwrap_err().is_cancelled()); + } + + #[test] + fn typo_filter_keeps_info_default_for_other_targets() { + let filter = format!("{},{}", LevelFilter::INFO, "blah_invalid"); + let parsed = EnvFilter::try_new(&filter).expect("filter parses"); + let rendered = parsed.to_string(); + assert!( + rendered.contains("info"), + "expected info default, got {rendered}" + ); + assert!( + rendered.contains("blah_invalid"), + "expected user override, got {rendered}", + ); + } + + #[test] + fn explicit_user_level_overrides_info_default() { + let filter = format!("{},{}", LevelFilter::INFO, "warn"); + let parsed = EnvFilter::try_new(&filter).expect("filter parses"); + assert_eq!(parsed.to_string(), "warn"); } } diff --git a/example.toml b/example.toml new file mode 100644 index 0000000..46632ce --- /dev/null +++ b/example.toml @@ -0,0 +1,95 @@ +[server] +# Addresses the XRPC server listens on. When using as an env var, comma-separated. +# +# Can also be specified via environment variable `BOBBIN_BIND`. +# +# Default value: ["127.0.0.1:8090", "[::1]:8090"] +#binds = ["127.0.0.1:8090", "[::1]:8090"] + +# The amount of time in seconds to allow in-flight requests to drain after sigterm +# before forcing the listener closed. +# +# Can also be specified via environment variable `BOBBIN_SHUTDOWN_GRACE_SECS`. +# +# Default value: 30 +#shutdown_grace_secs = 30 + +[hydrant] +# Base URL of the hydrant instance - the cursor-replayable /stream lives +# under this. Use `ws://` or `wss://` - `http://` and `https://` +# are rewritten to the corresponding ws scheme at connection time! +# +# Can also be specified via environment variable `BOBBIN_HYDRANT_URL`. +# +# Default value: "http://127.0.0.1:13010" +#url = "http://127.0.0.1:13010" + +# The cursor to request on first connect. Reconnects ignore this and resume +# strictly after the last cursor seen internally. +# +# Can also be specified via environment variable `BOBBIN_START_CURSOR`. +# +# Default value: 0 +#start_cursor = 0 + +[slingshot] +# Base URL of a slingshot instance. Used for record bodies and identity. +# +# Can also be specified via environment variable `BOBBIN_SLINGSHOT_URL`. +# +# Default value: "http://127.0.0.1:13011" +#url = "http://127.0.0.1:13011" + +[record_cache] +# Bound of bytes on the in-process record LRU. Records evict on a weighted +# LRU policy keyed on URI plus payload length. +# +# Can also be specified via environment variable `BOBBIN_RECORD_LRU_BYTES`. +# +# Default value: 67108864 +#lru_bytes = 67108864 + +[search] +# The heap size in bytes for the in-mem tantivy writer. Larger values +# trade RAM for fewer segment merges - the index itself lives in +# `RamDirectory` and is rebuilt from hydrant replay on every restart. +# +# Can also be specified via environment variable `BOBBIN_SEARCH_HEAP_BYTES`. +# +# Default value: 50000000 +#heap_bytes = 50000000 + +[knot] +# Whether to allow the knot proxy to dial private/loopback addresses. Off in +# production - on for local testing against a knotserver on localhost. +# +# Can also be specified via environment variable `BOBBIN_KNOT_ALLOW_PRIVATE`. +# +# Default value: false +#allow_private = false + +# Require https on knot hosts. Disable only when proxying to a local +# knot for development. +# +# Can also be specified via environment variable `BOBBIN_KNOT_REQUIRE_HTTPS`. +# +# Default value: true +#require_https = true + +[log] +# Log emitter format. `text` produces human-readable output for local +# development. `json` emits one structured object per line for log +# shippers in production. +# +# Can also be specified via environment variable `BOBBIN_LOG_FORMAT`. +# +# Default value: "text" +#format = "text" + +# `tracing-subscriber` env-filter directive. Defaults to `info` across +# every span - override with `BOBBIN_LOG=bobbin_xrpc=debug,info` etc. +# +# Can also be specified via environment variable `BOBBIN_LOG`. +# +# Default value: "info" +#filter = "info"