Skip to content

fix(image-normalizer): fail a dispatch whose stored image is gone as a terminal 422 - #1854

Open
sejori wants to merge 1 commit into
mainfrom
fix/image-unavailable-terminal
Open

sejori wants to merge 1 commit into
mainfrom
fix/image-unavailable-terminal

Conversation

@sejori

@sejori sejori commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

Problem

When the object behind a request's dw-img:// token is no longer in the store (the lifecycle rule removed it, see #1843 for how that happens), every dispatch fails identically: the engine cannot fetch the signed URL. That failure reached the batch daemon as a retriable 503, so the request retried to the attempt cap, hundreds of times over many hours, and every adaptive-concurrency limiter in the path counted the 503s as overload. On 2026-09-25 that drove the model's batch concurrency limit to 1 for seven hours while replicas sat idle.

Why message matching alone cannot fix it

Engines phrase a media failure differently and change the phrasing between versions: Python requests (404 Client Error: Not Found for url: …), aiohttp (404, message='Not Found', url='…'), current vLLM (Failed to fetch media from URL: HTTP 404 error, no URL), older vLLM (An exception occurred while loading IMAGE data at index 0: …), SGLang (Could not decode image: …). And Dynamo, our main gateway, sends the client only internal server error during processing with code: 500, keeping the diagnostic server-side. So onwards can never see the cause for Dynamo.

Change

dwctl, authoritative and engine-independent. After a dispatch returns an upstream failure (500/502/503/504) for a body that carried tokens, the image middleware probes the store for each token with a plain HEAD. If one is definitively missing, the response becomes a terminal 422 image_unavailable telling the caller to resubmit. A store error or a present object leaves the upstream failure untouched; capacity signals (429/529) are never reinterpreted. This adds ImageStore::is_present, a policy-free probe: exists may report a present-but-ageing object as absent so that ingest refreshes it (#1843), which would be a false positive here.

onwards, fast path. An embedded server error whose text says the engine could not fetch or decode a media input becomes a terminal 422 upstream_media_fetch_failed with no failover. The classifier (onwards::media_failure) keys on three independent signals rather than any engine's wording: a media word, a fetch/decode verb, and a standalone 4xx status bounded by non-alphanumerics so hashes and content-addressed paths never match. 408/429 from the media host and timeouts stay non-terminal. No regex dependency.

Fusillade treats 4xx as non-retriable, so both outcomes land in the batch error file and never feed the AIMD limiters.

Counters: dwctl_image_unavailable_total, onwards_upstream_media_failures_total{kind,upstream_status}.

Tests

  • dwctl (#[sqlx::test], real grants path): opaque 503 + object removed between signing and fetch → 422 image_unavailable; 500/502/503 with the object present → passed through unchanged; 529 with the object gone → untouched.
  • onwards: 10 classifier cases across engine phrasings plus negatives (capacity, template, context length, OOM, model-not-found); integration test that each phrasing, unary and streaming, yields one attempt and a 422 with providers that would otherwise fail over on 500, and that an opaque 500 still fails over to a 503.
  • cargo clippy -p onwards -p dwctl --all-features --no-deps clean; ZDR payload-logging guard clean.

🤖 Generated with Claude Code

Review in cubic

…a terminal 422

When the object behind a request's `dw-img://` token is no longer in the
store, every dispatch fails the same way: the engine cannot fetch the
signed URL. That failure reached the batch daemon as a retriable 503, so
the request retried to the attempt cap (hundreds of times over many
hours), and every adaptive-concurrency limiter in the path read the
stream of 503s as overload. On 2026-09-25 that drove the model's batch
concurrency to 1 for seven hours while capacity sat idle.

Two detectors, because inference engines do not agree on how they report
a media failure and some gateways do not report it at all:

- dwctl (authoritative, engine-independent): after a dispatch returns an
  upstream failure (500/502/503/504) for a body that carried tokens, the
  image middleware probes the store for each token with a plain HEAD. If
  one is definitively missing, the response becomes a 422
  `image_unavailable` telling the caller to resubmit. A store error or
  a present object leaves the upstream failure untouched; capacity
  signals (429/529) are never reinterpreted. `ImageStore::is_present`
  is the policy-free probe: `exists` may report a present-but-ageing
  object as absent so ingest refreshes it, which would be wrong here.

- onwards (fast path): an embedded server error whose text says the
  engine could not fetch or decode a media input becomes a terminal 422
  `upstream_media_fetch_failed` with no failover. The classifier keys on
  independent signals (a media word, a fetch/decode verb, a standalone
  4xx) rather than any engine's phrasing, and is tested against Python
  requests, aiohttp, current and older vLLM, SGLang and Dynamo wording.
  408/429 from the media host and timeouts stay non-terminal. Dynamo
  sends the client only "internal server error during processing", so
  it never reaches this path; the dwctl probe covers it.

Fusillade treats 4xx as non-retriable, so both land in the batch error
file. Counters: `dwctl_image_unavailable_total`,
`onwards_upstream_media_failures_total{kind,upstream_status}`.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Copilot AI lite review requested due to automatic review settings September 25, 2026 10:49
@cloudflare-workers-and-pages

Copy link
Copy Markdown

Deploying control-layer with  Cloudflare Pages  Cloudflare Pages

Latest commit: 72fea3a
Status: ✅  Deploy successful!
Preview URL: https://53fc786a.control-layer.pages.dev
Branch Preview URL: https://fix-image-unavailable-termin.control-layer.pages.dev

View logs

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Unresolved compatibility, detection, metadata, and classifier issues remain.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 1 High severity · 3 Medium severity · 1 Low severity

Open (5)
What changed in this PR

Adds terminal handling for unavailable stored images and upstream media-fetch failures, preventing futile retries and overload signals.

Changes:

  • Adds store-presence probing and 422 image_unavailable responses.
  • Adds engine-independent media failure classification.
  • Adds integration and classifier tests.
File Summary
onwards/​src/​media_failure.rs Media failure classifier and tests
onwards/​src/​lib.rs Integration coverage
onwards/​src/​handlers.rs Terminal media failure handling
onwards/​src/​errors.rs Media error response
dwctl/​src/​inference/​image_normalizer_middleware.rs Missing-image detection and tests
dwctl/​src/​image_normalizer/​store.rs Store presence probes
dwctl/​src/​image_normalizer/​mod.rs Normalizer presence API

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

/// present-but-ageing object as absent so ingest refreshes it; a caller
/// deciding whether a dispatch failed because the object is GONE needs
/// the plain answer.
async fn is_present(&self, token: ImageToken) -> Result<bool, StoreError>;
async fn missing_image_tokens(normalizer: &Arc<dyn ImageNormalizer>, tokens: &[ImageToken]) -> Vec<ImageToken> {
let mut seen = std::collections::HashSet::new();
let mut missing = Vec::new();
for &token in tokens.iter().filter(|t| seen.insert(**t)).take(MAX_PRESENCE_CHECKS) {
Comment on lines +75 to +78
const MEDIA_WORDS: &[&str] = &["image", "video", "audio", "media", "multimodal", "mm data"];
const FETCH_WORDS: &[&str] = &[
"fetch", "download", "load", "retriev", "for url", "url=", "url:", "url '",
];
if !MEDIA_WORDS.iter().any(|w| text.contains(w)) {
return None;
}
if DECODE_PHRASES.iter().any(|p| text.contains(p)) {
pub(crate) fn image_unavailable_response() -> Response {
let body = serde_json::json!({
"error": {
"message": "One or more image inputs of this request are no longer available in the image store: their stored copies have expired. Resubmit the request; the images will be re-uploaded automatically.",

@cubic-dev-ai cubic-dev-ai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

10 issues found across 7 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. If appropriate, use sub-agents to investigate and fix each issue separately.


<file name="onwards/src/errors.rs">

<violation number="1" location="onwards/src/errors.rs:216">
P2: This response labels every media failure as `image_url`, even though the classifier and message cover video and audio inputs. Pass the actual media parameter through the classification path, or omit `param` when the input field cannot be determined.</violation>
</file>

<file name="dwctl/src/inference/image_normalizer_middleware.rs">

<violation number="1" location="dwctl/src/inference/image_normalizer_middleware.rs:347">
P2: The terminal replacement drops the upstream response extensions, so failed dispatches lose `ServedBy` and serving metadata before request logging runs. Copy the original response extensions onto the 422 response before returning it.</violation>

<violation number="2" location="dwctl/src/inference/image_normalizer_middleware.rs:370">
P1: `missing_image_tokens` silently ignores every distinct token after the first 32, so a later missing image still produces a retriable 5xx instead of the promised terminal 422. Process all referenced tokens (with bounded concurrency if needed), or explicitly treat an incomplete probe as unknown rather than success.</violation>

<violation number="3" location="dwctl/src/inference/image_normalizer_middleware.rs:385">
P2: This promises automatic re-upload even when the request contains only a `dw-img` token whose stored bytes are gone. Tell callers to resubmit the original image inputs.</violation>
</file>

<file name="onwards/src/media_failure.rs">

<violation number="1" location="onwards/src/media_failure.rs:77">
P2: The substring fetch signal matches `payload` through `"load"`, so unrelated image validation errors can be converted into terminal 422 responses. Match `load` as a bounded verb or use explicit phrases instead of unrestricted substring matching.</violation>

<violation number="2" location="onwards/src/media_failure.rs:100">
P1: `classify_media_failure` makes every decode phrase terminal before honoring transient statuses and timeouts. A message such as `Could not decode image: HTTP 429 Too Many Requests` therefore stops failover with 422 instead of preserving the media-host rate limit; check 408/429 and timeout signals before the decode branch.</violation>

<violation number="3" location="onwards/src/media_failure.rs:110">
P2: `standalone_4xx` interprets arbitrary URL/query numbers as media-host statuses. A valid URL such as `...?width=404` can turn a statusless fetch error into terminal 422; extract statuses only from HTTP/error-status context or exclude URL spans before scanning.</violation>
</file>

<file name="onwards/src/handlers.rs">

<violation number="1" location="onwards/src/handlers.rs:1753">
P2: Direct HTTP 500/502/503/504 responses bypass `classify_media_failure`; only embedded errors in a 2xx response become terminal 422. Classify the non-2xx error body before status-based failover as well.</violation>

<violation number="2" location="onwards/src/handlers.rs:1772">
P2: This early terminal return skips the loop’s `serving_outcome` attachment, so analytics loses the resolved serving class for media-fetch failures. Set `serving_outcome` on this error before returning.</violation>
</file>

<file name="dwctl/src/image_normalizer/store.rs">

<violation number="1" location="dwctl/src/image_normalizer/store.rs:330">
P2: `GcsStore::is_present` performs the existing media `GET` instead of a metadata/HEAD probe. Failed dispatches can therefore trigger up to 32 object reads while checking image presence, amplifying store load during the very failures this path is meant to contain; implement a metadata-only presence request for GCS.</violation>
</file>

Reply with feedback, questions, or to request a fix.

Re-trigger cubic

async fn missing_image_tokens(normalizer: &Arc<dyn ImageNormalizer>, tokens: &[ImageToken]) -> Vec<ImageToken> {
let mut seen = std::collections::HashSet::new();
let mut missing = Vec::new();
for &token in tokens.iter().filter(|t| seen.insert(**t)).take(MAX_PRESENCE_CHECKS) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: missing_image_tokens silently ignores every distinct token after the first 32, so a later missing image still produces a retriable 5xx instead of the promised terminal 422. Process all referenced tokens (with bounded concurrency if needed), or explicitly treat an incomplete probe as unknown rather than success.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At dwctl/src/inference/image_normalizer_middleware.rs, line 370:

<comment>`missing_image_tokens` silently ignores every distinct token after the first 32, so a later missing image still produces a retriable 5xx instead of the promised terminal 422. Process all referenced tokens (with bounded concurrency if needed), or explicitly treat an incomplete probe as unknown rather than success.</comment>

<file context>
@@ -317,7 +318,77 @@ pub async fn image_normalizer_middleware(
+async fn missing_image_tokens(normalizer: &Arc<dyn ImageNormalizer>, tokens: &[ImageToken]) -> Vec<ImageToken> {
+    let mut seen = std::collections::HashSet::new();
+    let mut missing = Vec::new();
+    for &token in tokens.iter().filter(|t| seen.insert(**t)).take(MAX_PRESENCE_CHECKS) {
+        match normalizer.is_present(token).await {
+            Ok(true) => {}
</file context>

if !MEDIA_WORDS.iter().any(|w| text.contains(w)) {
return None;
}
if DECODE_PHRASES.iter().any(|p| text.contains(p)) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1: classify_media_failure makes every decode phrase terminal before honoring transient statuses and timeouts. A message such as Could not decode image: HTTP 429 Too Many Requests therefore stops failover with 422 instead of preserving the media-host rate limit; check 408/429 and timeout signals before the decode branch.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At onwards/src/media_failure.rs, line 100:

<comment>`classify_media_failure` makes every decode phrase terminal before honoring transient statuses and timeouts. A message such as `Could not decode image: HTTP 429 Too Many Requests` therefore stops failover with 422 instead of preserving the media-host rate limit; check 408/429 and timeout signals before the decode branch.</comment>

<file context>
@@ -0,0 +1,266 @@
+    if !MEDIA_WORDS.iter().any(|w| text.contains(w)) {
+        return None;
+    }
+    if DECODE_PHRASES.iter().any(|p| text.contains(p)) {
+        return Some(MediaFailure {
+            kind: MediaFailureKind::Decode,
</file context>

Comment thread onwards/src/errors.rs
body: Some(ErrorResponseBody {
message: message.to_string(),
r#type: "invalid_request_error".to_string(),
param: Some("image_url".to_string()),

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This response labels every media failure as image_url, even though the classifier and message cover video and audio inputs. Pass the actual media parameter through the classification path, or omit param when the input field cannot be determined.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At onwards/src/errors.rs, line 216:

<comment>This response labels every media failure as `image_url`, even though the classifier and message cover video and audio inputs. Pass the actual media parameter through the classification path, or omit `param` when the input field cannot be determined.</comment>

<file context>
@@ -197,6 +197,31 @@ impl OnwardsErrorResponse {
+            body: Some(ErrorResponseBody {
+                message: message.to_string(),
+                r#type: "invalid_request_error".to_string(),
+                param: Some("image_url".to_string()),
+                code: "upstream_media_fetch_failed".to_string(),
+            }),
</file context>

if !fetch_context {
return None;
}
let status = standalone_4xx(&text);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: standalone_4xx interprets arbitrary URL/query numbers as media-host statuses. A valid URL such as ...?width=404 can turn a statusless fetch error into terminal 422; extract statuses only from HTTP/error-status context or exclude URL spans before scanning.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At onwards/src/media_failure.rs, line 110:

<comment>`standalone_4xx` interprets arbitrary URL/query numbers as media-host statuses. A valid URL such as `...?width=404` can turn a statusless fetch error into terminal 422; extract statuses only from HTTP/error-status context or exclude URL spans before scanning.</comment>

<file context>
@@ -0,0 +1,266 @@
+    if !fetch_context {
+        return None;
+    }
+    let status = standalone_4xx(&text);
+    // A request timeout or rate limit from the media host is transient by
+    // definition: not ours to classify, whatever else the message says.
</file context>

Comment thread onwards/src/handlers.rs
"Upstream could not fetch or decode a media input of this request; failing it as 422 instead of retrying"
);
record_response_status(422);
return LoopAction::Done(Err(OnwardsErrorResponse::upstream_media_failure(failure.kind)));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This early terminal return skips the loop’s serving_outcome attachment, so analytics loses the resolved serving class for media-fetch failures. Set serving_outcome on this error before returning.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At onwards/src/handlers.rs, line 1772:

<comment>This early terminal return skips the loop’s `serving_outcome` attachment, so analytics loses the resolved serving class for media-fetch failures. Set `serving_outcome` on this error before returning.</comment>

<file context>
@@ -1741,6 +1741,37 @@ pub async fn target_message_handler<T: HttpClient>(
+                        "Upstream could not fetch or decode a media input of this request; failing it as 422 instead of retrying"
+                    );
+                    record_response_status(422);
+                    return LoopAction::Done(Err(OnwardsErrorResponse::upstream_media_failure(failure.kind)));
+                }
+
</file context>
Suggested change
return LoopAction::Done(Err(OnwardsErrorResponse::upstream_media_failure(failure.kind)));
let mut error = OnwardsErrorResponse::upstream_media_failure(failure.kind);
error.serving_outcome = Some(serving_resolution.outcome());
return LoopAction::Done(Err(error));


const MEDIA_WORDS: &[&str] = &["image", "video", "audio", "media", "multimodal", "mm data"];
const FETCH_WORDS: &[&str] = &[
"fetch", "download", "load", "retriev", "for url", "url=", "url:", "url '",

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: The substring fetch signal matches payload through "load", so unrelated image validation errors can be converted into terminal 422 responses. Match load as a bounded verb or use explicit phrases instead of unrestricted substring matching.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At onwards/src/media_failure.rs, line 77:

<comment>The substring fetch signal matches `payload` through `"load"`, so unrelated image validation errors can be converted into terminal 422 responses. Match `load` as a bounded verb or use explicit phrases instead of unrestricted substring matching.</comment>

<file context>
@@ -0,0 +1,266 @@
+
+const MEDIA_WORDS: &[&str] = &["image", "video", "audio", "media", "multimodal", "mm data"];
+const FETCH_WORDS: &[&str] = &[
+    "fetch", "download", "load", "retriev", "for url", "url=", "url:", "url '",
+];
+const DECODE_PHRASES: &[&str] = &[
</file context>


async fn is_present(&self, token: ImageToken) -> Result<bool, StoreError> {
// Same probe as `exists`; this backend applies no reuse policy.
self.exists(token).await

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: GcsStore::is_present performs the existing media GET instead of a metadata/HEAD probe. Failed dispatches can therefore trigger up to 32 object reads while checking image presence, amplifying store load during the very failures this path is meant to contain; implement a metadata-only presence request for GCS.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At dwctl/src/image_normalizer/store.rs, line 330:

<comment>`GcsStore::is_present` performs the existing media `GET` instead of a metadata/HEAD probe. Failed dispatches can therefore trigger up to 32 object reads while checking image presence, amplifying store load during the very failures this path is meant to contain; implement a metadata-only presence request for GCS.</comment>

<file context>
@@ -307,6 +325,11 @@ impl ImageStore for GcsStore {
 
+    async fn is_present(&self, token: ImageToken) -> Result<bool, StoreError> {
+        // Same probe as `exists`; this backend applies no reuse policy.
+        self.exists(token).await
+    }
+
</file context>

"upstream failed and stored image(s) this request references are gone; failing it as image_unavailable"
);
metrics::counter!("dwctl_image_unavailable_total").increment(1);
image_unavailable_response()

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: The terminal replacement drops the upstream response extensions, so failed dispatches lose ServedBy and serving metadata before request logging runs. Copy the original response extensions onto the 422 response before returning it.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At dwctl/src/inference/image_normalizer_middleware.rs, line 347:

<comment>The terminal replacement drops the upstream response extensions, so failed dispatches lose `ServedBy` and serving metadata before request logging runs. Copy the original response extensions onto the 422 response before returning it.</comment>

<file context>
@@ -317,7 +318,77 @@ pub async fn image_normalizer_middleware(
+        "upstream failed and stored image(s) this request references are gone; failing it as image_unavailable"
+    );
+    metrics::counter!("dwctl_image_unavailable_total").increment(1);
+    image_unavailable_response()
+}
+
</file context>
Suggested change
image_unavailable_response()
let mut unavailable = image_unavailable_response();
*unavailable.extensions_mut() = response.extensions().clone();
unavailable

Comment thread onwards/src/handlers.rs
// Engines that keep the diagnostic private (a bare "internal
// server error") do not reach here; the caller's own check on
// its stored media objects covers those.
if embedded >= 500

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: Direct HTTP 500/502/503/504 responses bypass classify_media_failure; only embedded errors in a 2xx response become terminal 422. Classify the non-2xx error body before status-based failover as well.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At onwards/src/handlers.rs, line 1753:

<comment>Direct HTTP 500/502/503/504 responses bypass `classify_media_failure`; only embedded errors in a 2xx response become terminal 422. Classify the non-2xx error body before status-based failover as well.</comment>

<file context>
@@ -1741,6 +1741,37 @@ pub async fn target_message_handler<T: HttpClient>(
+                // Engines that keep the diagnostic private (a bare "internal
+                // server error") do not reach here; the caller's own check on
+                // its stored media objects covers those.
+                if embedded >= 500
+                    && let Some(failure) = provider_error["message"]
+                        .as_str()
</file context>

pub(crate) fn image_unavailable_response() -> Response {
let body = serde_json::json!({
"error": {
"message": "One or more image inputs of this request are no longer available in the image store: their stored copies have expired. Resubmit the request; the images will be re-uploaded automatically.",

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This promises automatic re-upload even when the request contains only a dw-img token whose stored bytes are gone. Tell callers to resubmit the original image inputs.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At dwctl/src/inference/image_normalizer_middleware.rs, line 385:

<comment>This promises automatic re-upload even when the request contains only a `dw-img` token whose stored bytes are gone. Tell callers to resubmit the original image inputs.</comment>

<file context>
@@ -317,7 +318,77 @@ pub async fn image_normalizer_middleware(
+pub(crate) fn image_unavailable_response() -> Response {
+    let body = serde_json::json!({
+        "error": {
+            "message": "One or more image inputs of this request are no longer available in the image store: their stored copies have expired. Resubmit the request; the images will be re-uploaded automatically.",
+            "type": "invalid_request_error",
+            "param": "image_url",
</file context>
Suggested change
"message": "One or more image inputs of this request are no longer available in the image store: their stored copies have expired. Resubmit the request; the images will be re-uploaded automatically.",
"message": "One or more image inputs of this request are no longer available in the image store. Resubmit the request with the original image inputs.",

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants