diff --git a/deploy/helm/sie-cluster/files/alerts/sie-rules.yaml b/deploy/helm/sie-cluster/files/alerts/sie-rules.yaml index 8e90ef170..5571c884f 100644 --- a/deploy/helm/sie-cluster/files/alerts/sie-rules.yaml +++ b/deploy/helm/sie-cluster/files/alerts/sie-rules.yaml @@ -77,6 +77,24 @@ groups: summary: "SIE Gateway is down" description: "Gateway pod {{ $labels.pod }} is unreachable. Traffic cannot be routed." + - alert: SIERemoteFallbackPersistent + expr: | + max by (model) ( + sie_gateway_remote_serving_duration_seconds{namespace="__NAMESPACE__",service="__COLLECTOR_SERVICE__",endpoint="prometheus",producer_service="sie-gateway"} + and on (producer_instance, collector_generation, model) + ( + sum by (producer_instance, collector_generation, model) ( + increase(sie_gateway_remote_fallbacks_total{namespace="__NAMESPACE__",service="__COLLECTOR_SERVICE__",endpoint="prometheus",producer_service="sie-gateway",outcome="committed"}[5m]) + ) > 0 + ) + ) > __REMOTE_FALLBACK_PERSISTENCE_SECONDS__ + for: 1m + labels: + severity: warning + annotations: + summary: "SIE model {{ $labels.model }} continues to serve remotely" + description: "Committed fallback bridges for {{ $labels.model }} span more than __REMOTE_FALLBACK_PERSISTENCE_SECONDS__ seconds without observed local success. Check local worker readiness, model loading and admission pressure." + - alert: SIEHighErrorRate expr: | sum(rate(sie_gateway_requests_total{namespace="__NAMESPACE__",service="__COLLECTOR_SERVICE__",endpoint="prometheus",producer_service="sie-gateway",http_status_code=~"5.."}[5m])) diff --git a/deploy/helm/sie-cluster/files/dashboards/queue-routing.json b/deploy/helm/sie-cluster/files/dashboards/queue-routing.json index 494a877a1..84592530a 100644 --- a/deploy/helm/sie-cluster/files/dashboards/queue-routing.json +++ b/deploy/helm/sie-cluster/files/dashboards/queue-routing.json @@ -1738,6 +1738,84 @@ ], "title": "Gateway Queue Worker-Pool Backpressure and Shed", "type": "timeseries" + }, + { + "description": "Committed HTTP/first-output bridge responses and pre-output refusals. Later stream errors remain in generation telemetry.", + "fieldConfig": { + "defaults": { + "unit": "reqps" + } + }, + "gridPos": { + "x": 0, + "y": 125, + "w": 12, + "h": 8 + }, + "id": 119, + "options": { + "legend": { + "calcs": [ + "mean", + "max" + ], + "displayMode": "table", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "mode": "multi", + "sort": "desc" + } + }, + "targets": [ + { + "expr": "sum by (model, operation, fallback_reason, outcome) (rate(sie_gateway_remote_fallbacks_total{namespace=\"$namespace\", service=\"$collector\", endpoint=\"prometheus\", producer_service=\"sie-gateway\"}[$__rate_interval]))", + "legendFormat": "{{model}} {{operation}} {{fallback_reason}} {{outcome}}", + "refId": "A" + } + ], + "title": "Remote Fallback Responses", + "type": "timeseries" + }, + { + "description": "Observed seconds between first and latest committed bridge per gateway replica. Local success resets it. The persistence alert additionally requires recent committed activity from the same replica.", + "fieldConfig": { + "defaults": { + "unit": "s" + } + }, + "gridPos": { + "x": 12, + "y": 125, + "w": 12, + "h": 8 + }, + "id": 120, + "options": { + "legend": { + "calcs": [ + "mean", + "max" + ], + "displayMode": "table", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "mode": "multi", + "sort": "desc" + } + }, + "targets": [ + { + "expr": "max by (model) (sie_gateway_remote_serving_duration_seconds{namespace=\"$namespace\", service=\"$collector\", endpoint=\"prometheus\", producer_service=\"sie-gateway\"})", + "legendFormat": "{{model}}", + "refId": "A" + } + ], + "title": "Remote Serving Duration Since Local Success", + "type": "timeseries" } ], "refresh": "10s", diff --git a/deploy/helm/sie-cluster/templates/_otel-collector-config.tpl b/deploy/helm/sie-cluster/templates/_otel-collector-config.tpl index c45e40890..33e706cf3 100644 --- a/deploy/helm/sie-cluster/templates/_otel-collector-config.tpl +++ b/deploy/helm/sie-cluster/templates/_otel-collector-config.tpl @@ -89,7 +89,7 @@ processors: error_mode: propagate metrics: metric: - - 'not IsMatch(name, "^(sie[.]gateway[.]requests|sie[.]gateway[.]request[.]duration|sie[.]gateway[.]admission[.]decisions|sie[.]gateway[.]dispatches|sie[.]gateway[.]dispatch[.]duration|sie[.]gateway[.]config[.]applied_epoch|sie[.]gateway[.]config[.]operations|sie[.]gateway[.]config[.]bootstrap[.]degraded|sie[.]gateway[.]messaging[.]client[.]ready|sie[.]gateway[.]queue[.]publishes|sie[.]gateway[.]queue[.]publish[.]duration|sie[.]gateway[.]queue[.]publish[.]items|sie[.]gateway[.]queue[.]result_waits|sie[.]gateway[.]queue[.]result_wait[.]duration|sie[.]gateway[.]queue[.]result_chunks[.]received|sie[.]gateway[.]queue[.]result_chunk[.]bytes_received|sie[.]gateway[.]queue[.]result_chunk[.]rejections|sie[.]gateway[.]queue[.]result_chunk[.]transfers_completed|sie[.]gateway[.]queue[.]result_chunk[.]duplicates|sie[.]gateway[.]queue[.]result_chunk[.]retry_replacements|sie[.]gateway[.]queue[.]result_chunk[.]stale_retries|sie[.]gateway[.]queue[.]result_chunk[.]reserved_bytes|sie[.]gateway[.]queue[.]events|sie[.]gateway[.]queue[.]lane_admission[.]decisions|sie[.]gateway[.]queue[.]worker_pool[.]events|sie[.]gateway[.]provisioning[.]responses|sie[.]gateway[.]routing[.]unsupported_model_exclusions|sie[.]gateway[.]generation[.]events|sie[.]gateway[.]generation[.]ttft|sie[.]gateway[.]generation[.]tpot|sie[.]gateway[.]generation[.]tokens|sie[.]gateway[.]pool[.]pinned_model[.]loaded|sie[.]gateway[.]pending_demand|sie[.]gateway[.]lane[.]queue[.]depth|sie[.]gateway[.]lane[.]queue[.]snapshot[.]timestamp|sie[.]gateway[.]active_lease[.]gpus|sie[.]gateway[.]pool[.]warm_floor|sie[.]gateway[.]rejected[.]requests|sie[.]gateway[.]capacity[.]snapshot[.]timestamp|sie[.]gateway[.]key_snapshot[.]polls|sie[.]gateway[.]settlement[.]confirms|sie[.]gateway[.]settlement[.]confirm[.]duration)$")' + - 'not IsMatch(name, "^(sie[.]gateway[.]requests|sie[.]gateway[.]request[.]duration|sie[.]gateway[.]admission[.]decisions|sie[.]gateway[.]dispatches|sie[.]gateway[.]remote[.]fallbacks|sie[.]gateway[.]dispatch[.]duration|sie[.]gateway[.]remote[.]serving[.]duration|sie[.]gateway[.]config[.]applied_epoch|sie[.]gateway[.]config[.]operations|sie[.]gateway[.]config[.]bootstrap[.]degraded|sie[.]gateway[.]messaging[.]client[.]ready|sie[.]gateway[.]queue[.]publishes|sie[.]gateway[.]queue[.]publish[.]duration|sie[.]gateway[.]queue[.]publish[.]items|sie[.]gateway[.]queue[.]result_waits|sie[.]gateway[.]queue[.]result_wait[.]duration|sie[.]gateway[.]queue[.]result_chunks[.]received|sie[.]gateway[.]queue[.]result_chunk[.]bytes_received|sie[.]gateway[.]queue[.]result_chunk[.]rejections|sie[.]gateway[.]queue[.]result_chunk[.]transfers_completed|sie[.]gateway[.]queue[.]result_chunk[.]duplicates|sie[.]gateway[.]queue[.]result_chunk[.]retry_replacements|sie[.]gateway[.]queue[.]result_chunk[.]stale_retries|sie[.]gateway[.]queue[.]result_chunk[.]reserved_bytes|sie[.]gateway[.]queue[.]events|sie[.]gateway[.]queue[.]lane_admission[.]decisions|sie[.]gateway[.]queue[.]worker_pool[.]events|sie[.]gateway[.]provisioning[.]responses|sie[.]gateway[.]routing[.]unsupported_model_exclusions|sie[.]gateway[.]generation[.]events|sie[.]gateway[.]generation[.]ttft|sie[.]gateway[.]generation[.]tpot|sie[.]gateway[.]generation[.]tokens|sie[.]gateway[.]pool[.]pinned_model[.]loaded|sie[.]gateway[.]pending_demand|sie[.]gateway[.]lane[.]queue[.]depth|sie[.]gateway[.]lane[.]queue[.]snapshot[.]timestamp|sie[.]gateway[.]active_lease[.]gpus|sie[.]gateway[.]pool[.]warm_floor|sie[.]gateway[.]rejected[.]requests|sie[.]gateway[.]capacity[.]snapshot[.]timestamp|sie[.]gateway[.]key_snapshot[.]polls|sie[.]gateway[.]settlement[.]confirms|sie[.]gateway[.]settlement[.]confirm[.]duration)$")' - 'resource.attributes["service.name"] != "sie-gateway"' filter/prometheus_application_contract: error_mode: propagate @@ -120,8 +120,8 @@ processors: - context: metric statements: - set(description, "") - - 'set(unit, "{request}") where IsMatch(name, "^(sie[.]gateway[.]requests|sie[.]gateway[.]admission[.]decisions|sie[.]gateway[.]dispatches|sie[.]config[.]requests|sie[.]worker[.]ipc[.]requests|sie[.]worker[.]ipc[.]capacity|sie[.]worker[.]ipc[.]inflight|sie[.]worker[.]generation[.]inflight|sie[.]worker[.]generation[.]duplicate_prevented|sie[.]worker[.]generation[.]grammar[.]requests|sie[.]worker[.]upstream[.]refusals|sie[.]gateway[.]pending_demand|sie[.]gateway[.]rejected[.]requests)$")' - - 'set(unit, "s") where IsMatch(name, "^(sie[.]gateway[.]request[.]duration|sie[.]gateway[.]dispatch[.]duration|sie[.]gateway[.]queue[.]publish[.]duration|sie[.]gateway[.]queue[.]result_wait[.]duration|sie[.]gateway[.]generation[.]ttft|sie[.]gateway[.]generation[.]tpot|sie[.]dispatcher[.]invocation[.]duration|sie[.]dispatcher[.]sealed[.]stage[.]duration|sie[.]dispatcher[.]sealed[.]sandbox[.]seconds|sie[.]dispatcher[.]sealed[.]cold_start[.]seconds|sie[.]config[.]request[.]duration|sie[.]worker[.]queue[.]duration|sie[.]worker[.]scheduler[.]request_batch[.]dispatch_wait|sie[.]worker[.]scheduler[.]request_batch[.]total|sie[.]worker[.]scheduler[.]adaptive[.]wait|sie[.]worker[.]scheduler[.]adaptive[.]p50|sie[.]worker[.]ipc[.]request[.]duration|sie[.]worker[.]payload[.]fetch[.]duration|sie[.]worker[.]ipc[.]acquire[.]duration|sie[.]worker[.]shutdown[.]drain[.]duration|sie[.]worker[.]request[.]duration|sie[.]worker[.]inference[.]duration|sie[.]worker[.]model[.]load[.]duration|sie[.]worker[.]generation[.]worker_wait|sie[.]worker[.]generation[.]ttft|sie[.]worker[.]generation[.]tpot|sie[.]worker[.]generation[.]grammar[.]compile[.]duration|sie[.]worker[.]runtime[.]forward[.]duration|sie[.]worker[.]runtime[.]forward[.]permit[.]wait|sie[.]gateway[.]lane[.]queue[.]snapshot[.]timestamp|sie[.]gateway[.]capacity[.]snapshot[.]timestamp|sie[.]gateway[.]settlement[.]confirm[.]duration|sie[.]worker[.]work_item[.]age)$")' + - 'set(unit, "{request}") where IsMatch(name, "^(sie[.]gateway[.]requests|sie[.]gateway[.]admission[.]decisions|sie[.]gateway[.]dispatches|sie[.]gateway[.]remote[.]fallbacks|sie[.]config[.]requests|sie[.]worker[.]ipc[.]requests|sie[.]worker[.]ipc[.]capacity|sie[.]worker[.]ipc[.]inflight|sie[.]worker[.]generation[.]inflight|sie[.]worker[.]generation[.]duplicate_prevented|sie[.]worker[.]generation[.]grammar[.]requests|sie[.]worker[.]upstream[.]refusals|sie[.]gateway[.]pending_demand|sie[.]gateway[.]rejected[.]requests)$")' + - 'set(unit, "s") where IsMatch(name, "^(sie[.]gateway[.]request[.]duration|sie[.]gateway[.]dispatch[.]duration|sie[.]gateway[.]remote[.]serving[.]duration|sie[.]gateway[.]queue[.]publish[.]duration|sie[.]gateway[.]queue[.]result_wait[.]duration|sie[.]gateway[.]generation[.]ttft|sie[.]gateway[.]generation[.]tpot|sie[.]dispatcher[.]invocation[.]duration|sie[.]dispatcher[.]sealed[.]stage[.]duration|sie[.]dispatcher[.]sealed[.]sandbox[.]seconds|sie[.]dispatcher[.]sealed[.]cold_start[.]seconds|sie[.]config[.]request[.]duration|sie[.]worker[.]queue[.]duration|sie[.]worker[.]scheduler[.]request_batch[.]dispatch_wait|sie[.]worker[.]scheduler[.]request_batch[.]total|sie[.]worker[.]scheduler[.]adaptive[.]wait|sie[.]worker[.]scheduler[.]adaptive[.]p50|sie[.]worker[.]ipc[.]request[.]duration|sie[.]worker[.]payload[.]fetch[.]duration|sie[.]worker[.]ipc[.]acquire[.]duration|sie[.]worker[.]shutdown[.]drain[.]duration|sie[.]worker[.]request[.]duration|sie[.]worker[.]inference[.]duration|sie[.]worker[.]model[.]load[.]duration|sie[.]worker[.]generation[.]worker_wait|sie[.]worker[.]generation[.]ttft|sie[.]worker[.]generation[.]tpot|sie[.]worker[.]generation[.]grammar[.]compile[.]duration|sie[.]worker[.]runtime[.]forward[.]duration|sie[.]worker[.]runtime[.]forward[.]permit[.]wait|sie[.]gateway[.]lane[.]queue[.]snapshot[.]timestamp|sie[.]gateway[.]capacity[.]snapshot[.]timestamp|sie[.]gateway[.]settlement[.]confirm[.]duration|sie[.]worker[.]work_item[.]age)$")' - 'set(unit, "{epoch}") where IsMatch(name, "^(sie[.]gateway[.]config[.]applied_epoch|sie[.]config[.]epoch|sie[.]worker[.]config[.]epoch)$")' - 'set(unit, "{operation}") where IsMatch(name, "^(sie[.]gateway[.]config[.]operations|sie[.]config[.]publish|sie[.]config[.]store[.]writes|sie[.]worker[.]nats[.]operations)$")' - 'set(unit, "1") where IsMatch(name, "^(sie[.]gateway[.]config[.]bootstrap[.]degraded|sie[.]gateway[.]messaging[.]client[.]ready|sie[.]config[.]messaging[.]ready|sie[.]worker[.]batch[.]fill_ratio|sie[.]worker[.]config[.]degraded|sie[.]worker[.]saturated|sie[.]gateway[.]pool[.]pinned_model[.]loaded)$")' @@ -168,6 +168,8 @@ processors: - 'keep_keys(attributes, []) where IsMatch(metric.name, "^(sie[.]gateway[.]queue[.]result_chunks[.]received|sie[.]gateway[.]queue[.]result_chunk[.]bytes_received|sie[.]gateway[.]queue[.]result_chunk[.]transfers_completed|sie[.]gateway[.]queue[.]result_chunk[.]duplicates|sie[.]gateway[.]queue[.]result_chunk[.]retry_replacements|sie[.]gateway[.]queue[.]result_chunk[.]stale_retries|sie[.]gateway[.]queue[.]result_chunk[.]reserved_bytes)$")' - 'keep_keys(attributes, ["reason"]) where metric.name == "sie.gateway.queue.result_chunk.rejections"' - 'keep_keys(attributes, ["operation", "dispatch.path", "outcome", "fallback.reason", "lane"]) where IsMatch(metric.name, "^(sie[.]gateway[.]dispatches|sie[.]gateway[.]dispatch[.]duration)$")' + - 'keep_keys(attributes, ["operation", "model", "fallback.reason", "outcome"]) where metric.name == "sie.gateway.remote.fallbacks"' + - 'keep_keys(attributes, ["model"]) where metric.name == "sie.gateway.remote.serving.duration"' - 'keep_keys(attributes, []) where IsMatch(metric.name, "^(sie[.]gateway[.]config[.]applied_epoch|sie[.]gateway[.]config[.]bootstrap[.]degraded|sie[.]config[.]epoch|sie[.]gateway[.]capacity[.]snapshot[.]timestamp)$")' - 'keep_keys(attributes, ["transport"]) where IsMatch(metric.name, "^(sie[.]gateway[.]messaging[.]client[.]ready|sie[.]config[.]messaging[.]ready)$")' - 'keep_keys(attributes, ["event", "outcome"]) where metric.name == "sie.gateway.queue.events"' @@ -287,7 +289,7 @@ processors: error_mode: propagate metrics: metric: - - 'not IsMatch(name, "^(sie[.]gateway[.]requests|sie[.]gateway[.]request[.]duration|sie[.]gateway[.]admission[.]decisions|sie[.]gateway[.]dispatches|sie[.]gateway[.]dispatch[.]duration|sie[.]gateway[.]config[.]applied_epoch|sie[.]gateway[.]config[.]operations|sie[.]gateway[.]config[.]bootstrap[.]degraded|sie[.]gateway[.]messaging[.]client[.]ready|sie[.]gateway[.]queue[.]publishes|sie[.]gateway[.]queue[.]publish[.]duration|sie[.]gateway[.]queue[.]publish[.]items|sie[.]gateway[.]queue[.]result_waits|sie[.]gateway[.]queue[.]result_wait[.]duration|sie[.]gateway[.]queue[.]result_chunks[.]received|sie[.]gateway[.]queue[.]result_chunk[.]bytes_received|sie[.]gateway[.]queue[.]result_chunk[.]rejections|sie[.]gateway[.]queue[.]result_chunk[.]transfers_completed|sie[.]gateway[.]queue[.]result_chunk[.]duplicates|sie[.]gateway[.]queue[.]result_chunk[.]retry_replacements|sie[.]gateway[.]queue[.]result_chunk[.]stale_retries|sie[.]gateway[.]queue[.]result_chunk[.]reserved_bytes|sie[.]gateway[.]queue[.]events|sie[.]gateway[.]queue[.]lane_admission[.]decisions|sie[.]gateway[.]provisioning[.]responses|sie[.]gateway[.]routing[.]unsupported_model_exclusions|sie[.]gateway[.]generation[.]events|sie[.]gateway[.]generation[.]ttft|sie[.]gateway[.]generation[.]tpot|sie[.]gateway[.]generation[.]tokens|sie[.]gateway[.]pending_demand|sie[.]gateway[.]lane[.]queue[.]depth|sie[.]gateway[.]lane[.]queue[.]snapshot[.]timestamp|sie[.]gateway[.]active_lease[.]gpus|sie[.]gateway[.]pool[.]warm_floor|sie[.]gateway[.]rejected[.]requests|sie[.]gateway[.]key_snapshot[.]polls|sie[.]gateway[.]settlement[.]confirms|sie[.]gateway[.]settlement[.]confirm[.]duration)$")' + - 'not IsMatch(name, "^(sie[.]gateway[.]requests|sie[.]gateway[.]request[.]duration|sie[.]gateway[.]admission[.]decisions|sie[.]gateway[.]dispatches|sie[.]gateway[.]remote[.]fallbacks|sie[.]gateway[.]dispatch[.]duration|sie[.]gateway[.]remote[.]serving[.]duration|sie[.]gateway[.]config[.]applied_epoch|sie[.]gateway[.]config[.]operations|sie[.]gateway[.]config[.]bootstrap[.]degraded|sie[.]gateway[.]messaging[.]client[.]ready|sie[.]gateway[.]queue[.]publishes|sie[.]gateway[.]queue[.]publish[.]duration|sie[.]gateway[.]queue[.]publish[.]items|sie[.]gateway[.]queue[.]result_waits|sie[.]gateway[.]queue[.]result_wait[.]duration|sie[.]gateway[.]queue[.]result_chunks[.]received|sie[.]gateway[.]queue[.]result_chunk[.]bytes_received|sie[.]gateway[.]queue[.]result_chunk[.]rejections|sie[.]gateway[.]queue[.]result_chunk[.]transfers_completed|sie[.]gateway[.]queue[.]result_chunk[.]duplicates|sie[.]gateway[.]queue[.]result_chunk[.]retry_replacements|sie[.]gateway[.]queue[.]result_chunk[.]stale_retries|sie[.]gateway[.]queue[.]result_chunk[.]reserved_bytes|sie[.]gateway[.]queue[.]events|sie[.]gateway[.]queue[.]lane_admission[.]decisions|sie[.]gateway[.]provisioning[.]responses|sie[.]gateway[.]routing[.]unsupported_model_exclusions|sie[.]gateway[.]generation[.]events|sie[.]gateway[.]generation[.]ttft|sie[.]gateway[.]generation[.]tpot|sie[.]gateway[.]generation[.]tokens|sie[.]gateway[.]pending_demand|sie[.]gateway[.]lane[.]queue[.]depth|sie[.]gateway[.]lane[.]queue[.]snapshot[.]timestamp|sie[.]gateway[.]active_lease[.]gpus|sie[.]gateway[.]pool[.]warm_floor|sie[.]gateway[.]rejected[.]requests|sie[.]gateway[.]key_snapshot[.]polls|sie[.]gateway[.]settlement[.]confirms|sie[.]gateway[.]settlement[.]confirm[.]duration)$")' - 'resource.attributes["service.name"] != "sie-gateway"' filter/remote_application_contract: error_mode: propagate diff --git a/deploy/helm/sie-cluster/templates/prometheusrule.yaml b/deploy/helm/sie-cluster/templates/prometheusrule.yaml index 60f84058b..634a82fd8 100644 --- a/deploy/helm/sie-cluster/templates/prometheusrule.yaml +++ b/deploy/helm/sie-cluster/templates/prometheusrule.yaml @@ -1,4 +1,8 @@ {{- if or (index .Values "kube-prometheus-stack" "install") .Values.alertRules.enabled }} +{{- $remoteSeconds := .Values.alertRules.remoteFallbackPersistenceSeconds -}} +{{- if not (and (or (kindIs "int" $remoteSeconds) (kindIs "int64" $remoteSeconds) (kindIs "float64" $remoteSeconds)) (eq (float64 $remoteSeconds) (floor (float64 $remoteSeconds))) (gt (float64 $remoteSeconds) 0.0) (le (float64 $remoteSeconds) 86400.0)) -}} +{{- fail "alertRules.remoteFallbackPersistenceSeconds must be an integer from 1 to 86400" -}} +{{- end -}} {{- $fullname := include "sie-cluster.fullname" . -}} # Application rules below read the collector's Prometheus compatibility # exporter. The shared applicationPrometheusEnabled gate therefore also @@ -13,5 +17,5 @@ metadata: app.kubernetes.io/component: alert prometheus: kube-prometheus spec: -{{ .Files.Get "files/alerts/sie-rules.yaml" | replace "__NAMESPACE__" (include "sie-cluster.namespace" .) | replace "__FULLNAME__" $fullname | replace "__COLLECTOR_SERVICE__" (printf "%s-otel-collector" $fullname) | indent 2 }} +{{ .Files.Get "files/alerts/sie-rules.yaml" | replace "__NAMESPACE__" (include "sie-cluster.namespace" .) | replace "__FULLNAME__" $fullname | replace "__COLLECTOR_SERVICE__" (printf "%s-otel-collector" $fullname) | replace "__REMOTE_FALLBACK_PERSISTENCE_SECONDS__" (printf "%d" (int64 $remoteSeconds)) | indent 2 }} {{- end }} diff --git a/deploy/helm/sie-cluster/values.yaml b/deploy/helm/sie-cluster/values.yaml index af8c0c499..3183d95fe 100644 --- a/deploy/helm/sie-cluster/values.yaml +++ b/deploy/helm/sie-cluster/values.yaml @@ -1490,6 +1490,9 @@ serviceMonitor: # using an external Prometheus Operator. alertRules: enabled: false + # -- Alert after committed remote bridges span this many seconds without + # observed local success. Requires recent activity from that gateway replica. + remoteFallbackPersistenceSeconds: 600 # -- Grafana dashboard ConfigMaps dashboards: diff --git a/packages/sie_gateway/src/handlers/proxy.rs b/packages/sie_gateway/src/handlers/proxy.rs index 2e0c04e8b..539d4a697 100644 --- a/packages/sie_gateway/src/handlers/proxy.rs +++ b/packages/sie_gateway/src/handlers/proxy.rs @@ -1,6 +1,6 @@ use axum::body::{to_bytes, Body}; use axum::extract::{Request, State}; -use axum::http::{HeaderMap, HeaderName, HeaderValue, Method, StatusCode}; +use axum::http::{Extensions, HeaderMap, HeaderName, HeaderValue, Method, StatusCode}; use axum::response::{IntoResponse, Response}; use axum::Json; use base64::Engine; @@ -2714,6 +2714,7 @@ async fn proxy_request_inner( .model_registry .serving_execution_evidence(&bundle, &hash_pool, &dispatch_model); ServingDisclosure::record_evidence(req.extensions(), served_by.clone()); + FallbackAttempt::record_model(req.extensions(), &dispatch_model, endpoint); if let Some(RemoteFallbackOverride(plan)) = req.extensions().get::() { if plan.config_hash != bundle_config_hash || served_by.as_ref() != Some(&plan.served_by) @@ -3104,6 +3105,7 @@ async fn proxy_request_inner( // (encode / score) get the same text-appropriate 16 MiB cap as the // chat / embeddings paths. Extract accepts bounded binary media, so // its cap covers the maximum legal audio after JSON base64 expansion. + let request_extensions = req.extensions().clone(); let body_limit = native_request_body_limit(endpoint); let body_bytes = if let Some(body) = prepared_native_body { body @@ -3148,6 +3150,7 @@ async fn proxy_request_inner( batch_target, require_execution_authority_v1, prefetch_first_output, + &request_extensions, &physical_lane, ); // Scope an OTel context over the publish so the work-item envelope @@ -3613,6 +3616,7 @@ async fn queue_mode_proxy( batch_target: Option, require_execution_authority_v1: bool, prefetch_first_output: bool, + request_extensions: &Extensions, physical_lane: &PhysicalLane, ) -> Response { // Parse body once, extract items + params (avoids double parse) @@ -3670,6 +3674,7 @@ async fn queue_mode_proxy( if params.generate.as_ref().is_some_and(|params| params.stream) { return super::sse::build_sse_response(super::sse::SseParams { prefetch_first_output, + local_serving_model: FallbackAttempt::defer_local_stream(request_extensions), state, work_publisher: work_publisher_arc, physical_lane: physical_lane.clone(), @@ -7845,6 +7850,7 @@ async fn resolve_generation_route( .model_registry .serving_execution_evidence(&bundle, &hash_pool, dispatch_model); ServingDisclosure::record_evidence(ext, served_by.clone()); + FallbackAttempt::record_model(ext, dispatch_model, "generate"); if let Some(RemoteFallbackOverride(plan)) = bridge { if plan.config_hash != bundle_config_hash || served_by.as_ref() != Some(&plan.served_by) @@ -8659,6 +8665,7 @@ async fn proxy_chat_inner( .extensions .get::() .is_some_and(FallbackAttempt::active), + local_serving_model: FallbackAttempt::defer_local_stream(&parts.extensions), state: state.as_ref(), work_publisher: work_publisher_arc, physical_lane: physical_lane.clone(), @@ -9327,6 +9334,7 @@ async fn proxy_completions_inner(state: Arc, req: Request) -> Response .extensions .get::() .is_some_and(FallbackAttempt::active), + local_serving_model: FallbackAttempt::defer_local_stream(&parts.extensions), state: state.as_ref(), work_publisher: work_publisher_arc, physical_lane: physical_lane.clone(), diff --git a/packages/sie_gateway/src/handlers/serving_disclosure.rs b/packages/sie_gateway/src/handlers/serving_disclosure.rs index 27d3b7494..93b646453 100644 --- a/packages/sie_gateway/src/handlers/serving_disclosure.rs +++ b/packages/sie_gateway/src/handlers/serving_disclosure.rs @@ -12,6 +12,7 @@ use axum::extract::Request; use axum::http::{Extensions, HeaderMap, HeaderName, HeaderValue, StatusCode}; use axum::response::Response; +use crate::observability::metrics as telemetry; use crate::server::AppState; use crate::types::model::{FallbackTrigger, ServedBy}; @@ -95,7 +96,18 @@ impl ServingDisclosure { /// One request's original pre-acceptance refusal, retained through its bridge. /// It is not a retry counter: a request may install exactly one remote attempt. #[derive(Clone, Default)] -pub(crate) struct FallbackAttempt(Arc>>); +pub(crate) struct FallbackAttempt(Arc>); + +#[derive(Default)] +struct FallbackState { + original: Option, + observation: Option, +} + +struct ServingObservation { + model: String, + operation: &'static str, +} /// A compatibility facade owns response validation before a bridge commits. /// This extension is gateway-created and cannot be supplied by a caller. @@ -104,7 +116,7 @@ pub(crate) struct DeferredFallbackFinish; struct LocalRefusal { response: Response, - reason: &'static str, + trigger: FallbackTrigger, } impl FallbackAttempt { @@ -117,19 +129,56 @@ impl FallbackAttempt { attempt } + /// Only a resolved catalog route can name a telemetry model. Preserve the + /// original route while recursion selects a remote physical profile. + pub(crate) fn record_model(extensions: &Extensions, model: &str, operation: &str) { + let Some(attempt) = extensions.get::() else { + return; + }; + let operation = match operation { + "encode" => "encode", + "score" => "score", + "extract" => "extract", + "generate" => "generate", + _ => return, + }; + let mut state = attempt.0.lock().unwrap_or_else(PoisonError::into_inner); + if state.original.is_none() && state.observation.is_none() { + state.observation = Some(ServingObservation { + model: model.split(':').next().unwrap_or(model).to_string(), + operation, + }); + } + } + + /// Transfer local streaming observation to the output driver. HTTP 200 + /// alone does not prove a stream served any valid local output. + pub(crate) fn defer_local_stream(extensions: &Extensions) -> Option { + let disclosure = extensions.get::()?; + if !matches!( + *disclosure.0.lock().unwrap_or_else(PoisonError::into_inner), + Some(ServedBy::Local) + ) { + return None; + } + let attempt = extensions.get::()?; + let mut state = attempt.0.lock().unwrap_or_else(PoisonError::into_inner); + if state.original.is_some() { + return None; + } + state + .observation + .take() + .map(|observation| observation.model) + } + /// Retain the response before any local work has been accepted. pub(crate) fn begin(&self, response: Response, trigger: FallbackTrigger) -> bool { - let mut original = self.0.lock().unwrap_or_else(PoisonError::into_inner); - if original.is_some() { + let mut state = self.0.lock().unwrap_or_else(PoisonError::into_inner); + if state.original.is_some() { return false; } - let reason = match trigger { - FallbackTrigger::Provisioning => "provisioning", - FallbackTrigger::ModelLoading => "model_loading", - FallbackTrigger::Saturated => "saturated", - FallbackTrigger::Unhealthy => "unhealthy", - }; - *original = Some(LocalRefusal { response, reason }); + state.original = Some(LocalRefusal { response, trigger }); true } @@ -137,17 +186,39 @@ impl FallbackAttempt { self.0 .lock() .unwrap_or_else(PoisonError::into_inner) + .original .is_some() } /// Called after normal disclosure stamping and before HTTP success/output. /// A failed bridge preserves the original body and Retry-After verbatim. pub(crate) fn finish(&self, mut response: Response) -> Response { - let Some(mut original) = self.0.lock().unwrap_or_else(PoisonError::into_inner).take() - else { + let (original, observation) = { + let mut state = self.0.lock().unwrap_or_else(PoisonError::into_inner); + (state.original.take(), state.observation.take()) + }; + let Some(mut original) = original else { + if response.status().is_success() + && response + .headers() + .get(SERVED_BY_HEADER) + .is_some_and(|value| value == "local") + { + if let Some(observation) = observation { + telemetry::record_local_serving_success(&observation.model); + } + } return response; }; - let reason = HeaderValue::from_static(original.reason); + if let Some(observation) = observation { + telemetry::record_remote_fallback( + &observation.model, + observation.operation, + original.trigger, + response.status().is_success(), + ); + } + let reason = HeaderValue::from_static(original.trigger.as_str()); if response.status().is_success() { response .headers_mut() @@ -180,6 +251,38 @@ mod tests { use super::*; + #[test] + fn local_stream_observation_transfers_once_and_excludes_remote_routes_and_bridges() { + for (served_by, bridge, expected) in [ + (Some(ServedBy::Local), false, Some("acme/chat")), + (Some(ServedBy::Remote { upstream: None }), false, None), + (None, false, None), + (Some(ServedBy::Local), true, None), + ] { + let mut req = Request::new(Body::empty()); + ServingDisclosure::install(&mut req); + let attempt = FallbackAttempt::install(&mut req); + ServingDisclosure::record_evidence(req.extensions(), served_by); + FallbackAttempt::record_model(req.extensions(), "acme/chat:default", "generate"); + if bridge { + assert!(attempt.begin( + StatusCode::SERVICE_UNAVAILABLE.into_response(), + FallbackTrigger::Provisioning, + )); + } + assert_eq!( + FallbackAttempt::defer_local_stream(req.extensions()).as_deref(), + expected, + ); + if expected.is_some() { + assert!(FallbackAttempt::defer_local_stream(req.extensions()).is_none()); + assert!(attempt.0.lock().unwrap().observation.is_none()); + } else { + assert!(attempt.0.lock().unwrap().observation.is_some()); + } + } + } + async fn buffered_surface( gateway: &TestGateway, surface: &str, diff --git a/packages/sie_gateway/src/handlers/sse.rs b/packages/sie_gateway/src/handlers/sse.rs index d8332e3e9..7444d3415 100644 --- a/packages/sie_gateway/src/handlers/sse.rs +++ b/packages/sie_gateway/src/handlers/sse.rs @@ -94,6 +94,8 @@ pub enum SseEndpoint { pub struct SseParams<'a> { /// A bridge retains HTTP refusal authority until a valid event is ready. pub prefetch_first_output: bool, + /// A resolved local route resets remote persistence only on valid output. + pub local_serving_model: Option, pub state: &'a AppState, pub work_publisher: Arc, pub physical_lane: PhysicalLane, @@ -128,6 +130,7 @@ pub struct SseParams<'a> { pub async fn build_sse_response(params: SseParams<'_>) -> Response { let SseParams { prefetch_first_output, + local_serving_model, state, work_publisher, physical_lane, @@ -333,6 +336,7 @@ pub async fn build_sse_response(params: SseParams<'_>) -> Response { async move { run_sse_driver(SseDriverArgs { first_output, + local_serving_model, event_tx, chunk_rx, outcome_rx, @@ -418,6 +422,7 @@ type FirstOutputGate = Option, event_tx: tokio::sync::mpsc::Sender>, chunk_rx: broadcast::Receiver, outcome_rx: tokio::sync::oneshot::Receiver, @@ -564,12 +569,17 @@ async fn wait_for_terminal_durability( async fn run_sse_driver(args: SseDriverArgs) { let lifecycle = crate::observability::tracing::request_telemetry_enabled() .then(|| Lifecycle::generation_stream(opentelemetry::Context::current())); - run_sse_driver_with_lifecycle(args, lifecycle).await; + run_sse_driver_with_lifecycle(args, lifecycle, telemetry::record_local_serving_success).await; } -async fn run_sse_driver_with_lifecycle(args: SseDriverArgs, mut lifecycle: Option) { +async fn run_sse_driver_with_lifecycle( + args: SseDriverArgs, + mut lifecycle: Option, + record_local_success: impl Fn(&str), +) { let SseDriverArgs { mut first_output, + mut local_serving_model, event_tx, mut chunk_rx, outcome_rx, @@ -1258,6 +1268,14 @@ async fn run_sse_driver_with_lifecycle(args: SseDriverArgs, mut lifecycle: Optio .await; return; } + if chunk.error.is_none() + && !(chunk.done + && matches!(chunk.finish_reason.as_deref(), Some("cancelled" | "error"))) + { + if let Some(model) = local_serving_model.take() { + record_local_success(&model); + } + } } if is_terminal { @@ -2099,6 +2117,7 @@ mod tests { #[derive(Clone, Copy)] enum WorkerTerminalDelivery { + BeforeAnyDelta, BackloggedBeforeFirstPoll, AfterFirstDelta, } @@ -2151,7 +2170,9 @@ mod tests { endpoint: SseEndpoint, delivery: WorkerTerminalDelivery, error: ChunkError, - ) -> Vec { + ) -> (Vec, usize) { + use std::sync::atomic::{AtomicUsize, Ordering}; + use crate::queue::streaming::StreamCollector; use crate::state::demand_tracker::PhysicalLaneCatalog; @@ -2182,8 +2203,15 @@ mod tests { .expect("durability receiver is live"); let (event_tx, mut event_rx) = tokio::sync::mpsc::channel(8); + let local_successes = Arc::new(AtomicUsize::new(0)); + let recorded_successes = Arc::clone(&local_successes); + let record_local_success = move |model: &str| { + assert_eq!(model, "test/model"); + recorded_successes.fetch_add(1, Ordering::SeqCst); + }; let args = SseDriverArgs { first_output: None, + local_serving_model: Some("test/model".to_string()), event_tx, chunk_rx, outcome_rx, @@ -2215,6 +2243,7 @@ mod tests { Some(Lifecycle::generation_stream( opentelemetry::Context::current(), )), + record_local_success, ) .with_current_subscriber(), ); @@ -2249,10 +2278,12 @@ mod tests { .expect("outcome receiver is live"); driver.await.expect("driver task"); } else { - assert!(matches!( - collector.apply(_delta_chunk(41, "partial")), - ChunkApplied::Delta - )); + if !matches!(delivery, WorkerTerminalDelivery::BeforeAnyDelta) { + assert!(matches!( + collector.apply(_delta_chunk(41, "partial")), + ChunkApplied::Delta + )); + } let mut terminal = _terminal_chunk("error", None); terminal.seq = 42; terminal.usage = Some(UsageBlock { @@ -2276,6 +2307,7 @@ mod tests { Some(Lifecycle::generation_stream( opentelemetry::Context::current(), )), + record_local_success, ) .await; } @@ -2289,7 +2321,7 @@ mod tests { }; payloads.push(_event_data(event.expect("infallible event")).await); } - payloads + (payloads, local_successes.load(Ordering::SeqCst)) } #[tokio::test] @@ -2307,7 +2339,7 @@ mod tests { }, SseEndpoint::Generate, ] { - let payloads = run_driver_worker_error_race( + let (payloads, local_successes) = run_driver_worker_error_race( endpoint, delivery, ChunkError { @@ -2318,6 +2350,10 @@ mod tests { }, ) .await; + assert_eq!( + local_successes, 1, + "only valid delivered local output resets remote persistence", + ); assert_eq!( payloads .iter() @@ -2380,6 +2416,36 @@ mod tests { } } + #[tokio::test] + async fn live_driver_worker_error_before_any_output_preserves_remote_period() { + for endpoint in [ + SseEndpoint::Chat { + include_usage: false, + }, + SseEndpoint::Completion { + include_usage: false, + }, + SseEndpoint::Generate, + ] { + let (payloads, local_successes) = run_driver_worker_error_race( + endpoint, + WorkerTerminalDelivery::BeforeAnyDelta, + ChunkError { + code: "RESOURCE_EXHAUSTED".to_string(), + message: "scheduler full".to_string(), + param: None, + retry_after_s: None, + }, + ) + .await; + assert_eq!(local_successes, 0); + assert!(payloads + .iter() + .any(|payload| payload.contains("RESOURCE_EXHAUSTED"))); + assert_eq!(payloads.last().unwrap(), "[DONE]"); + } + } + #[tokio::test] async fn live_driver_surfaces_synthetic_only_outcome_when_tap_closes() { use crate::state::demand_tracker::PhysicalLaneCatalog; @@ -2439,6 +2505,7 @@ mod tests { run_sse_driver(SseDriverArgs { first_output: None, + local_serving_model: None, event_tx, chunk_rx, outcome_rx, @@ -2497,7 +2564,7 @@ mod tests { }, SseEndpoint::Generate, ] { - let payloads = run_driver_worker_error_race( + let (payloads, local_successes) = run_driver_worker_error_race( endpoint, WorkerTerminalDelivery::BackloggedBeforeFirstPoll, ChunkError { @@ -2508,6 +2575,7 @@ mod tests { }, ) .await; + assert_eq!(local_successes, 1); let error_payload = payloads .iter() .filter_map(|payload| serde_json::from_str::(payload).ok()) diff --git a/packages/sie_gateway/src/observability/metrics.rs b/packages/sie_gateway/src/observability/metrics.rs index 64ea40a88..75ca4a0fe 100644 --- a/packages/sie_gateway/src/observability/metrics.rs +++ b/packages/sie_gateway/src/observability/metrics.rs @@ -8,7 +8,7 @@ use std::collections::{BTreeSet, HashMap, HashSet}; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Mutex, OnceLock}; -use std::time::Duration; +use std::time::{Duration, Instant}; use opentelemetry::metrics::{Counter, Gauge, Histogram, ObservableGauge}; use opentelemetry::trace::SpanContext; @@ -16,12 +16,18 @@ use opentelemetry::{global, Context, KeyValue}; use crate::queue::lane_admission::LaneAdmissionOutcome; use crate::state::demand_tracker::{DemandTracker, PhysicalLane}; +use crate::types::model::FallbackTrigger; pub const REQUESTS_METRIC_NAME: &str = "sie.gateway.requests"; pub const REQUEST_DURATION_METRIC_NAME: &str = "sie.gateway.request.duration"; pub const ADMISSION_DECISIONS_METRIC_NAME: &str = "sie.gateway.admission.decisions"; pub const DISPATCHES_METRIC_NAME: &str = "sie.gateway.dispatches"; pub const DISPATCH_DURATION_METRIC_NAME: &str = "sie.gateway.dispatch.duration"; +pub const REMOTE_FALLBACKS_METRIC_NAME: &str = "sie.gateway.remote.fallbacks"; +pub const REMOTE_SERVING_DURATION_METRIC_NAME: &str = "sie.gateway.remote.serving.duration"; +pub const REMOTE_FALLBACK_MODEL_LIMIT: usize = 256; +const REMOTE_SERVING_MAX_IDLE: Duration = Duration::from_secs(300); +pub const REMOTE_FALLBACK_COUNTER_LIMIT: usize = (REMOTE_FALLBACK_MODEL_LIMIT + 1) * 7 * 4 * 2; pub const PENDING_DEMAND_METRIC_NAME: &str = "sie.gateway.pending_demand"; pub const LANE_QUEUE_DEPTH_METRIC_NAME: &str = "sie.gateway.lane.queue.depth"; pub const LANE_QUEUE_SNAPSHOT_TIMESTAMP_METRIC_NAME: &str = @@ -788,6 +794,52 @@ pub fn sanitize_label(value: &str) -> String { value.to_string() } +/// Retain only the first finite set of canonical catalog model labels. A +/// duration belongs to one exact model; overflow is counted but never tracked. +#[derive(Default)] +struct RemoteServingState { + models: HashMap>, +} + +struct RemoteServingPeriod { + first: Instant, + latest: Instant, +} + +impl RemoteServingState { + fn fallback(&mut self, model: &str, committed: bool, now: Instant) -> (String, Option) { + let model = sanitize_model_label(model); + if model == "other" + || (!self.models.contains_key(&model) + && self.models.len() >= REMOTE_FALLBACK_MODEL_LIMIT) + { + return ("other".to_string(), None); + } + let since = self.models.entry(model.clone()).or_default(); + let duration = committed.then(|| { + if since.as_ref().is_some_and(|period| { + now.saturating_duration_since(period.latest) > REMOTE_SERVING_MAX_IDLE + }) { + *since = None; + } + let period = since.get_or_insert(RemoteServingPeriod { + first: now, + latest: now, + }); + period.latest = now; + now.saturating_duration_since(period.first).as_secs_f64() + }); + (model, duration) + } + + fn local_success(&mut self, model: &str) -> Option { + let model = sanitize_model_label(model); + let since = self.models.get_mut(&model)?; + *since = None; + Some(model) + } +} + struct GatewayTelemetry { requests: Counter, request_duration: Histogram, @@ -796,6 +848,9 @@ struct GatewayTelemetry { dispatches: Counter, #[allow(dead_code)] // Managed composition API. dispatch_duration: Histogram, + remote_fallbacks: Counter, + remote_serving_duration: Gauge, + remote_serving_state: Mutex, pending_demand: Gauge, lane_queue_depth: Gauge, lane_queue_snapshot_timestamp: Gauge, @@ -888,6 +943,17 @@ impl GatewayTelemetry { .with_unit("s") .with_boundaries(REQUEST_LATENCY_BUCKETS.to_vec()) .build(), + remote_fallbacks: meter + .u64_counter(REMOTE_FALLBACKS_METRIC_NAME) + .with_description("Remote bridge responses committed or refused before output.") + .with_unit("{request}") + .build(), + remote_serving_duration: meter + .f64_gauge(REMOTE_SERVING_DURATION_METRIC_NAME) + .with_description("Seconds between committed bridges since the last local success.") + .with_unit("s") + .build(), + remote_serving_state: Mutex::new(RemoteServingState::default()), pending_demand: meter .f64_gauge(PENDING_DEMAND_METRIC_NAME) .with_description("Whether a physical worker lane has refreshable unmet demand.") @@ -1548,6 +1614,75 @@ fn telemetry() -> Option<&'static GatewayTelemetry> { Some(TELEMETRY.get_or_init(|| GatewayTelemetry::new(&global::meter("sie-gateway")))) } +/// One semantic observation at the response/first-output commitment boundary. +pub(crate) fn record_remote_fallback( + model: &str, + operation: &str, + reason: FallbackTrigger, + committed: bool, +) { + record_remote_fallback_to( + telemetry(), + model, + operation, + reason, + committed, + Instant::now, + ); +} + +fn record_remote_fallback_to( + telemetry: Option<&GatewayTelemetry>, + model: &str, + operation: &str, + reason: FallbackTrigger, + committed: bool, + now: impl FnOnce() -> Instant, +) { + let Some(telemetry) = telemetry else { + return; + }; + let mut state = telemetry + .remote_serving_state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let (model, duration) = state.fallback(model, committed, now()); + telemetry.remote_fallbacks.add( + 1, + &[ + KeyValue::new("model", model.clone()), + KeyValue::new("operation", bounded_operation(operation)), + KeyValue::new("fallback.reason", reason.as_str()), + KeyValue::new("outcome", if committed { "committed" } else { "refused" }), + ], + ); + if let Some(duration) = duration { + telemetry + .remote_serving_duration + .record(duration, &[KeyValue::new("model", model)]); + } +} + +pub(crate) fn record_local_serving_success(model: &str) { + record_local_serving_success_to(telemetry(), model); +} + +fn record_local_serving_success_to(telemetry: Option<&GatewayTelemetry>, model: &str) { + let Some(telemetry) = telemetry else { + return; + }; + if let Some(model) = telemetry + .remote_serving_state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .local_success(model) + { + telemetry + .remote_serving_duration + .record(0.0, &[KeyValue::new("model", model)]); + } +} + fn bounded_operation(operation: &str) -> &'static str { match operation { "encode" => "encode", @@ -2298,6 +2433,119 @@ mod tests { ] } + #[test] + fn remote_fallback_duration_tracks_commits_and_resets_on_local_success() { + let mut state = RemoteServingState::default(); + let now = Instant::now(); + assert_eq!( + state.fallback("acme/model", true, now), + ("acme/model".into(), Some(0.0)) + ); + assert_eq!( + state.fallback("acme/model", false, now + Duration::from_secs(20)), + ("acme/model".into(), None) + ); + assert_eq!( + state.fallback("acme/model", true, now + Duration::from_secs(42)), + ("acme/model".into(), Some(42.0)) + ); + assert_eq!(state.local_success("acme/model"), Some("acme/model".into())); + assert_eq!( + state.fallback("acme/model", true, now + Duration::from_secs(50)), + ("acme/model".into(), Some(0.0)) + ); + assert_eq!( + state.fallback("acme/model", false, now + Duration::from_secs(390)), + ("acme/model".into(), None) + ); + assert_eq!( + state.fallback("acme/model", true, now + Duration::from_secs(400)), + ("acme/model".into(), Some(0.0)) + ); + assert_eq!( + state.fallback("acme/model", true, now + Duration::from_secs(410)), + ("acme/model".into(), Some(10.0)) + ); + assert_eq!(state.local_success("unobserved/model"), None); + assert_eq!( + state.fallback("private invalid\nlabel", true, now), + ("other".into(), None) + ); + } + + #[test] + fn remote_fallback_full_declared_domain_exports_without_sdk_overflow() { + let (telemetry, exporter, provider) = metric_points(); + let now = Instant::now(); + for model in 0..=REMOTE_FALLBACK_MODEL_LIMIT { + for operation in [ + "encode", + "score", + "extract", + "embeddings", + "moderations", + "generate", + "other", + ] { + for reason in [ + FallbackTrigger::Provisioning, + FallbackTrigger::ModelLoading, + FallbackTrigger::Saturated, + FallbackTrigger::Unhealthy, + ] { + for committed in [false, true] { + record_remote_fallback_to( + Some(&telemetry), + &format!("acme/model-{model}"), + operation, + reason, + committed, + || now, + ); + } + } + } + } + assert_eq!( + telemetry.remote_serving_state.lock().unwrap().models.len(), + REMOTE_FALLBACK_MODEL_LIMIT + ); + record_local_serving_success_to(Some(&telemetry), "acme/model-0"); + provider.force_flush().unwrap(); + let resources = exporter.get_finished_metrics().unwrap(); + let metrics = resources + .iter() + .flat_map(|resource| resource.scope_metrics()) + .flat_map(|scope| scope.metrics()) + .collect::>(); + let counter = metrics + .iter() + .find(|metric| metric.name() == REMOTE_FALLBACKS_METRIC_NAME) + .unwrap(); + assert_eq!(counter.unit(), "{request}"); + let AggregatedMetrics::U64(MetricData::Sum(counter)) = counter.data() else { + panic!("expected fallback counter"); + }; + assert_eq!(counter.data_points().count(), REMOTE_FALLBACK_COUNTER_LIMIT); + assert!(counter.data_points().all(|point| point + .attributes() + .all(|attribute| attribute.key.as_str() != "otel.metric.overflow"))); + let gauge = metrics + .iter() + .find(|metric| metric.name() == REMOTE_SERVING_DURATION_METRIC_NAME) + .unwrap(); + assert_eq!(gauge.unit(), "s"); + let AggregatedMetrics::F64(MetricData::Gauge(gauge)) = gauge.data() else { + panic!("expected remote duration gauge"); + }; + assert_eq!(gauge.data_points().count(), REMOTE_FALLBACK_MODEL_LIMIT); + assert!(gauge + .data_points() + .all(|point| point.attributes().count() == 1)); + assert!(gauge.data_points().all(|point| point.value() == 0.0)); + provider.shutdown().unwrap(); + } + #[test] fn every_admission_outcome_is_declared_in_the_telemetry_contract() { let contract: serde_yaml::Value = serde_yaml::from_str(include_str!(concat!( diff --git a/packages/sie_gateway/src/observability/tracing.rs b/packages/sie_gateway/src/observability/tracing.rs index e8fe2929f..87224e28f 100644 --- a/packages/sie_gateway/src/observability/tracing.rs +++ b/packages/sie_gateway/src/observability/tracing.rs @@ -59,7 +59,9 @@ use crate::observability::metrics::{ RequestCompletionObservation, ACTIVE_LEASE_GPUS_METRIC_NAME, KEDA_SCALE_UP_REJECTION_REASON_CARDINALITY, LANE_QUEUE_DEPTH_METRIC_NAME, LANE_QUEUE_SNAPSHOT_TIMESTAMP_METRIC_NAME, PENDING_DEMAND_METRIC_NAME, - POOL_WARM_FLOOR_METRIC_NAME, REJECTED_REQUESTS_METRIC_NAME, + POOL_WARM_FLOOR_METRIC_NAME, REJECTED_REQUESTS_METRIC_NAME, REMOTE_FALLBACKS_METRIC_NAME, + REMOTE_FALLBACK_COUNTER_LIMIT, REMOTE_FALLBACK_MODEL_LIMIT, + REMOTE_SERVING_DURATION_METRIC_NAME, }; use crate::state::demand_tracker::MAX_CONFIGURED_PHYSICAL_LANES; @@ -504,8 +506,8 @@ fn init_metrics(endpoint: &str) -> bool { true } -/// Override the SDK's default 2,000-series ceiling for every KEDA-filtered -/// stream. Lane snapshots have one point per catalog member; the rejection +/// Override the SDK's default ceiling for declared KEDA and remote fallback +/// streams. Lane snapshots have one point per catalog member; the rejection /// counter has four scale-worthy reasons per member. These limits are exact, /// finite, and prevent a valid high-index lane from collapsing into the OTel /// overflow series (which PromQL's exact lane filters cannot see). @@ -517,13 +519,15 @@ pub(crate) fn keda_metric_cardinality_view(instrument: &Instrument) -> Option MAX_CONFIGURED_PHYSICAL_LANES, REJECTED_REQUESTS_METRIC_NAME => KEDA_REJECTED_REQUESTS_CARDINALITY_LIMIT, + REMOTE_FALLBACKS_METRIC_NAME => REMOTE_FALLBACK_COUNTER_LIMIT, + REMOTE_SERVING_DURATION_METRIC_NAME => REMOTE_FALLBACK_MODEL_LIMIT, _ => return None, }; Some( Stream::builder() .with_cardinality_limit(limit) .build() - .expect("constant KEDA cardinality limits must be valid"), + .expect("constant telemetry cardinality limits must be valid"), ) } diff --git a/packages/sie_gateway/src/types/model.rs b/packages/sie_gateway/src/types/model.rs index eb2ad2cf8..21c4c8f10 100644 --- a/packages/sie_gateway/src/types/model.rs +++ b/packages/sie_gateway/src/types/model.rs @@ -446,6 +446,17 @@ pub enum FallbackTrigger { Unhealthy, } +impl FallbackTrigger { + pub(crate) fn as_str(self) -> &'static str { + match self { + Self::Provisioning => "provisioning", + Self::ModelLoading => "model_loading", + Self::Saturated => "saturated", + Self::Unhealthy => "unhealthy", + } + } +} + impl TryFrom for RoutingConfig { type Error = String; fn try_from(fields: RoutingFields) -> Result { diff --git a/packages/sie_server/REMOTE_BACKENDS.md b/packages/sie_server/REMOTE_BACKENDS.md index 874b961e6..ec3b3a60d 100644 --- a/packages/sie_server/REMOTE_BACKENDS.md +++ b/packages/sie_server/REMOTE_BACKENDS.md @@ -440,3 +440,14 @@ remain in the stream. Explicit selectors and `X-SIE-Remote: forbid` retain their existing authority. Numerical fleet equivalence and coordinated threshold routing remain inactive. + +### Observing cluster fallback + +The queue-routing dashboard shows gateway fallback response rates by model, +operation, reason and commitment outcome, plus observed remote-serving duration. +Enable the chart's existing alert rules to receive `SIERemoteFallbackPersistent`; +set `alertRules.remoteFallbackPersistenceSeconds` to the desired threshold +(default 600 seconds). The rule requires recent committed remote activity from +the same gateway replica. Successful local serving clears that replica's +observed duration. See [the telemetry contract](../../telemetry/README.md#remote-fallback-observations) +for bounded-label and replica semantics. diff --git a/telemetry/README.md b/telemetry/README.md index 44cbc5479..e4a2645cd 100644 --- a/telemetry/README.md +++ b/telemetry/README.md @@ -709,3 +709,30 @@ histograms, timestamps, temporality or producer domain policies. Local Prometheu processing remains unchanged. The pinned collector regression sends raw duplicate keys through both rendered receiver branches, including reversed service claims, nested values and missing optional fields on successive points. + +### Remote fallback observations + +The gateway emits `sie.gateway.remote.fallbacks` once when a bridge response +commits before first output or restores its local refusal. Its bounded labels +are operation, canonical catalog model, fallback reason, and `committed` or +`refused`. A later streaming error is counted by the existing stream metrics. +`sie.gateway.remote.serving.duration` records seconds between the first and +latest committed bridges since that gateway last observed a successful local +response for the model. A local stream resets the period only after its first +valid output event; HTTP 200 before an error does not reset it. A gap longer than five minutes between committed +bridges starts a new period, so an idle model's next cold request cannot inherit +an old outage duration. Refused bridges do not extend the period. This is a +per-replica observation, not a fleet clock. + +The facade retains at most 256 exact model names for the process lifetime. +Additional models collapse to `model=other` on the counter and have no duration +series. Disabled telemetry constructs no point attributes. The collector +preserves only the declared attributes, exports through the existing OTLP and +Prometheus paths, and the queue-routing dashboard displays both instruments. + +`SIERemoteFallbackPersistent` requires recent committed activity from the same +producer instance and collector generation before comparing its duration with +`alertRules.remoteFallbackPersistenceSeconds` (600 by default, integer 1–86400). +The rule uses the maximum active replica duration per model. Local success +records zero; idle historical samples cannot sustain the alert. These metrics +are diagnostics and do not change KEDA control signals. diff --git a/telemetry/contract.yaml b/telemetry/contract.yaml index 059f7c19d..96a0814e3 100644 --- a/telemetry/contract.yaml +++ b/telemetry/contract.yaml @@ -399,6 +399,42 @@ cardinality_budgets: # Per-lane queue admission reuses the same deployment-owned lane tuple and # the same fail-closed rules as the KEDA control domain, but is not itself a # KEDA control signal: its Prometheus spelling is not a scaler API. + gateway_remote_fallback: + scope: per_process_per_instrument + source: canonical_catalog_route_before_bridge + admitted_exact_models_process_lifetime: 256 + counter_overflow_model: other + duration_overflow_models: omitted + model_label_maximum_length: 256 + total_counter_model_series: 257 + counter_series_limit: 14392 # 257 models * 7 operations * 4 reasons * 2 outcomes. + fixed_domains: + models: 257 + exact_models: 256 + operation: operation + reason: remote_fallback_reason + outcome: remote_fallback_outcome + instrument_limits: + fallback_responses: + metrics: [sie.gateway.remote.fallbacks] + formula: models * operation * reason * outcome + limit: 14392 + remote_duration: + metrics: [sie.gateway.remote.serving.duration] + formula: exact_models + limit: 256 + duration_series_limit: 256 + state_entries_limit: 256 + success_boundary: response_committed_before_first_output + duration_semantics: seconds_between_first_and_latest_committed_bridge_since_local_success_or_idle_reset + maximum_gap_between_committed_bridges_s: 300 + local_success_reset: true + stale_series_policy: alert_requires_recent_committed_bridge_activity + replicas: independent_observations_max_duration_and_sum_counter_rates + applies_to: + - sie.gateway.remote.fallbacks + - sie.gateway.remote.serving.duration + gateway_lane_admission_domain: source: SIE_GATEWAY_CONFIGURED_PHYSICAL_LANES tuple: [pool, machine_profile, bundle] @@ -788,6 +824,8 @@ enums: # `skipped` is the no-reservation branch, which issues no control-plane call # and therefore contributes no duration observation. settlement_confirm_outcome: [success, error, skipped] + remote_fallback_reason: [provisioning, model_loading, saturated, unhealthy] + remote_fallback_outcome: [committed, refused] gateway_dispatch_path: [i6pn, modal, other] dispatcher_path: [modal_remote, modal_stream, modal_spawn, generation_ws, other] dispatch_outcome: [success, error, timeout, cancelled, not_found, other] @@ -1062,6 +1100,13 @@ attribute_sets: outcome: dispatch_outcome fallback.reason: fallback_reason lane: lane + gateway_remote_fallback: + operation: operation + model: model + fallback.reason: remote_fallback_reason + outcome: remote_fallback_outcome + gateway_remote_serving: + model: model gateway_config_operation: operation: gateway_config_operation outcome: gateway_config_outcome @@ -1541,6 +1586,22 @@ metrics: buckets: gateway_request_seconds prometheus_name: sie_gateway_dispatch_duration_seconds export: [prometheus, otlp] + - name: sie.gateway.remote.fallbacks + owner: gateway + facade_event: gateway.remote_fallback_completed + type: counter + unit: "{request}" + attributes: gateway_remote_fallback + prometheus_name: sie_gateway_remote_fallbacks_total + export: [prometheus, otlp] + - name: sie.gateway.remote.serving.duration + owner: gateway + facade_event: gateway.remote_fallback_completed_or_local_response_committed + type: gauge + unit: s + attributes: gateway_remote_serving + prometheus_name: sie_gateway_remote_serving_duration_seconds + export: [prometheus, otlp] - name: sie.gateway.config.applied_epoch owner: gateway facade_event: gateway.config_state_changed diff --git a/tools/ci/tests/test_collector_metric_privacy.py b/tools/ci/tests/test_collector_metric_privacy.py index 88a350c70..26b0b5f28 100644 --- a/tools/ci/tests/test_collector_metric_privacy.py +++ b/tools/ci/tests/test_collector_metric_privacy.py @@ -81,22 +81,34 @@ def payload(receiver): add(rs.resource.attributes, key, SENTINEL) add(rs.resource.attributes, key, None) scope = rs.scope_metrics.add() - for suffix, kind in [("requests", "sum"), ("request.duration", "histogram")]: - metric = scope.metrics.add(name=f"{prefix}.{suffix}") + specs = [(f"{prefix}.requests", "sum"), (f"{prefix}.request.duration", "histogram")] + # Gateway-only diagnostics are also sent to the application receiver: + # its allowlist must drop them rather than authorize their service. + specs.extend( + [ + ("sie.gateway.remote.fallbacks", "sum"), + ("sie.gateway.remote.serving.duration", "gauge"), + ] + ) + for name, kind in specs: + metric = scope.metrics.add(name=name) data = getattr(metric, kind) - data.aggregation_temporality = AGGREGATION_TEMPORALITY_DELTA + if kind != "gauge": + data.aggregation_temporality = AGGREGATION_TEMPORALITY_DELTA if kind == "sum": data.is_monotonic = True for case in range(3): point = data.data_points.add(start_time_unix_nano=1_000, time_unix_nano=2_000 + case) if kind == "sum": point.as_int = 10 + case + elif kind == "gauge": + point.as_double = 10.0 + case else: point.count = 2 point.sum = 0.5 + case point.explicit_bounds.extend([1.0]) point.bucket_counts.extend([1, 1]) - add(point.attributes, "outcome", "success") + add(point.attributes, "outcome", "committed" if name == "sie.gateway.remote.fallbacks" else "success") add(point.attributes, "outcome", None) if case == 0: add(point.attributes, "operation", "encode") @@ -109,6 +121,12 @@ def payload(receiver): add(point.attributes, "http.status_code", 200) add(point.attributes, "http.status_code", None) add(point.attributes, "cloud.region", SENTINEL) + if name.startswith("sie.gateway.remote."): + add(point.attributes, "model", "acme/model") + add(point.attributes, "model", SENTINEL) + add(point.attributes, "model", None) + add(point.attributes, "fallback.reason", "provisioning") + add(point.attributes, "fallback.reason", None) return wire.SerializeToString() @@ -191,9 +209,10 @@ def test_remote_metric_maps_have_one_scalar_per_retained_key(tmp_path, receiver) } for scope in rs["scopeMetrics"]: for metric in scope["metrics"]: - kind = "sum" if "sum" in metric else "histogram" + kind = next(kind for kind in ["sum", "histogram", "gauge"] if kind in metric) data = metric[kind] - assert data["aggregationTemporality"] == 1 + if kind != "gauge": + assert data["aggregationTemporality"] == 1 for point in data["dataPoints"]: seen.append(point) case = int(point["timeUnixNano"]) - 2_000 @@ -202,17 +221,27 @@ def test_remote_metric_maps_have_one_scalar_per_retained_key(tmp_path, receiver) expected["operation"] = {"stringValue": "encode"} if receiver == "gateway": expected["http.status_code"] = {"intValue": "200"} + if metric["name"].startswith("sie.gateway.remote."): + expected.pop("http.status_code", None) + expected["model"] = {"stringValue": "acme/model"} + if kind == "gauge": + expected = {"model": expected["model"]} + else: + expected["outcome"] = {"stringValue": "committed"} + expected["fallback.reason"] = {"stringValue": "provisioning"} assert len(point["attributes"]) == len(expected) assert {a["key"]: a["value"] for a in point["attributes"]} == expected assert int(point["startTimeUnixNano"]) == 1_000 if kind == "sum": assert int(point["asInt"]) == 10 + case + elif kind == "gauge": + assert point["asDouble"] == 10.0 + case else: assert int(point["count"]) == 2 assert point["sum"] == 0.5 + case assert point["bucketCounts"] == ["1", "1"] assert point["explicitBounds"] == [1] - assert len(seen) == 6 + assert len(seen) == (12 if receiver == "gateway" else 6) finally: try: run("docker", "stop", "--time", "10", container) diff --git a/tools/ci/tests/test_helm_render.py b/tools/ci/tests/test_helm_render.py index e097b525a..64f9192fe 100644 --- a/tools/ci/tests/test_helm_render.py +++ b/tools/ci/tests/test_helm_render.py @@ -2547,3 +2547,41 @@ def test_without_a_remote_lane_sie_config_receives_no_upstream_names(tmp_path: P docs = rendered_documents(tmp_path, {"upstreams": upstreams_fixture()["values"], **L4_POOL}) assert container_env(docs, *CONFIG_SERVICE.split("/"))["SIE_UPSTREAM_NAMES"]["value"] == "" + + +@pytest.mark.parametrize("seconds", [1, 600, 86400]) +def test_remote_fallback_persistence_alert_uses_configured_threshold_and_replica_freshness( + tmp_path: Path, seconds: int +) -> None: + values = { + "alertRules": {"enabled": True, "remoteFallbackPersistenceSeconds": seconds}, + "observability": AUTOSCALING_VALUES["observability"], + } + result = render_template(tmp_path, values, "templates/prometheusrule.yaml") + assert result.returncode == 0, result.stderr + document = yaml.safe_load(result.stdout) + rule = next( + rule + for group in document["spec"]["groups"] + for rule in group["rules"] + if rule["alert"] == "SIERemoteFallbackPersistent" + ) + assert rule["expr"].strip().endswith(f"> {seconds}") + assert "and on (producer_instance, collector_generation, model)" in rule["expr"] + assert 'outcome="committed"' in rule["expr"] + assert "[5m]" in rule["expr"] + assert "__REMOTE_FALLBACK" not in result.stdout + + +@pytest.mark.parametrize("seconds", [0, -1, 86401, 1.5, "600", True]) +def test_remote_fallback_persistence_alert_refuses_invalid_thresholds(tmp_path: Path, seconds: object) -> None: + result = render_template( + tmp_path, + { + "alertRules": {"enabled": True, "remoteFallbackPersistenceSeconds": seconds}, + "observability": AUTOSCALING_VALUES["observability"], + }, + "templates/prometheusrule.yaml", + ) + assert result.returncode != 0 + assert "alertRules.remoteFallbackPersistenceSeconds" in result.stderr