diff --git a/Cargo.lock b/Cargo.lock index d2d0aa937..e3abc188b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6999,6 +6999,7 @@ dependencies = [ "aws-sdk-timestreamwrite", "aws-sdk-wafv2", "aws-types", + "fakecloud-core", "libc", "reqwest", "serde_json", diff --git a/crates/fakecloud-conformance/tests/ec2.rs b/crates/fakecloud-conformance/tests/ec2.rs index 2d0f88e1a..e64cb6823 100644 --- a/crates/fakecloud-conformance/tests/ec2.rs +++ b/crates/fakecloud-conformance/tests/ec2.rs @@ -13040,7 +13040,9 @@ fn xml_value(body: &str, tag: &str) -> String { body[start..end].to_string() } -/// An IPAM, an internet-registry association on it, and one registration. +/// An IPAM and a freshly created internet-registry association on it. The +/// association is still `pending-enable`, so it cannot publish registrations +/// yet. async fn make_ir_association(c: &aws_sdk_ec2::Client, q: &Ec2Query) -> String { let ipam = make_ipam(c).await; let body = q @@ -13056,6 +13058,30 @@ async fn make_ir_association(c: &aws_sdk_ec2::Client, q: &Ec2Query) -> String { xml_value(&body, "ipamInternetRegistryAssociationId") } +/// Enable an association against the registry's RPKI service, which is what +/// lets it publish Route Origin Authorizations. +async fn enable_ir_association(q: &Ec2Query, id: &str) { + q.call( + "EnableIpamInternetRegistryAssociation", + &[ + ("IpamInternetRegistryAssociationId", id), + ("RpkiVersion", "1"), + ("ServiceUri", "https://rpki.example/up-down"), + ("ChildHandle", "child"), + ("ParentHandle", "parent"), + ("ParentBpkiTa", "TA=="), + ], + ) + .await; +} + +/// An enabled association, ready to publish routing policy registrations. +async fn make_enabled_ir_association(c: &aws_sdk_ec2::Client, q: &Ec2Query) -> String { + let id = make_ir_association(c, q).await; + enable_ir_association(q, &id).await; + id +} + #[test_action("ec2", "CreateIpamInternetRegistryAssociation", checksum = "a4a63a5d")] #[test_action( "ec2", @@ -13099,6 +13125,34 @@ async fn ec2_ipam_internet_registry_association_lifecycle() { // entity-escaped. assert!(body.contains("child_handle="child""), "{body}"); + // The registrations have to be removed before the association goes: + // deleting one that still publishes ROAs would orphan them. + q.call( + "CreateIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "192.0.2.0/24"), + ("Asn.1", "64512"), + ], + ) + .await; + let (status, body) = q + .send( + "DeleteIpamInternetRegistryAssociation", + &[("IpamInternetRegistryAssociationId", &id)], + ) + .await; + assert_eq!(status, 400, "{body}"); + assert!(body.contains("DependencyViolation"), "{body}"); + q.call( + "DeleteIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "192.0.2.0/24"), + ], + ) + .await; + let body = q .call( "DeleteIpamInternetRegistryAssociation", @@ -13122,7 +13176,24 @@ async fn ec2_ipam_routing_policy_registration_lifecycle() { let s = TestServer::start().await; let c = s.ec2_client().await; let q = Ec2Query::new(&s); - let id = make_ir_association(&c, &q).await; + let pending = make_ir_association(&c, &q).await; + + // A registration publishes through the association's RPKI service, so an + // association that was never enabled cannot carry one. + let (status, _) = q + .send( + "CreateIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &pending), + ("Cidr", "10.0.0.0/16"), + ("Asn.1", "64512"), + ], + ) + .await; + assert_eq!(status, 400); + + let id = pending; + enable_ir_association(&q, &id).await; let body = q .call( @@ -13183,6 +13254,19 @@ async fn ec2_ipam_routing_policy_registration_lifecycle() { .await; assert!(body.contains("64513"), "{body}"); assert!(body.contains("update-complete"), "{body}"); + // Modify requires only `Asns`, so the members it leaves out survive it + // rather than being cleared out from under the published ROA. + assert!( + body.contains("24"), + "a modify that omits MaxLength must not clear it: {body}" + ); + let roas = q + .call( + "GetIpamRouteOriginAuthorizations", + &[("IpamInternetRegistryAssociationId", &id)], + ) + .await; + assert!(roas.contains("24"), "{roas}"); // Deltas are the audit trail, and survive the registration they describe. let body = q @@ -13214,12 +13298,19 @@ async fn ec2_ipam_routing_policy_registration_lifecycle() { ], ) .await; + // Both deltas have to be present before their order means anything: + // `None < Some(_)`, so comparing the finds directly would pass vacuously + // on a delta that went missing. + let at = |body: &str, delta: &str| { + body.find(delta) + .unwrap_or_else(|| panic!("delta {delta} missing from {body}")) + }; assert!( - forward.find(&first_delta) < forward.find(&second_delta), + at(&forward, &first_delta) < at(&forward, &second_delta), "forward is oldest first" ); assert!( - reverse.find(&second_delta) < reverse.find(&first_delta), + at(&reverse, &second_delta) < at(&reverse, &first_delta), "reverse is newest first" ); @@ -13238,6 +13329,160 @@ async fn ec2_ipam_routing_policy_registration_lifecycle() { "{one}" ); + // Modify is a partial update, so clearing a member needs a spelling of its + // own: an empty value is it, and the members the request does not name are + // still left alone. + q.call( + "ModifyIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "10.0.0.0/16"), + ("Asn.1", "64513"), + ("Description", "prod prefix"), + ], + ) + .await; + let body = q + .call( + "GetIpamRoutingPolicyRegistrations", + &[("IpamInternetRegistryAssociationId", &id)], + ) + .await; + assert!( + body.contains("prod prefix"), + "{body}" + ); + q.call( + "ModifyIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "10.0.0.0/16"), + ("Asn.1", "64513"), + ("Description", ""), + ], + ) + .await; + let body = q + .call( + "GetIpamRoutingPolicyRegistrations", + &[("IpamInternetRegistryAssociationId", &id)], + ) + .await; + assert!( + !body.contains(""), + "an empty Description clears it: {body}" + ); + assert!( + body.contains("24"), + "and the members the request does not name survive: {body}" + ); + + // `reverse` pages by delta id. Deltas are appended, so an offset cursor + // shifts under a caller who pages: a delta recorded between two pages + // would push items the first page already reported into the second. + for cidr in ["192.0.2.0/24", "198.51.100.0/24", "203.0.113.0/24"] { + q.call( + "CreateIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", cidr), + ("Asn.1", "64512"), + ], + ) + .await; + } + let first_page = q + .call( + "GetIpamRoutingPolicyRegistrationDeltas", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("ChronologicalOrder", "reverse"), + ("MaxResults", "5"), + ], + ) + .await; + let seen: Vec = first_page + .split("") + .skip(1) + .map(|s| s.split("").next().unwrap().to_string()) + .collect(); + assert_eq!(seen.len(), 5, "{first_page}"); + let token = xml_value(&first_page, "nextToken"); + // Two more deltas land between the two pages. + for cidr in ["172.16.0.0/16", "172.17.0.0/16"] { + q.call( + "CreateIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", cidr), + ("Asn.1", "64512"), + ], + ) + .await; + } + let second_page = q + .call( + "GetIpamRoutingPolicyRegistrationDeltas", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("ChronologicalOrder", "reverse"), + ("MaxResults", "5"), + ("NextToken", &token), + ], + ) + .await; + for delta in &seen { + assert!( + !second_page.contains(delta), + "the second page repeated {delta}: {second_page}" + ); + } + + // `MaxResults` that is not a number is rejected, rather than taken as "no + // limit" and answered with the whole set. + let (status, _) = q + .send( + "GetIpamRoutingPolicyRegistrationDeltas", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("MaxResults", "abc"), + ], + ) + .await; + assert_eq!(status, 400); + + // A client token reused with different parameters is not a retry: EC2 + // answers it with IdempotentParameterMismatch rather than reporting a + // removal of a CIDR that is still registered. + q.call( + "DeleteIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "192.0.2.0/24"), + ("ClientToken", "token-reused"), + ], + ) + .await; + let (status, body) = q + .send( + "DeleteIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "198.51.100.0/24"), + ("ClientToken", "token-reused"), + ], + ) + .await; + assert_eq!(status, 400, "{body}"); + assert!(body.contains("IdempotentParameterMismatch"), "{body}"); + let body = q + .call( + "GetIpamRoutingPolicyRegistrations", + &[("IpamInternetRegistryAssociationId", &id)], + ) + .await; + assert!(body.contains("198.51.100.0/24"), "{body}"); + q.call( "DeleteIpamRoutingPolicyRegistration", &[ @@ -13265,7 +13510,7 @@ async fn ec2_batch_modify_ipam_routing_policy_registrations() { let s = TestServer::start().await; let c = s.ec2_client().await; let q = Ec2Query::new(&s); - let id = make_ir_association(&c, &q).await; + let id = make_enabled_ir_association(&c, &q).await; let delta = r#"{"add":[{"cidr":"192.0.2.0/24","asns":["64512"],"maxLength":25}, {"cidr":"198.51.100.0/24","asns":["64513"]}]}"#; @@ -13309,6 +13554,100 @@ async fn ec2_batch_modify_ipam_routing_policy_registrations() { assert!(!body.contains("192.0.2.0/24"), "{body}"); assert!(body.contains("198.51.100.0/24"), "{body}"); + // A document is checked against the state it would itself leave, so it can + // remove a CIDR it adds: it is self-consistent, even though the CIDR is + // not registered when the document arrives. + q.call( + "BatchModifyIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ( + "DeltaJson", + r#"{"add":[{"cidr":"203.0.113.0/24","asns":["64512"]}], + "remove":["203.0.113.0/24"]}"#, + ), + ], + ) + .await; + let body = q + .call( + "GetIpamRoutingPolicyRegistrations", + &[("IpamInternetRegistryAssociationId", &id)], + ) + .await; + assert!(!body.contains("203.0.113.0/24"), "{body}"); + + // An `add` for a CIDR the association already holds is a partial update, + // so `null` is what clears a member: a member the entry leaves out keeps + // the value the registration already carries. + q.call( + "BatchModifyIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ( + "DeltaJson", + r#"{"add":[{"cidr":"198.51.100.0/24","asns":["64513"], + "maxLength":26,"description":"edge prefix"}]}"#, + ), + ], + ) + .await; + let body = q + .call( + "GetIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "198.51.100.0/24"), + ], + ) + .await; + assert!( + body.contains("edge prefix"), + "{body}" + ); + + q.call( + "BatchModifyIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ( + "DeltaJson", + r#"{"add":[{"cidr":"198.51.100.0/24","asns":["64513"],"description":null}]}"#, + ), + ], + ) + .await; + let body = q + .call( + "GetIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "198.51.100.0/24"), + ], + ) + .await; + assert!( + !body.contains(""), + "a null description clears it: {body}" + ); + assert!( + body.contains("26"), + "and an omitted member survives: {body}" + ); + + // Removing a CIDR the document neither holds nor adds is still a + // not-found: it removes nothing and must not be published as a delta. + let (status, _) = q + .send( + "BatchModifyIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("DeltaJson", r#"{"remove":["192.0.2.0/24"]}"#), + ], + ) + .await; + assert_eq!(status, 400); + // Malformed JSON is rejected rather than recorded as a delta. let (status, _) = q .send( @@ -13334,7 +13673,7 @@ async fn ec2_ipam_registry_views_derive_from_registrations() { let s = TestServer::start().await; let c = s.ec2_client().await; let q = Ec2Query::new(&s); - let id = make_ir_association(&c, &q).await; + let id = make_enabled_ir_association(&c, &q).await; q.call( "CreateIpamRoutingPolicyRegistration", @@ -13405,6 +13744,7 @@ async fn ec2_ipam_route_discovery_and_protection_findings() { ) .await; let id = xml_value(&body, "ipamInternetRegistryAssociationId"); + enable_ir_association(&q, &id).await; q.call( "CreateIpamRoutingPolicyRegistration", &[ diff --git a/crates/fakecloud-core/src/container_net.rs b/crates/fakecloud-core/src/container_net.rs index 6d099f498..f76c70c9b 100644 --- a/crates/fakecloud-core/src/container_net.rs +++ b/crates/fakecloud-core/src/container_net.rs @@ -82,22 +82,52 @@ pub fn cli_available(cli: &str) -> bool { /// Run the bounded ` info` liveness probe once (uncached). fn probe_cli(cli: &str) -> bool { - let child = std::process::Command::new(cli) - .arg("info") - .stdout(std::process::Stdio::null()) - .stderr(std::process::Stdio::null()) - .spawn(); + let child = spawn_bounded( + std::process::Command::new(cli) + .arg("info") + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()), + ); let Ok(mut child) = child else { return false; }; - wait_bounded(&mut child) && child.wait().map(|s| s.success()).unwrap_or(false) + wait_bounded_group(&mut child) && child.wait().map(|s| s.success()).unwrap_or(false) +} + +/// Spawn a container-CLI command in a process group of its own (Unix), so a +/// timed-out call can be torn down whole. `FAKECLOUD_CONTAINER_CLI` is +/// routinely a wrapper -- `sh -c 'exec docker "$@"'`, a `podman-remote` shim -- +/// which makes the real command a *grandchild*: it survives `Child::kill`, goes +/// on holding whatever pipes we handed it, and keeps running against a wedged +/// daemon forever. Its own group makes it reachable by a single signal. +/// Detaching these from terminal job control is fine: their lifetime is managed +/// by deadline here, not by the shell fakecloud was started from. +/// +/// stdin is /dev/null, and has to be: a new process group is a *background* +/// one, so a child that reads the controlling terminal -- a `sudo` or +/// credential-helper wrapper prompting, exactly the wrapper case above -- takes +/// SIGTTIN, which stops it rather than ending it. `try_wait` is WNOHANG without +/// WUNTRACED, so the loop below never sees a stopped child and the call burns +/// the whole deadline before being killed. These calls are non-interactive +/// anyway, so an immediate EOF is the right answer for them. +fn spawn_bounded(cmd: &mut std::process::Command) -> std::io::Result { + cmd.stdin(std::process::Stdio::null()); + #[cfg(unix)] + { + std::os::unix::process::CommandExt::process_group(cmd, 0); + } + cmd.spawn() } /// Wait for `child` up to [`CLI_PROBE_TIMEOUT`], killing it on expiry. Returns /// whether it exited on its own. Every container-CLI call goes through this: /// a liveness probe answering does not promise the next call will, and an /// unbounded one blocks the caller rather than just that command. -pub fn wait_bounded(child: &mut std::process::Child) -> bool { +/// +/// Only for a child from [`spawn_bounded`], which put it in a group of its own: +/// the expiry kill hits that whole group, so a wrapper CLI's grandchildren die +/// with it. +fn wait_bounded_group(child: &mut std::process::Child) -> bool { let deadline = std::time::Instant::now() + CLI_PROBE_TIMEOUT; loop { match child.try_wait() { @@ -107,7 +137,7 @@ pub fn wait_bounded(child: &mut std::process::Child) -> bool { } if std::time::Instant::now() >= deadline { // Daemon is wedged: kill the blocked call and report failure. - let _ = child.kill(); + kill_expired(child); let _ = child.wait(); return false; } @@ -115,47 +145,128 @@ pub fn wait_bounded(child: &mut std::process::Child) -> bool { } } +/// SIGKILL a timed-out child and its process group. [`spawn_bounded`] made the +/// child its own group leader, so the group id is the child's pid, and the +/// child is still unreaped here -- the pid cannot have been recycled and the +/// signal cannot stray onto an unrelated group. +#[cfg(unix)] +fn kill_expired(child: &mut std::process::Child) { + // SAFETY: `kill` with a negative pid targets the process group of that + // id; any pid value is safe to pass. + let _ = unsafe { libc::kill(-(child.id() as libc::pid_t), libc::SIGKILL) }; + let _ = child.kill(); +} + +/// Windows has no process-group signal (a job object would be needed), so the +/// direct child is as far as the kill reaches; [`run_bounded`] still bounds the +/// wait on its stdout reader so the caller can't be held by a surviving +/// grandchild. +#[cfg(not(unix))] +fn kill_expired(child: &mut std::process::Child) { + let _ = child.kill(); +} + +/// Whether the stdout reader thread ended before the call returned. +#[derive(Debug)] +enum ReaderState { + /// The reader returned; its thread is gone. + Finished, + /// The reader is still blocked on the pipe because a write end we could not + /// close is held outside the child's process group. The thread outlives the + /// call; the caller does not wait for it. + Abandoned, +} + +/// Floor on how long [`run_bounded`] waits for its stdout reader once the call +/// is over (it also gets whatever is left of the call's own budget). Both exits +/// close every write end we control -- the child exited, or its whole process +/// group was killed -- which ends the blocked `read_to_end` at once, so this +/// covers scheduling only. It exists so a write end held somewhere we cannot +/// reach costs the caller a few hundred milliseconds instead of blocking it for +/// good, which is what an unbounded join did. +const READER_DRAIN_GRACE: std::time::Duration = std::time::Duration::from_millis(500); + /// Run a container-CLI command and return its stdout, or `None` when it fails /// or outruns [`CLI_PROBE_TIMEOUT`]. pub fn bounded_output(cli: &str, args: &[&str]) -> Option { - let mut child = std::process::Command::new(cli) - .args(args) - .stdout(std::process::Stdio::piped()) - .stderr(std::process::Stdio::null()) - .spawn() - .ok()?; + run_bounded(cli, args).0 +} + +/// [`bounded_output`], plus whether its stdout reader finished -- so the +/// timeout path's "no reader left behind" guarantee is unit-testable instead of +/// only observable as a thread that never goes away. +fn run_bounded(cli: &str, args: &[&str]) -> (Option, ReaderState) { + let deadline = std::time::Instant::now() + CLI_PROBE_TIMEOUT; + let child = spawn_bounded( + std::process::Command::new(cli) + .args(args) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::null()), + ); + let Ok(mut child) = child else { + return (None, ReaderState::Finished); + }; + let Some(mut stdout) = child.stdout.take() else { + kill_expired(&mut child); + let _ = child.wait(); + return (None, ReaderState::Finished); + }; // Drain stdout while waiting. A child whose output outgrows the pipe // buffer blocks on write until someone reads it, so waiting for exit // first would deadlock until the deadline and then report the sweep as // failed -- `docker ps -a` across a busy host is exactly that much output. - let mut stdout = child.stdout.take()?; - let reader = std::thread::spawn(move || { + // + // The channel doubles as the reader's "I'm done" signal: the send is the + // last thing the thread does before dropping the pipe's read end, so a + // received buffer proves no reader is parked behind us. A `JoinHandle` + // can't say that without blocking, which on the timeout path is exactly + // what we must not do. + let (tx, rx) = std::sync::mpsc::channel(); + std::thread::spawn(move || { let mut buf = Vec::new(); let _ = std::io::Read::read_to_end(&mut stdout, &mut buf); - buf + let _ = tx.send(buf); }); - if !wait_bounded(&mut child) { - return None; - } - let status = child.wait().ok()?; - let buf = reader.join().ok()?; - status - .success() - .then(|| String::from_utf8_lossy(&buf).into_owned()) + // On expiry `wait_bounded_group` has killed the whole process group, so a + // wrapper CLI's grandchild releases the write end and the reader returns + // instead of blocking for the life of the process -- one leaked thread per + // call, on precisely the wedged-daemon path these bounds exist for. + let exited = wait_bounded_group(&mut child); + let status = child.wait().ok(); + // Whatever is left of the call's own budget, and never less than the grace: + // a prompt call can afford to wait out a reader thread the scheduler hasn't + // run yet, a timed-out one gets only the grace, and either way the caller is + // back within CLI_PROBE_TIMEOUT plus that grace. + let grace = deadline + .saturating_duration_since(std::time::Instant::now()) + .max(READER_DRAIN_GRACE); + let drained = rx.recv_timeout(grace).ok(); + let output = match (exited, status, &drained) { + (true, Some(status), Some(buf)) if status.success() => { + Some(String::from_utf8_lossy(buf).into_owned()) + } + _ => None, + }; + let reader = if drained.is_some() { + ReaderState::Finished + } else { + ReaderState::Abandoned + }; + (output, reader) } /// Run a container-CLI command for its effect only, bounded the same way. /// Returns whether it succeeded. -pub fn bounded_status(cli: &str, args: &[String]) -> bool { - let Ok(mut child) = std::process::Command::new(cli) - .args(args) - .stdout(std::process::Stdio::null()) - .stderr(std::process::Stdio::null()) - .spawn() - else { +pub fn bounded_status(cli: &str, args: &[&str]) -> bool { + let Ok(mut child) = spawn_bounded( + std::process::Command::new(cli) + .args(args) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()), + ) else { return false; }; - wait_bounded(&mut child) && child.wait().map(|s| s.success()).unwrap_or(false) + wait_bounded_group(&mut child) && child.wait().map(|s| s.success()).unwrap_or(false) } /// True if the given PID is a live process on this host. @@ -213,21 +324,25 @@ pub fn is_podman_binary(cli: &str) -> bool { /// Detect the Docker bridge gateway IP on Linux. Returns `None` if /// detection fails (caller falls back to the conventional `172.17.0.1`). +/// +/// Goes through [`bounded_output`] like every other container-CLI call here: +/// `network inspect` talks to the same daemon as the liveness probe, so a +/// wedged one blocks it on connect forever. This runs inside runtime +/// constructors on Linux, where an unbounded call hangs server startup outright +/// -- the exact failure [`CLI_PROBE_TIMEOUT`] exists to prevent. On timeout the +/// caller just takes the conventional fallback. pub fn detect_bridge_gateway(cli: &str) -> Option { - let output = std::process::Command::new(cli) - .args([ + let stdout = bounded_output( + cli, + &[ "network", "inspect", "bridge", "--format", "{{range .IPAM.Config}}{{.Gateway}}{{end}}", - ]) - .output() - .ok()?; - if !output.status.success() { - return None; - } - let gateway = String::from_utf8_lossy(&output.stdout).trim().to_string(); + ], + )?; + let gateway = stdout.trim().to_string(); if gateway.is_empty() || !gateway.contains('.') { return None; } @@ -715,13 +830,14 @@ mod bounded_cli_tests { #[test] fn a_hanging_cli_call_is_cut_off() { let start = std::time::Instant::now(); - let mut child = std::process::Command::new("sleep") - .arg("600") - .stdout(std::process::Stdio::null()) - .stderr(std::process::Stdio::null()) - .spawn() - .expect("sleep is available"); - assert!(!wait_bounded(&mut child)); + let mut child = spawn_bounded( + std::process::Command::new("sleep") + .arg("600") + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()), + ) + .expect("sleep is available"); + assert!(!wait_bounded_group(&mut child)); assert!( start.elapsed() < CLI_PROBE_TIMEOUT + std::time::Duration::from_secs(5), "the wait must end at the bound" @@ -755,4 +871,112 @@ mod bounded_cli_tests { assert!(bounded_status("true", &[])); assert!(!bounded_status("false", &[])); } + + /// A bounded call gets an empty stdin, never fakecloud's own. The process + /// group `spawn_bounded` creates is a *background* one, so a child that + /// reads the controlling terminal -- a `sudo`/credential-helper wrapper + /// prompting -- is stopped by SIGTTIN, which the WNOHANG `try_wait` loop + /// cannot see: the call would burn the whole deadline and be killed. + /// `cat` with no argument reads stdin to EOF, so it returns at once with + /// nothing only when stdin is /dev/null. + #[cfg(unix)] + #[test] + fn a_bounded_call_reads_an_empty_stdin() { + let start = std::time::Instant::now(); + let out = bounded_output("cat", &[]).expect("a call reading stdin must not time out"); + assert!(out.is_empty(), "stdin must be empty, got {out:?}"); + assert!( + start.elapsed() < CLI_PROBE_TIMEOUT, + "a call reading stdin must not reach the deadline" + ); + } + + /// The happy path must still hand back the child's output *and* collect the + /// reader, so the no-leak guarantee isn't bought by dropping output. + #[test] + fn a_prompt_cli_call_collects_its_reader() { + let (output, reader) = run_bounded("echo", &["abc123"]); + assert_eq!(output.as_deref().map(str::trim), Some("abc123")); + assert!( + matches!(reader, ReaderState::Finished), + "reader was {reader:?}, expected it collected" + ); + } + + /// The bridge-gateway probe talks to the same daemon as the liveness probe, + /// so it has to end at the same bound. It used to be a plain + /// `Command::output()`, which against a wedged daemon hung whichever runtime + /// constructor called it -- on Linux, every container-backed service at + /// server startup. Timing out is not an error here: the caller falls back to + /// the conventional `172.17.0.1`. + #[cfg(unix)] + #[test] + fn a_hanging_bridge_gateway_probe_is_cut_off() { + use std::os::unix::fs::PermissionsExt; + + let dir = std::env::temp_dir().join(format!("fc-gwtest-{}", std::process::id())); + std::fs::create_dir_all(&dir).unwrap(); + let script = dir.join("hangcli"); + std::fs::write(&script, "#!/bin/sh\nsleep 600\n").unwrap(); + std::fs::set_permissions(&script, std::fs::Permissions::from_mode(0o755)).unwrap(); + + let start = std::time::Instant::now(); + let gateway = detect_bridge_gateway(script.to_str().unwrap()); + let elapsed = start.elapsed(); + + std::fs::remove_dir_all(&dir).ok(); + assert_eq!(gateway, None, "a wedged daemon must report no gateway"); + assert!( + elapsed < CLI_PROBE_TIMEOUT + READER_DRAIN_GRACE + std::time::Duration::from_secs(5), + "the probe took {elapsed:?}, expected it bounded near {CLI_PROBE_TIMEOUT:?}" + ); + } + + /// The unavailable-CLI path keeps its shape: nothing to spawn, no gateway, + /// and the caller's fallback stands. + #[test] + fn a_missing_cli_reports_no_bridge_gateway() { + assert_eq!( + detect_bridge_gateway("definitely-not-a-real-cli-binary-xyz-123"), + None + ); + } + + /// A CLI that succeeds with no output -- an `inspect --format` over a bridge + /// with no IPAM config -- still means "no gateway", not an empty + /// `--add-host` value. Unchanged by the bounding; guarded so it stays that + /// way. + #[test] + fn an_empty_gateway_is_rejected() { + assert_eq!(detect_bridge_gateway("true"), None); + } + + /// `FAKECLOUD_CONTAINER_CLI` is routinely a wrapper (`sh -c 'exec docker + /// "$@"'`, a `podman-remote` shim), which makes the real command a + /// grandchild holding the stdout pipe. Killing only the direct child left + /// the reader's `read_to_end` blocked forever -- a thread parked for the + /// life of the process, once per call, on exactly the wedged-daemon path + /// these bounds were added for (the server reaper calls this at startup). + /// The reader reporting in is the evidence: EOF on that pipe is only + /// possible once every write end is closed, so a collected buffer proves + /// the grandchildren went down with the call. + #[cfg(unix)] + #[test] + fn a_timed_out_wrapper_call_leaves_no_reader_behind() { + let start = std::time::Instant::now(); + // A wrapper that outlives its own kill: the backgrounded sleep inherits + // the stdout pipe and is not the process we spawned. + let (output, reader) = run_bounded("sh", &["-c", "sleep 600 & sleep 600"]); + assert_eq!(output, None, "a wedged call must report failure"); + assert!( + matches!(reader, ReaderState::Finished), + "reader was {reader:?}: the stdout reader must not outlive the call" + ); + assert!( + start.elapsed() + < CLI_PROBE_TIMEOUT + READER_DRAIN_GRACE + std::time::Duration::from_secs(5), + "the call must still end at the bound, took {:?}", + start.elapsed() + ); + } } diff --git a/crates/fakecloud-ec2/src/service/ipam_registry.rs b/crates/fakecloud-ec2/src/service/ipam_registry.rs index 9325de18a..872ee933d 100644 --- a/crates/fakecloud-ec2/src/service/ipam_registry.rs +++ b/crates/fakecloud-ec2/src/service/ipam_registry.rs @@ -4,26 +4,130 @@ //! An association ties an IPAM to one Regional Internet Registry. Registrations //! hang off it, keyed by CIDR, and every change to them produces a delta: the //! deltas are the audit trail, so they outlive the registrations they describe. +//! +//! A registration publishes through the association's RPKI service, and +//! `EnableIpamInternetRegistryAssociation` is what establishes that service -- +//! the model says "after enabling, you can create Route Origin Authorizations +//! (ROAs)". So an association still in `pending-enable` publishes nothing. +//! +//! `ClientToken` is an idempotency token on every mutating operation here. The +//! association records the tokens it has served, keyed by operation, so a retry +//! replays the delta (or the association) the first call produced instead of +//! failing on the change that call already made. A record carries a +//! fingerprint of what the original call asked for, so a token reused with +//! different parameters is answered with `IdempotentParameterMismatch` rather +//! than a success for a change nobody made, and a timestamp, so the records do +//! not accumulate in the association forever. + +use std::collections::BTreeSet; use chrono::Utc; use fakecloud_aws::ec2query::{ec2_elem, ec2_list}; +use fakecloud_aws::xml::xml_escape; use fakecloud_core::service::{AwsRequest, AwsResponse, AwsServiceError}; use crate::service::Ec2Service; use crate::service_helpers::{ - gen_id, indexed_list, invalid_parameter_value, not_found, require, validate_enum, - validate_max_results, + filter_value_matches, gen_id, indexed_list, invalid_parameter_value, not_found, paginate, + parse_filters, require, validate_enum, validate_max_results, Filter, }; use crate::state::{ - Ec2State, IpamInternetRegistryAssociation, IpamRoutingPolicyRegistration, - IpamRoutingPolicyRegistrationDelta, Tag, + Ec2State, IpamIdempotencyRecord, IpamInternetRegistryAssociation, + IpamRoutingPolicyRegistration, IpamRoutingPolicyRegistrationDelta, Tag, }; const RIRS: &[&str] = &["ripe", "apnic", "arin", "lacnic"]; -fn mr(req: &AwsRequest) -> Result<(), AwsServiceError> { - validate_max_results(&req.query_params, 5, 1000) +/// `MaxResults` and the raw `NextToken` for the module's paginated reads. +/// `IpamMaxResults` carries `@range 5..1000`. +fn page_params(req: &AwsRequest) -> Result<(Option, Option), AwsServiceError> { + validate_max_results(&req.query_params, 5, 1000)?; + // `validate_max_results` only range-checks a value it can parse, so a + // `MaxResults` that is not a number has to be rejected here: dropping it + // would leave the read unpaginated and hand back the whole set, which is + // the opposite of what the caller asked for. A `NextToken` that is not a + // cursor is rejected the same way. + let max_results = + match req.query_params.get("MaxResults").filter(|v| !v.is_empty()) { + Some(v) => Some(v.parse::().map_err(|_| { + invalid_parameter_value(format!("Invalid value '{v}' for MaxResults")) + })?), + None => None, + }; + let next_token = req + .query_params + .get("NextToken") + .filter(|v| !v.is_empty()) + .cloned(); + Ok((max_results, next_token)) +} + +/// Reject a `NextToken` that is not the offset cursor [`paginate`] hands back, +/// rather than silently restarting the caller from the top. +fn validate_offset_token(token: Option<&str>) -> Result<(), AwsServiceError> { + match token { + Some(t) if t.parse::().is_err() => Err(invalid_parameter_value(format!( + "Invalid value '{t}' for NextToken" + ))), + _ => Ok(()), + } +} + +/// `MaxResults` and `NextToken` for the reads that page by offset. +fn pagination(req: &AwsRequest) -> Result<(Option, Option), AwsServiceError> { + let (max_results, next_token) = page_params(req)?; + validate_offset_token(next_token.as_deref())?; + Ok((max_results, next_token)) +} + +/// Page a rendered item list into the operation's set element, plus the +/// `nextToken` every one of these results models. +fn paged_response( + action: &'static str, + req: &AwsRequest, + wrapper: &str, + items: &[String], + page: (Option, Option), +) -> AwsResponse { + let (max_results, next_token) = page; + let (items, token) = paginate(items, next_token.as_deref(), max_results); + page_response(action, req, wrapper, &items, token) +} + +/// One already-paged set of rendered items, plus the `nextToken` that fetches +/// what did not fit. +fn page_response( + action: &'static str, + req: &AwsRequest, + wrapper: &str, + page: &[String], + token: Option, +) -> AwsResponse { + Ec2Service::respond( + action, + &req.request_id, + &format!( + "{}{}", + ec2_list(wrapper, page), + token.map(|t| ec2_elem("nextToken", &t)).unwrap_or_default() + ), + ) +} + +/// Apply a request's `Filter.N` set: values within a filter are OR'd and the +/// filters themselves are AND'd, which is what AWS does. `candidates` maps a +/// filter name to the values an item offers under it; `None` means the name is +/// not one this operation supports, and an unsupported name matches nothing -- +/// the same way the rest of the EC2 describes treat one. +fn matches_filters(filters: &[Filter], candidates: impl Fn(&str) -> Option>) -> bool { + filters.iter().all(|f| match candidates(&f.name) { + Some(values) => f + .values + .iter() + .any(|want| values.iter().any(|have| filter_value_matches(want, have))), + None => false, + }) } fn region_of(req: &AwsRequest) -> String { @@ -40,6 +144,171 @@ fn dry_run(req: &AwsRequest) -> bool { .is_some_and(|v| v.eq_ignore_ascii_case("true")) } +fn client_token(req: &AwsRequest) -> Option { + req.query_params + .get("ClientToken") + .filter(|v| !v.is_empty()) + .cloned() +} + +/// How long a served idempotency token keeps replaying. EC2 documents a +/// 24-hour idempotency window for these tokens, so a record older than that +/// can no longer serve a retry and is only weight in the snapshot. +const CLIENT_TOKEN_TTL_SECONDS: i64 = 24 * 60 * 60; + +/// How many records one association keeps at most. The window alone is not a +/// bound: `ClientToken` is an `@idempotencyToken`, so an SDK fills a fresh +/// UUID in on every call and a tight create/delete loop would leave hundreds +/// of thousands of live records inside the window. The oldest go first, which +/// is the order they stop being useful in. +const CLIENT_TOKEN_MAX_RECORDS: usize = 1000; + +/// Idempotency records are scoped to the operation, so a token a caller reuses +/// across two different calls cannot replay the other one's result. +fn token_key(action: &str, token: &str) -> String { + format!("{action}:{token}") +} + +/// A fingerprint of everything a mutating request asks for, so a retry can be +/// told apart from a token reused with different parameters. Every parameter +/// counts except the ones that do not describe the change: the protocol's own +/// envelope, the token itself, `DryRun` -- a dry run rehearses the same change, +/// and one carrying a served token replays -- and the SigV4 parameters a +/// presigned URL carries, which differ between two signings of the same call. +/// +/// Each pair is length-prefixed, so two different parameter sets cannot +/// flatten to the same string: a `DeltaJson` carrying `&` or `=` would +/// otherwise be able to impersonate another request's parameters and replay +/// its result. +fn request_fingerprint(req: &AwsRequest) -> String { + let mut pairs: Vec<(&str, &str)> = req + .query_params + .iter() + .filter(|(k, _)| { + !matches!(k.as_str(), "Action" | "Version" | "ClientToken" | "DryRun") + && !k.starts_with("X-Amz-") + }) + .map(|(k, v)| (k.as_str(), v.as_str())) + .collect(); + // `query_params` is a hash map, so its iteration order is not stable. + pairs.sort_unstable(); + let mut out = String::new(); + for (k, v) in pairs { + out.push_str(&format!("{}:{k}={}:{v};", k.len(), v.len())); + } + out +} + +/// AWS answers a token reused with different parameters with this rather than +/// the original result: the divergent call asked for something that was never +/// applied, and reporting success for it sends the caller on with a wrong +/// picture of the association. +fn idempotent_parameter_mismatch(token: &str) -> AwsServiceError { + AwsServiceError::aws_error( + http::StatusCode::BAD_REQUEST, + "IdempotentParameterMismatch", + format!("The client token '{token}' was already used with different parameters"), + ) +} + +/// Whether a record has fallen out of the idempotency window. A record whose +/// timestamp cannot be read cannot be aged, so it is treated as expired rather +/// than kept forever. +fn token_expired(record: &IpamIdempotencyRecord) -> bool { + match chrono::DateTime::parse_from_rfc3339(&record.recorded_at) { + Ok(t) => { + Utc::now() + .signed_duration_since(t.with_timezone(&Utc)) + .num_seconds() + > CLIENT_TOKEN_TTL_SECONDS + } + Err(_) => true, + } +} + +/// The record an earlier call under this idempotency token left, if the retry +/// asks for the same thing. A retry replays it rather than applying the change +/// a second time (or failing on the state the first call left behind); a token +/// reused with different parameters is not a retry at all. +fn replay_record<'a>( + a: &'a IpamInternetRegistryAssociation, + action: &str, + token: Option<&str>, + fingerprint: &str, +) -> Result, AwsServiceError> { + let Some(token) = token else { + return Ok(None); + }; + let Some(record) = a.client_tokens.get(&token_key(action, token)) else { + return Ok(None); + }; + if token_expired(record) { + return Ok(None); + } + if record.fingerprint != fingerprint { + return Err(idempotent_parameter_mismatch(token)); + } + Ok(Some(record)) +} + +/// The delta an earlier call under this idempotency token produced, if the +/// retry asks for the same thing. +fn replay_delta( + a: &IpamInternetRegistryAssociation, + action: &str, + token: Option<&str>, + fingerprint: &str, +) -> Result, AwsServiceError> { + let Some(record) = replay_record(a, action, token, fingerprint)? else { + return Ok(None); + }; + Ok(a.deltas + .iter() + .find(|d| d.delta_id == record.result_id) + .cloned()) +} + +fn record_client_token( + a: &mut IpamInternetRegistryAssociation, + action: &str, + token: Option<&str>, + fingerprint: &str, + result_id: &str, +) { + let Some(token) = token else { + return; + }; + prune_client_tokens(a); + a.client_tokens.insert( + token_key(action, token), + IpamIdempotencyRecord { + result_id: result_id.to_string(), + fingerprint: fingerprint.to_string(), + recorded_at: now_rfc3339(), + }, + ); +} + +/// Drop the records that can no longer serve a retry, so what the association +/// carries into every snapshot stays bounded: the aged-out ones first, then +/// the oldest survivors while the association is still at the cap. +fn prune_client_tokens(a: &mut IpamInternetRegistryAssociation) { + a.client_tokens.retain(|_, r| !token_expired(r)); + while a.client_tokens.len() >= CLIENT_TOKEN_MAX_RECORDS { + let oldest = a + .client_tokens + .iter() + .min_by(|(ak, ar), (bk, br)| ar.recorded_at.cmp(&br.recorded_at).then(ak.cmp(bk))) + .map(|(k, _)| k.clone()); + match oldest { + Some(key) => { + a.client_tokens.remove(&key); + } + None => break, + } + } +} + /// Parse an RFC 3339 time bound, rejecting a malformed one rather than letting /// a byte comparison silently filter everything out. fn parse_time_bound( @@ -85,6 +354,36 @@ fn get_association<'a>( .ok_or_else(|| association_not_found(id)) } +/// A registration is a ROA published through the association's RPKI service, +/// and that service only exists once `EnableIpamInternetRegistryAssociation` +/// has been called. Publishing through an association still in +/// `pending-enable` would make that operation decorative. +/// +/// Only the writes that publish are gated: `CreateIpamRoutingPolicyRegistration`, +/// `ModifyIpamRoutingPolicyRegistration`, and a batch document that adds. A +/// removal -- `DeleteIpamRoutingPolicyRegistration`, or a batch document that +/// only removes -- is deliberately not gated, because gating it would be a +/// trap with no way out: an association restored from a snapshot taken before +/// registrations were gated still carries them in `pending-enable`, and since +/// `DeleteIpamInternetRegistryAssociation` refuses while registrations remain, +/// gating the removal too would leave the association and its ROAs +/// undeletable. Un-publishing is also never the operation that needs the RPKI +/// service to exist. +fn require_enabled(a: &IpamInternetRegistryAssociation) -> Result<(), AwsServiceError> { + if a.state == "enable-complete" { + return Ok(()); + } + Err(AwsServiceError::aws_error( + http::StatusCode::BAD_REQUEST, + "IncorrectState", + format!( + "The internet registry association '{}' is in state '{}' and must be enabled before \ + it can publish routing policy registrations", + a.id, a.state + ), + )) +} + fn association_xml(a: &IpamInternetRegistryAssociation, owner: &str, tags: &[Tag]) -> String { let mut s = String::new(); s.push_str(&ec2_elem("ownerId", owner)); @@ -113,6 +412,31 @@ fn association_xml(a: &IpamInternetRegistryAssociation, owner: &str, tags: &[Tag s } +fn association_matches( + a: &IpamInternetRegistryAssociation, + owner: &str, + tags: &[Tag], + filters: &[Filter], +) -> bool { + matches_filters(filters, |name| match name { + "ipam-internet-registry-association-id" => Some(vec![a.id.clone()]), + "ipam-id" => Some(vec![a.ipam_id.clone()]), + "ipam-region" => Some(vec![a.region.clone()]), + "rir" => Some(vec![a.rir.clone()]), + "organization-handle" => Some(vec![a.organization_handle.clone()]), + "state" => Some(vec![a.state.clone()]), + "owner-id" => Some(vec![owner.to_string()]), + "tag-key" => Some(tags.iter().map(|t| t.key.clone()).collect()), + "tag-value" => Some(tags.iter().map(|t| t.value.clone()).collect()), + other => other.strip_prefix("tag:").map(|key| { + tags.iter() + .filter(|t| t.key == key) + .map(|t| t.value.clone()) + .collect() + }), + }) +} + fn delta_xml(d: &IpamRoutingPolicyRegistrationDelta) -> String { let mut s = String::new(); s.push_str(&ec2_elem("deltaId", &d.delta_id)); @@ -187,6 +511,18 @@ pub(crate) fn create_ipam_internet_registry_association( let rir = require(&req.query_params, "Rir")?; let organization_handle = require(&req.query_params, "OrganizationHandle")?; validate_enum(&req.query_params, "Rir", RIRS)?; + let token = client_token(req); + let fingerprint = request_fingerprint(req); + + let owner = req.account_id.clone(); + let region = region_of(req); + let mut accounts = svc.state.write(); + let state = accounts.get_or_create(&req.account_id); + if !state.ipams.contains_key(&ipam_id) { + return Err(not_found("InvalidIpamId.NotFound", &ipam_id)); + } + // A DryRun validates the request -- including that the IPAM exists, which + // is exactly the failure a dry run is for -- and changes nothing. if dry_run(req) { return Ok(Ec2Service::respond( "CreateIpamInternetRegistryAssociation", @@ -194,9 +530,39 @@ pub(crate) fn create_ipam_internet_registry_association( "", )); } + // A retry that carries the original token gets the original association + // back rather than a second one for the same registry. A call that reuses + // the token for a different IPAM, registry or handle is not a retry, and + // handing it the original association would report an association it never + // asked for. + if let Some(token) = &token { + if let Some(existing) = state + .ipam_ir_associations + .values() + .find(|a| a.client_token.as_deref() == Some(token.as_str())) + { + // An association restored from a snapshot that recorded no + // fingerprint offers nothing to compare against, so a retry on its + // token replays rather than failing on evidence never recorded. + if existing + .create_fingerprint + .as_deref() + .is_some_and(|recorded| recorded != fingerprint) + { + return Err(idempotent_parameter_mismatch(token)); + } + let tags = state.tags.get(&existing.id).cloned().unwrap_or_default(); + return Ok(Ec2Service::respond( + "CreateIpamInternetRegistryAssociation", + &req.request_id, + &format!( + "{}", + association_xml(existing, &owner, &tags) + ), + )); + } + } - let owner = req.account_id.clone(); - let region = region_of(req); let id = gen_id("ipam-ir-assoc"); let association = IpamInternetRegistryAssociation { id: id.clone(), @@ -211,13 +577,10 @@ pub(crate) fn create_ipam_internet_registry_association( child_request_xml: None, registrations: Default::default(), deltas: Vec::new(), + create_fingerprint: token.as_ref().map(|_| fingerprint), + client_token: token, + client_tokens: Default::default(), }; - - let mut accounts = svc.state.write(); - let state = accounts.get_or_create(&req.account_id); - if !state.ipams.contains_key(&association.ipam_id) { - return Err(not_found("InvalidIpamId.NotFound", &association.ipam_id)); - } let tags = { crate::service::tags::apply_tag_specifications( state, @@ -248,6 +611,8 @@ pub(crate) fn enable_ipam_internet_registry_association( let child_handle = require(&req.query_params, "ChildHandle")?; let parent_handle = require(&req.query_params, "ParentHandle")?; let parent_bpki_ta = require(&req.query_params, "ParentBpkiTa")?; + let token = client_token(req); + let fingerprint = request_fingerprint(req); let owner = req.account_id.clone(); let mut accounts = svc.state.write(); @@ -263,17 +628,47 @@ pub(crate) fn enable_ipam_internet_registry_association( "", )); } - // The child request is the RPKI provisioning document the registry needs; - // it is what the caller takes to the RIR to finish setup. - a.child_request_xml = Some(format!( - "\ - {parent_bpki_ta}\ - " - )); - a.state = "enable-complete".to_string(); + // A retry under the original token reports the association the first call + // enabled, leaving the child request it already issued alone. A reuse that + // carries a different service URI or handle is not a retry: it asks for a + // child request this association never issued. + let replaying = replay_record( + a, + "EnableIpamInternetRegistryAssociation", + token.as_deref(), + &fingerprint, + )? + .is_some(); + if !replaying { + // The child request is the RPKI provisioning document the registry + // needs; it is what the caller takes to the RIR to finish setup. Every + // value here is caller-supplied and lands in an attribute value or in + // element text, so it is entity-escaped as it goes in: the response + // escapes the blob as a whole, so an unescaped `&` or `"` would come + // back looking fine and only break when the RIR parses the document. + a.child_request_xml = Some(format!( + "\ + {}\ + ", + xml_escape(&rpki_version), + xml_escape(&service_uri), + xml_escape(&child_handle), + xml_escape(&parent_handle), + xml_escape(&parent_bpki_ta), + )); + a.state = "enable-complete".to_string(); + let association_id = a.id.clone(); + record_client_token( + a, + "EnableIpamInternetRegistryAssociation", + token.as_deref(), + &fingerprint, + &association_id, + ); + } let body = format!( "{}", association_xml(a, &owner, &tags) @@ -293,11 +688,30 @@ pub(crate) fn delete_ipam_internet_registry_association( let owner = req.account_id.clone(); let mut accounts = svc.state.write(); let state = accounts.get_or_create(&req.account_id); - if !state.ipam_ir_associations.contains_key(&id) { - return Err(association_not_found(&id)); + // The model is explicit that the registrations have to be removed before + // the association can go, and deleting one that still publishes ROAs would + // orphan them. EC2 reports a delete blocked by what depends on the + // resource as `DependencyViolation`; the model declares no errors of its + // own for this operation. + if !state + .ipam_ir_associations + .get(&id) + .ok_or_else(|| association_not_found(&id))? + .registrations + .is_empty() + { + return Err(AwsServiceError::aws_error( + http::StatusCode::BAD_REQUEST, + "DependencyViolation", + format!( + "The internet registry association '{id}' still has routing policy \ + registrations; remove them before deleting it" + ), + )); } // A DryRun validates the request -- including that the association exists - // -- and changes nothing, matching how the rest of EC2 treats one. + // and that nothing depends on it -- and changes nothing, matching how the + // rest of EC2 treats one. if dry_run(req) { return Ok(Ec2Service::respond( "DeleteIpamInternetRegistryAssociation", @@ -310,10 +724,8 @@ pub(crate) fn delete_ipam_internet_registry_association( .ipam_ir_associations .remove(&id) .ok_or_else(|| association_not_found(&id))?; - // The response reports the association in its terminal state; the - // registrations published through it go with it. + // The response reports the association in its terminal state. association.state = "delete-complete".to_string(); - association.registrations.clear(); state.tags.remove(&id); Ok(Ec2Service::respond( "DeleteIpamInternetRegistryAssociation", @@ -329,8 +741,9 @@ pub(crate) fn describe_ipam_internet_registry_associations( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - mr(req)?; + let page = pagination(req)?; let ids = indexed_list(&req.query_params, "IpamInternetRegistryAssociationId"); + let filters = parse_filters(&req.query_params); let owner = req.account_id.clone(); let accounts = svc.state.read(); let mut items = Vec::new(); @@ -340,93 +753,179 @@ pub(crate) fn describe_ipam_internet_registry_associations( continue; } let tags = state.tags.get(id).cloned().unwrap_or_default(); + if !association_matches(a, &owner, &tags, &filters) { + continue; + } items.push(association_xml(a, &owner, &tags)); } } - Ok(Ec2Service::respond( + Ok(paged_response( "DescribeIpamInternetRegistryAssociations", - &req.request_id, - &ec2_list("ipamInternetRegistryAssociationSet", &items), + req, + "ipamInternetRegistryAssociationSet", + &items, + page, )) } // ---- routing policy registrations ---- -/// Shared body for Create and Modify: both take the same registration fields -/// and report the delta the change produced. -fn upsert_registration( - svc: &Ec2Service, - req: &AwsRequest, - action: &'static str, -) -> Result { - let id = require(&req.query_params, "IpamInternetRegistryAssociationId")?; - let cidr = require(&req.query_params, "Cidr")?; - let asns = indexed_list(&req.query_params, "Asn"); - if asns.is_empty() { - return Err(invalid_parameter_value("Asns must not be empty")); +/// Which of the two single-CIDR registration writes is being applied. Create +/// rejects a CIDR that is already registered and Modify one that is not. +#[derive(Clone, Copy)] +enum RegistrationWrite { + Create, + Modify, +} + +impl RegistrationWrite { + /// The operation this write serves: it names the response and scopes the + /// idempotency records. + fn action(self) -> &'static str { + match self { + Self::Create => "CreateIpamRoutingPolicyRegistration", + Self::Modify => "ModifyIpamRoutingPolicyRegistration", + } } - let max_length = - match req.query_params.get("MaxLength").filter(|v| !v.is_empty()) { - Some(v) => Some(v.parse::().map_err(|_| { - invalid_parameter_value(format!("Invalid value '{v}' for MaxLength")) - })?), - None => None, - }; - // `IpamRoutingPolicyRegistrationMaxLength` carries `@range 0..48`, and the - // member documents that it must not be shorter than the CIDR's own prefix - // length -- a ROA that authorizes less than the prefix it covers announces - // nothing. - crate::service_helpers::validate_int_range(&req.query_params, "MaxLength", 0, 48)?; - if let (Some(m), Some(prefix_len)) = (max_length, cidr_prefix_len(&cidr)) { - if m < prefix_len { - return Err(invalid_parameter_value(format!( - "MaxLength must be greater than or equal to the prefix length of {cidr}" - ))); + + /// How the delta document spells the change. + fn delta_action(self) -> &'static str { + match self { + Self::Create => "create", + Self::Modify => "modify", } } +} - let mut accounts = svc.state.write(); - let state = accounts.get_or_create(&req.account_id); - let a = get_association(state, &id)?; - // A DryRun validates the request -- including that the association exists - // -- and changes nothing, matching how the rest of EC2 treats one. - if dry_run(req) { - return Ok(Ec2Service::respond(action, &req.request_id, "")); +/// What one request says about an optional member. Modify is a partial update, +/// so leaving a member out and clearing it cannot be the same thing: an +/// omitted member keeps the stored value, and a member the caller spells out +/// as empty (ec2Query) or `null` (a delta document) removes it. +#[derive(Clone, Debug, PartialEq)] +enum FieldUpdate { + Unchanged, + Clear, + Set(T), +} + +impl FieldUpdate { + /// What the member becomes, given what the registration already carries. + fn resolve(&self, previous: Option) -> Option { + match self { + Self::Unchanged => previous, + Self::Clear => None, + Self::Set(v) => Some(v.clone()), + } + } + + /// The value the request set, if it set one. + fn value(&self) -> Option<&T> { + match self { + Self::Set(v) => Some(v), + _ => None, + } } +} + +/// The fields one registration change carries. The single operations parse +/// them from indexed query parameters and a batch entry parses them from its +/// JSON object, so both reach the same validation and the same writer. +struct RegistrationChange { + cidr: String, + asns: Vec, + permit_more_specific_announcements: FieldUpdate, + max_length: FieldUpdate, + description: FieldUpdate, +} - let creating = action == "CreateIpamRoutingPolicyRegistration"; - if creating && a.registrations.contains_key(&cidr) { +/// `IpamRoutingPolicyRegistrationMaxLength` carries `@range 0..48`, and the +/// member documents that it must not be shorter than the CIDR's own prefix +/// length -- a ROA that authorizes less than the prefix it covers announces +/// nothing. Both bounds hold wherever the change came from. +fn validate_change(change: &RegistrationChange) -> Result<(), AwsServiceError> { + if change.asns.is_empty() { + return Err(invalid_parameter_value("Asns must not be empty")); + } + let Some(&m) = change.max_length.value() else { + return Ok(()); + }; + if !(0..=48).contains(&m) { + return Err(invalid_parameter_value( + "MaxLength must be between 0 and 48", + )); + } + if cidr_prefix_len(&change.cidr).is_some_and(|prefix_len| m < prefix_len) { return Err(invalid_parameter_value(format!( - "A routing policy registration already exists for {cidr}" + "MaxLength must be greater than or equal to the prefix length of {}", + change.cidr ))); } - if !creating && !a.registrations.contains_key(&cidr) { - return Err(not_found( - "InvalidIpamRoutingPolicyRegistration.NotFound", - &cidr, - )); + Ok(()) +} + +/// Whether the write is legal against what the association already holds. +/// Kept separate from applying it so a DryRun reaches the same verdict as the +/// real call. +fn check_write( + a: &IpamInternetRegistryAssociation, + cidr: &str, + write: RegistrationWrite, +) -> Result<(), AwsServiceError> { + match write { + RegistrationWrite::Create if a.registrations.contains_key(cidr) => { + Err(invalid_parameter_value(format!( + "A routing policy registration already exists for {cidr}" + ))) + } + RegistrationWrite::Create => Ok(()), + RegistrationWrite::Modify => require_registered(a, cidr), + } +} + +/// A change to a registration, and a delete of one, both need it to be there: +/// neither has anything to act on otherwise, and reporting success would tell +/// the caller a CIDR was changed or removed that never existed. +fn require_registered( + a: &IpamInternetRegistryAssociation, + cidr: &str, +) -> Result<(), AwsServiceError> { + if a.registrations.contains_key(cidr) { + return Ok(()); } + Err(not_found( + "InvalidIpamRoutingPolicyRegistration.NotFound", + cidr, + )) +} - let delta_json = serde_json::json!({ - "action": if creating { "create" } else { "modify" }, - "cidr": cidr, - "asns": asns, - "maxLength": max_length, - }) - .to_string(); - let delta_id = push_delta(a, delta_json); +/// Write one change into the association. A change that lands on a CIDR the +/// association already carries is a partial update: the model requires only +/// `Asns`, so a member the request leaves out keeps the value the registration +/// already carries, and only a member the request clears is removed. +fn apply_write( + a: &mut IpamInternetRegistryAssociation, + change: &RegistrationChange, + delta_id: &str, +) { + let previous = a.registrations.get(&change.cidr).cloned(); + let creating = previous.is_none(); a.registrations.insert( - cidr.clone(), + change.cidr.clone(), IpamRoutingPolicyRegistration { - cidr, - asns, - permit_more_specific_announcements: req - .query_params - .get("PermitMoreSpecificAnnouncements") - .map(|v| v.eq_ignore_ascii_case("true")), - max_length, - description: req.query_params.get("Description").cloned(), - latest_delta_id: delta_id.clone(), + cidr: change.cidr.clone(), + asns: change.asns.clone(), + permit_more_specific_announcements: change.permit_more_specific_announcements.resolve( + previous + .as_ref() + .and_then(|p| p.permit_more_specific_announcements), + ), + max_length: change + .max_length + .resolve(previous.as_ref().and_then(|p| p.max_length)), + description: change + .description + .resolve(previous.as_ref().and_then(|p| p.description.clone())), + latest_delta_id: delta_id.to_string(), state: if creating { "create-complete".to_string() } else { @@ -434,56 +933,267 @@ fn upsert_registration( }, }, ); +} + +/// Read one optional member out of an ec2Query request. An omitted member +/// leaves the stored value alone, because Modify is a partial update; a member +/// spelled with an empty value (`Description=`, `MaxLength=`) clears it. A +/// partial update needs some spelling that says "remove this", and the empty +/// value is the one the wire offers -- EC2 already distinguishes a +/// present-but-empty parameter from an absent one elsewhere (`DeleteTags` +/// deletes only the empty-value tag for `Tag.N.Value=`). +fn query_update( + req: &AwsRequest, + key: &str, + parse: impl Fn(&str) -> Result, +) -> Result, AwsServiceError> { + match req.query_params.get(key) { + None => Ok(FieldUpdate::Unchanged), + Some(v) if v.is_empty() => Ok(FieldUpdate::Clear), + Some(v) => parse(v).map(FieldUpdate::Set), + } +} + +fn change_from_request(req: &AwsRequest) -> Result { + let cidr = require(&req.query_params, "Cidr")?; + Ok(RegistrationChange { + cidr, + asns: indexed_list(&req.query_params, "Asn"), + permit_more_specific_announcements: query_update( + req, + "PermitMoreSpecificAnnouncements", + |v| Ok(v.eq_ignore_ascii_case("true")), + )?, + max_length: query_update(req, "MaxLength", |v| { + v.parse::() + .map_err(|_| invalid_parameter_value(format!("Invalid value '{v}' for MaxLength"))) + })?, + description: query_update(req, "Description", |v| Ok(v.to_string()))?, + }) +} + +/// Shared body for Create and Modify: both take the same registration fields +/// and report the delta the change produced. +fn upsert_registration( + svc: &Ec2Service, + req: &AwsRequest, + write: RegistrationWrite, +) -> Result { + let action = write.action(); + let id = require(&req.query_params, "IpamInternetRegistryAssociationId")?; + let change = change_from_request(req)?; + validate_change(&change)?; + let token = client_token(req); + let fingerprint = request_fingerprint(req); + + let mut accounts = svc.state.write(); + let state = accounts.get_or_create(&req.account_id); + let a = get_association(state, &id)?; + require_enabled(a)?; + // A retry replays the delta the first call produced instead of tripping + // over the registration that call already wrote. A DryRun carrying a + // served token replays too: it still changes nothing. + if let Some(delta) = replay_delta(a, action, token.as_deref(), &fingerprint)? { + return Ok(delta_response(action, req, &delta)); + } + check_write(a, &change.cidr, write)?; + // A DryRun validates the request -- including that the association exists, + // that it is enabled, and that the CIDR is in the state this operation + // needs -- and changes nothing, matching how the rest of EC2 treats one. + if dry_run(req) { + return Ok(Ec2Service::respond(action, &req.request_id, "")); + } + + let delta_id = push_delta(a, change_document(&change, write)); + record_client_token(a, action, token.as_deref(), &fingerprint, &delta_id); + apply_write(a, &change, &delta_id); let delta = a.deltas.last().expect("the delta was just pushed").clone(); Ok(delta_response(action, req, &delta)) } +/// The delta document one single-CIDR change records. It spells the optional +/// members the way a batch document spells them, so the audit trail says what +/// the call asked for: a member left out was left alone, and a `null` one was +/// cleared. +fn change_document(change: &RegistrationChange, write: RegistrationWrite) -> String { + let mut doc = serde_json::Map::new(); + doc.insert("action".to_string(), write.delta_action().into()); + doc.insert("cidr".to_string(), change.cidr.clone().into()); + doc.insert("asns".to_string(), change.asns.clone().into()); + delta_field(&mut doc, "maxLength", &change.max_length); + delta_field(&mut doc, "description", &change.description); + delta_field( + &mut doc, + "permitMoreSpecificAnnouncements", + &change.permit_more_specific_announcements, + ); + serde_json::Value::Object(doc).to_string() +} + +/// Spell one optional member into a delta document: an unchanged member is +/// absent from it, and a cleared one is `null`. +fn delta_field>( + doc: &mut serde_json::Map, + name: &str, + update: &FieldUpdate, +) { + let value = match update { + FieldUpdate::Unchanged => return, + FieldUpdate::Clear => serde_json::Value::Null, + FieldUpdate::Set(v) => v.clone().into(), + }; + doc.insert(name.to_string(), value); +} + pub(crate) fn create_ipam_routing_policy_registration( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - upsert_registration(svc, req, "CreateIpamRoutingPolicyRegistration") + upsert_registration(svc, req, RegistrationWrite::Create) } pub(crate) fn modify_ipam_routing_policy_registration( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - upsert_registration(svc, req, "ModifyIpamRoutingPolicyRegistration") + upsert_registration(svc, req, RegistrationWrite::Modify) } pub(crate) fn delete_ipam_routing_policy_registration( svc: &Ec2Service, req: &AwsRequest, ) -> Result { + const ACTION: &str = "DeleteIpamRoutingPolicyRegistration"; + let id = require(&req.query_params, "IpamInternetRegistryAssociationId")?; let cidr = require(&req.query_params, "Cidr")?; + let token = client_token(req); + let fingerprint = request_fingerprint(req); let mut accounts = svc.state.write(); let state = accounts.get_or_create(&req.account_id); let a = get_association(state, &id)?; + // A token reused for a different CIDR is not a retry: replaying the first + // delete's success would report a CIDR removed that is still registered, + // and the caller would then meet the `DependencyViolation` that CIDR + // raises when it deletes the association it was told was empty. + if let Some(delta) = replay_delta(a, ACTION, token.as_deref(), &fingerprint)? { + return Ok(delta_response(ACTION, req, &delta)); + } + // Removing a registration is deliberately not gated on the association + // being enabled -- see [`require_enabled`]. + // // A DryRun validates the request -- including that the association exists - // -- and changes nothing, matching how the rest of EC2 treats one. + // and that the CIDR is registered -- and changes nothing, matching how the + // rest of EC2 treats one. + require_registered(a, &cidr)?; if dry_run(req) { - return Ok(Ec2Service::respond( - "DeleteIpamRoutingPolicyRegistration", - &req.request_id, - "", - )); - } - if a.registrations.remove(&cidr).is_none() { - return Err(not_found( - "InvalidIpamRoutingPolicyRegistration.NotFound", - &cidr, - )); + return Ok(Ec2Service::respond(ACTION, &req.request_id, "")); } + a.registrations.remove(&cidr); let delta_json = serde_json::json!({ "action": "delete", "cidr": cidr }).to_string(); - push_delta(a, delta_json); + let delta_id = push_delta(a, delta_json); + record_client_token(a, ACTION, token.as_deref(), &fingerprint, &delta_id); let delta = a.deltas.last().expect("the delta was just pushed").clone(); - Ok(delta_response( - "DeleteIpamRoutingPolicyRegistration", - req, - &delta, - )) + Ok(delta_response(ACTION, req, &delta)) +} + +/// One `add` entry of a batch document. Every field is typed, and a field the +/// caller spelled as the wrong JSON type fails the request rather than being +/// dropped: the delta records the document as published, so an entry that is +/// quietly skipped leaves an audit trail claiming a registration nobody made. +fn batch_addition(entry: &serde_json::Value) -> Result { + let cidr = entry + .get("cidr") + .and_then(|v| v.as_str()) + .ok_or_else(|| { + invalid_parameter_value("Each DeltaJson 'add' entry must carry a string 'cidr'") + })? + .to_string(); + let change = RegistrationChange { + cidr, + asns: batch_asns(entry)?, + permit_more_specific_announcements: batch_field( + entry, + "permitMoreSpecificAnnouncements", + serde_json::Value::as_bool, + "a boolean", + )?, + max_length: batch_field(entry, "maxLength", serde_json::Value::as_i64, "an integer")?, + description: batch_field( + entry, + "description", + |v| v.as_str().map(str::to_string), + "a string", + )?, + }; + validate_change(&change)?; + Ok(change) +} + +/// Read one optional batch-entry field, rejecting a value of the wrong type. +/// A field the entry omits leaves the stored value alone -- an `add` for a +/// CIDR that is already registered is a partial update -- and an explicit +/// `null` clears it, which is the conventional way a JSON delta document says +/// "remove this". +fn batch_field( + entry: &serde_json::Value, + name: &str, + read: impl Fn(&serde_json::Value) -> Option, + expected: &str, +) -> Result, AwsServiceError> { + match entry.get(name) { + None => Ok(FieldUpdate::Unchanged), + Some(serde_json::Value::Null) => Ok(FieldUpdate::Clear), + Some(v) => read(v).map(FieldUpdate::Set).ok_or_else(|| { + invalid_parameter_value(format!("DeltaJson '{name}' must be {expected}")) + }), + } +} + +/// `AsnList` is a list of strings on the wire, but an ASN is a number and a +/// hand-written delta document spells it as one. Both spellings are accepted; +/// anything else fails rather than yielding a registration that publishes no +/// route origin authorization at all. +fn batch_asns(entry: &serde_json::Value) -> Result, AwsServiceError> { + let asns = entry + .get("asns") + .ok_or_else(|| invalid_parameter_value("Each DeltaJson 'add' entry must carry 'asns'"))? + .as_array() + .ok_or_else(|| invalid_parameter_value("DeltaJson 'asns' must be an array"))?; + asns.iter() + .map(|v| match v { + serde_json::Value::String(s) => Ok(s.clone()), + serde_json::Value::Number(n) => Ok(n.to_string()), + _ => Err(invalid_parameter_value( + "DeltaJson 'asns' entries must be ASNs written as strings or numbers", + )), + }) + .collect() +} + +/// The CIDRs a batch document removes: either bare strings or objects carrying +/// a `cidr`. +fn batch_removals(doc: &serde_json::Value) -> Result, AwsServiceError> { + let Some(entries) = doc.get("remove") else { + return Ok(Vec::new()); + }; + entries + .as_array() + .ok_or_else(|| invalid_parameter_value("DeltaJson 'remove' must be an array"))? + .iter() + .map(|entry| { + entry + .as_str() + .or_else(|| entry.get("cidr")?.as_str()) + .map(str::to_string) + .ok_or_else(|| { + invalid_parameter_value( + "Each DeltaJson 'remove' entry must be a CIDR string or carry a string \ + 'cidr'", + ) + }) + }) + .collect() } /// A batch of registration changes, described by a JSON document rather than @@ -492,80 +1202,81 @@ pub(crate) fn batch_modify_ipam_routing_policy_registrations( svc: &Ec2Service, req: &AwsRequest, ) -> Result { + const ACTION: &str = "BatchModifyIpamRoutingPolicyRegistrations"; let id = require(&req.query_params, "IpamInternetRegistryAssociationId")?; let delta_json = require(&req.query_params, "DeltaJson")?; let parsed: serde_json::Value = serde_json::from_str(&delta_json) .map_err(|_| invalid_parameter_value("DeltaJson is not valid JSON"))?; + let token = client_token(req); + let fingerprint = request_fingerprint(req); + + // The document lists the registrations to add and the CIDRs to remove. + // Every entry is parsed and validated before anything is written: the + // delta records the whole document as published, so one mistyped entry has + // to fail the request rather than leave an audit trail claiming changes + // that were never applied. + let additions: Vec = match parsed.get("add") { + Some(add) => add + .as_array() + .ok_or_else(|| invalid_parameter_value("DeltaJson 'add' must be an array"))? + .iter() + .map(batch_addition) + .collect::>()?, + None => Vec::new(), + }; + let removals = batch_removals(&parsed)?; let mut accounts = svc.state.write(); let state = accounts.get_or_create(&req.account_id); let a = get_association(state, &id)?; - // A DryRun validates the request -- including that the association exists + // Only a document that publishes needs the RPKI service; one that only + // removes registrations does not -- see [`require_enabled`]. + if !additions.is_empty() { + require_enabled(a)?; + } + if let Some(delta) = replay_delta(a, ACTION, token.as_deref(), &fingerprint)? { + return Ok(delta_response(ACTION, req, &delta)); + } + // The document is checked against the state it would itself leave, in the + // order it is applied below -- additions first, removals after -- so a + // document that adds a CIDR and removes it again is self-consistent rather + // than a not-found against the state it has not been applied to yet. An + // `add` entry may create or update, so it needs no check of its own, but a + // `remove` for a CIDR that neither exists nor is added by the document + // removes nothing and must not be reported as published. + let mut registered: BTreeSet<&str> = a.registrations.keys().map(String::as_str).collect(); + registered.extend(additions.iter().map(|change| change.cidr.as_str())); + for cidr in &removals { + if !registered.remove(cidr.as_str()) { + return Err(not_found( + "InvalidIpamRoutingPolicyRegistration.NotFound", + cidr, + )); + } + } + // A DryRun validates the request -- including every entry of the document // -- and changes nothing, matching how the rest of EC2 treats one. if dry_run(req) { - return Ok(Ec2Service::respond( - "BatchModifyIpamRoutingPolicyRegistrations", - &req.request_id, - "", - )); + return Ok(Ec2Service::respond(ACTION, &req.request_id, "")); } let delta_id = push_delta(a, delta_json.clone()); - - // The document lists the registrations to add and the CIDRs to remove. - if let Some(additions) = parsed.get("add").and_then(|v| v.as_array()) { - for entry in additions { - let Some(cidr) = entry.get("cidr").and_then(|v| v.as_str()) else { - continue; - }; - let asns: Vec = entry - .get("asns") - .and_then(|v| v.as_array()) - .map(|a| { - a.iter() - .filter_map(|v| v.as_str().map(str::to_string)) - .collect() - }) - .unwrap_or_default(); - a.registrations.insert( - cidr.to_string(), - IpamRoutingPolicyRegistration { - cidr: cidr.to_string(), - asns, - permit_more_specific_announcements: entry - .get("permitMoreSpecificAnnouncements") - .and_then(|v| v.as_bool()), - max_length: entry.get("maxLength").and_then(|v| v.as_i64()), - description: entry - .get("description") - .and_then(|v| v.as_str()) - .map(str::to_string), - latest_delta_id: delta_id.clone(), - state: "create-complete".to_string(), - }, - ); - } + record_client_token(a, ACTION, token.as_deref(), &fingerprint, &delta_id); + for change in &additions { + apply_write(a, change, &delta_id); } - if let Some(removals) = parsed.get("remove").and_then(|v| v.as_array()) { - for entry in removals { - if let Some(cidr) = entry.as_str().or_else(|| entry.get("cidr")?.as_str()) { - a.registrations.remove(cidr); - } - } + for cidr in &removals { + a.registrations.remove(cidr); } let delta = a.deltas.last().expect("the delta was just pushed").clone(); - Ok(delta_response( - "BatchModifyIpamRoutingPolicyRegistrations", - req, - &delta, - )) + Ok(delta_response(ACTION, req, &delta)) } pub(crate) fn get_ipam_routing_policy_registrations( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - mr(req)?; + let page = pagination(req)?; let id = require(&req.query_params, "IpamInternetRegistryAssociationId")?; let cidr = req.query_params.get("Cidr").filter(|v| !v.is_empty()); let accounts = svc.state.read(); @@ -578,10 +1289,12 @@ pub(crate) fn get_ipam_routing_policy_registrations( .filter(|r| cidr.is_none_or(|c| &r.cidr == c)) .map(registration_xml) .collect(); - Ok(Ec2Service::respond( + Ok(paged_response( "GetIpamRoutingPolicyRegistrations", - &req.request_id, - &ec2_list("ipamRoutingPolicyRegistrationSet", &items), + req, + "ipamRoutingPolicyRegistrationSet", + &items, + page, )) } @@ -589,7 +1302,7 @@ pub(crate) fn get_ipam_routing_policy_registration_deltas( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - mr(req)?; + let (max_results, next_token) = page_params(req)?; let id = require(&req.query_params, "IpamInternetRegistryAssociationId")?; validate_enum( &req.query_params, @@ -617,19 +1330,46 @@ pub(crate) fn get_ipam_routing_policy_registration_deltas( .filter(|d| end.is_none_or(|e| delta_time(d).is_none_or(|t| t <= e))) .collect(); // Deltas are stored oldest first; `reverse` reports newest first. - if req + let reverse = req .query_params .get("ChronologicalOrder") .map(String::as_str) - == Some("reverse") - { + == Some("reverse"); + if reverse { deltas.reverse(); } - let items: Vec = deltas.into_iter().map(delta_xml).collect(); - Ok(Ec2Service::respond( + // Deltas are appended, so under `reverse` an offset cursor moves: every + // delta recorded between two pages shifts the index of everything the + // first page already reported, and the second page repeats items the + // caller has seen. `reverse` therefore pages by the delta id the next page + // starts at, which does not move. Forward order is stable under appends + // and keeps the shared offset cursor. + let page_start = match next_token.as_deref() { + None => 0, + Some(t) if reverse => deltas + .iter() + .position(|d| d.delta_id == t) + .ok_or_else(|| invalid_parameter_value(format!("Invalid value '{t}' for NextToken")))?, + Some(t) => { + validate_offset_token(Some(t))?; + t.parse::().unwrap_or(0).min(deltas.len()) + } + }; + let page_end = max_results.map_or(deltas.len(), |n| (page_start + n).min(deltas.len())); + let token = (page_end < deltas.len()).then(|| { + if reverse { + deltas[page_end].delta_id.clone() + } else { + page_end.to_string() + } + }); + let items: Vec = deltas.into_iter().map(delta_xml).collect(); + Ok(page_response( "GetIpamRoutingPolicyRegistrationDeltas", - &req.request_id, - &ec2_list("ipamRoutingPolicyRegistrationDeltaSet", &items), + req, + "ipamRoutingPolicyRegistrationDeltaSet", + &items[page_start..page_end], + token, )) } @@ -639,7 +1379,7 @@ pub(crate) fn get_ipam_route_origin_authorizations( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - mr(req)?; + let page = pagination(req)?; let id = require(&req.query_params, "IpamInternetRegistryAssociationId")?; let cidr = req.query_params.get("Cidr").filter(|v| !v.is_empty()); let accounts = svc.state.read(); @@ -661,10 +1401,12 @@ pub(crate) fn get_ipam_route_origin_authorizations( items.push(s); } } - Ok(Ec2Service::respond( + Ok(paged_response( "GetIpamRouteOriginAuthorizations", - &req.request_id, - &ec2_list("ipamRouteOriginAuthorizationSet", &items), + req, + "ipamRouteOriginAuthorizationSet", + &items, + page, )) } @@ -674,8 +1416,9 @@ pub(crate) fn get_ipam_internet_registry_association_asns( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - mr(req)?; + let page = pagination(req)?; let id = require(&req.query_params, "IpamInternetRegistryAssociationId")?; + let filters = parse_filters(&req.query_params); let accounts = svc.state.read(); let a = accounts .get(&req.account_id) @@ -688,12 +1431,20 @@ pub(crate) fn get_ipam_internet_registry_association_asns( let now = now_rfc3339(); let items: Vec = asns .into_iter() + .filter(|asn| { + matches_filters(&filters, |name| match name { + "asn" => Some(vec![asn.to_string()]), + _ => None, + }) + }) .map(|asn| ec2_elem("asn", asn) + &ec2_elem("lastObservedAt", &now)) .collect(); - Ok(Ec2Service::respond( + Ok(paged_response( "GetIpamInternetRegistryAssociationAsns", - &req.request_id, - &ec2_list("ipamInternetRegistryAssociationAsnSet", &items), + req, + "ipamInternetRegistryAssociationAsnSet", + &items, + page, )) } @@ -701,8 +1452,9 @@ pub(crate) fn get_ipam_internet_registry_association_cidrs( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - mr(req)?; + let page = pagination(req)?; let id = require(&req.query_params, "IpamInternetRegistryAssociationId")?; + let filters = parse_filters(&req.query_params); let accounts = svc.state.read(); let a = accounts .get(&req.account_id) @@ -713,12 +1465,20 @@ pub(crate) fn get_ipam_internet_registry_association_cidrs( let items: Vec = a .registrations .keys() + .filter(|cidr| { + matches_filters(&filters, |name| match name { + "cidr" => Some(vec![cidr.to_string()]), + _ => None, + }) + }) .map(|cidr| ec2_elem("cidr", cidr) + &ec2_elem("lastObservedAt", &now)) .collect(); - Ok(Ec2Service::respond( + Ok(paged_response( "GetIpamInternetRegistryAssociationCidrs", - &req.request_id, - &ec2_list("ipamInternetRegistryAssociationCidrSet", &items), + req, + "ipamInternetRegistryAssociationCidrSet", + &items, + page, )) } @@ -729,9 +1489,10 @@ pub(crate) fn get_ipam_discovered_routes( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - mr(req)?; + let page = pagination(req)?; let discovery_id = require(&req.query_params, "IpamResourceDiscoveryId")?; let resource_region = require(&req.query_params, "ResourceRegion")?; + let filters = parse_filters(&req.query_params); let owner = req.account_id.clone(); let accounts = svc.state.read(); let state = accounts @@ -752,6 +1513,18 @@ pub(crate) fn get_ipam_discovered_routes( } for r in a.registrations.values() { let asn = r.asns.first().cloned().unwrap_or_default(); + let keep = matches_filters(&filters, |name| match name { + "ipam-resource-discovery-id" => Some(vec![discovery_id.clone()]), + "resource-region" => Some(vec![resource_region.clone()]), + "resource-owner-id" => Some(vec![owner.clone()]), + "cidr" => Some(vec![r.cidr.clone()]), + "asn" => Some(r.asns.clone()), + "state" => Some(vec!["advertised".to_string()]), + _ => None, + }); + if !keep { + continue; + } items.push(format!( "{}{}{}{}{}{}{}", ec2_elem("ipamResourceDiscoveryId", &discovery_id), @@ -764,10 +1537,12 @@ pub(crate) fn get_ipam_discovered_routes( )); } } - Ok(Ec2Service::respond( + Ok(paged_response( "GetIpamDiscoveredRoutes", - &req.request_id, - &ec2_list("ipamDiscoveredRouteSet", &items), + req, + "ipamDiscoveredRouteSet", + &items, + page, )) } @@ -778,8 +1553,9 @@ pub(crate) fn get_ipam_route_protection_findings( svc: &Ec2Service, req: &AwsRequest, ) -> Result { - mr(req)?; + let (max_results, next_token) = pagination(req)?; let ipam_id = require(&req.query_params, "IpamId")?; + let filters = parse_filters(&req.query_params); let owner = req.account_id.clone(); let accounts = svc.state.read(); let state = accounts @@ -805,6 +1581,18 @@ pub(crate) fn get_ipam_route_protection_findings( } else { ("valid", "strict") }; + let keep = matches_filters(&filters, |name| match name { + "resource-owner-id" => Some(vec![owner.clone()]), + "resource-region" => Some(vec![a.region.clone()]), + "cidr" => Some(vec![r.cidr.clone()]), + "asn" => Some(r.asns.clone()), + "rpki-status" => Some(vec![status.to_string()]), + "rpki-strength" => Some(vec![strength.to_string()]), + _ => None, + }); + if !keep { + continue; + } // A finding's `roaSet` holds `IpamRouteOriginAuthorization`, whose // prefix member is `prefix`. The `cidr` spelling belongs to // `IpamRouteOriginAuthorizationInfo`, the shape @@ -837,13 +1625,15 @@ pub(crate) fn get_ipam_route_protection_findings( items.push(finding); } } + let (page, token) = paginate(&items, next_token.as_deref(), max_results); Ok(Ec2Service::respond( "GetIpamRouteProtectionFindings", &req.request_id, &format!( - "{}{}", + "{}{}{}", ec2_elem("ipamId", &ipam_id), - ec2_list("routeProtectionFindingSet", &items) + ec2_list("routeProtectionFindingSet", &page), + token.map(|t| ec2_elem("nextToken", &t)).unwrap_or_default() ), )) } @@ -873,7 +1663,8 @@ mod tests { ); } - fn make_association(svc: &Ec2Service) -> String { + /// Create an association, without enabling it: it cannot publish yet. + fn make_pending_association(svc: &Ec2Service) -> String { seed_ipam(svc); let b = body( create_ipam_internet_registry_association( @@ -898,6 +1689,34 @@ mod tests { .to_string() } + fn enable(svc: &Ec2Service, id: &str, service_uri: &str, child_handle: &str) -> String { + body( + enable_ipam_internet_registry_association( + svc, + &req( + "EnableIpamInternetRegistryAssociation", + &[ + ("IpamInternetRegistryAssociationId", id), + ("RpkiVersion", "1"), + ("ServiceUri", service_uri), + ("ChildHandle", child_handle), + ("ParentHandle", "parent"), + ("ParentBpkiTa", "TA=="), + ], + ), + ) + .unwrap(), + ) + } + + /// An association that has been enabled against the registry, which is + /// what a registration needs. + fn make_association(svc: &Ec2Service) -> String { + let id = make_pending_association(svc); + enable(svc, &id, "https://rpki.example/up-down", "child"); + id + } + fn register(svc: &Ec2Service, id: &str, cidr: &str, max_length: Option<&str>) { let mut params: Vec<(&str, &str)> = vec![ ("IpamInternetRegistryAssociationId", id), @@ -914,6 +1733,90 @@ mod tests { .unwrap(); } + fn registrations(svc: &Ec2Service, id: &str, params: &[(&str, &str)]) -> String { + let mut all: Vec<(&str, &str)> = vec![("IpamInternetRegistryAssociationId", id)]; + all.extend_from_slice(params); + body( + get_ipam_routing_policy_registrations( + svc, + &req("GetIpamRoutingPolicyRegistrations", &all), + ) + .unwrap(), + ) + } + + fn stored_child_request(svc: &Ec2Service, id: &str) -> String { + svc.state + .read() + .get("000000000000") + .unwrap() + .ipam_ir_associations + .get(id) + .unwrap() + .child_request_xml + .clone() + .unwrap() + } + + fn batch(svc: &Ec2Service, id: &str, delta_json: &str) -> Result { + batch_modify_ipam_routing_policy_registrations( + svc, + &req( + "BatchModifyIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", id), + ("DeltaJson", delta_json), + ], + ), + ) + } + + fn deltas(svc: &Ec2Service, id: &str, params: &[(&str, &str)]) -> String { + let mut all: Vec<(&str, &str)> = vec![("IpamInternetRegistryAssociationId", id)]; + all.extend_from_slice(params); + body( + get_ipam_routing_policy_registration_deltas( + svc, + &req("GetIpamRoutingPolicyRegistrationDeltas", &all), + ) + .unwrap(), + ) + } + + fn elements(body: &str, tag: &str) -> Vec { + body.split(&format!("<{tag}>")) + .skip(1) + .map(|s| { + s.split(&format!("")) + .next() + .unwrap_or_default() + .to_string() + }) + .collect() + } + + fn with_association( + svc: &Ec2Service, + id: &str, + f: impl FnOnce(&mut IpamInternetRegistryAssociation) -> T, + ) -> T { + let mut accounts = svc.state.write(); + let state = accounts.get_or_create("000000000000"); + f(state.ipam_ir_associations.get_mut(id).unwrap()) + } + + /// Push every recorded idempotency token out of the window, the way the + /// clock does to a record nobody retried in time. + fn age_client_tokens(svc: &Ec2Service, id: &str) { + let stale = (Utc::now() - chrono::Duration::seconds(CLIENT_TOKEN_TTL_SECONDS + 60)) + .to_rfc3339_opts(chrono::SecondsFormat::Millis, true); + with_association(svc, id, |a| { + for record in a.client_tokens.values_mut() { + record.recorded_at = stale.clone(); + } + }); + } + /// A finding's `roaSet` carries `IpamRouteOriginAuthorization`, whose /// prefix member is `prefix`; `cidr` belongs to a different shape and an /// SDK discards it. @@ -935,11 +1838,36 @@ mod tests { !b.contains(""), "cidr is the wrong member name here: {b}" ); - // `IpamRpkiStrength` is `strict | permissive` — nothing else. + // `IpamRpkiStrength` is `strict | permissive` -- nothing else. assert!(b.contains("strict"), "{b}"); assert!(!b.contains("strong"), "{b}"); } + /// The child request is a document the caller hands to the RIR, so every + /// value interpolated into it has to be entity-escaped. The response + /// escapes the blob as a whole, so an unescaped `&` or `"` would look fine + /// on the wire and only break when the registry parses the document. + #[test] + fn the_child_request_escapes_every_interpolated_value() { + let svc = Ec2Service::new(); + let id = make_pending_association(&svc); + enable( + &svc, + &id, + "https://rpki.example/up-down?src=a&v=2", + "ch\"ild", + ); + + let doc = stored_child_request(&svc, &id); + assert!( + doc.contains("service_uri=\"https://rpki.example/up-down?src=a&v=2\""), + "{doc}" + ); + assert!(doc.contains("child_handle=\"ch"ild\""), "{doc}"); + // The attribute never closes early, so the document stays parseable. + assert_eq!(doc.matches('"').count(), 8, "{doc}"); + } + /// A DryRun validates the request, so it cannot report success for an /// association that does not exist. #[test] @@ -990,17 +1918,109 @@ mod tests { ), ) .unwrap(); + assert!(registrations(&svc, &id, &[]).contains("192.0.2.0/24")); + } + + /// A dry run reaches the same verdict as the real call, so it runs after + /// every existence check rather than before them: otherwise it reports + /// success and the call it was rehearsing fails. + #[test] + fn a_dry_run_reaches_the_same_verdict_as_the_real_call() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + register(&svc, &id, "192.0.2.0/24", None); + + // Creating an already-registered CIDR is a conflict, dry run or not. + let err = err_of(create_ipam_routing_policy_registration( + &svc, + &req( + "CreateIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "192.0.2.0/24"), + ("Asn.1", "64512"), + ("DryRun", "true"), + ], + ), + )); + assert_eq!(err.code(), "InvalidParameterValue"); + + // And modifying one that was never registered is still a not-found. + let err = err_of(modify_ipam_routing_policy_registration( + &svc, + &req( + "ModifyIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "198.51.100.0/24"), + ("Asn.1", "64512"), + ("DryRun", "true"), + ], + ), + )); + assert_eq!(err.code(), "InvalidIpamRoutingPolicyRegistration.NotFound"); + + // A dry-run create naming an IPAM that does not exist is the very + // failure a dry run exists to surface. + let err = err_of(create_ipam_internet_registry_association( + &svc, + &req( + "CreateIpamInternetRegistryAssociation", + &[ + ("IpamId", "ipam-ghost"), + ("Rir", "arin"), + ("OrganizationHandle", "ORG-1"), + ("DryRun", "true"), + ], + ), + )); + assert_eq!(err.code(), "InvalidIpamId.NotFound"); + } + + /// The model requires the registrations to be removed before the + /// association goes, so a delete that would orphan published ROAs is + /// refused -- on a dry run exactly as for real. + #[test] + fn deleting_an_association_requires_its_registrations_to_be_gone() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + register(&svc, &id, "192.0.2.0/24", None); + + for dry in ["false", "true"] { + let err = err_of(delete_ipam_internet_registry_association( + &svc, + &req( + "DeleteIpamInternetRegistryAssociation", + &[("IpamInternetRegistryAssociationId", &id), ("DryRun", dry)], + ), + )); + assert_eq!(err.code(), "DependencyViolation", "DryRun={dry}"); + } + // The association and its registration are both still there. + assert!(registrations(&svc, &id, &[]).contains("192.0.2.0/24")); + + delete_ipam_routing_policy_registration( + &svc, + &req( + "DeleteIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "192.0.2.0/24"), + ], + ), + ) + .unwrap(); let b = body( - get_ipam_routing_policy_registrations( + delete_ipam_internet_registry_association( &svc, &req( - "GetIpamRoutingPolicyRegistrations", + "DeleteIpamInternetRegistryAssociation", &[("IpamInternetRegistryAssociationId", &id)], ), ) .unwrap(), ); - assert!(b.contains("192.0.2.0/24"), "{b}"); + assert!(b.contains("delete-complete"), "{b}"); } /// `MaxLength` carries `@range 0..48` and must cover at least the prefix. @@ -1026,41 +2046,783 @@ mod tests { register(&svc, &id, "192.0.2.0/24", Some("32")); } - /// A time bound is compared as an instant, so a delta recorded in the same - /// second as the bound is not silently dropped, and a malformed bound is - /// rejected rather than filtering everything out. + /// Modify takes only `Asns` as required, so the members it leaves out keep + /// the values the registration already carries instead of being cleared. #[test] - fn delta_time_bounds_compare_instants() { + fn modify_is_a_partial_update() { let svc = Ec2Service::new(); let id = make_association(&svc); - register(&svc, &id, "192.0.2.0/24", None); + create_ipam_routing_policy_registration( + &svc, + &req( + "CreateIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "10.0.0.0/16"), + ("Asn.1", "64512"), + ("MaxLength", "24"), + ("Description", "prod prefix"), + ("PermitMoreSpecificAnnouncements", "true"), + ], + ), + ) + .unwrap(); - // A whole-second bound at the epoch start still includes the delta. + modify_ipam_routing_policy_registration( + &svc, + &req( + "ModifyIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "10.0.0.0/16"), + ("Asn.1", "64513"), + ], + ), + ) + .unwrap(); + + let b = registrations(&svc, &id, &[]); + assert!(b.contains("64513"), "{b}"); + assert!(b.contains("24"), "{b}"); + assert!(b.contains("prod prefix"), "{b}"); + assert!( + b.contains("true"), + "{b}" + ); + assert!(b.contains("update-complete"), "{b}"); + + // The ROAs the registration publishes keep the max length too. let b = body( - get_ipam_routing_policy_registration_deltas( + get_ipam_route_origin_authorizations( &svc, &req( - "GetIpamRoutingPolicyRegistrationDeltas", - &[ - ("IpamInternetRegistryAssociationId", &id), - ("StartTime", "2000-01-01T00:00:00Z"), - ], + "GetIpamRouteOriginAuthorizations", + &[("IpamInternetRegistryAssociationId", &id)], ), ) .unwrap(), ); - assert!(b.contains(""), "{b}"); + assert!(b.contains("24"), "{b}"); + } - let err = err_of(get_ipam_routing_policy_registration_deltas( + /// A registration publishes through the association's RPKI service, which + /// only exists once the association has been enabled. + #[test] + fn a_registration_needs_an_enabled_association() { + let svc = Ec2Service::new(); + let id = make_pending_association(&svc); + + let err = err_of(create_ipam_routing_policy_registration( &svc, &req( - "GetIpamRoutingPolicyRegistrationDeltas", + "CreateIpamRoutingPolicyRegistration", &[ ("IpamInternetRegistryAssociationId", &id), - ("StartTime", "banana"), + ("Cidr", "192.0.2.0/24"), + ("Asn.1", "64512"), ], ), )); - assert_eq!(err.code(), "InvalidParameterValue"); + assert_eq!(err.code(), "IncorrectState"); + + let err = err_of(batch( + &svc, + &id, + r#"{"add":[{"cidr":"192.0.2.0/24","asns":["64512"]}]}"#, + )); + assert_eq!(err.code(), "IncorrectState"); + + // Enabling it opens the association up. + enable(&svc, &id, "https://rpki.example/up-down", "child"); + register(&svc, &id, "192.0.2.0/24", None); + assert!(registrations(&svc, &id, &[]).contains("192.0.2.0/24")); + } + + /// An ASN is a number, and a hand-written delta document spells it as one. + /// Dropping those ASNs would publish a registration authorizing nobody. + #[test] + fn a_batch_accepts_asns_written_as_numbers() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + batch( + &svc, + &id, + r#"{"add":[{"cidr":"192.0.2.0/24","asns":[64512,"64513"]}]}"#, + ) + .unwrap(); + + let b = registrations(&svc, &id, &[]); + assert!(b.contains("64512"), "{b}"); + assert!(b.contains("64513"), "{b}"); + + // And the registration publishes the route origin authorizations that + // make the finding `valid` rather than `unknown`. + let b = body( + get_ipam_route_protection_findings( + &svc, + &req("GetIpamRouteProtectionFindings", &[("IpamId", "ipam-1")]), + ) + .unwrap(), + ); + assert!(b.contains("valid"), "{b}"); + } + + /// The delta records the whole document as published, so an entry the + /// caller mistyped fails the request instead of leaving an audit trail + /// claiming registrations that were never applied. + #[test] + fn a_batch_entry_that_cannot_be_applied_fails_the_whole_request() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + for bad in [ + r#"{"add":[{"cidr":"192.0.2.0/24","asns":["64512"]},{"asns":["64513"]}]}"#, + r#"{"add":[{"cidr":"192.0.2.0/24","asns":["64512"]},{"cidr":198,"asns":["64513"]}]}"#, + r#"{"add":[{"cidr":"10.0.0.0/16","asns":["64512"],"maxLength":200}]}"#, + r#"{"add":[{"cidr":"10.0.0.0/16","asns":["64512"],"maxLength":8}]}"#, + r#"{"add":[{"cidr":"10.0.0.0/16","asns":[]}]}"#, + r#"{"add":[{"cidr":"10.0.0.0/16"}]}"#, + r#"{"add":[{"cidr":"10.0.0.0/16","asns":["64512"],"maxLength":"24"}]}"#, + r#"{"remove":[{"asn":"64512"}]}"#, + ] { + let err = err_of(batch(&svc, &id, bad)); + assert_eq!(err.code(), "InvalidParameterValue", "{bad}"); + } + // Removing a CIDR that was never registered is a not-found, not a + // delta claiming a removal that did not happen. + let err = err_of(batch(&svc, &id, r#"{"remove":["203.0.113.0/24"]}"#)); + assert_eq!(err.code(), "InvalidIpamRoutingPolicyRegistration.NotFound"); + + // Nothing was applied and no delta was recorded. + let b = registrations(&svc, &id, &[]); + assert!(!b.contains("192.0.2.0/24"), "{b}"); + let b = body( + get_ipam_routing_policy_registration_deltas( + &svc, + &req( + "GetIpamRoutingPolicyRegistrationDeltas", + &[("IpamInternetRegistryAssociationId", &id)], + ), + ) + .unwrap(), + ); + assert!(!b.contains(""), "{b}"); + } + + /// A batch `add` is documented to create or update, so a second document + /// for the same CIDR updates it -- partially, the way Modify does. + #[test] + fn a_batch_add_updates_a_cidr_it_already_registered() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + batch( + &svc, + &id, + r#"{"add":[{"cidr":"10.0.0.0/16","asns":["64512"],"maxLength":24}]}"#, + ) + .unwrap(); + batch( + &svc, + &id, + r#"{"add":[{"cidr":"10.0.0.0/16","asns":[64513]}]}"#, + ) + .unwrap(); + + let b = registrations(&svc, &id, &[]); + assert_eq!(b.matches("10.0.0.0/16").count(), 1, "{b}"); + assert!(b.contains("64513"), "{b}"); + assert!(b.contains("24"), "{b}"); + assert!(b.contains("update-complete"), "{b}"); + } + + /// A retry under the original `ClientToken` replays the first call's + /// result instead of failing on the change that call already made. + #[test] + fn a_client_token_replays_the_original_result() { + let svc = Ec2Service::new(); + seed_ipam(&svc); + let create = |token: &str| { + body( + create_ipam_internet_registry_association( + &svc, + &req( + "CreateIpamInternetRegistryAssociation", + &[ + ("IpamId", "ipam-1"), + ("Rir", "arin"), + ("OrganizationHandle", "ORG-1"), + ("ClientToken", token), + ], + ), + ) + .unwrap(), + ) + }; + let first = create("token-a"); + assert_eq!(first, create("token-a"), "the retry replays the original"); + assert_ne!(first, create("token-b"), "a new token is a new association"); + + let id = make_association(&svc); + let register_once = |token: &str| { + body( + create_ipam_routing_policy_registration( + &svc, + &req( + "CreateIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "192.0.2.0/24"), + ("Asn.1", "64512"), + ("ClientToken", token), + ], + ), + ) + .unwrap(), + ) + }; + let first = register_once("token-c"); + assert_eq!(first, register_once("token-c")); + // The replay minted no second delta and no second registration. + let b = registrations(&svc, &id, &[]); + assert_eq!(b.matches("192.0.2.0/24").count(), 1, "{b}"); + let b = body( + get_ipam_routing_policy_registration_deltas( + &svc, + &req( + "GetIpamRoutingPolicyRegistrationDeltas", + &[("IpamInternetRegistryAssociationId", &id)], + ), + ) + .unwrap(), + ); + assert_eq!(b.matches("").count(), 1, "{b}"); + } + + /// A time bound is compared as an instant, so a delta recorded in the same + /// second as the bound is not silently dropped, and a malformed bound is + /// rejected rather than filtering everything out. + #[test] + fn delta_time_bounds_compare_instants() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + register(&svc, &id, "192.0.2.0/24", None); + + // A whole-second bound at the epoch start still includes the delta. + let b = body( + get_ipam_routing_policy_registration_deltas( + &svc, + &req( + "GetIpamRoutingPolicyRegistrationDeltas", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("StartTime", "2000-01-01T00:00:00Z"), + ], + ), + ) + .unwrap(), + ); + assert!(b.contains(""), "{b}"); + + let err = err_of(get_ipam_routing_policy_registration_deltas( + &svc, + &req( + "GetIpamRoutingPolicyRegistrationDeltas", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("StartTime", "banana"), + ], + ), + )); + assert_eq!(err.code(), "InvalidParameterValue"); + } + + /// `MaxResults` bounds a page and the `nextToken` it returns fetches the + /// rest, rather than every read handing back the whole set. + #[test] + fn reads_page_and_round_trip_the_next_token() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + for i in 0..7 { + register(&svc, &id, &format!("10.{i}.0.0/16"), None); + } + + let first = registrations(&svc, &id, &[("MaxResults", "5")]); + assert_eq!(first.matches("").count(), 5, "{first}"); + let token = first + .split("") + .nth(1) + .unwrap_or_else(|| panic!("no nextToken in {first}")) + .split("") + .next() + .unwrap() + .to_string(); + + let second = registrations(&svc, &id, &[("MaxResults", "5"), ("NextToken", &token)]); + assert_eq!(second.matches("").count(), 2, "{second}"); + assert!(!second.contains(""), "{second}"); + + // The deltas each registration produced page the same way. + let b = body( + get_ipam_routing_policy_registration_deltas( + &svc, + &req( + "GetIpamRoutingPolicyRegistrationDeltas", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("MaxResults", "5"), + ], + ), + ) + .unwrap(), + ); + assert_eq!(b.matches("").count(), 5, "{b}"); + assert!(b.contains(""), "{b}"); + + // `IpamMaxResults` is `@range 5..1000`. + let err = err_of(get_ipam_routing_policy_registrations( + &svc, + &req( + "GetIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("MaxResults", "1"), + ], + ), + )); + assert_eq!(err.code(), "InvalidParameterValue"); + } + + /// The `Filters` these operations model narrow the result rather than + /// being accepted and discarded. + #[test] + fn filters_narrow_the_results() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + register(&svc, &id, "192.0.2.0/24", None); + + let describe = |params: &[(&str, &str)]| { + body( + describe_ipam_internet_registry_associations( + &svc, + &req("DescribeIpamInternetRegistryAssociations", params), + ) + .unwrap(), + ) + }; + assert!(describe(&[("Filter.1.Name", "rir"), ("Filter.1.Value.1", "arin")]).contains(&id)); + assert!(!describe(&[("Filter.1.Name", "rir"), ("Filter.1.Value.1", "ripe")]).contains(&id)); + assert!( + !describe(&[("Filter.1.Name", "nonsense"), ("Filter.1.Value.1", "arin")]).contains(&id), + "an unknown filter name matches nothing" + ); + + // The per-CIDR and per-ASN views filter too. + let cidrs = body( + get_ipam_internet_registry_association_cidrs( + &svc, + &req( + "GetIpamInternetRegistryAssociationCidrs", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Filter.1.Name", "cidr"), + ("Filter.1.Value.1", "198.51.100.0/24"), + ], + ), + ) + .unwrap(), + ); + assert!(!cidrs.contains("192.0.2.0/24"), "{cidrs}"); + let asns = body( + get_ipam_internet_registry_association_asns( + &svc, + &req( + "GetIpamInternetRegistryAssociationAsns", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Filter.1.Name", "asn"), + ("Filter.1.Value.1", "64512"), + ], + ), + ) + .unwrap(), + ); + assert!(asns.contains("64512"), "{asns}"); + } + + /// A token reused with different parameters is not a retry. Replaying the + /// original result there reports success for a change nobody made: the + /// worst case is a delete of a CIDR that is still registered, which the + /// caller then meets again as a `DependencyViolation` on the association + /// it was told was empty. + #[test] + fn a_client_token_reused_with_different_parameters_is_rejected() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + register(&svc, &id, "192.0.2.0/24", None); + register(&svc, &id, "198.51.100.0/24", None); + + let delete = |cidr: &str| { + delete_ipam_routing_policy_registration( + &svc, + &req( + "DeleteIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", cidr), + ("ClientToken", "token-delete"), + ], + ), + ) + }; + let first = body(delete("192.0.2.0/24").unwrap()); + let err = err_of(delete("198.51.100.0/24")); + assert_eq!(err.code(), "IdempotentParameterMismatch"); + let b = registrations(&svc, &id, &[]); + assert!(b.contains("198.51.100.0/24"), "nothing was removed: {b}"); + // The same call under the same token still replays. + assert_eq!(first, body(delete("192.0.2.0/24").unwrap())); + + // A create that reuses a token for a different registry handle is not + // a retry either. + let create = |handle: &str| { + create_ipam_internet_registry_association( + &svc, + &req( + "CreateIpamInternetRegistryAssociation", + &[ + ("IpamId", "ipam-1"), + ("Rir", "arin"), + ("OrganizationHandle", handle), + ("ClientToken", "token-create"), + ], + ), + ) + }; + create("ORG-1").unwrap(); + assert_eq!( + err_of(create("ORG-2")).code(), + "IdempotentParameterMismatch" + ); + + // And so is a batch document that changed under a reused token. + let batch_once = |doc: &str| { + batch_modify_ipam_routing_policy_registrations( + &svc, + &req( + "BatchModifyIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("DeltaJson", doc), + ("ClientToken", "token-batch"), + ], + ), + ) + }; + batch_once(r#"{"add":[{"cidr":"203.0.113.0/24","asns":["64512"]}]}"#).unwrap(); + let err = err_of(batch_once(r#"{"remove":["198.51.100.0/24"]}"#)); + assert_eq!(err.code(), "IdempotentParameterMismatch"); + let b = registrations(&svc, &id, &[]); + assert!(b.contains("198.51.100.0/24"), "{b}"); + } + + /// `ClientToken` is an `@idempotencyToken`, so an SDK fills a fresh one in + /// on every call rather than only on retries: the records have to age out + /// and stay capped, or a create/delete loop would leave one dead record + /// per call in the association and in every snapshot of it, forever. + #[test] + fn client_token_records_age_out_and_stay_bounded() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + register(&svc, &id, "192.0.2.0/24", None); + + let modify = |token: &str| { + body( + modify_ipam_routing_policy_registration( + &svc, + &req( + "ModifyIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "192.0.2.0/24"), + ("Asn.1", "64512"), + ("ClientToken", token), + ], + ), + ) + .unwrap(), + ) + }; + let first = modify("token-m"); + assert_eq!(first, modify("token-m"), "a retry in the window replays"); + + age_client_tokens(&svc, &id); + assert_ne!( + first, + modify("token-m"), + "an aged-out token no longer replays" + ); + + // Every call mints its own token, so the cap is what bounds the map. + for i in 0..CLIENT_TOKEN_MAX_RECORDS + 50 { + modify(&format!("token-{i}")); + } + let kept = with_association(&svc, &id, |a| a.client_tokens.len()); + assert!(kept <= CLIENT_TOKEN_MAX_RECORDS, "{kept} records kept"); + } + + /// A batch document is checked against the state it would itself leave, so + /// it can remove a CIDR it adds; the removal of a CIDR the document + /// neither holds nor adds is still a not-found. + #[test] + fn a_batch_can_remove_a_cidr_it_adds() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + + batch( + &svc, + &id, + r#"{"add":[{"cidr":"10.0.0.0/16","asns":["64512"]}],"remove":["10.0.0.0/16"]}"#, + ) + .unwrap(); + let b = registrations(&svc, &id, &[]); + assert!(!b.contains("10.0.0.0/16"), "{b}"); + // The document was published, so it left the audit trail behind. + assert_eq!(elements(&deltas(&svc, &id, &[]), "deltaId").len(), 1); + + let err = err_of(batch( + &svc, + &id, + r#"{"add":[{"cidr":"10.0.0.0/16","asns":["64512"]}],"remove":["203.0.113.0/24"]}"#, + )); + assert_eq!(err.code(), "InvalidIpamRoutingPolicyRegistration.NotFound"); + } + + /// Modify is a partial update, so an omitted member keeps its value -- + /// which leaves the caller needing a spelling that says "remove this". An + /// empty value clears the member on the wire, and `null` clears it in a + /// delta document. + #[test] + fn optional_members_can_be_cleared() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + create_ipam_routing_policy_registration( + &svc, + &req( + "CreateIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "10.0.0.0/16"), + ("Asn.1", "64512"), + ("MaxLength", "24"), + ("Description", "prod prefix"), + ("PermitMoreSpecificAnnouncements", "true"), + ], + ), + ) + .unwrap(); + + modify_ipam_routing_policy_registration( + &svc, + &req( + "ModifyIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "10.0.0.0/16"), + ("Asn.1", "64512"), + ("MaxLength", ""), + ("Description", ""), + ("PermitMoreSpecificAnnouncements", ""), + ], + ), + ) + .unwrap(); + + let b = registrations(&svc, &id, &[]); + assert!(!b.contains(""), "{b}"); + assert!(!b.contains(""), "{b}"); + assert!(!b.contains(""), "{b}"); + // The delta says what was asked for: a cleared member is `null`, which + // is how a batch document spells the same change. + let recorded = with_association(&svc, &id, |a| { + a.deltas + .last() + .expect("a delta was recorded") + .delta_json + .clone() + }); + assert!(recorded.contains(r#""maxLength":null"#), "{recorded}"); + assert!(recorded.contains(r#""description":null"#), "{recorded}"); + + // A batch entry clears with `null` and leaves an omitted member alone. + batch( + &svc, + &id, + r#"{"add":[{"cidr":"192.0.2.0/24","asns":["64512"],"maxLength":25, + "description":"edge prefix"}]}"#, + ) + .unwrap(); + batch( + &svc, + &id, + r#"{"add":[{"cidr":"192.0.2.0/24","asns":["64512"],"description":null}]}"#, + ) + .unwrap(); + let b = registrations(&svc, &id, &[("Cidr", "192.0.2.0/24")]); + assert!(!b.contains(""), "{b}"); + assert!(b.contains("25"), "{b}"); + } + + /// Deltas are appended, so an offset into the reversed list moves under a + /// caller who pages: `reverse` pages by delta id instead, and the second + /// page cannot repeat what the first already reported. + #[test] + fn reverse_delta_pages_do_not_repeat_when_deltas_are_appended() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + for i in 0..7 { + register(&svc, &id, &format!("10.{i}.0.0/16"), None); + } + + let reverse = [("ChronologicalOrder", "reverse"), ("MaxResults", "5")]; + let first = deltas(&svc, &id, &reverse); + let seen = elements(&first, "deltaId"); + assert_eq!(seen.len(), 5, "{first}"); + let token = elements(&first, "nextToken") + .pop() + .unwrap_or_else(|| panic!("no nextToken in {first}")); + + // Two more deltas land between the two pages. + register(&svc, &id, "10.7.0.0/16", None); + register(&svc, &id, "10.8.0.0/16", None); + + let mut params = reverse.to_vec(); + params.push(("NextToken", &token)); + let second = deltas(&svc, &id, ¶ms); + let rest = elements(&second, "deltaId"); + assert_eq!(rest.len(), 2, "{second}"); + assert!( + rest.iter().all(|d| !seen.contains(d)), + "a second page must not repeat the first: {first} {second}" + ); + assert!(!second.contains(""), "{second}"); + + // A `reverse` token that names no delta is rejected rather than + // silently restarting the caller from the newest one. + let bad: Vec<(&str, &str)> = vec![ + ("IpamInternetRegistryAssociationId", &id), + ("ChronologicalOrder", "reverse"), + ("NextToken", "ipam-delta-nope"), + ]; + let err = err_of(get_ipam_routing_policy_registration_deltas( + &svc, + &req("GetIpamRoutingPolicyRegistrationDeltas", &bad), + )); + assert_eq!(err.code(), "InvalidParameterValue"); + } + + /// A `MaxResults` that is not a number is rejected the way a `NextToken` + /// that is not a cursor is: taking it as "no limit" would hand back the + /// whole set unpaginated. + #[test] + fn a_max_results_that_is_not_a_number_is_rejected() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + for i in 0..7 { + register(&svc, &id, &format!("10.{i}.0.0/16"), None); + } + + let err = err_of(get_ipam_routing_policy_registrations( + &svc, + &req( + "GetIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("MaxResults", "abc"), + ], + ), + )); + assert_eq!(err.code(), "InvalidParameterValue"); + + let err = err_of(get_ipam_routing_policy_registration_deltas( + &svc, + &req( + "GetIpamRoutingPolicyRegistrationDeltas", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("MaxResults", "abc"), + ], + ), + )); + assert_eq!(err.code(), "InvalidParameterValue"); + + let err = err_of(get_ipam_routing_policy_registrations( + &svc, + &req( + "GetIpamRoutingPolicyRegistrations", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("NextToken", "abc"), + ], + ), + )); + assert_eq!(err.code(), "InvalidParameterValue"); + } + + /// Publishing needs the association's RPKI service; un-publishing does + /// not. An association restored from a snapshot taken before registrations + /// were gated still carries them in `pending-enable`, and the association + /// cannot be deleted while they remain -- gating removals too would strand + /// it and its ROAs for good. + #[test] + fn removing_a_registration_does_not_need_an_enabled_association() { + let svc = Ec2Service::new(); + let id = make_association(&svc); + register(&svc, &id, "192.0.2.0/24", None); + register(&svc, &id, "198.51.100.0/24", None); + with_association(&svc, &id, |a| a.state = "pending-enable".to_string()); + + let err = err_of(create_ipam_routing_policy_registration( + &svc, + &req( + "CreateIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "203.0.113.0/24"), + ("Asn.1", "64512"), + ], + ), + )); + assert_eq!(err.code(), "IncorrectState"); + let err = err_of(batch( + &svc, + &id, + r#"{"add":[{"cidr":"203.0.113.0/24","asns":["64512"]}],"remove":["192.0.2.0/24"]}"#, + )); + assert_eq!(err.code(), "IncorrectState"); + + // Both spellings of a removal go through, so the association can be + // emptied and then deleted. + delete_ipam_routing_policy_registration( + &svc, + &req( + "DeleteIpamRoutingPolicyRegistration", + &[ + ("IpamInternetRegistryAssociationId", &id), + ("Cidr", "192.0.2.0/24"), + ], + ), + ) + .unwrap(); + batch(&svc, &id, r#"{"remove":["198.51.100.0/24"]}"#).unwrap(); + let b = body( + delete_ipam_internet_registry_association( + &svc, + &req( + "DeleteIpamInternetRegistryAssociation", + &[("IpamInternetRegistryAssociationId", &id)], + ), + ) + .unwrap(), + ); + assert!(b.contains("delete-complete"), "{b}"); } } diff --git a/crates/fakecloud-ec2/src/state.rs b/crates/fakecloud-ec2/src/state.rs index 017459ccb..42d31abcf 100644 --- a/crates/fakecloud-ec2/src/state.rs +++ b/crates/fakecloud-ec2/src/state.rs @@ -1184,6 +1184,42 @@ pub struct IpamInternetRegistryAssociation { /// registration changes, so they outlive the registrations themselves. #[serde(default)] pub deltas: Vec, + /// The `ClientToken` the association was created with, if any. A create + /// that retries under the same token replays this association instead of + /// minting a second one against the same registry. + #[serde(default)] + pub client_token: Option, + /// The parameters the create that minted the association asked for, so a + /// retry under `client_token` that asks for something else is answered + /// with `IdempotentParameterMismatch` rather than this association. + #[serde(default)] + pub create_fingerprint: Option, + /// The idempotency tokens this association has already served, keyed by + /// `{action}:{token}` so a token reused across two operations cannot + /// replay the other one's result. + #[serde(default)] + pub client_tokens: BTreeMap, +} + +/// One idempotency token an association has already served. +/// +/// `ClientToken` is an `@idempotencyToken`, so every SDK fills a fresh UUID in +/// on every call rather than only on retries: these records are written far +/// more often than they are read, and they are cloned with the association +/// into every snapshot. They are therefore aged out and capped -- see +/// `CLIENT_TOKEN_TTL_SECONDS` in the handler module. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct IpamIdempotencyRecord { + /// The delta the original call produced, or the association's own id for + /// Enable, which produces no delta. + pub result_id: String, + /// A fingerprint of the parameters the original call carried. A token + /// reused with different parameters is not a retry. + #[serde(default)] + pub fingerprint: String, + /// When the record was written, RFC 3339, so it can be aged out. + #[serde(default)] + pub recorded_at: String, } /// One CIDR's route origin authorization within an association. diff --git a/crates/fakecloud-server/src/reaper.rs b/crates/fakecloud-server/src/reaper.rs index 9ff385e57..7f1e62b52 100644 --- a/crates/fakecloud-server/src/reaper.rs +++ b/crates/fakecloud-server/src/reaper.rs @@ -27,9 +27,7 @@ pub fn reap_stale_containers() { return; }; - let reaped = reap_orphans(&cli, &["ps", "-a"], |id| { - vec!["rm".to_string(), "-f".to_string(), id.to_string()] - }); + let reaped = reap_orphans(&cli, &["ps", "-a"], &["rm", "-f"]); if reaped > 0 { tracing::info!(count = reaped, "reaped orphaned backing containers"); } @@ -39,20 +37,18 @@ pub fn reap_stale_containers() { // so prune networks *after* containers. `network rm` is a no-op for an // already-gone network, so a partial container reap above doesn't wedge // this pass. - let reaped_networks = reap_orphans(&cli, &["network", "ls"], |id| { - vec!["network".to_string(), "rm".to_string(), id.to_string()] - }); + let reaped_networks = reap_orphans(&cli, &["network", "ls"], &["network", "rm"]); if reaped_networks > 0 { tracing::info!(count = reaped_networks, "reaped orphaned backing networks"); } } /// List objects carrying the `fakecloud-instance` label via -/// ` --filter label=fakecloud-instance`, then run the -/// `remove_argv(id)` command for every object whose owning PID is no longer +/// ` --filter label=fakecloud-instance`, then run +/// ` ` for every object whose owning PID is no longer /// alive (skipping the current process and live owners). Returns the number /// removed. Shared by the container and network reap passes. -fn reap_orphans(cli: &str, list_args: &[&str], remove_argv: impl Fn(&str) -> Vec) -> usize { +fn reap_orphans(cli: &str, list_args: &[&str], remove_args: &[&str]) -> usize { let mut args: Vec<&str> = list_args.to_vec(); args.extend_from_slice(&[ "--filter", @@ -80,7 +76,9 @@ fn reap_orphans(cli: &str, list_args: &[&str], remove_argv: impl Fn(&str) -> Vec ) { continue; } - let removed = fakecloud_core::container_net::bounded_status(cli, &remove_argv(id)); + let mut remove_argv = remove_args.to_vec(); + remove_argv.push(id); + let removed = fakecloud_core::container_net::bounded_status(cli, &remove_argv); if removed { reaped += 1; } diff --git a/crates/fakecloud-testkit/Cargo.toml b/crates/fakecloud-testkit/Cargo.toml index 81edc72da..3a68652a1 100644 --- a/crates/fakecloud-testkit/Cargo.toml +++ b/crates/fakecloud-testkit/Cargo.toml @@ -71,6 +71,11 @@ sdk-clients = [ ] [dependencies] +# For the bounded container-CLI helpers in `fakecloud_core::container_net`. +# The harness sweeps leftover containers on drop and must not hang the test +# run doing it; sharing core's helpers keeps that bound in one place instead +# of a copy that drifts (and re-grows the leaks core has already fixed). +fakecloud-core = { workspace = true } tokio = { workspace = true } reqwest = { workspace = true } tempfile = { workspace = true } diff --git a/crates/fakecloud-testkit/src/lib.rs b/crates/fakecloud-testkit/src/lib.rs index 57cf5e948..6408c9e82 100644 --- a/crates/fakecloud-testkit/src/lib.rs +++ b/crates/fakecloud-testkit/src/lib.rs @@ -36,6 +36,7 @@ use std::time::Duration; use aws_config::BehaviorVersion; use aws_credential_types::Credentials; use aws_types::region::Region; +use fakecloud_core::container_net::{bounded_output, bounded_status, cli_available}; /// A test server that spawns fakecloud on a random port. pub struct TestServer { @@ -466,6 +467,10 @@ fn graceful_kill(child: &mut Child) { /// support entirely (`FAKECLOUD_CONTAINER_CLI=false`). Most e2e tests /// never touch the lambda/rds/elasticache runtimes, so the sweep is /// pure overhead on the drop path for those. +/// +/// Every CLI call goes through `fakecloud_core::container_net`'s bounded +/// helpers: a wedged daemon blocks on connect forever, and an unbounded call +/// here would hang the whole test run rather than the one container sweep. fn sweep_instance_containers(cli: &str, pid: u32) { if cli.is_empty() || cli == "false" { return; @@ -476,71 +481,7 @@ fn sweep_instance_containers(cli: &str, pid: u32) { return; }; for id in ids.split_whitespace() { - bounded_status(cli, &["rm", "-f", id]); - } -} - -/// How long any container-CLI call in the harness may take. A healthy daemon -/// answers immediately; a wedged one (stale `DOCKER_HOST`, Docker Desktop mid -/// start, a broken socket) blocks on connect forever, and an unbounded call -/// here hangs the whole test run rather than the one container sweep. -const CLI_TIMEOUT: Duration = Duration::from_secs(10); - -/// Run a container-CLI command, returning its stdout, or `None` when it fails -/// or outruns [`CLI_TIMEOUT`]. -fn bounded_output(cli: &str, args: &[&str]) -> Option { - let mut child = Command::new(cli) - .args(args) - .stdout(Stdio::piped()) - .stderr(Stdio::null()) - .spawn() - .ok()?; - // Drain stdout while waiting: a child that fills the pipe buffer blocks on - // write, so waiting for exit first would deadlock until the deadline. - let mut stdout = child.stdout.take()?; - let reader = std::thread::spawn(move || { - let mut buf = Vec::new(); - let _ = std::io::Read::read_to_end(&mut stdout, &mut buf); - buf - }); - if !wait_bounded(&mut child) { - return None; - } - let status = child.wait().ok()?; - let buf = reader.join().ok()?; - status - .success() - .then(|| String::from_utf8_lossy(&buf).into_owned()) -} - -/// Run a container-CLI command for its effect only, bounded the same way. -fn bounded_status(cli: &str, args: &[&str]) { - if let Ok(mut child) = Command::new(cli) - .args(args) - .stdout(Stdio::null()) - .stderr(Stdio::null()) - .spawn() - { - wait_bounded(&mut child); - } -} - -/// Wait for `child` up to [`CLI_TIMEOUT`], killing it on expiry. Returns -/// whether it exited on its own. -fn wait_bounded(child: &mut std::process::Child) -> bool { - let deadline = std::time::Instant::now() + CLI_TIMEOUT; - loop { - match child.try_wait() { - Ok(Some(_)) => return true, - Ok(None) => {} - Err(_) => return false, - } - if std::time::Instant::now() >= deadline { - let _ = child.kill(); - let _ = child.wait(); - return false; - } - std::thread::sleep(Duration::from_millis(25)); + let _ = bounded_status(cli, &["rm", "-f", id]); } } @@ -644,6 +585,15 @@ fn find_binary() -> String { ); } +/// Pick the container CLI the harness tells the server to use, falling back to +/// `docker` when neither answers (the server re-probes and disables container +/// support itself). +/// +/// Liveness comes from `fakecloud_core::container_net::cli_available`, which +/// bounds the ` info` probe: an unreachable or wedged daemon (stale +/// `DOCKER_HOST`, Docker Desktop mid start, a broken socket) can leave the CLI +/// blocked on connect *forever*, and without that bound a conformance `*_probe` +/// that only calls `TestServer::start()` would never return. fn detect_container_cli() -> String { if cli_available("docker") { "docker".to_string() @@ -654,41 +604,6 @@ fn detect_container_cli() -> String { } } -/// True when ` info` succeeds within a bounded window. A healthy daemon -/// answers in well under a second; an unreachable or wedged daemon (stale -/// `DOCKER_HOST`, Docker Desktop mid start, a broken socket) can leave the CLI -/// blocked on connect *forever*. Without this bound, `detect_container_cli` -/// hangs the whole test harness — a conformance `*_probe` that only calls -/// `TestServer::start()` would never return. Mirrors -/// `fakecloud_core::container_net::cli_available` (testkit can't depend on -/// core, so the bounded pattern is duplicated here on purpose). -fn cli_available(cli: &str) -> bool { - static CACHE: std::sync::OnceLock>> = - std::sync::OnceLock::new(); - let cache = CACHE.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new())); - if let Some(&cached) = cache.lock().unwrap().get(cli) { - return cached; - } - let result = probe_cli(cli); - cache.lock().unwrap().insert(cli.to_string(), result); - result -} - -fn probe_cli(cli: &str) -> bool { - let Ok(mut child) = Command::new(cli) - .arg("info") - .stdout(Stdio::null()) - .stderr(Stdio::null()) - .spawn() - else { - return false; - }; - if !wait_bounded(&mut child) { - return false; - } - child.wait().map(|s| s.success()).unwrap_or(false) -} - /// Prefix that `fakecloud-server` prints before the bound port on stdout. /// Must stay in sync with `PORT_HANDSHAKE_PREFIX` in /// `crates/fakecloud-server/src/main.rs`. @@ -1057,37 +972,55 @@ mod handshake_tests { } #[cfg(test)] -mod bounded_cli_tests { +mod container_sweep_tests { use super::*; + use fakecloud_core::container_net::CLI_PROBE_TIMEOUT; - /// A container CLI that never returns must not hang the harness. `sleep` - /// stands in for a wedged daemon: the bound has to cut it off. + /// Tests that disable container support (`FAKECLOUD_CONTAINER_CLI=false`) + /// must not pay for a subprocess on every server drop. #[test] - fn a_hanging_cli_call_is_cut_off() { + fn a_disabled_cli_skips_the_sweep() { let start = std::time::Instant::now(); - let mut child = Command::new("sleep") - .arg("600") - .stdout(Stdio::null()) - .stderr(Stdio::null()) - .spawn() - .expect("sleep is available"); - assert!( - !wait_bounded(&mut child), - "a hung call must not report success" - ); + sweep_instance_containers("false", std::process::id()); + sweep_instance_containers("", std::process::id()); assert!( - start.elapsed() < CLI_TIMEOUT + Duration::from_secs(5), - "the wait must end at the bound, not run on" + start.elapsed() < Duration::from_secs(1), + "a disabled CLI must not spawn anything, took {:?}", + start.elapsed() ); } #[test] - fn a_prompt_cli_call_returns_its_output() { - assert_eq!( - bounded_output("echo", &["container-id"]) - .as_deref() - .map(str::trim), - Some("container-id") + fn a_missing_cli_leaves_the_sweep_a_no_op() { + // Nothing to sweep and nothing to hang on: an absent binary fails to + // spawn, so the sweep returns instead of panicking on the drop path. + sweep_instance_containers("definitely-not-a-real-cli-binary-xyz-123", 1); + } + + /// The sweep runs from `Drop`, once per test server. A container CLI that + /// never returns (a wedged daemon) must end at the shared bound rather than + /// hang the whole test run -- the reason the harness routes every call + /// through `fakecloud_core::container_net`'s bounded helpers instead of + /// keeping its own copy. + #[cfg(unix)] + #[test] + fn a_hanging_cli_cannot_hang_the_sweep() { + use std::os::unix::fs::PermissionsExt; + + let dir = std::env::temp_dir().join(format!("fc-sweeptest-{}", std::process::id())); + std::fs::create_dir_all(&dir).unwrap(); + let script = dir.join("hangcli"); + std::fs::write(&script, "#!/bin/sh\nsleep 600\n").unwrap(); + std::fs::set_permissions(&script, std::fs::Permissions::from_mode(0o755)).unwrap(); + + let start = std::time::Instant::now(); + sweep_instance_containers(script.to_str().unwrap(), std::process::id()); + let elapsed = start.elapsed(); + + std::fs::remove_dir_all(&dir).ok(); + assert!( + elapsed < CLI_PROBE_TIMEOUT + Duration::from_secs(5), + "the sweep took {elapsed:?}, expected it bounded near {CLI_PROBE_TIMEOUT:?}" ); } }