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