From 5e1e998db3a2be167e8286bbd17178e37e57ad23 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 11:39:30 -0300 Subject: [PATCH 1/6] fix(ecs,lambda): fall back to the cached image when a pull fails A bare `docker pull` always contacts the registry, even for an image already in the local cache, so any registry error failed the launch. Anonymous pulls from public.ecr.aws are rate limited per source IP, and a burst of task launches -- or a shared CI runner address -- gets `429 Too Many Requests`. The ECS task then stopped with TaskFailedToStart and every Batch job on it went FAILED. That is the cause of the batch_real_execution e2e flakes seen on nearly every CI run: array, depends_on, timeout, retry and submit tests all launch alpine containers concurrently. Reproduced by running the suite with a container CLI that answers every pull with a 429: the old pull path fails 6 of 7 tests with the 429 in statusReason, the new one passes all 7. The ECS agent's default ECS_IMAGE_PULL_BEHAVIOR uses the cached image when a pull fails and retries pulls with backoff. The new fakecloud_core::container_image::pull_image does the same and is used by the ECS task runtime (so Batch) and by Lambda PackageType=Image function starts and prewarms: - pull succeeds -> use it - pull fails, image cached -> use the cached image, log the pull error - pull rate limited, nothing cached -> retry with backoff (5 attempts) - any other failure with nothing cached -> fail at once with the registry's error, so a missing image still fails fast The Batch e2e wait helper now asserts the expected terminal status itself and, on a mismatch, prints each job's status, statusReason and exit code, including every array child, so a CI failure names its cause. --- crates/fakecloud-core/src/container_image.rs | 244 ++++++++++++++++++ crates/fakecloud-core/src/lib.rs | 1 + .../tests/batch_real_execution.rs | 62 ++++- .../src/runtime/task_lifecycle.rs | 17 +- crates/fakecloud-lambda/src/runtime/docker.rs | 46 ++-- website/content/docs/services/ecs.md | 2 +- website/content/docs/services/lambda.md | 2 +- 7 files changed, 322 insertions(+), 52 deletions(-) create mode 100644 crates/fakecloud-core/src/container_image.rs diff --git a/crates/fakecloud-core/src/container_image.rs b/crates/fakecloud-core/src/container_image.rs new file mode 100644 index 000000000..cde5b6cfc --- /dev/null +++ b/crates/fakecloud-core/src/container_image.rs @@ -0,0 +1,244 @@ +//! Image pulls for the runtimes that launch user-supplied images (ECS, and +//! Batch through it; Lambda `PackageType=Image` functions). +//! +//! A bare `docker pull` always contacts the registry, even when the image is +//! already in the local cache, so any registry hiccup fails the launch. The +//! common one is rate limiting: anonymous pulls from `public.ecr.aws` are +//! capped per source IP, and a burst of task launches -- or several processes +//! sharing one NAT address -- gets `429 Too Many Requests` back. +//! +//! The ECS container agent handles this the way this module does. Under its +//! default `ECS_IMAGE_PULL_BEHAVIOR`, a failed pull falls back to the image +//! cached on the instance, and pulls are retried with backoff before giving +//! up. A pull that fails with nothing cached still fails, with the registry's +//! own error. + +use std::path::Path; +use std::time::Duration; + +use tokio::process::Command; + +/// Pull attempts made while the registry keeps rate limiting, before the +/// last error is returned. +const MAX_PULL_ATTEMPTS: u32 = 5; + +/// Delay before the first retry; doubles on each further retry. +const BASE_RETRY_DELAY: Duration = Duration::from_secs(1); + +/// How an image became available for a launch. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum PulledImage { + /// The registry served the image. + Pulled, + /// The pull failed but the image was already cached locally, so the + /// cached copy is used. Carries the pull's error for logging. + Cached { pull_error: String }, +} + +/// Pull `reference` with the container `cli`, falling back to a locally +/// cached copy when the pull fails and retrying while the registry is rate +/// limiting. `docker_config` is exported as `DOCKER_CONFIG` for the pull so +/// registry credentials resolve. +/// +/// Returns the pull's stderr as the error when the image is neither pullable +/// nor cached. +pub async fn pull_image( + cli: &str, + docker_config: Option<&Path>, + reference: &str, +) -> Result { + pull_image_with(cli, docker_config, reference, BASE_RETRY_DELAY).await +} + +async fn pull_image_with( + cli: &str, + docker_config: Option<&Path>, + reference: &str, + base_delay: Duration, +) -> Result { + let mut delay = base_delay; + let mut attempt = 1; + loop { + let mut cmd = Command::new(cli); + if let Some(p) = docker_config { + cmd.env("DOCKER_CONFIG", p); + } + let out = cmd + .args(["pull", reference]) + .output() + .await + .map_err(|e| format!("{cli} pull: {e}"))?; + if out.status.success() { + return Ok(PulledImage::Pulled); + } + let pull_error = String::from_utf8_lossy(&out.stderr).trim().to_string(); + + if image_cached(cli, reference).await { + tracing::warn!( + image = %reference, + error = %pull_error, + "image pull failed; using the locally cached image" + ); + return Ok(PulledImage::Cached { pull_error }); + } + if attempt >= MAX_PULL_ATTEMPTS || !is_rate_limited(&pull_error) { + return Err(pull_error); + } + tracing::info!( + image = %reference, + attempt, + retry_in_ms = delay.as_millis() as u64, + "image pull rate limited; retrying" + ); + tokio::time::sleep(delay).await; + delay *= 2; + attempt += 1; + } +} + +/// Whether `reference` resolves to an image in the local cache. +async fn image_cached(cli: &str, reference: &str) -> bool { + Command::new(cli) + .args(["image", "inspect", reference]) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .status() + .await + .map(|s| s.success()) + .unwrap_or(false) +} + +/// Whether a pull error is the registry throttling the caller. Docker reports +/// the status line (`429 Too Many Requests`) or the registry error code +/// (`toomanyrequests`); ECR Public's throttle message is `Rate exceeded`. +fn is_rate_limited(stderr: &str) -> bool { + let lower = stderr.to_ascii_lowercase(); + lower.contains("toomanyrequests") + || lower.contains("too many requests") + || lower.contains("rate exceeded") +} + +#[cfg(all(test, unix))] +mod tests { + use super::*; + use std::os::unix::fs::PermissionsExt; + + /// A stand-in container CLI. `pull` fails with `pull_stderr` for the first + /// `pull_failures` calls and succeeds after; `image inspect` succeeds only + /// when `cached`. Every invocation is appended to `calls.log`. + struct FakeCli { + dir: tempfile::TempDir, + } + + impl FakeCli { + fn new(pull_failures: u32, pull_stderr: &str, cached: bool) -> Self { + let dir = tempfile::tempdir().unwrap(); + let script = format!( + r#"#!/bin/sh +d="{dir}" +echo "$*" >> "$d/calls.log" +case "$1" in + pull) + n=$(cat "$d/pulls" 2>/dev/null || echo 0) + n=$((n + 1)) + echo "$n" > "$d/pulls" + if [ "$n" -le {pull_failures} ]; then + echo '{pull_stderr}' >&2 + exit 1 + fi + exit 0 ;; + image) + [ "{cached}" = "true" ] && exit 0 + echo 'Error: No such image' >&2 + exit 1 ;; +esac +exit 2 +"#, + dir = dir.path().display(), + ); + let path = dir.path().join("cli"); + std::fs::write(&path, script).unwrap(); + std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o755)).unwrap(); + Self { dir } + } + + fn cli(&self) -> String { + self.dir.path().join("cli").display().to_string() + } + + fn calls(&self) -> Vec { + std::fs::read_to_string(self.dir.path().join("calls.log")) + .unwrap_or_default() + .lines() + .map(String::from) + .collect() + } + + async fn pull(&self) -> Result { + pull_image_with(&self.cli(), None, "alpine:3.20", Duration::from_millis(1)).await + } + } + + const THROTTLED: &str = "Error response from daemon: unexpected status from HEAD request to https://public.ecr.aws/v2/docker/library/alpine/manifests/3.20: 429 Too Many Requests"; + + #[tokio::test] + async fn a_successful_pull_needs_no_cache_check() { + let cli = FakeCli::new(0, "", false); + assert_eq!(cli.pull().await, Ok(PulledImage::Pulled)); + assert_eq!(cli.calls(), ["pull alpine:3.20"]); + } + + #[tokio::test] + async fn a_failed_pull_uses_the_cached_image() { + let cli = FakeCli::new(u32::MAX, THROTTLED, true); + let got = cli.pull().await; + assert_eq!( + got, + Ok(PulledImage::Cached { + pull_error: THROTTLED.to_string() + }) + ); + assert_eq!( + cli.calls(), + ["pull alpine:3.20", "image inspect alpine:3.20"], + "a cached image is used at once, without retrying the pull" + ); + } + + #[tokio::test] + async fn a_rate_limited_pull_with_nothing_cached_is_retried() { + let cli = FakeCli::new(2, THROTTLED, false); + assert_eq!(cli.pull().await, Ok(PulledImage::Pulled)); + let pulls = cli.calls().iter().filter(|c| c.starts_with("pull")).count(); + assert_eq!(pulls, 3); + } + + #[tokio::test] + async fn retries_stop_after_the_attempt_cap() { + let cli = FakeCli::new(u32::MAX, THROTTLED, false); + assert_eq!(cli.pull().await, Err(THROTTLED.to_string())); + let pulls = cli.calls().iter().filter(|c| c.starts_with("pull")).count(); + assert_eq!(pulls, MAX_PULL_ATTEMPTS as usize); + } + + #[tokio::test] + async fn a_missing_image_fails_without_retrying() { + let missing = + "Error response from daemon: manifest for alpine:nope not found: manifest unknown"; + let cli = FakeCli::new(u32::MAX, missing, false); + assert_eq!(cli.pull().await, Err(missing.to_string())); + let pulls = cli.calls().iter().filter(|c| c.starts_with("pull")).count(); + assert_eq!(pulls, 1, "only a throttled pull is worth retrying"); + } + + #[test] + fn rate_limit_detection_covers_each_registry_form() { + assert!(is_rate_limited(THROTTLED)); + assert!(is_rate_limited( + "toomanyrequests: You have reached your pull rate limit." + )); + assert!(is_rate_limited("Error: Rate exceeded")); + assert!(!is_rate_limited("manifest unknown")); + assert!(!is_rate_limited("pull access denied for foo")); + } +} diff --git a/crates/fakecloud-core/src/lib.rs b/crates/fakecloud-core/src/lib.rs index 5a7302a03..d7d5d6558 100644 --- a/crates/fakecloud-core/src/lib.rs +++ b/crates/fakecloud-core/src/lib.rs @@ -1,6 +1,7 @@ pub mod auth; pub mod auth_message; pub mod cfn_template; +pub mod container_image; pub mod container_net; pub mod delivery; pub mod dispatch; diff --git a/crates/fakecloud-e2e/tests/batch_real_execution.rs b/crates/fakecloud-e2e/tests/batch_real_execution.rs index 2aebe25a9..2f3f07bc6 100644 --- a/crates/fakecloud-e2e/tests/batch_real_execution.rs +++ b/crates/fakecloud-e2e/tests/batch_real_execution.rs @@ -97,7 +97,12 @@ async fn run_job(batch: &aws_sdk_batch::Client, name: &str, command: Vec<&str>) job.job_id().unwrap().to_string() } -async fn wait_terminal(batch: &aws_sdk_batch::Client, job_id: &str) -> String { +/// Poll `job_id` to a terminal status and assert it is `want`. On a mismatch +/// the panic carries each job's status, `statusReason` and exit code -- for an +/// array parent, every child's too -- so a CI failure names its cause (an +/// image pull the registry refused, a container that exited non-zero) rather +/// than just the wrong status. +async fn expect_terminal(batch: &aws_sdk_batch::Client, job_id: &str, want: &str) { for _ in 0..120 { let d = batch.describe_jobs().jobs(job_id).send().await.unwrap(); let status = d.jobs()[0] @@ -105,11 +110,48 @@ async fn wait_terminal(batch: &aws_sdk_batch::Client, job_id: &str) -> String { .map(|s| s.as_str().to_string()) .unwrap_or_default(); if status == "SUCCEEDED" || status == "FAILED" { - return status; + if status != want { + panic!( + "job {job_id} ended {status}, expected {want}\n{}", + describe_for_diagnostics(batch, job_id).await + ); + } + return; } tokio::time::sleep(Duration::from_secs(1)).await; } - panic!("job {job_id} never reached a terminal status"); + panic!( + "job {job_id} never reached a terminal status\n{}", + describe_for_diagnostics(batch, job_id).await + ); +} + +async fn describe_for_diagnostics(batch: &aws_sdk_batch::Client, job_id: &str) -> String { + let mut ids = vec![job_id.to_string()]; + let d = batch.describe_jobs().jobs(job_id).send().await.unwrap(); + if let Some(size) = d.jobs()[0].array_properties().and_then(|a| a.size()) { + ids.extend((0..size).map(|i| format!("{job_id}:{i}"))); + } + let d = batch + .describe_jobs() + .set_jobs(Some(ids)) + .send() + .await + .unwrap(); + d.jobs() + .iter() + .map(|j| { + format!( + " {} status={:?} statusReason={:?} exitCode={:?} containerReason={:?}", + j.job_id().unwrap_or_default(), + j.status(), + j.status_reason(), + j.container().and_then(|c| c.exit_code()), + j.container().and_then(|c| c.reason()), + ) + }) + .collect::>() + .join("\n") } #[tokio::test] @@ -120,7 +162,7 @@ async fn submit_job_runs_real_container_and_succeeds() { let s = TestServer::start().await; let batch = aws_sdk_batch::Client::new(&s.aws_config().await); let job_id = run_job(&batch, "ok", vec!["sh", "-c", "exit 0"]).await; - assert_eq!(wait_terminal(&batch, &job_id).await, "SUCCEEDED"); + expect_terminal(&batch, &job_id, "SUCCEEDED").await; let d = batch.describe_jobs().jobs(&job_id).send().await.unwrap(); assert_eq!(d.jobs()[0].container().and_then(|c| c.exit_code()), Some(0)); } @@ -133,7 +175,7 @@ async fn submit_job_failing_container_fails_the_job() { let s = TestServer::start().await; let batch = aws_sdk_batch::Client::new(&s.aws_config().await); let job_id = run_job(&batch, "bad", vec!["sh", "-c", "exit 7"]).await; - assert_eq!(wait_terminal(&batch, &job_id).await, "FAILED"); + expect_terminal(&batch, &job_id, "FAILED").await; let d = batch.describe_jobs().jobs(&job_id).send().await.unwrap(); assert_eq!(d.jobs()[0].container().and_then(|c| c.exit_code()), Some(7)); } @@ -173,8 +215,8 @@ async fn depends_on_job_waits_for_its_dependency() { "B must wait for A, was {b_early}" ); - assert_eq!(wait_terminal(&batch, &a).await, "SUCCEEDED"); - assert_eq!(wait_terminal(&batch, &b).await, "SUCCEEDED"); + expect_terminal(&batch, &a, "SUCCEEDED").await; + expect_terminal(&batch, &b, "SUCCEEDED").await; } #[tokio::test] @@ -204,7 +246,7 @@ async fn array_job_runs_every_child_and_parent_succeeds() { .unwrap() .to_string(); - assert_eq!(wait_terminal(&batch, &parent).await, "SUCCEEDED"); + expect_terminal(&batch, &parent, "SUCCEEDED").await; let d = batch.describe_jobs().jobs(&parent).send().await.unwrap(); let summary = d.jobs()[0] .array_properties() @@ -238,7 +280,7 @@ async fn retry_strategy_reattempts_a_failing_job() { .unwrap() .to_string(); - assert_eq!(wait_terminal(&batch, &job_id).await, "FAILED"); + expect_terminal(&batch, &job_id, "FAILED").await; // Two attempts were made: one recorded retry + the final. let d = batch.describe_jobs().jobs(&job_id).send().await.unwrap(); assert_eq!(d.jobs()[0].attempts().len(), 1); @@ -270,7 +312,7 @@ async fn timeout_fails_an_overrunning_job() { .unwrap() .to_string(); - assert_eq!(wait_terminal(&batch, &job_id).await, "FAILED"); + expect_terminal(&batch, &job_id, "FAILED").await; let d = batch.describe_jobs().jobs(&job_id).send().await.unwrap(); assert!(d.jobs()[0] .status_reason() diff --git a/crates/fakecloud-ecs/src/runtime/task_lifecycle.rs b/crates/fakecloud-ecs/src/runtime/task_lifecycle.rs index 1d1498598..f8043d06a 100644 --- a/crates/fakecloud-ecs/src/runtime/task_lifecycle.rs +++ b/crates/fakecloud-ecs/src/runtime/task_lifecycle.rs @@ -112,16 +112,13 @@ impl EcsRuntime { self.server_port, ); let pull_uri = local_pull_uri.as_deref().unwrap_or(&rp.plan.image); - let pull_out = self - .cli_command() - .args(["pull", pull_uri]) - .output() - .await - .map_err(|e| RuntimeError::ImagePull(e.to_string()))?; - if !pull_out.status.success() { - let err = String::from_utf8_lossy(&pull_out.stderr).to_string(); - return Err(RuntimeError::ImagePull(err)); - } + fakecloud_core::container_image::pull_image( + &self.cli, + self.docker_config_path().as_deref(), + pull_uri, + ) + .await + .map_err(RuntimeError::ImagePull)?; // Retag the local pull URI to the AWS URI so `docker run` finds // the image under the user-facing name. Digest-pinned refs // can't be `docker tag` targets, so we fall through and run diff --git a/crates/fakecloud-lambda/src/runtime/docker.rs b/crates/fakecloud-lambda/src/runtime/docker.rs index 7df322232..a5d33906a 100644 --- a/crates/fakecloud-lambda/src/runtime/docker.rs +++ b/crates/fakecloud-lambda/src/runtime/docker.rs @@ -124,21 +124,13 @@ impl DockerBackend { ); let pull_uri = local_pull_uri.as_deref().unwrap_or(image); - let mut pull_cmd = tokio::process::Command::new(&self.cli); - if let Some(p) = self.docker_config_path() { - pull_cmd.env("DOCKER_CONFIG", p); - } - let pull_out = pull_cmd - .args(["pull", pull_uri]) - .output() - .await - .map_err(|e| RuntimeError::ContainerStartFailed(format!("docker pull: {e}")))?; - if !pull_out.status.success() { - return Err(RuntimeError::ContainerStartFailed(format!( - "docker pull failed: {}", - String::from_utf8_lossy(&pull_out.stderr) - ))); - } + fakecloud_core::container_image::pull_image( + &self.cli, + self.docker_config_path().as_deref(), + pull_uri, + ) + .await + .map_err(|e| RuntimeError::ContainerStartFailed(format!("docker pull failed: {e}")))?; // Retag the local pull URI to the AWS URI so `docker create` // finds the image under the user-visible name. Digest-pinned // refs can't be `docker tag` targets, so fall through and @@ -503,21 +495,15 @@ impl LambdaBackend for DockerBackend { ); let pull_uri = local_uri.as_deref().unwrap_or(image); - let mut cmd = tokio::process::Command::new(&self.cli); - if let Some(p) = self.docker_config_path() { - cmd.env("DOCKER_CONFIG", p); - } - let out = cmd - .args(["pull", pull_uri]) - .output() - .await - .map_err(|e| RuntimeError::ContainerStartFailed(format!("docker pull: {e}")))?; - if !out.status.success() { - return Err(RuntimeError::ContainerStartFailed(format!( - "docker pull failed for {pull_uri}: {}", - String::from_utf8_lossy(&out.stderr) - ))); - } + fakecloud_core::container_image::pull_image( + &self.cli, + self.docker_config_path().as_deref(), + pull_uri, + ) + .await + .map_err(|e| { + RuntimeError::ContainerStartFailed(format!("docker pull failed for {pull_uri}: {e}")) + })?; Ok(()) } } diff --git a/website/content/docs/services/ecs.md b/website/content/docs/services/ecs.md index 96d61f331..96f3c4551 100644 --- a/website/content/docs/services/ecs.md +++ b/website/content/docs/services/ecs.md @@ -54,7 +54,7 @@ Task-definition families track revisions monotonically; `DeleteTaskDefinitions` `RunTask` records the task synchronously and kicks off a background docker execution per spawned task: -1. `docker pull ` (timestamps captured on the task: `pullStartedAt` / `pullStoppedAt`) +1. `docker pull ` (timestamps captured on the task: `pullStartedAt` / `pullStoppedAt`). As with the ECS agent's default `ECS_IMAGE_PULL_BEHAVIOR`, a failed pull falls back to the image already cached locally, and a pull the registry rate limits (`429 Too Many Requests`, common for anonymous `public.ecr.aws` pulls from a shared IP) is retried with backoff. A pull that fails with nothing cached stops the task with `TaskFailedToStart`. 2. `docker run -d ` (container ID recorded on the task's container) 3. `docker wait ` (blocks on container exit; exit code → `containers[].exitCode`) 4. `docker logs ` (captured stdout/stderr stored on the task + exposed via the introspection endpoint) diff --git a/website/content/docs/services/lambda.md b/website/content/docs/services/lambda.md index eae66965f..db8a90e94 100644 --- a/website/content/docs/services/lambda.md +++ b/website/content/docs/services/lambda.md @@ -79,7 +79,7 @@ Lambda is a target for most event-producing services: ## Gotchas - **Requires a Docker socket.** Lambda needs access to `/var/run/docker.sock` to start and stop containers. Only use in environments you trust — Docker socket access is effectively host-level privilege. -- **First invocation of a runtime pulls the image.** Expect a slower first run while the Lambda runtime image downloads. Subsequent invocations are fast. +- **First invocation of a runtime pulls the image.** Expect a slower first run while the Lambda runtime image downloads. Subsequent invocations are fast. For `PackageType=Image` functions, a failed pull of the function's image falls back to a copy already cached locally, and a pull the registry rate limits is retried with backoff. - **Cold vs. warm containers.** fakecloud reuses containers between invocations for the same function. Force a cold start via `/_fakecloud/lambda/{name}/evict-container`. ## Source From 2055db64b8811e2621831c7b2e747d55038015de Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 11:41:51 -0300 Subject: [PATCH 2/6] fix(core): only a transient pull failure falls back to the cached image A refused pull -- an image deleted from the registry, or one a repository policy denies -- must fail the launch even when an earlier launch left a copy in the local cache. Neither Fargate nor Lambda has a per-host image cache, so on AWS such a pull never succeeds. Restrict both the cache fallback and the retry to transient failures: throttling, registry 5xx, and network timeouts or dropped connections. --- crates/fakecloud-core/src/container_image.rs | 138 +++++++++++++------ website/content/docs/services/ecs.md | 2 +- website/content/docs/services/lambda.md | 2 +- 3 files changed, 101 insertions(+), 41 deletions(-) diff --git a/crates/fakecloud-core/src/container_image.rs b/crates/fakecloud-core/src/container_image.rs index cde5b6cfc..6144acb4d 100644 --- a/crates/fakecloud-core/src/container_image.rs +++ b/crates/fakecloud-core/src/container_image.rs @@ -2,24 +2,28 @@ //! Batch through it; Lambda `PackageType=Image` functions). //! //! A bare `docker pull` always contacts the registry, even when the image is -//! already in the local cache, so any registry hiccup fails the launch. The -//! common one is rate limiting: anonymous pulls from `public.ecr.aws` are -//! capped per source IP, and a burst of task launches -- or several processes -//! sharing one NAT address -- gets `429 Too Many Requests` back. +//! already in the local cache, so a momentary registry failure fails the +//! launch. The common one is rate limiting: anonymous pulls from +//! `public.ecr.aws` are capped per source IP, and a burst of task launches -- +//! or several processes sharing one NAT address -- gets +//! `429 Too Many Requests` back. //! -//! The ECS container agent handles this the way this module does. Under its -//! default `ECS_IMAGE_PULL_BEHAVIOR`, a failed pull falls back to the image -//! cached on the instance, and pulls are retried with backoff before giving -//! up. A pull that fails with nothing cached still fails, with the registry's -//! own error. +//! A transient failure (throttling, a registry 5xx, a network timeout) is +//! retried with backoff, and falls back to the image already cached locally +//! instead of failing the launch. A refusal is final: an image that no longer +//! exists or a pull the registry denies fails at once with the registry's own +//! error, even when a stale copy is cached. Otherwise an image deleted from +//! ECR, or one a repository policy denies, would keep launching from the +//! local copy -- which neither Fargate nor Lambda, having no per-host image +//! cache, ever does. use std::path::Path; use std::time::Duration; use tokio::process::Command; -/// Pull attempts made while the registry keeps rate limiting, before the -/// last error is returned. +/// Pull attempts made while the registry keeps failing transiently, before +/// the last error is returned. const MAX_PULL_ATTEMPTS: u32 = 5; /// Delay before the first retry; doubles on each further retry. @@ -30,18 +34,18 @@ const BASE_RETRY_DELAY: Duration = Duration::from_secs(1); pub enum PulledImage { /// The registry served the image. Pulled, - /// The pull failed but the image was already cached locally, so the - /// cached copy is used. Carries the pull's error for logging. + /// The pull failed transiently but the image was already cached locally, + /// so the cached copy is used. Carries the pull's error for logging. Cached { pull_error: String }, } -/// Pull `reference` with the container `cli`, falling back to a locally -/// cached copy when the pull fails and retrying while the registry is rate -/// limiting. `docker_config` is exported as `DOCKER_CONFIG` for the pull so -/// registry credentials resolve. +/// Pull `reference` with the container `cli`. A transient registry failure +/// falls back to a locally cached copy, or is retried with backoff when +/// nothing is cached; any other failure is returned at once. `docker_config` +/// is exported as `DOCKER_CONFIG` for the pull so registry credentials +/// resolve. /// -/// Returns the pull's stderr as the error when the image is neither pullable -/// nor cached. +/// Returns the pull's stderr as the error. pub async fn pull_image( cli: &str, docker_config: Option<&Path>, @@ -72,23 +76,25 @@ async fn pull_image_with( return Ok(PulledImage::Pulled); } let pull_error = String::from_utf8_lossy(&out.stderr).trim().to_string(); - + if !is_transient(&pull_error) { + return Err(pull_error); + } if image_cached(cli, reference).await { tracing::warn!( image = %reference, error = %pull_error, - "image pull failed; using the locally cached image" + "image pull failed transiently; using the locally cached image" ); return Ok(PulledImage::Cached { pull_error }); } - if attempt >= MAX_PULL_ATTEMPTS || !is_rate_limited(&pull_error) { + if attempt >= MAX_PULL_ATTEMPTS { return Err(pull_error); } tracing::info!( image = %reference, attempt, retry_in_ms = delay.as_millis() as u64, - "image pull rate limited; retrying" + "image pull failed transiently; retrying" ); tokio::time::sleep(delay).await; delay *= 2; @@ -108,14 +114,28 @@ async fn image_cached(cli: &str, reference: &str) -> bool { .unwrap_or(false) } -/// Whether a pull error is the registry throttling the caller. Docker reports -/// the status line (`429 Too Many Requests`) or the registry error code -/// (`toomanyrequests`); ECR Public's throttle message is `Rate exceeded`. -fn is_rate_limited(stderr: &str) -> bool { +/// Whether a pull error is one a later attempt could succeed past: the +/// registry throttling the caller (the `429 Too Many Requests` status line, +/// the `toomanyrequests` error code, ECR Public's `Rate exceeded`), a +/// registry-side 5xx, or the network timing out or dropping the connection. +/// Not-found and access-denied responses are not transient. +fn is_transient(stderr: &str) -> bool { + const MARKERS: [&str; 12] = [ + "toomanyrequests", + "too many requests", + "rate exceeded", + "500 internal server error", + "502 bad gateway", + "503 service unavailable", + "504 gateway timeout", + "i/o timeout", + "tls handshake timeout", + "connection reset by peer", + "context deadline exceeded", + "request canceled while waiting for connection", + ]; let lower = stderr.to_ascii_lowercase(); - lower.contains("toomanyrequests") - || lower.contains("too many requests") - || lower.contains("rate exceeded") + MARKERS.iter().any(|m| lower.contains(m)) } #[cfg(all(test, unix))] @@ -189,7 +209,7 @@ exit 2 } #[tokio::test] - async fn a_failed_pull_uses_the_cached_image() { + async fn a_throttled_pull_uses_the_cached_image() { let cli = FakeCli::new(u32::MAX, THROTTLED, true); let got = cli.pull().await; assert_eq!( @@ -227,18 +247,58 @@ exit 2 "Error response from daemon: manifest for alpine:nope not found: manifest unknown"; let cli = FakeCli::new(u32::MAX, missing, false); assert_eq!(cli.pull().await, Err(missing.to_string())); - let pulls = cli.calls().iter().filter(|c| c.starts_with("pull")).count(); - assert_eq!(pulls, 1, "only a throttled pull is worth retrying"); + assert_eq!( + cli.calls(), + ["pull alpine:3.20"], + "a refused pull is neither retried nor checked against the cache" + ); + } + + #[tokio::test] + async fn a_refused_pull_fails_even_with_a_stale_cached_copy() { + // The image was deleted from the registry, or a policy now denies the + // pull. A copy cached by an earlier launch must not be used. + for refused in [ + "Error response from daemon: manifest for alpine:3.20 not found: manifest unknown", + "Error response from daemon: pull access denied for alpine, repository does not exist or may require authorization: denied", + ] { + let cli = FakeCli::new(u32::MAX, refused, true); + assert_eq!(cli.pull().await, Err(refused.to_string())); + assert_eq!(cli.calls(), ["pull alpine:3.20"]); + } + } + + #[tokio::test] + async fn a_registry_server_error_uses_the_cached_image() { + let unavailable = + "Error response from daemon: received unexpected HTTP status: 503 Service Unavailable"; + let cli = FakeCli::new(u32::MAX, unavailable, true); + assert_eq!( + cli.pull().await, + Ok(PulledImage::Cached { + pull_error: unavailable.to_string() + }) + ); } #[test] - fn rate_limit_detection_covers_each_registry_form() { - assert!(is_rate_limited(THROTTLED)); - assert!(is_rate_limited( + fn transient_detection_separates_retryable_from_refused() { + assert!(is_transient(THROTTLED)); + assert!(is_transient( "toomanyrequests: You have reached your pull rate limit." )); - assert!(is_rate_limited("Error: Rate exceeded")); - assert!(!is_rate_limited("manifest unknown")); - assert!(!is_rate_limited("pull access denied for foo")); + assert!(is_transient("Error: Rate exceeded")); + assert!(is_transient( + "received unexpected HTTP status: 502 Bad Gateway" + )); + assert!(is_transient( + "Get \"https://public.ecr.aws/v2/\": net/http: TLS handshake timeout" + )); + assert!(is_transient( + "read tcp 10.0.0.2:4431->1.2.3.4:443: read: connection reset by peer" + )); + assert!(!is_transient("manifest unknown")); + assert!(!is_transient("pull access denied for foo")); + assert!(!is_transient("unauthorized: authentication required")); } } diff --git a/website/content/docs/services/ecs.md b/website/content/docs/services/ecs.md index 96f3c4551..19480bfb3 100644 --- a/website/content/docs/services/ecs.md +++ b/website/content/docs/services/ecs.md @@ -54,7 +54,7 @@ Task-definition families track revisions monotonically; `DeleteTaskDefinitions` `RunTask` records the task synchronously and kicks off a background docker execution per spawned task: -1. `docker pull ` (timestamps captured on the task: `pullStartedAt` / `pullStoppedAt`). As with the ECS agent's default `ECS_IMAGE_PULL_BEHAVIOR`, a failed pull falls back to the image already cached locally, and a pull the registry rate limits (`429 Too Many Requests`, common for anonymous `public.ecr.aws` pulls from a shared IP) is retried with backoff. A pull that fails with nothing cached stops the task with `TaskFailedToStart`. +1. `docker pull ` (timestamps captured on the task: `pullStartedAt` / `pullStoppedAt`). A transient registry failure (rate limiting such as `429 Too Many Requests`, common for anonymous `public.ecr.aws` pulls from a shared IP; a registry 5xx; a network timeout) is retried with backoff and falls back to a copy of the image already cached locally. A refused pull (the image or tag no longer exists, or access is denied) stops the task with `TaskFailedToStart` even when a stale copy is cached, as does a transient failure with nothing cached. 2. `docker run -d ` (container ID recorded on the task's container) 3. `docker wait ` (blocks on container exit; exit code → `containers[].exitCode`) 4. `docker logs ` (captured stdout/stderr stored on the task + exposed via the introspection endpoint) diff --git a/website/content/docs/services/lambda.md b/website/content/docs/services/lambda.md index db8a90e94..24aa0277f 100644 --- a/website/content/docs/services/lambda.md +++ b/website/content/docs/services/lambda.md @@ -79,7 +79,7 @@ Lambda is a target for most event-producing services: ## Gotchas - **Requires a Docker socket.** Lambda needs access to `/var/run/docker.sock` to start and stop containers. Only use in environments you trust — Docker socket access is effectively host-level privilege. -- **First invocation of a runtime pulls the image.** Expect a slower first run while the Lambda runtime image downloads. Subsequent invocations are fast. For `PackageType=Image` functions, a failed pull of the function's image falls back to a copy already cached locally, and a pull the registry rate limits is retried with backoff. +- **First invocation of a runtime pulls the image.** Expect a slower first run while the Lambda runtime image downloads. Subsequent invocations are fast. For `PackageType=Image` functions, a transient failure pulling the function's image (rate limiting, a registry 5xx, a network timeout) is retried with backoff and falls back to a copy already cached locally; a deleted or access-denied image fails the start. - **Cold vs. warm containers.** fakecloud reuses containers between invocations for the same function. Force a cold start via `/_fakecloud/lambda/{name}/evict-container`. ## Source From 60eda4dab59030ac70de9c0a22fe38658010cab4 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 11:53:09 -0300 Subject: [PATCH 3/6] fix(core): classify pull errors so an image name cannot fake a transient failure - Check refusal wording (not found, denied, unauthorized, forbidden) first, on the unmodified message, so a refused pull of a repository named like a throttling code (toomanyrequests) is not read as transient and does not launch a stale cached copy. A repository name cannot contain a space, so it can only spell the one-word markers. - Treat any 5xx status as transient, parsed from the status position, instead of four fixed reason phrases. - The Batch e2e diagnostics no longer index an empty DescribeJobs result. --- crates/fakecloud-core/src/container_image.rs | 86 +++++++++++++++++-- .../tests/batch_real_execution.rs | 5 +- 2 files changed, 81 insertions(+), 10 deletions(-) diff --git a/crates/fakecloud-core/src/container_image.rs b/crates/fakecloud-core/src/container_image.rs index 6144acb4d..6739ecf60 100644 --- a/crates/fakecloud-core/src/container_image.rs +++ b/crates/fakecloud-core/src/container_image.rs @@ -116,18 +116,29 @@ async fn image_cached(cli: &str, reference: &str) -> bool { /// Whether a pull error is one a later attempt could succeed past: the /// registry throttling the caller (the `429 Too Many Requests` status line, -/// the `toomanyrequests` error code, ECR Public's `Rate exceeded`), a -/// registry-side 5xx, or the network timing out or dropping the connection. -/// Not-found and access-denied responses are not transient. +/// the `toomanyrequests` error code, ECR Public's `Rate exceeded`), any 5xx +/// from the registry, or the network timing out or dropping the connection. +/// +/// A refusal (not found, access denied, unauthorized) is never transient, and +/// is checked first so it wins over transient wording in the same message. +/// That matters because the message quotes the image name, which the user +/// chooses: a repository named `toomanyrequests` whose pull is refused must +/// not be read as throttled and launch a stale cached copy. A repository +/// name cannot contain a space, so it can only spell the one-word markers, +/// and a refusal always carries refusal wording of its own. fn is_transient(stderr: &str) -> bool { - const MARKERS: [&str; 12] = [ + const REFUSED: [&str; 6] = [ + "manifest unknown", + "not found", + "denied", + "unauthorized", + "forbidden", + "does not exist", + ]; + const TRANSIENT: [&str; 8] = [ "toomanyrequests", "too many requests", "rate exceeded", - "500 internal server error", - "502 bad gateway", - "503 service unavailable", - "504 gateway timeout", "i/o timeout", "tls handshake timeout", "connection reset by peer", @@ -135,7 +146,30 @@ fn is_transient(stderr: &str) -> bool { "request canceled while waiting for connection", ]; let lower = stderr.to_ascii_lowercase(); - MARKERS.iter().any(|m| lower.contains(m)) + if REFUSED.iter().any(|m| lower.contains(m)) { + return false; + } + TRANSIENT.iter().any(|m| lower.contains(m)) || has_server_error_status(&lower) +} + +/// Whether the message carries a 5xx HTTP status. Docker and Podman quote the +/// status as `: 503 Service Unavailable`, `status: 500`, or +/// `status code 502`; a three-digit number in that position from 500 to 599 +/// counts, whatever reason phrase follows. +fn has_server_error_status(message: &str) -> bool { + let bytes = message.as_bytes(); + ["status code ", "status: ", "status ", ": "] + .iter() + .flat_map(|prefix| message.match_indices(prefix).map(|(i, p)| i + p.len())) + .any(|start| { + let code = &bytes[start..bytes.len().min(start + 3)]; + code.len() == 3 + && code[0] == b'5' + && code.iter().all(u8::is_ascii_digit) + && !bytes + .get(start + 3) + .is_some_and(|c| c.is_ascii_alphanumeric()) + }) } #[cfg(all(test, unix))] @@ -281,6 +315,16 @@ exit 2 ); } + #[tokio::test] + async fn a_refused_pull_of_a_repository_named_like_a_marker_is_still_refused() { + // The message quotes the image name. A repository spelled like a + // throttling code must not make a refused pull look transient. + let refused = "Error response from daemon: manifest for toomanyrequests:latest not found: manifest unknown: manifest unknown"; + let cli = FakeCli::new(u32::MAX, refused, true); + assert_eq!(cli.pull().await, Err(refused.to_string())); + assert_eq!(cli.calls(), ["pull alpine:3.20"]); + } + #[test] fn transient_detection_separates_retryable_from_refused() { assert!(is_transient(THROTTLED)); @@ -300,5 +344,29 @@ exit 2 assert!(!is_transient("manifest unknown")); assert!(!is_transient("pull access denied for foo")); assert!(!is_transient("unauthorized: authentication required")); + assert!(!is_transient( + "pull access denied for toomanyrequests, repository does not exist or may require authorization" + )); + } + + #[test] + fn any_5xx_status_is_transient_whatever_its_reason_phrase() { + for msg in [ + "received unexpected HTTP status: 500 Internal Server Error", + "unexpected status from GET request to https://r.example/v2/: 507 Insufficient Storage", + "unexpected status code 520", + "error pulling image: status: 599", + "unexpected status from HEAD request to https://r.example/v2/a/manifests/1: 503", + ] { + assert!(is_transient(msg), "{msg}"); + } + for msg in [ + // A registry port or a 4xx is not a server error. + "Get \"http://127.0.0.1:5000/v2/\": dial tcp 127.0.0.1:5000: connect: connection refused", + "unexpected status code 400 Bad Request", + "status: 5001", + ] { + assert!(!is_transient(msg), "{msg}"); + } } } diff --git a/crates/fakecloud-e2e/tests/batch_real_execution.rs b/crates/fakecloud-e2e/tests/batch_real_execution.rs index 2f3f07bc6..99df00bb7 100644 --- a/crates/fakecloud-e2e/tests/batch_real_execution.rs +++ b/crates/fakecloud-e2e/tests/batch_real_execution.rs @@ -129,7 +129,10 @@ async fn expect_terminal(batch: &aws_sdk_batch::Client, job_id: &str, want: &str async fn describe_for_diagnostics(batch: &aws_sdk_batch::Client, job_id: &str) -> String { let mut ids = vec![job_id.to_string()]; let d = batch.describe_jobs().jobs(job_id).send().await.unwrap(); - if let Some(size) = d.jobs()[0].array_properties().and_then(|a| a.size()) { + let Some(job) = d.jobs().first() else { + return format!(" {job_id} not returned by DescribeJobs"); + }; + if let Some(size) = job.array_properties().and_then(|a| a.size()) { ids.extend((0..size).map(|i| format!("{job_id}:{i}"))); } let d = batch From 476097d27d7b3e46c07e0ab64c02afafcd12c9f8 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 11:56:04 -0300 Subject: [PATCH 4/6] fix(core): classify a pull error without its image name, and check the cache on the same daemon - Remove the image's own name (the reference, and each trailing path of its repository) from the error before matching, whole names only. A throttled pull of a repository named like a refusal (acme/access-denied-page) is transient again, and a refused pull of one named like a throttle code stays refused. - Run the cache check with the pull's DOCKER_CONFIG. The config selects the Docker context, so without it the check could consult a different daemon than the one that pulls and runs the image. --- crates/fakecloud-core/src/container_image.rs | 159 +++++++++++++++---- 1 file changed, 131 insertions(+), 28 deletions(-) diff --git a/crates/fakecloud-core/src/container_image.rs b/crates/fakecloud-core/src/container_image.rs index 6739ecf60..d43006b94 100644 --- a/crates/fakecloud-core/src/container_image.rs +++ b/crates/fakecloud-core/src/container_image.rs @@ -76,10 +76,10 @@ async fn pull_image_with( return Ok(PulledImage::Pulled); } let pull_error = String::from_utf8_lossy(&out.stderr).trim().to_string(); - if !is_transient(&pull_error) { + if !is_transient(&pull_error, reference) { return Err(pull_error); } - if image_cached(cli, reference).await { + if image_cached(cli, docker_config, reference).await { tracing::warn!( image = %reference, error = %pull_error, @@ -102,10 +102,16 @@ async fn pull_image_with( } } -/// Whether `reference` resolves to an image in the local cache. -async fn image_cached(cli: &str, reference: &str) -> bool { - Command::new(cli) - .args(["image", "inspect", reference]) +/// Whether `reference` resolves to an image in the local cache. Runs with +/// the same `DOCKER_CONFIG` as the pull: the config also selects the Docker +/// context, so without it the check could consult a different daemon than +/// the one that pulls and later runs the image. +async fn image_cached(cli: &str, docker_config: Option<&Path>, reference: &str) -> bool { + let mut cmd = Command::new(cli); + if let Some(p) = docker_config { + cmd.env("DOCKER_CONFIG", p); + } + cmd.args(["image", "inspect", reference]) .stdout(std::process::Stdio::null()) .stderr(std::process::Stdio::null()) .status() @@ -120,13 +126,12 @@ async fn image_cached(cli: &str, reference: &str) -> bool { /// from the registry, or the network timing out or dropping the connection. /// /// A refusal (not found, access denied, unauthorized) is never transient, and -/// is checked first so it wins over transient wording in the same message. -/// That matters because the message quotes the image name, which the user -/// chooses: a repository named `toomanyrequests` whose pull is refused must -/// not be read as throttled and launch a stale cached copy. A repository -/// name cannot contain a space, so it can only spell the one-word markers, -/// and a refusal always carries refusal wording of its own. -fn is_transient(stderr: &str) -> bool { +/// wins over transient wording in the same message. Both are matched only +/// after the image's own name is removed from the message: the name is chosen +/// by the user and quoted in the error, so a repository spelled like either +/// kind of marker (`toomanyrequests/app`, `acme/access-denied-page`) must not +/// flip the classification. +fn is_transient(stderr: &str, reference: &str) -> bool { const REFUSED: [&str; 6] = [ "manifest unknown", "not found", @@ -145,11 +150,64 @@ fn is_transient(stderr: &str) -> bool { "context deadline exceeded", "request canceled while waiting for connection", ]; - let lower = stderr.to_ascii_lowercase(); - if REFUSED.iter().any(|m| lower.contains(m)) { + let message = without_image_name(&stderr.to_ascii_lowercase(), reference); + if REFUSED.iter().any(|m| message.contains(m)) { return false; } - TRANSIENT.iter().any(|m| lower.contains(m)) || has_server_error_status(&lower) + TRANSIENT.iter().any(|m| message.contains(m)) || has_server_error_status(&message) +} + +/// `message` (lowercase) with each way an error can quote the image blanked +/// out: the reference as given, and every trailing path of its repository +/// (`public.ecr.aws/acme/app`, `acme/app`, `app`) -- registries put the +/// repository path in URLs, and Podman expands short names to +/// `docker.io/library/`. Only whole names are removed, bounded by +/// characters a name cannot contain, so a short repository like `d` never +/// cuts letters out of the surrounding words. +fn without_image_name(message: &str, reference: &str) -> String { + let reference = reference.to_ascii_lowercase(); + let untagged = reference.split('@').next().unwrap_or(&reference); + // A `:` after the last `/` starts the tag; one before it is a registry port. + let repository = match (untagged.rfind(':'), untagged.rfind('/')) { + (Some(colon), Some(slash)) if colon < slash => untagged, + (Some(colon), _) => &untagged[..colon], + (None, _) => untagged, + }; + let mut names = vec![reference.as_str(), repository]; + names.extend( + repository + .match_indices('/') + .map(|(i, _)| &repository[i + 1..]), + ); + names.sort_by_key(|n| std::cmp::Reverse(n.len())); + + let mut out = message.to_string(); + for name in names.into_iter().filter(|n| !n.is_empty()) { + out = remove_whole(&out, name); + } + out +} + +/// `haystack` with every occurrence of `name` that is not part of a longer +/// name (flanked by a letter, digit, `.`, `_` or `-`) replaced by a space. +fn remove_whole(haystack: &str, name: &str) -> String { + let is_name_char = |c: char| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'); + let mut out = String::with_capacity(haystack.len()); + let mut rest = haystack; + while let Some(i) = rest.find(name) { + let end = i + name.len(); + let before = rest[..i].chars().next_back(); + let after = rest[end..].chars().next(); + if before.is_some_and(is_name_char) || after.is_some_and(is_name_char) { + out.push_str(&rest[..end]); + } else { + out.push_str(&rest[..i]); + out.push(' '); + } + rest = &rest[end..]; + } + out.push_str(rest); + out } /// Whether the message carries a 5xx HTTP status. Docker and Podman quote the @@ -233,6 +291,10 @@ exit 2 } } + fn is_transient_for_test(stderr: &str) -> bool { + is_transient(stderr, "alpine:3.20") + } + const THROTTLED: &str = "Error response from daemon: unexpected status from HEAD request to https://public.ecr.aws/v2/docker/library/alpine/manifests/3.20: 429 Too Many Requests"; #[tokio::test] @@ -327,24 +389,26 @@ exit 2 #[test] fn transient_detection_separates_retryable_from_refused() { - assert!(is_transient(THROTTLED)); - assert!(is_transient( + assert!(is_transient_for_test(THROTTLED)); + assert!(is_transient_for_test( "toomanyrequests: You have reached your pull rate limit." )); - assert!(is_transient("Error: Rate exceeded")); - assert!(is_transient( + assert!(is_transient_for_test("Error: Rate exceeded")); + assert!(is_transient_for_test( "received unexpected HTTP status: 502 Bad Gateway" )); - assert!(is_transient( + assert!(is_transient_for_test( "Get \"https://public.ecr.aws/v2/\": net/http: TLS handshake timeout" )); - assert!(is_transient( + assert!(is_transient_for_test( "read tcp 10.0.0.2:4431->1.2.3.4:443: read: connection reset by peer" )); - assert!(!is_transient("manifest unknown")); - assert!(!is_transient("pull access denied for foo")); - assert!(!is_transient("unauthorized: authentication required")); - assert!(!is_transient( + assert!(!is_transient_for_test("manifest unknown")); + assert!(!is_transient_for_test("pull access denied for foo")); + assert!(!is_transient_for_test( + "unauthorized: authentication required" + )); + assert!(!is_transient_for_test( "pull access denied for toomanyrequests, repository does not exist or may require authorization" )); } @@ -358,7 +422,7 @@ exit 2 "error pulling image: status: 599", "unexpected status from HEAD request to https://r.example/v2/a/manifests/1: 503", ] { - assert!(is_transient(msg), "{msg}"); + assert!(is_transient_for_test(msg), "{msg}"); } for msg in [ // A registry port or a 4xx is not a server error. @@ -366,7 +430,46 @@ exit 2 "unexpected status code 400 Bad Request", "status: 5001", ] { - assert!(!is_transient(msg), "{msg}"); + assert!(!is_transient_for_test(msg), "{msg}"); + } + } + + #[test] + fn a_throttled_pull_of_a_repository_named_like_a_refusal_is_still_transient() { + let reference = "public.ecr.aws/acme/access-denied-page:1"; + for msg in [ + "Error response from daemon: unexpected status from HEAD request to https://public.ecr.aws/v2/acme/access-denied-page/manifests/1: 429 Too Many Requests", + "Error response from daemon: toomanyrequests: Rate exceeded for public.ecr.aws/acme/access-denied-page:1", + ] { + assert!(is_transient(msg, reference), "{msg}"); } + // Its genuine refusals are still refusals. + assert!(!is_transient( + "Error response from daemon: manifest for public.ecr.aws/acme/access-denied-page:1 not found: manifest unknown", + reference + )); + } + + #[test] + fn only_whole_names_are_removed() { + // A one-letter repository must not cut the `d` out of `denied`. + assert!(!is_transient( + "Error response from daemon: pull access denied for d, repository does not exist", + "d" + )); + assert_eq!( + without_image_name( + "pull access denied for docker.io/library/alpine", + "alpine:3.20" + ), + "pull access denied for docker.io/library/ " + ); + assert_eq!( + without_image_name( + "get https://127.0.0.1:5000/v2/team/app/manifests/v1", + "127.0.0.1:5000/team/app:v1" + ), + "get https://127.0.0.1:5000/v2/ /manifests/v1" + ); } } From d1f53bcddc5364bbe16900d82e4fac8cd4c6c7f3 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 11:58:24 -0300 Subject: [PATCH 5/6] fix(core): also remove the registry host from a pull error before classifying it URLs in the error quote the registry host on its own, so a host named like a refusal (denied.example) would have hidden a 429 from it. --- crates/fakecloud-core/src/container_image.rs | 22 ++++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/crates/fakecloud-core/src/container_image.rs b/crates/fakecloud-core/src/container_image.rs index d43006b94..eefd807da 100644 --- a/crates/fakecloud-core/src/container_image.rs +++ b/crates/fakecloud-core/src/container_image.rs @@ -158,9 +158,9 @@ fn is_transient(stderr: &str, reference: &str) -> bool { } /// `message` (lowercase) with each way an error can quote the image blanked -/// out: the reference as given, and every trailing path of its repository -/// (`public.ecr.aws/acme/app`, `acme/app`, `app`) -- registries put the -/// repository path in URLs, and Podman expands short names to +/// out: the reference as given, every trailing path of its repository +/// (`public.ecr.aws/acme/app`, `acme/app`, `app`), and its registry host -- +/// registries put the host and repository path in URLs, and Podman expands short names to /// `docker.io/library/`. Only whole names are removed, bounded by /// characters a name cannot contain, so a short repository like `d` never /// cuts letters out of the surrounding words. @@ -179,6 +179,12 @@ fn without_image_name(message: &str, reference: &str) -> String { .match_indices('/') .map(|(i, _)| &repository[i + 1..]), ); + // The registry host, which URLs in the error quote on its own. + if let Some((host, _)) = repository.split_once('/') { + if host.contains(['.', ':']) || host == "localhost" { + names.push(host); + } + } names.sort_by_key(|n| std::cmp::Reverse(n.len())); let mut out = message.to_string(); @@ -450,6 +456,14 @@ exit 2 )); } + #[test] + fn a_registry_host_named_like_a_refusal_does_not_hide_a_throttle() { + assert!(is_transient( + "Error response from daemon: unexpected status from HEAD request to https://denied.example/v2/app/manifests/1: 429 Too Many Requests", + "denied.example/app:1" + )); + } + #[test] fn only_whole_names_are_removed() { // A one-letter repository must not cut the `d` out of `denied`. @@ -469,7 +483,7 @@ exit 2 "get https://127.0.0.1:5000/v2/team/app/manifests/v1", "127.0.0.1:5000/team/app:v1" ), - "get https://127.0.0.1:5000/v2/ /manifests/v1" + "get https:// /v2/ /manifests/v1" ); } } From 147a0b12b5f6460b46a980d8d6643e23471ae8d6 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 12:06:00 -0300 Subject: [PATCH 6/6] test(core): drive the image-name tests through pull_image with the named reference The fake CLI pull always used alpine:3.20, so the marker-named repository test never exercised name removal. The fake now takes the reference, and a pull-level test covers a throttled pull of acme/access-denied-page falling back to the cache, which fails if the name is not removed first. --- crates/fakecloud-core/src/container_image.rs | 36 ++++++++++++++++++-- 1 file changed, 33 insertions(+), 3 deletions(-) diff --git a/crates/fakecloud-core/src/container_image.rs b/crates/fakecloud-core/src/container_image.rs index eefd807da..20c00fc3f 100644 --- a/crates/fakecloud-core/src/container_image.rs +++ b/crates/fakecloud-core/src/container_image.rs @@ -293,7 +293,11 @@ exit 2 } async fn pull(&self) -> Result { - pull_image_with(&self.cli(), None, "alpine:3.20", Duration::from_millis(1)).await + self.pull_ref("alpine:3.20").await + } + + async fn pull_ref(&self, reference: &str) -> Result { + pull_image_with(&self.cli(), None, reference, Duration::from_millis(1)).await } } @@ -389,8 +393,34 @@ exit 2 // throttling code must not make a refused pull look transient. let refused = "Error response from daemon: manifest for toomanyrequests:latest not found: manifest unknown: manifest unknown"; let cli = FakeCli::new(u32::MAX, refused, true); - assert_eq!(cli.pull().await, Err(refused.to_string())); - assert_eq!(cli.calls(), ["pull alpine:3.20"]); + assert_eq!( + cli.pull_ref("toomanyrequests:latest").await, + Err(refused.to_string()) + ); + assert_eq!(cli.calls(), ["pull toomanyrequests:latest"]); + } + + #[tokio::test] + async fn a_throttled_pull_of_a_repository_named_like_a_refusal_uses_the_cache() { + // Only removing the image name from the message keeps `denied` in the + // repository from reading as a refusal; without it this pull would + // fail instead of falling back to the cached copy. + let reference = "public.ecr.aws/acme/access-denied-page:1"; + let throttled = "Error response from daemon: unexpected status from HEAD request to https://public.ecr.aws/v2/acme/access-denied-page/manifests/1: 429 Too Many Requests"; + let cli = FakeCli::new(u32::MAX, throttled, true); + assert_eq!( + cli.pull_ref(reference).await, + Ok(PulledImage::Cached { + pull_error: throttled.to_string() + }) + ); + assert_eq!( + cli.calls(), + [ + format!("pull {reference}"), + format!("image inspect {reference}") + ] + ); } #[test]