diff --git a/CHANGELOG.md b/CHANGELOG.md index b9365d0..f1587e5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,12 +10,26 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added - Optional `mimalloc` feature that uses mimalloc as the global allocator. Enabled by default in the Docker images to reduce memory growth in long-running deployments (#68) +- Opt-in access-audit log (`telemetry.audit.enabled`, disabled by default): one JSON line per HTTP request on stdout + with the user and subject verified by an authenticating reverse proxy (`telemetry.audit.user-header` and + `telemetry.audit.subject-header`, default `X-Forwarded-Email` and `X-Forwarded-User`), source address, method, + path, DICOM coordinates, status and duration. With auditing enabled, the C-MOVE completion log line also carries the + Study Instance UID (#62) +- With auditing enabled, line breaks inside a regular log message are escaped, so request data carried by a log message + cannot forge a line on the stdout stream that also carries the audit records (#62) +- Trusted relays for the access-audit log (`telemetry.audit.trusted-relays`, `telemetry.audit.on-behalf-of-header`): + explicitly trusted callers can name the end user they act for, recorded as `on_behalf_of` next to the caller's own + identity (a recorded claim, never used for authorization); claims from other callers or malformed claims are + recorded as `on_behalf_of_rejected` without the claimed value (#62) +- `request_id` in the access-audit record, taken from `X-Request-Id`, for correlation with proxy access logs (#62) ### Fixed ### Changed - Updated `dicom-rs` dependency to 0.10.0 +- ANSI colors in the log output are disabled when stdout is not a terminal (e.g. in containers), so collected logs stay + machine-parseable. `NO_COLOR` is still honoured on terminals (#62) ## [0.3.1] - 2026-09-14 diff --git a/Cargo.lock b/Cargo.lock index a0db6fc..73e9661 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1711,6 +1711,7 @@ dependencies = [ "axum-extra", "axum-streams", "bytes", + "chrono", "config", "dicom", "dicom-json", diff --git a/Cargo.toml b/Cargo.toml index 635c746..6e5d0a9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -32,6 +32,8 @@ serde = { version = "1.0.228", features = ["derive"] } serde_json = "1.0.145" # Logging tracing = "0.1.41" +# Timestamps for the audit log; already transitive via dicom-core. +chrono = { version = "0.4.42", default-features = false, features = ["clock"] } tracing-subscriber = { version = "0.3.20", features = ["env-filter"] } # Convenient error handling thiserror = "2.0.17" @@ -66,6 +68,8 @@ aws-credential-types = { version = "1.2.13", optional = true } testcontainers = "0.27.3" dicom-test-files = "0.4.0" dicom-web = "0.6.0" +# `ServiceExt::oneshot` for middleware tests +tower = { version = "0.5.2", features = ["util"] } [lints.rust] unsafe_code = "forbid" diff --git a/docs/topics/configuration.md b/docs/topics/configuration.md index 547443f..af18a3f 100644 --- a/docs/topics/configuration.md +++ b/docs/topics/configuration.md @@ -81,6 +81,8 @@ aets: telemetry: sentry: https://sentry.local/dsn level: INFO + audit: + enabled: false ``` @@ -98,8 +100,177 @@ telemetry:
  • TRACE
  • + + Structured access-audit logging, disabled by default. + See Access Audit Config. + +
    + +## Access Audit Config {id="access-audit-config"} + +%product% performs no authentication itself. When it runs behind an authenticating reverse proxy, it can write one +access-audit record per HTTP request, naming the user the proxy verified and the DICOM resources that were accessed. +Read the trust model before relying on the identity fields. +All settings are optional; with enabled: false (the default) nothing changes. + +```yaml +telemetry: + audit: + enabled: true + user-header: X-Forwarded-Email + subject-header: X-Forwarded-User + trusted-relays: [] + on-behalf-of-header: X-On-Behalf-Of +``` + + + + Enables the access-audit log (default false). + Every HTTP request then emits one self-contained JSON line on stdout. + Delivery is fail-open: records pass through a bounded buffer to a dedicated writer thread, so a slow log + consumer never blocks requests. When the buffer is full, records are dropped and counted; a warning with the + running count is logged for the first drop and then for every 100th. + On a graceful shutdown (server.http.graceful-shutdown, enabled by default) %product% waits up to + 5 seconds for buffered records to be written. Without a graceful shutdown, buffered records are lost. + + + The request header carrying the user verified by the proxy, recorded as user + (default X-Forwarded-Email). + This is also the identity that trusted-relays is matched against. + + + The request header carrying the subject identifier verified by the proxy, recorded as subject + (default X-Forwarded-User). + + + Identities, exactly as the proxy asserts them in user-header, that may name the end user they act + for (default: empty, nobody may). + Use this for services that call %product% with their own credentials on behalf of a signed-in user. + Matching is ASCII case-insensitive. Empty entries are rejected at startup. + + + The request header in which a trusted relay names the end user (default X-On-Behalf-Of). + It is only read when trusted-relays is not empty, and must differ from + user-header and subject-header. + +Header names are validated when the configuration is loaded; an invalid name, or a header that carries credentials +(Authorization, Proxy-Authorization, Cookie, +X-Forwarded-Access-Token), stops %product% at startup with an error naming the offending key. The credential +check is a guard against an obvious misconfiguration, not an exhaustive list: never point an identity header at a header +that carries a secret. + +### Audit Record + +```json +{"audit":"http-access","ts":"2026-08-17T17:16:55Z","user":"jane.doe@example.org","subject":"3f2c9a4e","source":"192.0.2.10","method":"GET","path":"/aets/PACS/studies/1.2.3.4","aet":"PACS","study":"1.2.3.4","status":200,"duration_ms":4886,"request_id":"0f8c2b7e"} +``` + +| Field | Description | +|-------------------------|-------------------------------------------------------------------------------------------| +| `audit` | Always `http-access`. | +| `ts` | Time the response head was produced (UTC, RFC 3339, second precision). | +| `user`, `subject` | Values of `user-header` and `subject-header`: 1 to 320 bytes of UTF-8 without control characters. Omitted if absent, sent more than once or malformed (never truncated). | +| `on_behalf_of` | The end user named by a trusted relay: a claim, not a verified identity (see below). | +| `on_behalf_of_rejected` | Why an on-behalf-of header was ignored: `untrusted-caller` or `invalid`. | +| `source` | Leftmost entry of `X-Forwarded-For`, at most 64 bytes. | +| `method`, `path` | Request method, and path with query string (at most 8 KiB). | +| `aet`, `study`, `series`, `instance` | DICOM coordinates from the request path, at most 256 bytes each. If any path parameter cannot be decoded, all four are omitted; `path` is still recorded. | +| `status`, `duration_ms` | Status and elapsed time when the response head was produced, including `408` for timed-out requests. | +| `user_agent` | `User-Agent` header, at most 512 bytes. | +| `request_id` | `X-Request-Id` header, if sent once with 1 to 128 printable ASCII characters (no spaces). | + +Absent values are omitted. Values longer than their limit are cut at a character boundary and end in `…`; the limits +include the marker. +With auditing enabled, the log line of a completed C-MOVE also carries the study_uid, and line breaks inside +a regular log message are escaped (\n, \r): audit records and the regular log share stdout, and a +message that carries request data (a percent-decoded path segment can contain a newline) must never start a line of its +own that looks like an audit record. With auditing disabled, the log output is unchanged. The guarantee is about +\n and \r; a consumer that also splits on Unicode line separators (U+2028, U+2029, U+0085) is not +covered. + +What the record does and does not show: + +- `source` and `request_id` are copied from the incoming request. Both are whatever the client sent, unless the + proxies in front of %product% overwrite them. +- `status` and `duration_ms` are taken when the response head is produced. A streamed retrieve that fails after that + is still recorded with the status of its head. +- A client that disconnects before the response head is produced, or a request whose handler panics, can leave no + record at all. +- A record without `user` behind an authenticating proxy means the identity header did not arrive. Alert on such + records: besides misconfiguration, a client can make some proxies drop headers they inject by listing them as + hop-by-hop headers in `Connection`; the reverse proxy in Go's standard library, for example, removes every header + named there. This erases the identity; it cannot replace it with a chosen one. + +### Trust Model {id="audit-trust-model"} + +The identity fields are only as trustworthy as the deployment around %product%. All of the following must hold: + + + +
  • The authenticating proxy removes any user-header and subject-header a client + sends and sets them itself from the verified session. It must not remove the + on-behalf-of-header: a relay is itself a client of the proxy, and removing the header would switch + the feature off.
  • +
  • Every trusted relay sets the on-behalf-of-header itself, overwriting any existing value, to the + end user of its own verified session, and never forwards a copy it received from its own clients. Otherwise any + user of the relay can name someone else.
  • +
  • %product% is reachable only through the proxy, for example by binding server.http.interface to + 127.0.0.1 next to a sidecar proxy, or with a firewall or network policy. Anyone who can reach the + port directly can send any header.
  • +
    +
    + +A relay is trusted because of the identity the proxy verified for it, never because of a header it sends itself. +An on-behalf-of claim is recorded as on_behalf_of only if the request's user is listed in +trusted-relays, the header occurs exactly once, and its value is 1 to 320 bytes of UTF-8 without +whitespace or control characters. Otherwise the record carries on_behalf_of_rejected +(untrusted-caller or invalid) and the claimed value is not logged. + +on_behalf_of is an unverifiable claim made by an authenticated relay, recorded next to the relay's own +identity in user. %product% cannot check it, and it never grants or restricts access. + +### Example: oauth2-proxy + +In oauth2-proxy v7.15.3 (pkg/middleware/headers.go, pkg/apis/options/legacy_options.go), +pass_user_headers (enabled by default) removes client-supplied X-Forwarded-User, +X-Forwarded-Email, X-Forwarded-Groups, X-Forwarded-Preferred-Username and +X-Forwarded-Access-Token from the request before setting them from the session, which makes the default +user-header and subject-header suitable. The X-Auth-Request-* headers are response +headers only: they are neither set on nor removed from the upstream request, so they must not be used as identity +headers. Verify that the proxy version you run behaves the same. + +An oauth2-proxy configuration (excerpt) in front of %product%, accepting both interactive users and services that +present their own bearer token: + +```toml +provider = "oidc" +oidc_issuer_url = "https://idp.example.org" +upstreams = ["http://127.0.0.1:8080/"] +email_domains = ["*"] +# Default: sets X-Forwarded-User/-Email from the session, replacing client copies +pass_user_headers = true +# Lets services call with a bearer token issued by the same provider +skip_jwt_bearer_tokens = true +``` + +For a service, the proxy fills the identity headers from the claims of its token, so with the default +user-header the relay's token needs an e-mail claim. +With the matching %product% configuration, requests from viewer@example.org may name the end user in +X-On-Behalf-Of: + +```yaml +server: + http: + interface: 127.0.0.1 +telemetry: + audit: + enabled: true + trusted-relays: + - viewer@example.org +``` + ## Global Server Config ```yaml diff --git a/src/audit.rs b/src/audit.rs new file mode 100644 index 0000000..72156df --- /dev/null +++ b/src/audit.rs @@ -0,0 +1,1238 @@ +//! Structured access-audit logging. +//! +//! When enabled (`telemetry.audit.enabled: true`), every HTTP request emits +//! one self-contained JSON line on stdout describing WHO accessed WHAT: +//! +//! ```json +//! {"audit":"http-access","ts":"2026-08-17T17:16:55Z","user":"jane.doe@example.org", +//! "subject":"3f2c9a4e-…","source":"192.0.2.10","method":"GET", +//! "path":"/aets/PACS/studies/1.2.3.4","aet":"PACS", +//! "study":"1.2.3.4","status":200,"duration_ms":4886} +//! ``` +//! +//! Identity is read from the request headers named by +//! `telemetry.audit.user-header` (default `X-Forwarded-Email`) and +//! `telemetry.audit.subject-header` (default `X-Forwarded-User`). +//! DICOM-RST itself performs no authentication (see #15/#42): these fields +//! are TRUSTWORTHY ONLY when an authenticating proxy replaces any client +//! copies of these headers with values from its verified session and is the +//! only way to reach DICOM-RST. Whether a given proxy does so, and the rest +//! of the trust model, is documented in the "Access Audit Config" section +//! of `docs/topics/configuration.md`. A header that occurs more than once is +//! ambiguous and treated as absent, as is a value that is not 1..=320 bytes +//! of UTF-8 without control characters (never truncated: a shortened +//! identity could equal someone else's). The record is emitted regardless — +//! an absent identity is itself audit-relevant. +//! +//! A caller that acts for someone else (e.g. a backend service fetching +//! images for a signed-in user) can name that end user in the header set by +//! `telemetry.audit.on-behalf-of-header` (default `X-On-Behalf-Of`). The +//! claim is recorded as `on_behalf_of` only if the caller's own verified +//! identity (`user`) is listed in `telemetry.audit.trusted-relays` (ASCII +//! case-insensitive), the header occurs exactly once, and its value is +//! 1..=320 bytes of UTF-8 without whitespace or control characters. +//! Otherwise `on_behalf_of_rejected` says why (`"untrusted-caller"` or +//! `"invalid"`) and the claimed value is NOT recorded. With no trusted +//! relays configured (the default) the header is not read at all and +//! neither field ever appears. `on_behalf_of` is an unverifiable claim by an +//! authenticated relay, recorded beside the relay's own identity; it is +//! never used for authorization. +//! +//! `source` is the leftmost `X-Forwarded-For` entry and `request_id` the +//! incoming `X-Request-Id` (1..=128 printable ASCII characters without +//! space, sent once): both are client-asserted unless the proxy chain +//! overwrites them. `request_id` lets a record be correlated with the +//! access log of the proxy or ingress that set it. `path` (8 KiB), +//! `user_agent` (512 bytes), `source` (64 bytes) and each DICOM coordinate +//! (256 bytes) are capped; a cut value ends in `…` and, marker included, +//! stays within its cap. If a path parameter cannot be decoded, the record carries +//! no DICOM coordinates at all; `path` is still recorded. +//! +//! `status` and `duration_ms` are taken when the response head is produced, +//! so a streamed retrieve that fails mid-body is recorded with the status of +//! its head. A client that disconnects before the head, or a handler panic, +//! can leave no record. +//! +//! Delivery is FAIL-OPEN by design: records flow through a bounded channel +//! to a dedicated writer thread (not a Tokio task, so a stalled stdout never +//! ties up a runtime worker); when the buffer is full the record is dropped +//! and counted, and a warning with the running count is logged for the first +//! drop and every 100th — a slow disk or collector never blocks request +//! handling. Deployments with stricter requirements should alert on the drop +//! warnings. After a graceful shutdown, [`AuditWriter::finish`] waits a +//! bounded time for buffered records to be written; without a graceful +//! shutdown they are lost. + +use std::io::Write; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::Arc; +use std::thread::JoinHandle; +use std::time::{Duration, Instant}; + +use axum::extract::{RawPathParams, Request, State}; +use axum::http::{HeaderMap, HeaderName, HeaderValue}; +use axum::middleware::Next; +use axum::response::Response; +use axum::RequestExt; +use chrono::{SecondsFormat, Utc}; +use serde::Serialize; +use tokio::sync::mpsc; +use tracing::warn; + +use crate::config::AuditConfig; + +/// One audit record per HTTP request. +#[derive(Debug, Serialize)] +pub struct AuditRecord { + /// Discriminator for log pipelines; always `"http-access"` for now. + pub audit: &'static str, + /// Wall-clock request completion time (UTC, RFC 3339, second precision). + pub ts: String, + /// Proxy-verified user from `telemetry.audit.user-header`, if present. + #[serde(skip_serializing_if = "Option::is_none")] + pub user: Option, + /// Proxy-verified subject from `telemetry.audit.subject-header`, if present. + #[serde(skip_serializing_if = "Option::is_none")] + pub subject: Option, + /// End user named by a trusted relay (see the module docs). + #[serde(skip_serializing_if = "Option::is_none")] + pub on_behalf_of: Option, + /// Why an on-behalf-of header was ignored. The ignored value itself is + /// never recorded: it came from a caller that may not make the claim. + #[serde(skip_serializing_if = "Option::is_none")] + pub on_behalf_of_rejected: Option, + /// First `X-Forwarded-For` entry, if present (at most 64 bytes, marker + /// included). + #[serde(skip_serializing_if = "Option::is_none")] + pub source: Option, + pub method: String, + /// Full request path and query (at most 8 KiB, marker included). QIDO + /// match parameters are part of "which data was accessed" and are + /// deliberately included. + pub path: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub aet: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub study: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub series: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub instance: Option, + pub status: u16, + pub duration_ms: u128, + /// `User-Agent`, if present (at most 512 bytes, marker included). + #[serde(skip_serializing_if = "Option::is_none")] + pub user_agent: Option, + /// `X-Request-Id`, if sent once as 1..=128 printable ASCII characters + /// without space. + #[serde(skip_serializing_if = "Option::is_none")] + pub request_id: Option, +} + +/// Why an on-behalf-of header was not honoured. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "kebab-case")] +pub enum OnBehalfOfRejection { + /// The caller is unidentified or not a configured trusted relay. + UntrustedCaller, + /// A trusted relay sent the header more than once or a malformed value. + Invalid, +} + +/// Outcome of evaluating the on-behalf-of header for one request. +#[derive(Debug)] +enum Delegation { + /// Nothing to decide: no trusted relays configured, or no header sent. + NotClaimed, + Honoured(String), + Rejected(OnBehalfOfRejection), +} + +impl Delegation { + fn into_fields(self) -> (Option, Option) { + match self { + Self::NotClaimed => (None, None), + Self::Honoured(end_user) => (Some(end_user), None), + Self::Rejected(reason) => (None, Some(reason)), + } + } +} + +/// Longest accepted identity: room for the longest e-mail address +/// (64-octet local part, `@`, 255-octet domain). +const MAX_IDENTITY_LEN: usize = 320; + +const X_REQUEST_ID: HeaderName = HeaderName::from_static("x-request-id"); +const MAX_REQUEST_ID_LEN: usize = 128; + +// Caps on values copied from the request as-is, so that one oversized +// request cannot produce an oversized audit line. Longer values are cut at +// a character boundary and end in `TRUNCATED`, the marker included in the +// cap. +const MAX_PATH_LEN: usize = 8 * 1024; +const MAX_USER_AGENT_LEN: usize = 512; +const MAX_SOURCE_LEN: usize = 64; +/// Per DICOM coordinate (`aet`, `study`, `series`, `instance`): far above +/// any valid AE title (16) or UID (64), so only garbage is ever cut. +const MAX_COORDINATE_LEN: usize = 256; +const TRUNCATED: &str = "…"; + +/// The human-readable log's writer while auditing is enabled: stdout, with +/// every line break inside one formatted log event escaped (`\n`, `\r`). +/// +/// Audit records and the regular log share stdout, and a log message may +/// carry request data (a percent-decoded path segment can hold a newline). +/// Unescaped, such a message could start a line of its own that looks like +/// an audit record. tracing-subscriber formats each event into a buffer and +/// hands it to the writer with one `write_all` (`fmt_layer.rs`, 0.3.20), so +/// only the event's final line break is kept; the test +/// `a_newline_ending_a_format_argument_is_escaped_too` fails if a future +/// version ever streams an event in pieces. The guarantee is about `\n` and +/// `\r`: consumers that also split on Unicode line separators (U+2028, +/// U+2029, U+0085) are not covered. Not used while auditing is disabled, +/// which leaves the log output unchanged. +pub struct LineSafeStdout; + +impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for LineSafeStdout { + type Writer = LineSafe; + + fn make_writer(&'a self) -> Self::Writer { + LineSafe(std::io::stdout()) + } +} + +/// See [`LineSafeStdout`]. +pub struct LineSafe(pub W); + +impl Write for LineSafe { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + let (body, end): (&[u8], &[u8]) = match buf.split_last() { + Some((b'\n', body)) => (body, b"\n"), + _ => (buf, b""), + }; + let mut escaped = Vec::with_capacity(buf.len() + 8); + for &byte in body { + match byte { + b'\n' => escaped.extend_from_slice(br"\n"), + b'\r' => escaped.extend_from_slice(br"\r"), + other => escaped.push(other), + } + } + escaped.extend_from_slice(end); + self.0.write_all(&escaped)?; + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + self.0.flush() + } +} + +/// Cloneable handle to the audit writer. Only exists while auditing is +/// enabled. +#[derive(Clone)] +pub struct AuditSink { + tx: mpsc::Sender, + config: Arc, +} + +/// Records dropped because the buffer was full (fail-open pressure valve). +static DROPPED: AtomicU64 = AtomicU64::new(0); + +const BUFFER: usize = 1024; + +/// Start the stdout writer thread and return a sink feeding it, or `None` +/// when auditing is disabled. Without a sink the middleware must not be +/// installed at all, so the request path is exactly that of a build without +/// auditing. +/// +/// # Errors +/// Returns an error if the writer thread cannot be spawned. +pub fn start(config: &AuditConfig) -> std::io::Result> { + if !config.enabled { + return Ok(None); + } + let (tx, rx) = mpsc::channel::(BUFFER); + // `Stdout::write_all` takes the lock once per call, and the regular log + // also writes each event with one `write_all` (see `LineSafeStdout`), so + // audit records and log lines never interleave mid-line. + let writer = spawn_writer(rx, std::io::stdout())?; + let sink = AuditSink { + tx, + config: Arc::new(config.clone()), + }; + Ok(Some((sink, writer))) +} + +/// Handle to the thread that writes audit records. +pub struct AuditWriter { + thread: JoinHandle<()>, +} + +impl AuditWriter { + /// Waits at most `timeout` for the writer to finish, which it does once + /// every [`AuditSink`] clone is dropped and all buffered records are + /// written. Returns whether it finished cleanly. + pub fn finish(self, timeout: Duration) -> bool { + let deadline = Instant::now() + timeout; + while !self.thread.is_finished() { + if Instant::now() >= deadline { + return false; + } + std::thread::sleep(Duration::from_millis(10)); + } + self.thread.join().is_ok() + } +} + +fn spawn_writer(mut rx: mpsc::Receiver, mut out: W) -> std::io::Result +where + W: Write + Send + 'static, +{ + let thread = std::thread::Builder::new() + .name("audit-writer".to_owned()) + .spawn(move || { + while let Some(record) = rx.blocking_recv() { + match serde_json::to_string(&record) { + Ok(mut line) => { + line.push('\n'); + let _ = out.write_all(line.as_bytes()); + } + Err(err) => warn!("failed to serialize audit record: {err}"), + } + } + let _ = out.flush(); + })?; + Ok(AuditWriter { thread }) +} + +impl AuditSink { + /// An enabled sink whose records are handed to the caller instead of + /// being written to stdout. + #[cfg(test)] + fn with_receiver(config: AuditConfig) -> (Self, mpsc::Receiver) { + let (tx, rx) = mpsc::channel(BUFFER); + let sink = Self { + tx, + config: Arc::new(config), + }; + (sink, rx) + } + + fn emit(&self, record: AuditRecord) { + if self.tx.try_send(record).is_err() { + let dropped = DROPPED.fetch_add(1, Ordering::Relaxed) + 1; + // Every drop is a warning-worthy event, but do not spam a + // saturated system: log the first and then every 100th. + if dropped == 1 || dropped.is_multiple_of(100) { + warn!("audit buffer full: {dropped} record(s) dropped so far (fail-open)"); + } + } + } +} + +/// Axum middleware producing one [`AuditRecord`] per request. +/// +/// Attach with `axum::middleware::from_fn_with_state(sink, audit::middleware)` +/// OUTSIDE the timeout layer, so timed-out requests are recorded with their +/// 408 as well. +/// +/// The middleware only observes: it never answers a request itself. Path +/// parameters are therefore read without a rejecting extractor — a request +/// whose parameters cannot be decoded is passed on unchanged and still +/// audited, with its path but without any of the DICOM coordinates. +pub async fn middleware( + State(sink): State, + mut request: Request, + next: Next, +) -> Response { + let params = request.extract_parts::().await.ok(); + + let mut aet = None; + let mut study = None; + let mut series = None; + let mut instance = None; + for (name, value) in params.iter().flatten() { + match name { + "aet" => aet = Some(bounded(value.to_owned(), MAX_COORDINATE_LEN)), + "study" => study = Some(bounded(value.to_owned(), MAX_COORDINATE_LEN)), + "series" => series = Some(bounded(value.to_owned(), MAX_COORDINATE_LEN)), + "instance" => instance = Some(bounded(value.to_owned(), MAX_COORDINATE_LEN)), + _ => {} + } + } + + // Extract everything BEFORE the await, inside a block that ends first: + // a closure borrowing `&Request` held across `next.run().await` makes + // the future `!Send` (axum's `Body` is `!Sync`), failing the middleware + // `Service` bound with a famously opaque error. + let (user, subject, delegation, source, user_agent, request_id) = { + let headers = request.headers(); + let get = |name: &str| { + headers + .get(name) + .and_then(|value| value.to_str().ok()) + .map(str::to_owned) + }; + let user = identity(headers, &sink.config.user_header); + let delegation = delegation_for(&sink.config, headers, user.as_deref()); + ( + user, + identity(headers, &sink.config.subject_header), + delegation, + get("x-forwarded-for").map(|forwarded| { + let first = forwarded.split(',').next().unwrap_or_default(); + bounded(first.trim().to_owned(), MAX_SOURCE_LEN) + }), + get("user-agent").map(|user_agent| bounded(user_agent, MAX_USER_AGENT_LEN)), + request_id(headers), + ) + }; + let method = request.method().to_string(); + let path = bounded(request.uri().to_string(), MAX_PATH_LEN); + + let (on_behalf_of, on_behalf_of_rejected) = delegation.into_fields(); + + let started = std::time::Instant::now(); + let response = next.run(request).await; + + sink.emit(AuditRecord { + audit: "http-access", + ts: Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + user, + subject, + on_behalf_of, + on_behalf_of_rejected, + source, + method, + path, + aet, + study, + series, + instance, + status: response.status().as_u16(), + duration_ms: started.elapsed().as_millis(), + user_agent, + request_id, + }); + + response +} + +/// `value` if it fits in `max` bytes; otherwise cut at a character boundary +/// and followed by [`TRUNCATED`], the whole at most `max` bytes. A cap too +/// small for the marker gets the cut value alone (every cap here is far +/// larger than the marker). +fn bounded(mut value: String, max: usize) -> String { + if value.len() > max { + let marker = if max >= TRUNCATED.len() { + TRUNCATED + } else { + "" + }; + let mut end = max - marker.len(); + while !value.is_char_boundary(end) { + end -= 1; + } + value.truncate(end); + value.push_str(marker); + } + value +} + +/// The value of `name`, if the header occurs exactly once: a repeated header +/// is ambiguous and must not decide who the caller is. +fn single<'h>(headers: &'h HeaderMap, name: &HeaderName) -> Option<&'h HeaderValue> { + let mut values = headers.get_all(name).iter(); + let value = values.next()?; + values.next().is_none().then_some(value) +} + +/// A proxy-asserted identity: the header occurs exactly once and its value +/// is an [`identity_value`]. +fn identity(headers: &HeaderMap, name: &HeaderName) -> Option { + single(headers, name) + .and_then(identity_value) + .map(str::to_owned) +} + +/// 1..=320 bytes of UTF-8 without control characters. Anything else is +/// rejected as a whole, never truncated: a shortened identity could equal +/// someone else's. +fn identity_value(value: &HeaderValue) -> Option<&str> { + let value = std::str::from_utf8(value.as_bytes()).ok()?; + let well_formed = + (1..=MAX_IDENTITY_LEN).contains(&value.len()) && !value.chars().any(char::is_control); + well_formed.then_some(value) +} + +/// Decides whether `caller` (the verified `user`) may name the end user it +/// acts for. See the module docs for the rules. +fn delegation_for(config: &AuditConfig, headers: &HeaderMap, caller: Option<&str>) -> Delegation { + // Without trusted relays the header is not even looked at. + if config.trusted_relays.is_empty() { + return Delegation::NotClaimed; + } + let mut values = headers.get_all(&config.on_behalf_of_header).iter(); + let Some(value) = values.next() else { + return Delegation::NotClaimed; + }; + if !caller.is_some_and(|caller| config.trusted_relays.contains(caller)) { + return Delegation::Rejected(OnBehalfOfRejection::UntrustedCaller); + } + if values.next().is_some() { + return Delegation::Rejected(OnBehalfOfRejection::Invalid); + } + end_user(value).map_or( + Delegation::Rejected(OnBehalfOfRejection::Invalid), + Delegation::Honoured, + ) +} + +/// An end-user identity as named by a trusted relay: an [`identity_value`] +/// that also contains no whitespace. +fn end_user(value: &HeaderValue) -> Option { + identity_value(value) + .filter(|value| !value.chars().any(char::is_whitespace)) + .map(str::to_owned) +} + +/// A correlation id: exactly one value of 1..=128 printable ASCII +/// characters without space. +fn request_id(headers: &HeaderMap) -> Option { + single(headers, &X_REQUEST_ID) + .map(HeaderValue::as_bytes) + .filter(|id| (1..=MAX_REQUEST_ID_LEN).contains(&id.len())) + .filter(|id| id.iter().all(u8::is_ascii_graphic)) + .and_then(|id| std::str::from_utf8(id).ok()) + .map(str::to_owned) +} + +#[cfg(test)] +mod tests { + use super::*; + use axum::body::Body; + use axum::routing::get; + use axum::Router; + use serde_json::json; + use tower::ServiceExt; + + const RELAY: &str = "relay@example.org"; + const END_USER: &str = "jane.doe@example.org"; + + fn enabled() -> AuditConfig { + AuditConfig { + enabled: true, + ..AuditConfig::default() + } + } + + /// Parsed like the real configuration, so relays are normalised. + fn with_relays(relays: &[&str]) -> AuditConfig { + serde_json::from_value(json!({ "enabled": true, "trusted-relays": relays })) + .expect("valid audit config") + } + + fn from_relay() -> axum::http::request::Builder { + get_study().header("x-forwarded-email", RELAY) + } + + /// The record as a JSON object, minus the fields that vary between runs. + fn stable_json(record: &AuditRecord) -> serde_json::Value { + let mut json = serde_json::to_value(record).expect("serialize"); + let object = json.as_object_mut().expect("object"); + object.remove("ts"); + object.remove("duration_ms"); + json + } + + fn app() -> Router { + Router::new().route("/aets/{aet}/studies/{study}", get(|| async { "ok" })) + } + + /// Sends one request through a router carrying the audit middleware and + /// returns the record it produced, if any. + async fn audit(config: AuditConfig, request: Request) -> Option { + let (sink, mut rx) = AuditSink::with_receiver(config); + let app = app().layer(axum::middleware::from_fn_with_state(sink, middleware)); + let response = app.oneshot(request).await.expect("infallible"); + assert_eq!(response.status(), 200); + rx.try_recv().ok() + } + + fn get_study() -> axum::http::request::Builder { + Request::get("/aets/PACS/studies/1.2.3.4") + } + + #[tokio::test] + async fn records_identity_from_forwarded_headers() { + let request = get_study() + .header("x-forwarded-email", "jane.doe@example.org") + .header("x-forwarded-user", "3f2c9a4e") + .body(Body::empty()) + .expect("request"); + let record = audit(enabled(), request).await.expect("record"); + assert_eq!(record.user.as_deref(), Some("jane.doe@example.org")); + assert_eq!(record.subject.as_deref(), Some("3f2c9a4e")); + assert_eq!(record.aet.as_deref(), Some("PACS")); + assert_eq!(record.study.as_deref(), Some("1.2.3.4")); + } + + #[tokio::test] + async fn ignores_x_auth_request_headers() { + // Response headers in oauth2-proxy's reverse-proxy mode: never set on + // the upstream request and never stripped from it, so a client can + // send them at will. + let request = get_study() + .header("x-auth-request-email", "forged@example.com") + .header("x-auth-request-user", "forged") + .body(Body::empty()) + .expect("request"); + let record = audit(enabled(), request).await.expect("record"); + assert_eq!(record.user, None); + assert_eq!(record.subject, None); + } + + #[tokio::test] + async fn identity_header_is_configurable() { + let config = AuditConfig { + user_header: HeaderName::from_static("x-forwarded-preferred-username"), + ..enabled() + }; + let request = get_study() + .header("x-forwarded-preferred-username", "jdoe") + .header("x-forwarded-email", "jane.doe@example.org") + .body(Body::empty()) + .expect("request"); + let record = audit(config, request).await.expect("record"); + assert_eq!(record.user.as_deref(), Some("jdoe")); + } + + #[tokio::test] + async fn records_utf8_identity() { + let request = get_study() + .header( + "x-forwarded-email", + HeaderValue::from_bytes("jürgen.müller@example.org".as_bytes()).expect("header"), + ) + .header( + "x-forwarded-user", + HeaderValue::from_bytes("Jürgen Müller".as_bytes()).expect("header"), + ) + .body(Body::empty()) + .expect("request"); + let record = audit(enabled(), request).await.expect("record"); + assert_eq!(record.user.as_deref(), Some("jürgen.müller@example.org")); + assert_eq!(record.subject.as_deref(), Some("Jürgen Müller")); + } + + #[tokio::test] + async fn malformed_identity_is_absent_not_truncated() { + let longest = "a".repeat(MAX_IDENTITY_LEN); + let request = get_study() + .header("x-forwarded-email", longest.as_str()) + .body(Body::empty()) + .expect("request"); + let record = audit(enabled(), request).await.expect("record"); + assert_eq!(record.user.as_deref(), Some(longest.as_str())); + + let too_long = "a".repeat(MAX_IDENTITY_LEN + 1); + let malformed: [&[u8]; 4] = [ + b"", + too_long.as_bytes(), + b"jane\tdoe@example.org", + b"jane\xFFdoe@example.org", + ]; + for value in malformed { + let request = get_study() + .header( + "x-forwarded-email", + HeaderValue::from_bytes(value).expect("header"), + ) + .body(Body::empty()) + .expect("request"); + let record = audit(enabled(), request).await.expect("record"); + assert_eq!(record.user, None, "{value:?}"); + } + } + + #[tokio::test] + async fn repeated_identity_header_is_treated_as_absent() { + let request = get_study() + .header("x-forwarded-email", "jane.doe@example.org") + .header("x-forwarded-email", "john.doe@example.org") + .body(Body::empty()) + .expect("request"); + let record = audit(enabled(), request).await.expect("record"); + assert_eq!(record.user, None); + } + + #[test] + fn disabled_audit_has_no_sink() { + // No sink means no writer thread and no middleware: nothing is emitted. + assert!(start(&AuditConfig::default()).expect("start").is_none()); + } + + /// An in-memory `Write` target shared with the writer thread. + #[derive(Clone, Default)] + struct SharedBuffer(Arc>>); + + impl Write for SharedBuffer { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.lock().expect("lock").extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + fn sample_record() -> AuditRecord { + AuditRecord { + audit: "http-access", + ts: "2026-08-17T00:00:00Z".to_owned(), + user: Some(END_USER.to_owned()), + subject: None, + on_behalf_of: None, + on_behalf_of_rejected: None, + source: None, + method: "GET".to_owned(), + path: "/aets".to_owned(), + aet: None, + study: None, + series: None, + instance: None, + status: 200, + duration_ms: 3, + user_agent: None, + request_id: None, + } + } + + #[test] + fn writer_thread_flushes_buffered_records_on_finish() { + // A plain #[test]: the writer must not need a Tokio runtime. + let (tx, rx) = mpsc::channel(BUFFER); + let out = SharedBuffer::default(); + let writer = spawn_writer(rx, out.clone()).expect("writer thread"); + for _ in 0..3 { + tx.try_send(sample_record()).expect("buffer has room"); + } + drop(tx); + assert!(writer.finish(Duration::from_secs(5))); + + let written = String::from_utf8(out.0.lock().expect("lock").clone()).expect("UTF-8"); + assert_eq!(written.lines().count(), 3, "{written}"); + for line in written.lines() { + let json: serde_json::Value = serde_json::from_str(line).expect("JSON line"); + assert_eq!(json["audit"], "http-access"); + } + } + + #[test] + fn writer_finish_is_bounded() { + let (tx, rx) = mpsc::channel::(BUFFER); + let writer = spawn_writer(rx, SharedBuffer::default()).expect("writer thread"); + // A sink that outlives the deadline keeps the writer running. + let late_drop = std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(500)); + drop(tx); + }); + assert!(!writer.finish(Duration::from_millis(50))); + late_drop.join().expect("dropper"); + } + + #[tokio::test] + async fn undecodable_path_parameter_passes_through_and_is_audited() { + let request = || { + Request::get("/aets/%FF/studies/1.2.3.4") + .body(Body::empty()) + .expect("request") + }; + let unaudited = app().oneshot(request()).await.expect("infallible"); + assert_eq!(unaudited.status(), 200); + + let record = audit(enabled(), request()).await.expect("record"); + assert_eq!(record.status, 200); + assert_eq!(record.aet, None); + assert_eq!(record.study, None); + } + + #[tokio::test] + async fn trusted_relay_names_end_user() { + let request = from_relay() + .header("x-on-behalf-of", END_USER) + .body(Body::empty()) + .expect("request"); + let record = audit(with_relays(&[RELAY]), request).await.expect("record"); + assert_eq!(record.user.as_deref(), Some(RELAY)); + assert_eq!(record.on_behalf_of.as_deref(), Some(END_USER)); + assert_eq!(record.on_behalf_of_rejected, None); + } + + #[tokio::test] + async fn trusted_relay_without_on_behalf_of_header_records_neither_field() { + let request = from_relay().body(Body::empty()).expect("request"); + let record = audit(with_relays(&[RELAY]), request).await.expect("record"); + assert_eq!(record.user.as_deref(), Some(RELAY)); + assert_eq!(record.on_behalf_of, None); + assert_eq!(record.on_behalf_of_rejected, None); + let line = serde_json::to_string(&record).expect("serialize"); + assert!(!line.contains("on_behalf_of"), "{line}"); + } + + #[tokio::test] + async fn relay_is_matched_through_a_custom_user_header() { + let config: AuditConfig = serde_json::from_value(json!({ + "enabled": true, + "user-header": "X-Forwarded-Preferred-Username", + "trusted-relays": ["viewer-service"], + })) + .expect("valid audit config"); + + let request = get_study() + .header("x-forwarded-preferred-username", "viewer-service") + .header("x-forwarded-email", "other@example.org") + .header("x-on-behalf-of", END_USER) + .body(Body::empty()) + .expect("request"); + let record = audit(config.clone(), request).await.expect("record"); + assert_eq!(record.user.as_deref(), Some("viewer-service")); + assert_eq!(record.on_behalf_of.as_deref(), Some(END_USER)); + + // The relay name in the default user header does not count. + let request = get_study() + .header("x-forwarded-preferred-username", "jdoe") + .header("x-forwarded-email", "viewer-service") + .header("x-on-behalf-of", END_USER) + .body(Body::empty()) + .expect("request"); + let record = audit(config, request).await.expect("record"); + assert_eq!(record.on_behalf_of, None); + assert_eq!( + record.on_behalf_of_rejected, + Some(OnBehalfOfRejection::UntrustedCaller) + ); + } + + #[tokio::test] + async fn records_408_from_a_timeout_layer_inside_the_audit_layer() { + use axum::http::StatusCode; + use tower_http::timeout::TimeoutLayer; + + let (sink, mut rx) = AuditSink::with_receiver(enabled()); + let slow = || async { + tokio::time::sleep(Duration::from_secs(30)).await; + "late" + }; + let app = Router::new() + .route("/aets/{aet}/studies/{study}", get(slow)) + .layer(TimeoutLayer::with_status_code( + StatusCode::REQUEST_TIMEOUT, + Duration::from_millis(20), + )) + .layer(axum::middleware::from_fn_with_state(sink, middleware)); + let request = get_study().body(Body::empty()).expect("request"); + let response = app.oneshot(request).await.expect("infallible"); + assert_eq!(response.status(), StatusCode::REQUEST_TIMEOUT); + + let record = rx.try_recv().expect("record"); + assert_eq!(record.status, 408); + assert_eq!(record.study.as_deref(), Some("1.2.3.4")); + } + + #[tokio::test] + async fn relay_match_is_ascii_case_insensitive() { + let request = get_study() + .header("x-forwarded-email", "RELAY@Example.ORG") + .header("x-on-behalf-of", END_USER) + .body(Body::empty()) + .expect("request"); + let record = audit(with_relays(&["Relay@example.org"]), request) + .await + .expect("record"); + assert_eq!(record.on_behalf_of.as_deref(), Some(END_USER)); + } + + #[tokio::test] + async fn untrusted_caller_is_rejected_without_recording_the_claim() { + let request = get_study() + .header("x-forwarded-email", "mallory@example.com") + .header("x-on-behalf-of", END_USER) + .body(Body::empty()) + .expect("request"); + let record = audit(with_relays(&[RELAY]), request).await.expect("record"); + assert_eq!(record.user.as_deref(), Some("mallory@example.com")); + assert_eq!(record.on_behalf_of, None); + assert_eq!( + record.on_behalf_of_rejected, + Some(OnBehalfOfRejection::UntrustedCaller) + ); + let line = serde_json::to_string(&record).expect("serialize"); + assert!(!line.contains(END_USER), "claimed value leaked: {line}"); + assert!(line.contains(r#""on_behalf_of_rejected":"untrusted-caller""#)); + } + + #[tokio::test] + async fn unidentified_caller_is_rejected() { + let request = get_study() + .header("x-on-behalf-of", END_USER) + .body(Body::empty()) + .expect("request"); + let record = audit(with_relays(&[RELAY]), request).await.expect("record"); + assert_eq!(record.user, None); + assert_eq!(record.on_behalf_of, None); + assert_eq!( + record.on_behalf_of_rejected, + Some(OnBehalfOfRejection::UntrustedCaller) + ); + } + + #[tokio::test] + async fn relay_identity_must_come_from_the_user_header() { + // The relay's address in any other header (here the subject header) + // does not make the caller a relay. + let request = get_study() + .header("x-forwarded-user", RELAY) + .header("x-on-behalf-of", END_USER) + .body(Body::empty()) + .expect("request"); + let record = audit(with_relays(&[RELAY]), request).await.expect("record"); + assert_eq!( + record.on_behalf_of_rejected, + Some(OnBehalfOfRejection::UntrustedCaller) + ); + } + + #[tokio::test] + async fn without_trusted_relays_the_header_changes_nothing() { + let plain = from_relay().body(Body::empty()).expect("request"); + let claimed = from_relay() + .header("x-on-behalf-of", END_USER) + .body(Body::empty()) + .expect("request"); + let plain = audit(enabled(), plain).await.expect("record"); + let claimed = audit(enabled(), claimed).await.expect("record"); + assert_eq!(stable_json(&claimed), stable_json(&plain)); + let line = serde_json::to_string(&claimed).expect("serialize"); + assert!(!line.contains("on_behalf_of"), "{line}"); + } + + #[tokio::test] + async fn repeated_on_behalf_of_header_is_invalid() { + let request = from_relay() + .header("x-on-behalf-of", END_USER) + .header("x-on-behalf-of", "john.doe@example.org") + .body(Body::empty()) + .expect("request"); + let record = audit(with_relays(&[RELAY]), request).await.expect("record"); + assert_eq!(record.on_behalf_of, None); + assert_eq!( + record.on_behalf_of_rejected, + Some(OnBehalfOfRejection::Invalid) + ); + } + + #[tokio::test] + async fn honoured_value_is_json_escaped_in_the_audit_line() { + // A trusted relay's value is recorded verbatim, but as data inside + // the JSON line: quotes and backslashes can neither close the field + // nor forge another key. + let forged = r#"a"b\c","audit":"forged"#; + let request = from_relay() + .header("x-on-behalf-of", forged) + .body(Body::empty()) + .expect("request"); + let record = audit(with_relays(&[RELAY]), request).await.expect("record"); + assert_eq!(record.on_behalf_of.as_deref(), Some(forged)); + let line = serde_json::to_string(&record).expect("serialize"); + let parsed: serde_json::Value = serde_json::from_str(&line).expect("one JSON object"); + assert_eq!(parsed["audit"], "http-access", "{line}"); + assert_eq!(parsed["on_behalf_of"], forged, "{line}"); + } + + #[tokio::test] + async fn malformed_on_behalf_of_values_are_invalid() { + let longest = "a".repeat(MAX_IDENTITY_LEN); + let too_long = "a".repeat(MAX_IDENTITY_LEN + 1); + let malformed: [&[u8]; 6] = [ + b"", + too_long.as_bytes(), + b"jane doe@example.org", + b"jane\tdoe@example.org", + "jane\u{85}doe@example.org".as_bytes(), + b"jane\xFFdoe@example.org", + ]; + for value in malformed { + let request = from_relay() + .header( + "x-on-behalf-of", + HeaderValue::from_bytes(value).expect("header"), + ) + .body(Body::empty()) + .expect("request"); + let record = audit(with_relays(&[RELAY]), request).await.expect("record"); + assert_eq!(record.on_behalf_of, None, "{value:?}"); + assert_eq!( + record.on_behalf_of_rejected, + Some(OnBehalfOfRejection::Invalid), + "{value:?}" + ); + } + + for value in [longest.as_str(), "jürgen@example.org"] { + let request = from_relay() + .header( + "x-on-behalf-of", + HeaderValue::from_bytes(value.as_bytes()).expect("header"), + ) + .body(Body::empty()) + .expect("request"); + let record = audit(with_relays(&[RELAY]), request).await.expect("record"); + assert_eq!(record.on_behalf_of.as_deref(), Some(value)); + } + } + + #[tokio::test] + async fn records_well_formed_request_id() { + let longest = "f".repeat(MAX_REQUEST_ID_LEN); + for id in ["0f8c2b7e-5d1a-4c3b-9e6f-2a1d0c9b8a7e", longest.as_str()] { + let request = get_study() + .header("x-request-id", id) + .body(Body::empty()) + .expect("request"); + let record = audit(enabled(), request).await.expect("record"); + assert_eq!(record.request_id.as_deref(), Some(id)); + } + } + + #[tokio::test] + async fn omits_malformed_request_id() { + let too_long = "f".repeat(MAX_REQUEST_ID_LEN + 1); + let malformed: [&[&[u8]]; 5] = [ + &[b""], + &[too_long.as_bytes()], + &[b"abc def"], + &["abc\u{e9}".as_bytes()], + &[b"abc", b"def"], + ]; + for values in malformed { + let mut request = get_study(); + for value in values { + request = request.header( + "x-request-id", + HeaderValue::from_bytes(value).expect("header"), + ); + } + let request = request.body(Body::empty()).expect("request"); + let record = audit(enabled(), request).await.expect("record"); + assert_eq!(record.request_id, None, "{values:?}"); + } + } + + #[test] + fn line_safe_keeps_one_event_on_one_line() { + // Interior line breaks are escaped; the event's final one is kept. + let cases: [(&[u8], &[u8]); 5] = [ + (b"plain\n", b"plain\n"), + (b"a\nb\n", b"a\\nb\n"), + (b"a\r\nb", b"a\\r\\nb"), + (b"\n", b"\n"), + (b"", b""), + ]; + for (input, expected) in cases { + let mut out = LineSafe(Vec::new()); + let written = out.write(input).expect("write"); + assert_eq!(written, input.len()); + assert_eq!( + out.0, + expected, + "{:?} -> {:?}", + String::from_utf8_lossy(input), + String::from_utf8_lossy(&out.0) + ); + } + } + + #[test] + fn a_log_message_cannot_forge_an_audit_line() { + #[derive(Clone, Default)] + struct Shared(Arc>>); + impl Write for Shared { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.lock().expect("lock").extend_from_slice(buf); + Ok(buf.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + let sink = Shared::default(); + let make = { + let sink = sink.clone(); + move || LineSafe(sink.clone()) + }; + let subscriber = tracing_subscriber::fmt() + .compact() + .with_ansi(false) + .with_writer(make) + .finish(); + let forged = "1.2.3\n{\"audit\":\"http-access\",\"user\":\"someone@example.org\"}"; + tracing::subscriber::with_default(subscriber, || { + tracing::info!("Requesting {} from S3", forged); + }); + let output = String::from_utf8(sink.0.lock().expect("lock").clone()).expect("utf-8"); + assert_eq!(output.lines().count(), 1, "{output}"); + assert!( + !output.lines().any(|line| line.starts_with('{')), + "{output}" + ); + assert!(output.ends_with('\n'), "{output}"); + } + + #[test] + fn a_newline_ending_a_format_argument_is_escaped_too() { + // Pins that the formatter hands the writer one whole event: were an + // event ever streamed in pieces, a piece ending in a newline would + // keep it and the second argument would start a line of its own. + #[derive(Clone, Default)] + struct Shared(Arc>>); + impl Write for Shared { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.lock().expect("lock").extend_from_slice(buf); + Ok(buf.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + let sink = Shared::default(); + let make = { + let sink = sink.clone(); + move || LineSafe(sink.clone()) + }; + let subscriber = tracing_subscriber::fmt() + .compact() + .with_ansi(false) + .with_writer(make) + .finish(); + let forged = r#"{"audit":"http-access","user":"someone@example.org"}"#; + tracing::subscriber::with_default(subscriber, || { + tracing::info!("{}{}", "1.2.3\n", forged); + tracing::info!(prefix = "1.2.3\n", "{forged}"); + }); + let output = String::from_utf8(sink.0.lock().expect("lock").clone()).expect("utf-8"); + assert_eq!(output.lines().count(), 2, "{output}"); + assert!( + !output.lines().any(|line| line.starts_with('{')), + "{output}" + ); + } + + #[test] + fn bounded_cuts_at_a_character_boundary() { + assert_eq!(bounded("abcde".to_owned(), 5), "abcde"); + // The 3-byte marker counts against the cap. + assert_eq!(bounded("abcdef".to_owned(), 5), "ab…"); + // "ü" is two bytes: cutting at 2 would split it. + assert_eq!(bounded("aüüb".to_owned(), 5), "a…"); + assert_eq!(bounded("aüüb".to_owned(), 6), "aüüb"); + } + + #[test] + fn bounded_never_exceeds_its_cap() { + let value = "aü߀x".repeat(10); + for max in 0..value.len() { + let cut = bounded(value.clone(), max); + assert!(cut.len() <= max, "{max}: {} bytes", cut.len()); + assert!( + value.starts_with(cut.trim_end_matches(TRUNCATED)), + "{max}: {cut}" + ); + if max >= TRUNCATED.len() { + assert!(cut.ends_with(TRUNCATED), "{max}: {cut}"); + } + } + } + + #[tokio::test] + async fn copied_request_values_are_bounded() { + let long_study = "1".repeat(MAX_PATH_LEN); + let long_agent = "a".repeat(MAX_USER_AGENT_LEN + 1); + let long_source = format!("{}, 192.0.2.10", "f".repeat(MAX_SOURCE_LEN + 1)); + let request = Request::get(format!("/aets/PACS/studies/{long_study}")) + .header("user-agent", long_agent.as_str()) + .header("x-forwarded-for", long_source.as_str()) + .body(Body::empty()) + .expect("request"); + let record = audit(enabled(), request).await.expect("record"); + + let path_prefix = format!("/aets/PACS/studies/{long_study}"); + let keep = |max: usize| max - TRUNCATED.len(); + assert_eq!( + record.path, + format!("{}{TRUNCATED}", &path_prefix[..keep(MAX_PATH_LEN)]) + ); + assert_eq!(record.path.len(), MAX_PATH_LEN); + assert_eq!( + record.study, + Some(format!( + "{}{TRUNCATED}", + &long_study[..keep(MAX_COORDINATE_LEN)] + )) + ); + assert_eq!(record.aet.as_deref(), Some("PACS")); + assert_eq!( + record.user_agent, + Some(format!( + "{}{TRUNCATED}", + &long_agent[..keep(MAX_USER_AGENT_LEN)] + )) + ); + assert_eq!( + record.source, + Some(format!("{}{TRUNCATED}", "f".repeat(keep(MAX_SOURCE_LEN)))) + ); + + let at_cap = "a".repeat(MAX_USER_AGENT_LEN); + let request = get_study() + .header("user-agent", at_cap.as_str()) + .header("x-forwarded-for", "192.0.2.10, 198.51.100.7") + .body(Body::empty()) + .expect("request"); + let record = audit(enabled(), request).await.expect("record"); + assert_eq!(record.user_agent, Some(at_cap)); + assert_eq!(record.source.as_deref(), Some("192.0.2.10")); + assert_eq!(record.path, "/aets/PACS/studies/1.2.3.4"); + } + + #[test] + fn record_serializes_without_absent_fields() { + let record = AuditRecord { + audit: "http-access", + ts: "2026-08-17T00:00:00Z".to_owned(), + user: None, + subject: None, + on_behalf_of: None, + on_behalf_of_rejected: None, + source: None, + method: "GET".to_owned(), + path: "/aets".to_owned(), + aet: None, + study: None, + series: None, + instance: None, + status: 200, + duration_ms: 3, + user_agent: None, + request_id: None, + }; + let json = serde_json::to_string(&record).expect("serialize"); + assert!( + !json.contains("user"), + "absent fields must be omitted: {json}" + ); + assert!(json.contains("\"audit\":\"http-access\"")); + } +} diff --git a/src/backend/dimse/cmove/movescu.rs b/src/backend/dimse/cmove/movescu.rs index dbac541..7979ac2 100644 --- a/src/backend/dimse/cmove/movescu.rs +++ b/src/backend/dimse/cmove/movescu.rs @@ -15,16 +15,34 @@ use tracing::{error, info, instrument, trace}; pub struct MoveServiceClassUser { pool: AssociationPool, timeout: Duration, + /// Name the study on the completion log line. Set only with + /// `telemetry.audit.enabled`, so default log output is unchanged. + log_study_uid: bool, } impl MoveServiceClassUser { - pub const fn new(pool: AssociationPool, timeout: Duration) -> Self { - Self { pool, timeout } + pub const fn new(pool: AssociationPool, timeout: Duration, log_study_uid: bool) -> Self { + Self { + pool, + timeout, + log_study_uid, + } } #[instrument(skip_all, name = "MOVE-SCU")] #[allow(clippy::significant_drop_tightening)] pub async fn invoke(&self, request: CompositeMoveRequest) -> Result<(), MoveError> { + // When auditing, surface WHICH study the C-MOVE concerns — the audit + // trail needs more than "a move happened". + let study_uid = self.log_study_uid.then(|| { + request + .identifier + .element(tags::STUDY_INSTANCE_UID) + .ok() + .and_then(|element| element.to_str().ok()) + .map(|uid| uid.trim_end_matches('\0').to_owned()) + .unwrap_or_default() + }); let mut association = self .pool .get(PresentationParameter { @@ -54,7 +72,13 @@ impl MoveServiceClassUser { match status_type { StatusType::Success => { - info!("C-MOVE completed successfully"); + if let Some(study_uid) = &study_uid { + // Debug-formatted: the UID is percent-decoded from the + // request URL, so control characters must stay escaped. + info!(study_uid = ?study_uid, "C-MOVE completed successfully"); + } else { + info!("C-MOVE completed successfully"); + } break; } StatusType::Pending => { diff --git a/src/backend/dimse/wado.rs b/src/backend/dimse/wado.rs index c0d8a47..72015cb 100644 --- a/src/backend/dimse/wado.rs +++ b/src/backend/dimse/wado.rs @@ -125,8 +125,9 @@ impl DimseWadoService { mediator: MoveMediator, timeout: Duration, config: WadoConfig, + log_study_uid: bool, ) -> Self { - let movescu = MoveServiceClassUser::new(pool, timeout); + let movescu = MoveServiceClassUser::new(pool, timeout, log_study_uid); Self { movescu: Arc::new(movescu), mediator, diff --git a/src/backend/mod.rs b/src/backend/mod.rs index 973e37c..1765327 100644 --- a/src/backend/mod.rs +++ b/src/backend/mod.rs @@ -68,6 +68,7 @@ where state.mediator, Duration::from_millis(ae_config.wado.timeout), ae_config.wado.clone(), + state.config.telemetry.audit.enabled, ))), stow: Some(Box::new(DimseStowService::new( pool.to_owned(), diff --git a/src/config/mod.rs b/src/config/mod.rs index d7047ef..0ec1085 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -1,6 +1,7 @@ use crate::types::AE; use crate::DEFAULT_AET; +use axum::http::{header, HeaderName}; use serde::de::Error; use serde::{Deserialize, Deserializer}; use std::net::IpAddr; @@ -340,6 +341,9 @@ pub struct TelemetryConfig { pub sentry: Option, #[serde(deserialize_with = "deserialize_log_level")] pub level: tracing::Level, + /// Structured access-audit logging — see [`crate::audit`]. + #[serde(default)] + pub audit: AuditConfig, } impl Default for TelemetryConfig { @@ -347,8 +351,169 @@ impl Default for TelemetryConfig { Self { sentry: None, level: tracing::Level::INFO, + audit: AuditConfig::default(), + } + } +} + +/// Configuration for the structured access-audit log ([`crate::audit`]). +/// +/// Disabled by default: enabling it emits one JSON line per HTTP request on +/// stdout, carrying the identity an authenticating reverse proxy forwards +/// plus the DICOM resource coordinates. Delivery is fail-open (bounded +/// buffer, drops are counted and logged). +/// +/// Parsed and validated once when the configuration is loaded, so an invalid +/// setting is a startup error rather than a per-request surprise. +#[derive(Debug, Clone, Deserialize)] +#[serde(try_from = "RawAuditConfig")] +pub struct AuditConfig { + pub enabled: bool, + /// Request header carrying the proxy-verified user, e.g. an e-mail address. + pub user_header: HeaderName, + /// Request header carrying the proxy-verified subject identifier. + pub subject_header: HeaderName, + /// Callers (as identified by `user_header`) that may name the end user + /// they act for in `on_behalf_of_header`. Empty: nobody may. + pub trusted_relays: TrustedRelays, + /// Request header in which a trusted relay names the end user. + pub on_behalf_of_header: HeaderName, +} + +impl AuditConfig { + /// oauth2-proxy sets this from the verified session (`pass_user_headers`). + pub const DEFAULT_USER_HEADER: HeaderName = HeaderName::from_static("x-forwarded-email"); + /// oauth2-proxy sets this from the verified session (`pass_user_headers`). + pub const DEFAULT_SUBJECT_HEADER: HeaderName = HeaderName::from_static("x-forwarded-user"); + pub const DEFAULT_ON_BEHALF_OF_HEADER: HeaderName = HeaderName::from_static("x-on-behalf-of"); +} + +impl Default for AuditConfig { + fn default() -> Self { + Self { + enabled: false, + user_header: Self::DEFAULT_USER_HEADER, + subject_header: Self::DEFAULT_SUBJECT_HEADER, + trusted_relays: TrustedRelays::default(), + on_behalf_of_header: Self::DEFAULT_ON_BEHALF_OF_HEADER, + } + } +} + +/// Identities trusted to act on behalf of an end user, as they appear in +/// the user header. Trimmed, non-empty and ASCII-lowercased at load time; +/// matching is ASCII case-insensitive. +#[derive(Debug, Clone, Default)] +pub struct TrustedRelays(Vec); + +impl TrustedRelays { + fn parse(relays: Vec) -> Result { + relays + .into_iter() + .map(|relay| { + let relay = relay.trim(); + if relay.is_empty() { + Err(AuditConfigError::EmptyTrustedRelay) + } else { + Ok(relay.to_ascii_lowercase()) + } + }) + .collect::>() + .map(Self) + } + + pub const fn is_empty(&self) -> bool { + self.0.is_empty() + } + + pub fn contains(&self, identity: &str) -> bool { + self.0 + .iter() + .any(|relay| relay.eq_ignore_ascii_case(identity)) + } +} + +/// [`AuditConfig`] as written in the configuration, before validation. +#[derive(Debug, Default, Deserialize)] +#[serde(default, rename_all = "kebab-case")] +struct RawAuditConfig { + enabled: bool, + user_header: Option, + subject_header: Option, + trusted_relays: Vec, + on_behalf_of_header: Option, +} + +#[derive(Debug, thiserror::Error)] +pub enum AuditConfigError { + #[error("telemetry.audit.{key}: {value:?} is not a valid HTTP header name")] + InvalidHeaderName { key: &'static str, value: String }, + #[error("telemetry.audit.{key}: {name} carries credentials and must not be recorded")] + CredentialHeader { key: &'static str, name: HeaderName }, + #[error("telemetry.audit.trusted-relays: entries must not be empty")] + EmptyTrustedRelay, + #[error("telemetry.audit.on-behalf-of-header must differ from user-header and subject-header")] + OnBehalfOfHeaderCollision, +} + +impl TryFrom for AuditConfig { + type Error = AuditConfigError; + + fn try_from(raw: RawAuditConfig) -> Result { + let config = Self { + enabled: raw.enabled, + user_header: parse_header_name( + "user-header", + raw.user_header, + Self::DEFAULT_USER_HEADER, + )?, + subject_header: parse_header_name( + "subject-header", + raw.subject_header, + Self::DEFAULT_SUBJECT_HEADER, + )?, + trusted_relays: TrustedRelays::parse(raw.trusted_relays)?, + on_behalf_of_header: parse_header_name( + "on-behalf-of-header", + raw.on_behalf_of_header, + Self::DEFAULT_ON_BEHALF_OF_HEADER, + )?, + }; + // The relay's own identity header cannot double as the end user's. + if config.on_behalf_of_header == config.user_header + || config.on_behalf_of_header == config.subject_header + { + return Err(AuditConfigError::OnBehalfOfHeaderCollision); } + Ok(config) + } +} + +/// Headers that carry credentials and must never become an audit field. A +/// best-effort guard against an obvious misconfiguration, not an exhaustive +/// list: `X-Forwarded-Access-Token` is the one oauth2-proxy sets itself. +const CREDENTIAL_HEADERS: [HeaderName; 4] = [ + header::AUTHORIZATION, + header::PROXY_AUTHORIZATION, + header::COOKIE, + HeaderName::from_static("x-forwarded-access-token"), +]; + +fn parse_header_name( + key: &'static str, + value: Option, + default: HeaderName, +) -> Result { + let Some(value) = value else { + return Ok(default); + }; + // `HeaderName` is lowercase, so the comparison is case-insensitive. + let name = HeaderName::from_bytes(value.as_bytes()) + .map_err(|_| AuditConfigError::InvalidHeaderName { key, value })?; + if CREDENTIAL_HEADERS.contains(&name) { + return Err(AuditConfigError::CredentialHeader { key, name }); } + Ok(name) } /// Deserializer for [`tracing::Level`] as it does not implement [Deserialize] @@ -361,3 +526,110 @@ where tracing::Level::from_str(&value) .map_err(|_| Error::unknown_variant(&value, &["TRACE", "DEBUG", "INFO", "WARN", "ERROR"])) } + +#[cfg(test)] +mod tests { + use super::*; + + /// Loads a configuration the same way [`AppConfig::new`] does, minus the + /// file system and environment. + fn load(yaml: &str) -> Result { + use config::{Config, File, FileFormat}; + Config::builder() + .add_source(File::from_str(yaml, FileFormat::Yaml)) + .build()? + .try_deserialize() + } + + #[test] + fn audit_defaults_without_audit_section() { + let config = load("telemetry:\n level: INFO\n").expect("valid config"); + let audit = config.telemetry.audit; + assert!(!audit.enabled); + assert_eq!(audit.user_header, "x-forwarded-email"); + assert_eq!(audit.subject_header, "x-forwarded-user"); + assert!(audit.trusted_relays.is_empty()); + assert_eq!(audit.on_behalf_of_header, "x-on-behalf-of"); + } + + #[test] + fn audit_defaults_without_new_keys() { + let yaml = "telemetry:\n level: INFO\n audit:\n enabled: true\n"; + let audit = load(yaml).expect("valid config").telemetry.audit; + assert!(audit.enabled); + assert_eq!(audit.user_header, AuditConfig::DEFAULT_USER_HEADER); + assert_eq!(audit.subject_header, AuditConfig::DEFAULT_SUBJECT_HEADER); + assert!(audit.trusted_relays.is_empty()); + assert_eq!( + audit.on_behalf_of_header, + AuditConfig::DEFAULT_ON_BEHALF_OF_HEADER + ); + } + + #[test] + fn audit_trusted_relays_are_normalised() { + let yaml = "telemetry:\n level: INFO\n audit:\n enabled: true\n \ + trusted-relays:\n - \" Relay@Example.ORG \"\n \ + on-behalf-of-header: X-Acting-For\n"; + let audit = load(yaml).expect("valid config").telemetry.audit; + assert!(audit.trusted_relays.contains("relay@example.org")); + assert!(audit.trusted_relays.contains("RELAY@EXAMPLE.ORG")); + assert!(!audit.trusted_relays.contains("other@example.org")); + assert_eq!(audit.on_behalf_of_header, "x-acting-for"); + } + + #[test] + fn empty_trusted_relay_is_a_load_error() { + let yaml = "telemetry:\n level: INFO\n audit:\n trusted-relays:\n - \" \"\n"; + let error = load(yaml).expect_err("empty relay must be rejected"); + assert!(error.to_string().contains("trusted-relays"), "{error}"); + } + + #[test] + fn credential_headers_are_a_load_error() { + for key in ["user-header", "subject-header", "on-behalf-of-header"] { + for name in [ + "Authorization", + "proxy-authorization", + "COOKIE", + "X-Forwarded-Access-Token", + ] { + let yaml = format!("telemetry:\n level: INFO\n audit:\n {key}: {name}\n"); + let error = load(&yaml).expect_err("credential header must be rejected"); + assert!( + error + .to_string() + .contains(&format!("{key}: {}", name.to_ascii_lowercase())), + "{key}={name}: {error}" + ); + } + } + } + + #[test] + fn on_behalf_of_header_must_not_be_an_identity_header() { + let yaml = "telemetry:\n level: INFO\n audit:\n \ + on-behalf-of-header: X-Forwarded-Email\n"; + let error = load(yaml).expect_err("header collision must be rejected"); + assert!(error.to_string().contains("on-behalf-of-header"), "{error}"); + } + + #[test] + fn audit_identity_headers_are_configurable() { + let yaml = "telemetry:\n level: INFO\n audit:\n enabled: true\n \ + user-header: X-Forwarded-Preferred-Username\n subject-header: X-Subject\n"; + let audit = load(yaml).expect("valid config").telemetry.audit; + assert_eq!(audit.user_header, "x-forwarded-preferred-username"); + assert_eq!(audit.subject_header, "x-subject"); + } + + #[test] + fn invalid_audit_header_name_is_a_load_error() { + let yaml = "telemetry:\n level: INFO\n audit:\n user-header: \"X Bad Header\"\n"; + let error = load(yaml).expect_err("invalid header name must be rejected"); + assert!( + error.to_string().contains("user-header"), + "error should name the key: {error}" + ); + } +} diff --git a/src/main.rs b/src/main.rs index 5d9b0cc..b2d607d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,6 +1,7 @@ #![allow(clippy::multiple_crate_versions)] pub(crate) mod api; +pub(crate) mod audit; pub(crate) mod backend; pub(crate) mod config; pub(crate) mod rendering; @@ -17,6 +18,8 @@ use axum::extract::{DefaultBodyLimit, Request}; use axum::http::StatusCode; use axum::response::Response; use axum::ServiceExt; +use std::ffi::OsStr; +use std::io::IsTerminal; use std::net::SocketAddr; use std::time::Duration; use tokio::net::TcpListener; @@ -27,6 +30,7 @@ use tower_http::normalize_path::NormalizePathLayer; use tower_http::timeout::TimeoutLayer; use tower_http::trace; use tracing::{error, info, level_filters::LevelFilter, Level}; +use tracing_subscriber::fmt::writer::BoxMakeWriter; use tracing_subscriber::layer::SubscriberExt; use tracing_subscriber::util::SubscriberInitExt; use tracing_subscriber::EnvFilter; @@ -45,14 +49,32 @@ pub const IMPLEMENTATION_VERSION_NAME: &str = concat!("DICOM-RST ", env!("CARGO_ pub const DEFAULT_AET: &str = "DICOM-RST"; -fn init_logger(level: tracing::Level) { +/// How long to wait after shutdown for buffered audit records to be written. +const AUDIT_FLUSH_TIMEOUT: Duration = Duration::from_secs(5); + +fn init_logger(level: tracing::Level, escape_line_breaks: bool) { + // With auditing on, the log shares stdout with the audit records, so a + // message must never span lines (see `audit::LineSafeStdout`). Otherwise + // the default writer, i.e. unchanged output. + let writer = if escape_line_breaks { + BoxMakeWriter::new(audit::LineSafeStdout) + } else { + BoxMakeWriter::new(std::io::stdout) + }; tracing_subscriber::registry() .with( tracing_subscriber::fmt::layer() .compact() + // ANSI escapes belong on terminals, not in collected pod + // logs (they garble downstream log pipelines). + .with_ansi(use_ansi( + std::io::stdout().is_terminal(), + std::env::var_os("NO_COLOR").as_deref(), + )) .with_file(false) .with_line_number(false) - .with_target(false), + .with_target(false) + .with_writer(writer), ) .with( EnvFilter::builder() @@ -63,6 +85,13 @@ fn init_logger(level: tracing::Level) { .init(); } +/// ANSI colors only on a terminal, and not when `NO_COLOR` is set to a +/// non-empty value (), which tracing-subscriber would +/// otherwise honour by default. +fn use_ansi(stdout_is_terminal: bool, no_color: Option<&OsStr>) -> bool { + stdout_is_terminal && no_color.is_none_or(OsStr::is_empty) +} + #[derive(Clone)] pub struct AppState { pub config: AppConfig, @@ -90,25 +119,36 @@ fn init_sentry(config: &AppConfig) -> sentry::ClientInitGuard { fn main() -> Result<(), Box> { let config = AppConfig::new()?; - init_logger(config.telemetry.level); + init_logger(config.telemetry.level, config.telemetry.audit.enabled); // Manually create the Tokio runtime because the Sentry client needs to be created *before* the // Tokio runtime, which prevents us from using the #[tokio::main] macro. // See https://docs.sentry.io/platforms/rust/#async-main-function let _sentry = init_sentry(&config); + // The audit writer is a plain thread, independent of the runtime. + let (audit_sink, audit_writer) = audit::start(&config.telemetry.audit)?.unzip(); + tokio::runtime::Builder::new_multi_thread() .enable_all() .build()? .block_on(async move { - if let Err(error) = run(config).await { + if let Err(error) = run(config, audit_sink).await { error!("Failed to start application due to error: {error}"); } }); + + // The runtime and with it every audit sink are gone: let the writer + // drain what is still buffered. + if let Some(writer) = audit_writer { + if !writer.finish(AUDIT_FLUSH_TIMEOUT) { + error!("Audit log writer did not finish; buffered audit records may be lost"); + } + } Ok(()) } -async fn run(config: AppConfig) -> anyhow::Result<()> { +async fn run(config: AppConfig, audit_sink: Option) -> anyhow::Result<()> { let mediator = MoveMediator::new(&config); let pools = AssociationPools::new(&config); @@ -151,8 +191,17 @@ async fn run(config: AppConfig) -> anyhow::Result<()> { .layer(TimeoutLayer::with_status_code( StatusCode::REQUEST_TIMEOUT, Duration::from_secs(config.server.http.request_timeout), - )) - .with_state(app_state); + )); + // Outside the timeout layer, so timed-out requests are audited with + // their 408 as well. Not installed at all unless telemetry.audit.enabled. + let app = match audit_sink { + Some(sink) => app.layer(axum::middleware::from_fn_with_state( + sink, + audit::middleware, + )), + None => app, + } + .with_state(app_state); let app = NormalizePathLayer::trim_trailing_slash().layer(app); let service = ServiceExt::::into_make_service(app); @@ -211,3 +260,17 @@ async fn add_common_headers(req: Request, next: axum::middleware::Next) -> Respo headers.insert("Server", axum::http::HeaderValue::from_static(server_name)); response } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn ansi_only_on_a_terminal_without_no_color() { + assert!(use_ansi(true, None)); + assert!(use_ansi(true, Some(OsStr::new("")))); + assert!(!use_ansi(true, Some(OsStr::new("1")))); + assert!(!use_ansi(false, None)); + assert!(!use_ansi(false, Some(OsStr::new("1")))); + } +} diff --git a/tests/audit.rs b/tests/audit.rs new file mode 100644 index 0000000..30a780c --- /dev/null +++ b/tests/audit.rs @@ -0,0 +1,105 @@ +// Uses only part of the shared helpers. +#[allow(dead_code)] +mod common; + +use common::spawn_dicomrst; +use std::fmt::Write as _; +use std::time::Duration; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpStream; + +/// A server that needs no PACS: the request below only lists the AETs. +fn config(audit: &str) -> String { + format!( + " + telemetry: + level: INFO +{audit} + server: + http: + interface: 127.0.0.1 + port: 0 + dimse: + - aet: DICOM-RST + interface: 127.0.0.1 + port: 0 + aets: + - aet: PACS + host: 127.0.0.1 + port: 104 + backend: DIMSE + " + ) +} + +/// Sends one plain HTTP/1.1 GET request and returns the raw response. +async fn get(port: u16, path: &str, headers: &[(&str, &str)]) -> anyhow::Result { + let mut request = format!("GET {path} HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n"); + for (name, value) in headers { + write!(request, "{name}: {value}\r\n")?; + } + request.push_str("\r\n"); + + let mut stream = TcpStream::connect(("127.0.0.1", port)).await?; + stream.write_all(request.as_bytes()).await?; + let mut response = String::new(); + stream.read_to_string(&mut response).await?; + Ok(response) +} + +fn audit_lines(logs: &[String]) -> Vec { + logs.iter() + .filter_map(|line| serde_json::from_str::(line).ok()) + .filter(|json| json["audit"] == "http-access") + .collect() +} + +#[tokio::test] +async fn enabled_audit_writes_a_json_line_per_request() -> anyhow::Result<()> { + let config = config(" audit:\n enabled: true"); + let mut server = spawn_dicomrst(&config).await?; + + let response = get( + server.http_port(), + "/aets", + &[("X-Forwarded-Email", "jane.doe@example.org")], + ) + .await?; + assert!(response.starts_with("HTTP/1.1 200"), "{response}"); + + let logs = server.collect_logs(Duration::from_secs(1)).await; + let records = audit_lines(&logs); + assert_eq!(records.len(), 1, "{logs:#?}"); + assert_eq!(records[0]["user"], "jane.doe@example.org"); + assert_eq!(records[0]["method"], "GET"); + assert_eq!(records[0]["path"], "/aets"); + assert_eq!(records[0]["status"], 200); + Ok(()) +} + +#[tokio::test] +async fn disabled_audit_writes_no_audit_line() -> anyhow::Result<()> { + let config = config(""); + let mut server = spawn_dicomrst(&config).await?; + + let response = get( + server.http_port(), + "/aets", + &[("X-Forwarded-Email", "jane.doe@example.org")], + ) + .await?; + assert!(response.starts_with("HTTP/1.1 200"), "{response}"); + + let logs = server.collect_logs(Duration::from_secs(1)).await; + assert!( + logs.iter() + .any(|line| line.contains("finished processing request")), + "the request itself should be logged: {logs:#?}" + ); + assert!(audit_lines(&logs).is_empty(), "{logs:#?}"); + assert!( + !logs.iter().any(|line| line.contains("http-access")), + "{logs:#?}" + ); + Ok(()) +} diff --git a/tests/common/mod.rs b/tests/common/mod.rs index 7e74ab9..0a2aea1 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -88,6 +88,11 @@ impl ServerProcess { .context("Timed out waiting for DICOM-RST to start")? } + /// The port the HTTP server is listening on. + pub const fn http_port(&self) -> u16 { + self.http_port + } + /// Collects log lines from the server's stdout until no new line arrives within /// `quiet_period`. pub async fn collect_logs(&mut self, quiet_period: Duration) -> Vec { @@ -132,7 +137,7 @@ pub async fn with_test_server( let client = DicomWebClient::with_single_url(&format!( "http://localhost:{}/aets/ORTHANC", - server.http_port + server.http_port() )); test(client, &mut server).await?;