This is an automated email from the ASF dual-hosted git repository.
hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/master by this push:
new bc9ef0c3f fix(connectors): warn when the runtime API is exposed
without a key (#3804)
bc9ef0c3f is described below
commit bc9ef0c3f938b82ac0fe3c3a89f64d0242660edd
Author: Maxim Levkov <[email protected]>
AuthorDate: Fri Sep 4 13:02:41 2026 -0700
fix(connectors): warn when the runtime API is exposed without a key (#3804)
---
core/connectors/runtime/README.md | 86 +++-
core/connectors/runtime/config.toml | 9 +-
core/connectors/runtime/example_config/config.toml | 9 +-
core/connectors/runtime/src/api/config.rs | 111 +++-
core/connectors/runtime/src/api/mod.rs | 573 ++++++++++++++++++++-
5 files changed, 780 insertions(+), 8 deletions(-)
diff --git a/core/connectors/runtime/README.md
b/core/connectors/runtime/README.md
index 590d0a8f8..23803704d 100644
--- a/core/connectors/runtime/README.md
+++ b/core/connectors/runtime/README.md
@@ -197,10 +197,17 @@ Connector runtime has an optional HTTP API that can be
enabled by setting the `e
```toml
[http] # Optional HTTP API configuration
enabled = true
+# Loopback on purpose: the configuration endpoints return plugin credentials in
+# plaintext and also accept writes. Set api_key in the same edit if you move
+# this off loopback, and http.tls unless cleartext is acceptable.
address = "127.0.0.1:8081"
-api_key = "" # Optional API key for authentication to be passed as `api-key`
header
+api_key = "" # Optional API key for authentication to be passed as `api-key`
header; empty disables authentication
[http.cors] # Optional CORS configuration for HTTP API
+# Enabling this with the shipped allowed_origins = ["*"] lets any page the
+# operator visits read these endpoints cross-origin, whatever address is bound.
+# Pin the origins you own, or set api_key, before turning it on.
+# "*" is honored only as the first entry; anywhere else it panics at startup.
enabled = false
allowed_methods = ["GET", "POST", "PUT", "DELETE"]
allowed_origins = ["*"]
@@ -219,29 +226,104 @@ cert_file = "core/certs/iggy_cert.pem"
key_file = "core/certs/iggy_key.pem"
```
+> [!IMPORTANT]
+> **Treat this API as privileged. It reads and it writes.**
+>
+> The configuration endpoints return plugin configuration exactly as stored,
+> credentials included - a database connection string, an S3 secret key, a
+> webhook signing secret. Nothing redacts them on the way out. The runtime
+> masks the `api-key` in its own logs and redacts the state-store headers when
+> it serializes them, but no such path exists for plugin configuration.
+>
+> The exposure is not limited to disclosure. Publishing a configuration with
+> `POST /{sinks,sources}/{key}/configs` and then calling `POST .../restart` is
+> enough to repoint a connector at a destination of the caller's choosing,
+> because `restart` re-reads the stored configuration and starts the connector
+> from it - on the local provider that is whatever version was published last,
+> with no activation step in between. The runtime then forwards the operator's
+> topic data using its own Iggy credentials. `PUT .../configs/active` and
+> `DELETE .../configs` sit behind the same key.
+>
+> A rewritten plugin `path` is not loaded by the restart. `start_connector`
+> reuses the container `dlopen`ed at boot and only re-runs the plugin's init
+> with the new configuration, so a hostile path sits in the stored config until
+> the next time the runtime process starts, which makes it deferred code
+> execution rather than immediate.
+>
+> `api_key` is empty by default, which means authentication is **off** by
+> default. Only `/` and `/health` are exempt once it is set, so everything
above
+> sits behind that one empty string, and the loopback default `address` is what
+> confines it to local processes.
+>
+> Three edits take that containment away:
+>
+> - **Moving `address` off loopback.** Set `api_key` in the same edit. The
+> runtime warns at startup when the address resolves beyond loopback with no
+> key configured, but nothing prevents it.
+> - **Enabling `[http.cors]` with the shipped `allowed_origins = ["*"]`.** That
+> becomes `AllowOrigin::any()`, and the CORS layer wraps *outside*
+> authentication, so tower-http answers the preflight before
`resolve_api_key`
+> runs. A browser is a local process, so with a wildcard origin and no key
any
+> page the operator visits gets a green preflight for a configuration `POST`
+> and can then read *and rewrite* configuration cross-origin - the whole
+> publish-then-restart repoint above, not just disclosure. The rewrite half
+> needs the method on `allowed_methods` and `content-type` on
+> `allowed_headers`, both of which the shipped block grants.
+>
+> Setting `api_key` closes it, since an attacker's page cannot supply the
+> header. Pinning `allowed_origins` to origins the operator owns closes the
+> cross-origin *read*, and no startup warning fires for that case. One fires
+> for `null`, which is narrower than `*` but not owned by anyone: a browser
+> sends it from a sandboxed iframe, a `data:` URL and a `file://` page, all
of
+> which an attacker can produce, so pinning it buys nothing. No origin
setting
+> closes the CORS-simple route below; `api_key` does.
+> - **Leaving `http.tls.enabled = false`.** It ships disabled, so the `api-key`
+> header and the credential-bearing responses both travel in cleartext.
Enable
+> TLS alongside `api_key` whenever this API leaves loopback. The runtime
warns
+> at startup whenever the address resolves beyond loopback with TLS off, and
+> keeps warning after `api_key` is set, because the key crosses in the clear
+> too. Terminating TLS at an ingress or a service mesh is a valid answer to
+> that warning; the runtime cannot see it, so the line stays.
+>
+> And one that needs no edit at all. `POST .../restart` carries no body and no
+> content type, which makes it a CORS-simple request: a page can issue it with
+> `mode: 'no-cors'` and the browser sends it whatever `[http.cors]` says,
+> because CORS gates reading a response rather than issuing a request. So on
the
+> shipped keyless default, any page the operator visits can restart any
+> connector, and no startup warning covers it because it is true of the
defaults
+> rather than of an edit. Chrome's private network access blocks the
+> public-origin case; Firefox and Safari do not, a page served from a local
+> origin bypasses it everywhere, and setting `allow_private_network = true`
+> hands back the case Chrome would otherwise block. `api_key` is what closes
+> this one.
+
Currently, it does expose the following endpoints:
- `GET /`: welcome message.
- `GET /health`: health status of the runtime.
- `GET /stats`: runtime statistics including process info, memory/CPU usage,
and connector status.
-- `GET /metrics`: Prometheus-formatted metrics (when `http.metrics.enabled` is
`true`).
+- `GET {http.metrics.endpoint}` (default `/metrics`): Prometheus-formatted
metrics, when `http.metrics.enabled` is `true`.
- `GET /sinks`: list of sinks.
- `GET /sinks/{key}`: sink details.
- `GET /sinks/{key}/configs`: list of configuration versions for the sink.
- `POST /sinks/{key}/configs`: add a new configuration version for the sink.
+- `DELETE /sinks/{key}/configs`: delete one configuration version for the sink
- the `version` query parameter, or on the local provider the active version
when it is omitted.
- `GET /sinks/{key}/configs/{version}`: configuration details for a specific
version.
- `GET /sinks/{key}/configs/active`: active configuration details.
- `PUT /sinks/{key}/configs/active`: activate a specific configuration version
for the sink.
- `GET /sinks/{key}/configs/plugin`: sink plugin config, including the
optional `format` query parameter to specify the config format.
+- `POST /sinks/{key}/restart`: stop the sink and start it again from its
highest stored configuration version, which on the local provider is not
necessarily the active one
([#3848](https://github.com/apache/iggy/issues/3848)).
- `GET /sinks/{key}/transforms`: sink transforms to be applied to the fields.
- `GET /sources`: list of sources.
- `GET /sources/{key}`: source details.
- `GET /sources/{key}/configs`: list of configuration versions for the source.
- `POST /sources/{key}/configs`: add a new configuration version for the
source.
+- `DELETE /sources/{key}/configs`: delete one configuration version for the
source - the `version` query parameter, or on the local provider the active
version when it is omitted.
- `GET /sources/{key}/configs/{version}`: configuration details for a specific
version.
- `GET /sources/{key}/configs/active`: active configuration details.
- `PUT /sources/{key}/configs/active`: activate a specific configuration
version for the source.
- `GET /sources/{key}/configs/plugin`: source plugin config, including the
optional `format` query parameter to specify the config format.
+- `POST /sources/{key}/restart`: stop the source and start it again from its
highest stored configuration version, which on the local provider is not
necessarily the active one
([#3848](https://github.com/apache/iggy/issues/3848)).
- `GET /sources/{key}/transforms`: source transforms to be applied to the
fields.
## Telemetry
diff --git a/core/connectors/runtime/config.toml
b/core/connectors/runtime/config.toml
index d52909122..9a3e3befb 100644
--- a/core/connectors/runtime/config.toml
+++ b/core/connectors/runtime/config.toml
@@ -17,10 +17,17 @@
[http] # Optional HTTP API configuration
enabled = true
+# Loopback on purpose: the configuration endpoints return plugin credentials in
+# plaintext and also accept writes. Set api_key in the same edit if you move
+# this off loopback, and http.tls unless cleartext is acceptable.
address = "127.0.0.1:8081"
-api_key = "" # Optional API key for authentication to be passed as `api-key`
header
+api_key = "" # Optional API key for authentication to be passed as `api-key`
header; empty disables authentication
[http.cors] # Optional CORS configuration for HTTP API
+# Enabling this with the shipped allowed_origins = ["*"] lets any page the
+# operator visits read these endpoints cross-origin, whatever address is bound.
+# Pin the origins you own, or set api_key, before turning it on.
+# "*" is honored only as the first entry; anywhere else it panics at startup.
enabled = false
allowed_methods = ["GET", "POST", "PUT", "DELETE"]
allowed_origins = ["*"]
diff --git a/core/connectors/runtime/example_config/config.toml
b/core/connectors/runtime/example_config/config.toml
index 29afad4ce..a6690b801 100644
--- a/core/connectors/runtime/example_config/config.toml
+++ b/core/connectors/runtime/example_config/config.toml
@@ -17,10 +17,17 @@
[http] # Optional HTTP API configuration
enabled = true
+# Loopback on purpose: the configuration endpoints return plugin credentials in
+# plaintext and also accept writes. Set api_key in the same edit if you move
+# this off loopback, and http.tls unless cleartext is acceptable.
address = "127.0.0.1:8081"
-api_key = "" # Optional API key for authentication to be passed as `api-key`
header
+api_key = "" # Optional API key for authentication to be passed as `api-key`
header; empty disables authentication
[http.cors] # Optional CORS configuration for HTTP API
+# Enabling this with the shipped allowed_origins = ["*"] lets any page the
+# operator visits read these endpoints cross-origin, whatever address is bound.
+# Pin the origins you own, or set api_key, before turning it on.
+# "*" is honored only as the first entry; anywhere else it panics at startup.
enabled = false
allowed_methods = ["GET", "POST", "PUT", "DELETE"]
allowed_origins = ["*"]
diff --git a/core/connectors/runtime/src/api/config.rs
b/core/connectors/runtime/src/api/config.rs
index 1221bd3e8..78f5125a3 100644
--- a/core/connectors/runtime/src/api/config.rs
+++ b/core/connectors/runtime/src/api/config.rs
@@ -73,6 +73,44 @@ pub struct HttpCorsConfig {
pub allow_private_network: bool,
}
+impl HttpCorsConfig {
+ /// Whether `configure_cors` will turn this into `AllowOrigin::any()`.
+ ///
+ /// Only the first entry decides, because that is what the mapping below
+ /// reads. A `"*"` in any later position never reaches a served request:
+ /// `AllowOrigin::list` panics on a wildcard, so such a config takes the
+ /// process down at startup rather than allowing anything.
+ ///
+ /// Compared untrimmed, because `configure_cors` compares untrimmed: `" *"`
+ /// becomes a list entry no `Origin` matches, so it allows nothing and
there
+ /// is nothing to warn about. `core/server/src/http.rs` trims before the
same
+ /// comparison; porting that here means trimming in both places at once,
+ /// since trimming only this one would warn about a closed config and
+ /// trimming only the mapping would open one silently.
+ pub fn allows_any_origin(&self) -> bool {
+ self.allowed_origins
+ .first()
+ .is_some_and(|origin| origin == "*")
+ }
+
+ /// Whether the resulting policy admits an origin nobody owns, assuming one
+ /// gets built: a wildcard past the first position panics `configure_cors`,
+ /// so that shape warns and then dies.
+ ///
+ /// `*` is one. `null` is the other: a browser sends it from a sandboxed
+ /// iframe, a `data:` URL and a `file://` page, so listing it hands the
+ /// cross-origin read to whoever gets the operator to open a page.
+ /// `AllowOrigin::list` echoes any listed value back on a match, so a
+ /// `null` anywhere counts, not only first. Untrimmed for the reason above.
+ ///
+ /// Separate from `allows_any_origin` because `configure_cors` has to keep
+ /// mapping `["null"]` to a one-entry list rather than widening it, and
+ /// defined in terms of it so the two cannot disagree about `*`.
+ pub fn allows_unowned_origin(&self) -> bool {
+ self.allows_any_origin() || self.allowed_origins.iter().any(|origin|
origin == "null")
+ }
+}
+
#[derive(Debug, Default, Deserialize, Serialize, Clone, ConfigEnv)]
pub struct HttpTlsConfig {
pub enabled: bool,
@@ -111,7 +149,7 @@ pub fn map_connector_config(
pub fn configure_cors(config: &HttpCorsConfig) -> CorsLayer {
let allowed_origins = match &config.allowed_origins {
origins if origins.is_empty() => AllowOrigin::default(),
- origins if origins.first().unwrap() == "*" => AllowOrigin::any(),
+ _ if config.allows_any_origin() => AllowOrigin::any(),
origins => AllowOrigin::list(origins.iter().map(|s|
s.parse().unwrap())),
};
@@ -223,3 +261,74 @@ impl std::fmt::Display for HttpCorsConfig {
)
}
}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ fn cors(allowed_origins: &[&str]) -> HttpCorsConfig {
+ HttpCorsConfig {
+ allowed_origins: allowed_origins.iter().map(|s|
(*s).to_owned()).collect(),
+ ..HttpCorsConfig::default()
+ }
+ }
+
+ #[test]
+ fn given_a_leading_wildcard_when_classified_should_allow_any_origin() {
+ assert!(cors(&["*"]).allows_any_origin());
+ assert!(cors(&["*"]).allows_unowned_origin());
+ }
+
+ #[test]
+ fn given_pinned_or_absent_origins_when_classified_should_allow_none() {
+ assert!(
+ !cors(&[]).allows_any_origin(),
+ "an empty list emits no header"
+ );
+ assert!(!cors(&[]).allows_unowned_origin());
+ assert!(!cors(&["https://console.example"]).allows_any_origin());
+ assert!(!cors(&["https://console.example"]).allows_unowned_origin());
+ assert!(
+ !cors(&["https://console.example", "*"]).allows_any_origin(),
+ "only the first entry reaches `AllowOrigin::any`; a later `*`
makes \
+ `AllowOrigin::list` panic, so that config never serves a request"
+ );
+ }
+
+ #[test]
+ fn given_the_classified_origin_shapes_when_built_should_produce_a_layer() {
+ // The predicates above describe what `configure_cors` does with each
+ // shape, so a shape that cannot build would make the warning the last
+ // line an operator sees before the process dies.
+ for origins in [
+ &["*"][..],
+ &[][..],
+ &["https://console.example"][..],
+ &["null"][..],
+ ] {
+ let _layer = configure_cors(&cors(origins));
+ }
+ }
+
+ #[test]
+ #[should_panic(expected = "Wildcard origin")]
+ fn
given_a_trailing_wildcard_when_built_should_panic_rather_than_allow_nothing() {
+ let _layer = configure_cors(&cors(&["https://console.example", "*"]));
+ }
+
+ #[test]
+ fn given_a_null_origin_when_classified_should_report_an_unowned_one() {
+ assert!(
+ cors(&["null"]).allows_unowned_origin(),
+ "a sandboxed iframe, a data: URL and a file:// page all send
Origin: null"
+ );
+ assert!(
+ cors(&["https://console.example", "null"]).allows_unowned_origin(),
+ "`AllowOrigin::list` echoes any listed value, so position does not
matter"
+ );
+ assert!(
+ !cors(&["null"]).allows_any_origin(),
+ "`configure_cors` still has to build a one-entry list, not widen
to any()"
+ );
+ }
+}
diff --git a/core/connectors/runtime/src/api/mod.rs
b/core/connectors/runtime/src/api/mod.rs
index 38cd8f6ed..62950fa6e 100644
--- a/core/connectors/runtime/src/api/mod.rs
+++ b/core/connectors/runtime/src/api/mod.rs
@@ -22,9 +22,14 @@ use axum::{Json, Router, extract::State, middleware,
routing::get};
use axum_server::tls_rustls::RustlsConfig;
use config::{HttpConfig, configure_cors};
use iggy_connector_sdk::api::ConnectorRuntimeStats;
-use std::{net::SocketAddr, path::PathBuf, sync::Arc};
-use tokio::spawn;
-use tracing::{error, info};
+use secrecy::ExposeSecret;
+use std::{
+ net::{IpAddr, SocketAddr},
+ path::PathBuf,
+ sync::Arc,
+};
+use tokio::{net::lookup_host, spawn};
+use tracing::{error, info, warn};
mod auth;
pub mod config;
@@ -34,6 +39,12 @@ mod sink;
mod source;
const NAME: &str = env!("CARGO_PKG_NAME");
+/// Where the full exposure surface is written down, so the startup warnings
+/// can name one thing each instead of restating it. A URL rather than a repo
+/// path: the published image carries the binary and its licences, nothing
else,
+/// so the reader most likely to see these lines has no source tree.
+const EXPOSURE_DOC: &str =
+
"https://github.com/apache/iggy/blob/master/core/connectors/runtime/README.md";
pub async fn init(config: &HttpConfig, context: Arc<RuntimeContext>) {
if !config.enabled {
@@ -41,6 +52,8 @@ pub async fn init(config: &HttpConfig, context:
Arc<RuntimeContext>) {
return;
}
+ warn_on_weak_containment(config).await;
+
let mut system_router = Router::new().route("/stats", get(get_stats));
if config.metrics.enabled {
@@ -121,6 +134,83 @@ pub async fn init(config: &HttpConfig, context:
Arc<RuntimeContext>) {
});
}
+/// Warns once for each configuration change that leaves this API without a
+/// compensating control: no key beyond loopback, no key with an unowned CORS
+/// origin, no TLS beyond loopback.
+///
+/// Separate warnings rather than one, because the three compose independently
+/// and an operator who closes one has not necessarily closed the others. Each
+/// carries a pointer instead of its own remediation: the README holds the full
+/// exposure surface, including the parts no startup warning can help with
+/// because they are already true of the shipped defaults.
+async fn warn_on_weak_containment(config: &HttpConfig) {
+ let unauthenticated = config.api_key.expose_secret().is_empty();
+ let beyond_loopback = resolves_beyond_loopback(&config.address).await;
+
+ if unauthenticated && beyond_loopback {
+ warn!(
+ "{NAME} HTTP API on {} has no api_key: anyone who reaches it can
read and rewrite every connector configuration. {EXPOSURE_DOC}",
+ config.address
+ );
+ }
+
+ // Loopback does not contain this one. A browser is a local process and the
+ // CORS layer wraps outside authentication, so the shipped
+ // `allowed_origins = ["*"]` lets any page the operator visits read these
+ // endpoints cross-origin. Pinning origins the operator owns closes that
+ // read, and an empty list emits no header at all, so neither is warned
+ // about. `null` is not such an origin, so a list holding it still is. No
+ // origin setting closes `POST .../restart`, which is CORS-simple and
+ // reaches the handler regardless; that one is true of the defaults, so the
+ // README carries it rather than a warning.
+ if unauthenticated && config.cors.enabled &&
config.cors.allows_unowned_origin() {
+ warn!(
+ "{NAME} HTTP API has http.cors open to * or null with no api_key:
any page the operator visits can read and rewrite its config. {EXPOSURE_DOC}"
+ );
+ }
+
+ if beyond_loopback && !config.tls.enabled {
+ warn!(
+ "{NAME} HTTP API on {} has http.tls disabled: the api-key header
and the credentials in its responses cross in cleartext. {EXPOSURE_DOC}",
+ config.address
+ );
+ }
+}
+
+/// Whether `address` resolves to anything outside loopback.
+///
+/// Resolves rather than parses because `address` is a free-form `String` that
+/// takes a hostname, as `[iggy] address` does in the same file. Not because
the
+/// default needs it: the embedded `config.toml` is the first figment layer, so
+/// the effective default is `127.0.0.1:8081` and would parse. An address that
+/// cannot resolve counts as exposed, since it is about to fail the bind anyway
+/// and staying quiet about one we could not classify is the wrong direction to
+/// be wrong in.
+///
+/// Classification only. Do not bind what this resolves: `TcpListener::bind`
+/// walks every resolved address and takes the first that works, so collapsing
+/// to one would drop the `localhost` -> `[::1, 127.0.0.1]` fallback on hosts
+/// with IPv6 disabled. The cost is resolving twice at startup, which is the
+/// trade for keeping that fallback.
+async fn resolves_beyond_loopback(address: &str) -> bool {
+ let Ok(resolved) = lookup_host(address).await else {
+ return true;
+ };
+ let addresses: Vec<SocketAddr> = resolved.collect();
+ // Empty is reported as exposed rather than confined: `all` over nothing is
+ // vacuously true, which would quietly invert the policy above.
+ addresses.is_empty()
+ || !addresses.iter().all(|address| match address.ip() {
+ // `Ipv6Addr::is_loopback` matches only `::1`, so an IPv4-mapped
+ // `::ffff:127.0.0.1` binds to 127.0.0.1 and would still read as
+ // exposed.
+ IpAddr::V6(v6) => v6
+ .to_ipv4_mapped()
+ .map_or(v6.is_loopback(), |v4| v4.is_loopback()),
+ ip => ip.is_loopback(),
+ })
+}
+
async fn get_metrics(State(context): State<Arc<RuntimeContext>>) -> String {
context.metrics.get_formatted_output()
}
@@ -128,3 +218,480 @@ async fn get_metrics(State(context):
State<Arc<RuntimeContext>>) -> String {
async fn get_stats(State(context): State<Arc<RuntimeContext>>) ->
Json<ConnectorRuntimeStats> {
Json(stats::get_runtime_stats(&context).await)
}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::configs::connectors::create_connectors_config_provider;
+ use crate::configs::runtime::{ConnectorsConfig, LocalConnectorsConfig};
+ use crate::manager::sink::SinkManager;
+ use crate::manager::source::SourceManager;
+ use crate::metrics::Metrics;
+ use crate::state::FileStateFactory;
+ use crate::stream::IggyClients;
+ use iggy::prelude::IggyClient;
+ use iggy_common::IggyTimestamp;
+ use secrecy::SecretString;
+ use std::fmt::Write as _;
+ use std::sync::Mutex;
+ use tempfile::TempDir;
+ use tracing::Level;
+ use tracing::field::{Field, Visit};
+ use tracing::subscriber::DefaultGuard;
+ use tracing_subscriber::Layer as _;
+ use tracing_subscriber::filter::LevelFilter;
+ use tracing_subscriber::layer::{Context as LayerContext, SubscriberExt};
+
+ /// Reserved for documentation by RFC 5737, so the bind fails and the tests
+ /// that use it reach the warning without listening anywhere. Hosts running
+ /// `net.ipv4.ip_nonlocal_bind=1`, which keepalived and haproxy boxes set,
+ /// bind it regardless; `routable_address_binds` excuses those.
+ const UNASSIGNABLE_ROUTABLE_ADDRESS: &str = "192.0.2.1:8081";
+
+ /// The pointer has to be fetchable by whoever reads the warning, and the
+ /// published image ships no source tree. Asserting the const against
itself
+ /// elsewhere cannot catch a value that stops being a URL.
+ #[test]
+ fn given_the_exposure_pointer_when_read_should_be_a_fetchable_url() {
+ assert!(
+ EXPOSURE_DOC.starts_with("https://"),
+ "a repo-relative path does not exist for an operator running the
image"
+ );
+ }
+ const EPHEMERAL_LOOPBACK_ADDRESS: &str = "127.0.0.1:0";
+
+ type Captured = Arc<Mutex<Vec<(Level, String)>>>;
+
+ fn config(address: &str, api_key: &str) -> HttpConfig {
+ HttpConfig {
+ address: address.to_owned(),
+ api_key: SecretString::from(api_key.to_owned()),
+ ..HttpConfig::default()
+ }
+ }
+
+ /// Captures events for the current thread only, for as long as the guard
+ /// lives. Not a global subscriber: that slot is process-wide, and a shared
+ /// buffer would leave negative assertions hostage to the rest of the
binary.
+ ///
+ /// `#[tokio::test]` builds a current-thread runtime, so a task spawned by
+ /// the test body sees this subscriber. Under a multi-thread flavour the
+ /// capture would come back empty and these tests would fail, not pass.
+ ///
+ /// Filtered rather than checking the level in `on_event`: a layer with no
+ /// filter reports no `max_level_hint`, which pushes the global max level
to
+ /// TRACE and stops every callsite in the binary short-circuiting.
+ fn capture_events() -> (DefaultGuard, Captured) {
+ let captured: Captured = Arc::new(Mutex::new(Vec::new()));
+ let layer = CaptureEvents {
+ captured: Arc::clone(&captured),
+ }
+ .with_filter(LevelFilter::INFO);
+ let guard =
tracing::subscriber::set_default(tracing_subscriber::registry().with(layer));
+ (guard, captured)
+ }
+
+ fn warnings(captured: &Captured) -> Vec<String> {
+ captured
+ .lock()
+ .expect("the capture mutex is only held to push a line")
+ .iter()
+ .filter(|(level, _)| *level == Level::WARN)
+ .map(|(_, message)| message.clone())
+ .collect()
+ }
+
+ /// Whether `init` got as far as serving. The positive control for tests
+ /// whose real assertion is that something was not warned about.
+ fn started_serving(captured: &Captured) -> bool {
+ logged_info(captured, "Started")
+ }
+
+ fn logged_info(captured: &Captured, needle: &str) -> bool {
+ captured
+ .lock()
+ .expect("the capture mutex is only held to push a line")
+ .iter()
+ .any(|(level, message)| *level == Level::INFO &&
message.contains(needle))
+ }
+
+ /// Whether this host binds an address RFC 5737 reserves for documentation.
+ ///
+ /// `net.ipv4.ip_nonlocal_bind=1`, which keepalived and haproxy boxes set,
+ /// makes it succeed. The two tests that read a failed bind as their
+ /// ordering control stand those two assertions down there rather than
going
+ /// red over a sysctl. They do not skip: a skipped test reports `ok` with
no
+ /// signal, and the warning itself is required on every host.
+ fn routable_address_binds() -> bool {
+ std::net::TcpListener::bind(UNASSIGNABLE_ROUTABLE_ADDRESS).is_ok()
+ }
+
+ struct CaptureEvents {
+ captured: Captured,
+ }
+
+ impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for
CaptureEvents {
+ fn on_event(&self, event: &tracing::Event<'_>, _context:
LayerContext<'_, S>) {
+ let mut recorded = Recorded(String::new());
+ event.record(&mut recorded);
+ self.captured
+ .lock()
+ .expect("the capture mutex is only held to push a line")
+ .push((*event.metadata().level(), recorded.0));
+ }
+ }
+
+ /// Unconditional on purpose: singling out the `message` field would add a
+ /// branch whose other side nothing here takes.
+ struct Recorded(String);
+
+ impl Visit for Recorded {
+ fn record_debug(&mut self, _field: &Field, value: &dyn
std::fmt::Debug) {
+ let _ = write!(self.0, "{value:?} ");
+ }
+ }
+
+ /// The cheapest context `init` will accept. Nothing here reaches Iggy.
+ ///
+ /// `api_key` is a parameter because the guard reads `config.api_key` while
+ /// the middleware enforces `context.api_key`; a test that set only one
+ /// would exercise a state the runtime cannot reach.
+ async fn context(api_key: &str) -> (Arc<RuntimeContext>, TempDir) {
+ let directory = tempfile::tempdir().expect("a temp dir must be
available");
+ let config_provider =
+
create_connectors_config_provider(&ConnectorsConfig::Local(LocalConnectorsConfig
{
+ config_dir: directory.path().display().to_string(),
+ }))
+ .await
+ .expect("an empty config dir must initialize with no connectors");
+
+ let context = RuntimeContext {
+ sinks: SinkManager::new(vec![]),
+ sources: SourceManager::new(vec![]),
+ api_key: SecretString::from(api_key.to_owned()),
+ config_provider: Arc::from(config_provider),
+ metrics: Arc::new(Metrics::init()),
+ start_time: IggyTimestamp::now(),
+ iggy_clients: Arc::new(IggyClients {
+ producer: IggyClient::default(),
+ consumer: IggyClient::default(),
+ }),
+ state_factory: Arc::new(FileStateFactory::new(
+ directory.path().display().to_string(),
+ )),
+ };
+ (Arc::new(context), directory)
+ }
+
+ #[tokio::test]
+ async fn
given_loopback_addresses_when_classified_should_report_contained() {
+ assert!(!resolves_beyond_loopback("127.0.0.1:8081").await);
+ assert!(!resolves_beyond_loopback("[::1]:8081").await);
+ assert!(
+ !resolves_beyond_loopback("localhost:8081").await,
+ "`address` accepts a hostname, and parsing alone would misjudge
one"
+ );
+ assert!(
+ !resolves_beyond_loopback("[::ffff:127.0.0.1]:8081").await,
+ "an IPv4-mapped address binds to 127.0.0.1, and
`Ipv6Addr::is_loopback` \
+ alone matches only `::1`"
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_routable_or_unresolvable_addresses_when_classified_should_report_exposed()
{
+ assert!(
+ resolves_beyond_loopback("0.0.0.0:8081").await,
+ "binding every interface to reach the API from outside a container
\
+ is the case this exists to catch"
+ );
+ assert!(resolves_beyond_loopback("192.0.2.10:8081").await);
+ // About to fail the bind regardless, so staying quiet about an address
+ // we cannot classify is the wrong direction to be wrong in.
+ assert!(resolves_beyond_loopback("not a valid address").await);
+ }
+
+ #[tokio::test]
+ async fn
given_no_key_and_a_routable_address_when_initialized_should_warn_before_binding()
{
+ let bind_is_refused = !routable_address_binds();
+ // Fixture first: `create_connectors_config_provider` logs, and
anything
+ // it warns about later would otherwise land in this buffer.
+ let (context, _directory) = context("").await;
+ let (_capture, captured) = capture_events();
+ let config = config(UNASSIGNABLE_ROUTABLE_ADDRESS, "");
+
+ // `init` panics when the bind fails, which is what makes this the
+ // ordering test: the warning has to already be out by then, or an
+ // operator whose bind fails never learns the API was unauthenticated.
+ let bind_failed = tokio::spawn(async move { init(&config,
context).await })
+ .await
+ .is_err();
+
+ if bind_is_refused {
+ assert!(
+ bind_failed,
+ "a documentation-range address must not bind here"
+ );
+ assert!(
+ !started_serving(&captured),
+ "the bind must not have completed, or this proves nothing
about ordering"
+ );
+ }
+ // Keyless and untrusting on a routable address, so the address warning
+ // and the cleartext one both fire. Asserted with `all` and a count
+ // rather than `any`: either message could stop naming the address and
+ // the other would still satisfy an `any`.
+ let warnings = warnings(&captured);
+ assert_eq!(
+ warnings.len(),
+ 2,
+ "init must consult the guard for both the missing key and the
missing \
+ TLS: {warnings:?}"
+ );
+ assert!(
+ warnings
+ .iter()
+ .all(|warning|
warning.contains(UNASSIGNABLE_ROUTABLE_ADDRESS)),
+ "each names the address it is exposing: {warnings:?}"
+ );
+ assert!(
+ warnings
+ .iter()
+ .all(|warning| warning.contains(EXPOSURE_DOC)),
+ "each carries the pointer instead of its own remediation:
{warnings:?}"
+ );
+ assert!(
+ warnings
+ .iter()
+ .any(|warning| warning.contains("read and rewrite")),
+ "the exposure is not read-only, and saying so is the point:
{warnings:?}"
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_loopback_address_when_initialized_should_serve_without_warning() {
+ let (context, _directory) = context("").await;
+ let (_capture, captured) = capture_events();
+
+ init(&config(EPHEMERAL_LOOPBACK_ADDRESS, ""), context).await;
+
+ // Positive control first: without it the assertion below passes for
any
+ // reason `init` returns early, including `enabled` ever defaulting to
+ // false, and the only in-`init` loopback coverage disappears silently.
+ assert!(
+ started_serving(&captured),
+ "init must reach the listener, or the assertion below proves
nothing"
+ );
+ // Any warning at all, not one matching this address: matching on the
+ // address goes vacuous the moment the message is reworded.
+ assert!(
+ warnings(&captured).is_empty(),
+ "the shipped posture is loopback with no key; warning about it
would \
+ teach operators to ignore the ones that matter: {:?}",
+ warnings(&captured)
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_wildcard_cors_and_no_key_when_initialized_should_warn_despite_loopback() {
+ let (context, _directory) = context("").await;
+ let (_capture, captured) = capture_events();
+ let mut config = config(EPHEMERAL_LOOPBACK_ADDRESS, "");
+ config.cors.enabled = true;
+ config.cors.allowed_origins = vec!["*".to_owned()];
+
+ init(&config, context).await;
+
+ assert!(started_serving(&captured));
+ assert!(
+ warnings(&captured)
+ .iter()
+ .any(|warning| warning.contains("http.cors")),
+ "loopback does not contain CORS: a browser is a local process and
the \
+ layer wraps outside authentication"
+ );
+ assert!(
+ warnings(&captured)
+ .iter()
+ .all(|warning| warning.contains(EXPOSURE_DOC)),
+ "each warning carries the pointer instead of its own remediation:
{:?}",
+ warnings(&captured)
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_owned_cors_origins_when_initialized_should_not_warn_about_cors() {
+ let (context, _directory) = context("").await;
+ let (_capture, captured) = capture_events();
+ let mut config = config(EPHEMERAL_LOOPBACK_ADDRESS, "");
+ config.cors.enabled = true;
+ config.cors.allowed_origins =
vec!["https://console.example".to_owned()];
+
+ init(&config, context).await;
+
+ assert!(
+ started_serving(&captured),
+ "init must reach the listener, or the assertion below proves
nothing"
+ );
+ assert!(
+ warnings(&captured).is_empty(),
+ "an operator who pinned origins they own closed the cross-origin
read, \
+ and warning at them is what teaches everyone to ignore the rest:
{:?}",
+ warnings(&captured)
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_a_null_cors_origin_when_initialized_should_warn_despite_the_pin() {
+ let (context, _directory) = context("").await;
+ let (_capture, captured) = capture_events();
+ let mut config = config(EPHEMERAL_LOOPBACK_ADDRESS, "");
+ config.cors.enabled = true;
+ config.cors.allowed_origins = vec!["null".to_owned()];
+
+ init(&config, context).await;
+
+ assert!(started_serving(&captured));
+ assert!(
+ warnings(&captured)
+ .iter()
+ .any(|warning| warning.contains("http.cors")),
+ "`null` is a pinned origin that nobody owns, so pinning it buys
nothing \
+ a wildcard would not have given away"
+ );
+ assert!(
+ warnings(&captured)
+ .iter()
+ .all(|warning| warning.contains(EXPOSURE_DOC)),
+ "each warning carries the pointer instead of its own remediation:
{:?}",
+ warnings(&captured)
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_a_key_but_no_tls_beyond_loopback_when_initialized_should_still_warn() {
+ let bind_is_refused = !routable_address_binds();
+ let (context, _directory) = context("configured").await;
+ let (_capture, captured) = capture_events();
+ let config = config(UNASSIGNABLE_ROUTABLE_ADDRESS, "configured");
+
+ let bind_failed = tokio::spawn(async move { init(&config,
context).await })
+ .await
+ .is_err();
+
+ if bind_is_refused {
+ assert!(
+ bind_failed,
+ "a documentation-range address must not bind here"
+ );
+ }
+ let warnings = warnings(&captured);
+ assert!(
+ warnings.iter().any(|warning| warning.contains("http.tls")),
+ "setting a key does not stop the key and the credential-bearing \
+ responses crossing the network in cleartext"
+ );
+ assert_eq!(
+ warnings.len(),
+ 1,
+ "the key that was set must not still be reported as missing, and a
\
+ count cannot go vacuous the way matching on the old wording
would: \
+ {warnings:?}"
+ );
+ assert!(
+ warnings
+ .iter()
+ .all(|warning| warning.contains(EXPOSURE_DOC)),
+ "each warning carries the pointer instead of its own remediation:
{warnings:?}"
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_wildcard_cors_and_a_key_when_initialized_should_not_warn_about_cors() {
+ let (context, _directory) = context("configured").await;
+ let (_capture, captured) = capture_events();
+ let mut config = config(EPHEMERAL_LOOPBACK_ADDRESS, "configured");
+ config.cors.enabled = true;
+ config.cors.allowed_origins = vec!["*".to_owned()];
+
+ init(&config, context).await;
+
+ assert!(started_serving(&captured));
+ assert!(
+ warnings(&captured).is_empty(),
+ "an attacker's page cannot supply the api-key header, so a key
closes \
+ the cross-origin path a wildcard opens: {:?}",
+ warnings(&captured)
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_cors_switched_off_when_initialized_should_not_warn_about_its_origins() {
+ let (context, _directory) = context("").await;
+ let (_capture, captured) = capture_events();
+ // The shipped block verbatim: a wildcard list that is not switched on.
+ // Warning here would fire on every stock deployment.
+ let mut config = config(EPHEMERAL_LOOPBACK_ADDRESS, "");
+ config.cors.allowed_origins = vec!["*".to_owned()];
+
+ init(&config, context).await;
+
+ assert!(started_serving(&captured));
+ assert!(
+ warnings(&captured).is_empty(),
+ "no CORS layer is installed unless http.cors is enabled, so the
origin \
+ list allows nothing on its own: {:?}",
+ warnings(&captured)
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_tls_enabled_beyond_loopback_when_initialized_should_not_warn_about_cleartext()
{
+ let (context, _directory) = context("configured").await;
+ let (_capture, captured) = capture_events();
+ let mut config = config(UNASSIGNABLE_ROUTABLE_ADDRESS, "configured");
+ config.tls.enabled = true;
+
+ // The TLS branch loads a certificate this config does not name, so
+ // `init` panics there. That is after the warning block, which is what
+ // makes the panic the positive control.
+ let reached_the_certificate = tokio::spawn(async move { init(&config,
context).await })
+ .await
+ .is_err();
+
+ assert!(
+ reached_the_certificate,
+ "init must get past the warning block, or the assertion below
proves nothing"
+ );
+ assert!(
+ warnings(&captured).is_empty(),
+ "TLS on the listener is exactly what the cleartext warning asks
for: {:?}",
+ warnings(&captured)
+ );
+ }
+
+ #[tokio::test]
+ async fn given_a_disabled_api_when_initialized_should_warn_about_nothing()
{
+ let (context, _directory) = context("").await;
+ let (_capture, captured) = capture_events();
+ // Routable, keyless and untrusting in every direction, but switched
off.
+ let mut config = config(UNASSIGNABLE_ROUTABLE_ADDRESS, "");
+ config.enabled = false;
+ config.cors.enabled = true;
+
+ init(&config, context).await;
+
+ // Positive control: without it this passes for any reason `init`
returns
+ // early, including one that never reaches the guard at all.
+ assert!(
+ logged_info(&captured, "HTTP API is disabled"),
+ "init must take the disabled path, or the assertion below proves
nothing"
+ );
+ assert!(
+ warnings(&captured).is_empty(),
+ "an API that is not listening exposes nothing, and warning about
one \
+ is the false positive that teaches operators to ignore the rest:
{:?}",
+ warnings(&captured)
+ );
+ }
+}