use std::{
collections::HashSet,
future::{Future, ready},
sync::Arc,
time::Duration,
};
use dynamo_kv_router::{
protocols::{TokensWithHashes, WorkerConfigLike, WorkerWithDpRank},
selector::{WorkerInputs, WorkerSelector},
};
use dynamo_runtime::{
error::{DynamoError, ErrorType, match_error_chain},
metrics::frontend_perf::{STAGE_ROUTE, StageGuard},
pipeline::{
AsyncEngine, AsyncEngineContext, AsyncEngineContextProvider, Error, ManyOut, PushRouter,
ResponseStream, RouterMode, SingleIn, async_trait,
network::egress::route_span::{
get_route_trace_context, record_route_error, record_route_span_start, wrap_route_span,
},
},
protocols::annotated::Annotated,
};
use futures::stream::{self, StreamExt};
use tracing::Instrument;
use crate::{
kv_router::{
KvRouter, metrics::RouterRequestMetrics, scheduler::DefaultWorkerSelector,
to_worker_selection_session_context,
},
local_model::runtime_config::ModelRuntimeConfig,
lora::{LoadEstimator, LoraFilter},
preprocessor::PreprocessedRequest,
protocols::common::{
FinishReason,
extensions::SessionAffinityId,
llm_backend::LLMEngineOutput,
timing::{RequestPhase, RoutingData, WORKER_TYPE_DECODE, WORKER_TYPE_PREFILL},
},
session_affinity::{
AffinityAcquire, AffinityCoordinator, AffinityTarget, SessionAffinityMode, affinity_id,
explicit_target, invalid_argument,
},
};
mod builtin;
mod cancellation;
mod kv;
mod kv_selection;
mod occupancy;
mod request_guard;
use builtin::BuiltinWorkerSelector;
use cancellation::{CleanupBudget, DispatchCancellation, StagedKv, await_with_cleanup_policy};
use kv_selection::{RoutingRequestParts, SelectionOptions, WorkerSelection};
use occupancy::HostedOccupancy;
pub(crate) use request_guard::prompt_private_blocks;
use request_guard::{KvRequestCleanup, LoraLoadGuard, RequestGuard};
const OUTPUT_REPLAY_ID_ANNOTATION_KEY: &str = "output_replay_id";
const OUTPUT_REPLAY_CONSUMER_RUNTIME_KEY: &str = "output_replay_consumer";
const DRAIN_TIMEOUT: Duration = Duration::from_secs(5);
pub(crate) fn is_cancelled(error: &Error) -> bool {
match_error_chain(error.as_ref(), &[ErrorType::Cancelled], &[])
}
fn route_target(worker: WorkerWithDpRank) -> AffinityTarget {
AffinityTarget::new(worker.worker_id, Some(worker.dp_rank))
}
fn monitor_response_stream<Sel>(
mut response_stream: ManyOut<Annotated<LLMEngineOutput>>,
context: Arc<dyn AsyncEngineContext>,
mut guard: RequestGuard<Sel>,
) -> impl futures::Stream<Item = Annotated<LLMEngineOutput>> + Send
where
Sel: WorkerSelector<ModelRuntimeConfig> + Send + 'static,
{
async_stream::stream! {
let stopped = context.stopped();
tokio::pin!(stopped);
let mut drainable_terminal = false;
let mut pending_terminal: Option<Annotated<LLMEngineOutput>> = None;
let drain_deadline = tokio::time::sleep(Duration::ZERO);
tokio::pin!(drain_deadline);
let completed = loop {
tokio::select! {
biased;
_ = &mut stopped => {
tracing::debug!(request_id = context.id(), "Request cancelled, ending stream");
drop(pending_terminal.take());
break false;
}
item = response_stream.next() => {
let Some(item) = item else {
if drainable_terminal {
guard.record_migration_failure(None);
}
break !drainable_terminal;
};
let outcome = classify_response_item(&item);
guard.on_item(&item).await;
match outcome {
ResponseItemOutcome::Failed => {
drop(pending_terminal.take());
guard.record_migration_failure(item.error.clone());
guard.abort().await;
yield item;
break false;
}
ResponseItemOutcome::DrainableTerminal => {
if !drainable_terminal {
drainable_terminal = true;
drain_deadline.as_mut().reset(tokio::time::Instant::now() + DRAIN_TIMEOUT);
}
if let Some(previous) = pending_terminal.replace(item) {
yield previous;
}
if tokio::time::Instant::now() >= drain_deadline.deadline() {
guard.record_migration_failure(None);
break false;
}
}
ResponseItemOutcome::Healthy => {
drainable_terminal = false;
if let Some(previous) = pending_terminal.take() {
yield previous;
}
yield item;
}
}
}
_ = &mut drain_deadline, if drainable_terminal => {
tracing::debug!(
request_id = context.id(),
"Terminal frame was not followed by an error within {DRAIN_TIMEOUT:?}, ending stream"
);
guard.record_migration_failure(None);
break false;
}
}
};
if completed {
guard.finish().await;
} else {
guard.abort().await;
}
if let Some(pending) = pending_terminal.take() {
yield pending;
}
}
}
fn into_monitored_response<Sel>(
response_stream: ManyOut<Annotated<LLMEngineOutput>>,
guard: RequestGuard<Sel>,
) -> ManyOut<Annotated<LLMEngineOutput>>
where
Sel: WorkerSelector<ModelRuntimeConfig> + Send + 'static,
{
let stream_context = response_stream.context();
let wrapped_stream = Box::pin(monitor_response_stream(
response_stream,
stream_context.clone(),
guard,
));
ResponseStream::new(wrapped_stream, stream_context)
}
enum RoutingPolicy<Sel>
where
Sel: WorkerSelector<ModelRuntimeConfig> + Send + 'static,
{
Kv(Arc<KvRouter<Sel>>),
Builtin(BuiltinWorkerSelector),
Direct,
DeviceAwareWeighted,
}
struct LoraRouting {
filter: Arc<LoraFilter>,
load_estimator: Arc<LoadEstimator>,
selector: BuiltinWorkerSelector,
}
struct LoraSelection {
target: u64,
allowed_fallback: HashSet<u64>,
load_guard: LoraLoadGuard,
}
struct HostedSelection {
initial_worker: u64,
target_constraint: Option<AffinityTarget>,
occupancy_reservation: Option<dynamo_runtime::pipeline::OccupancyReservation>,
candidate_count: usize,
selected_occupancy: Option<u64>,
device_aware_telemetry: Option<DeviceAwareTelemetry>,
}
struct DeviceAwareTelemetry {
is_cpu: bool,
embedding_cache_hit: bool,
request_cache_keys: usize,
}
pub struct RoutingHost<Sel = DefaultWorkerSelector>
where
Sel: WorkerSelector<ModelRuntimeConfig> + Send + 'static,
{
inner: PushRouter<PreprocessedRequest, Annotated<LLMEngineOutput>>,
policy: RoutingPolicy<Sel>,
request_metrics: Arc<RouterRequestMetrics>,
affinity: Option<AffinityCoordinator>,
session_affinity_mode: SessionAffinityMode,
hosted_occupancy: Option<HostedOccupancy>,
lora: Option<LoraRouting>,
#[allow(dead_code)]
routing_context: Option<Arc<crate::kv_router::RoutingLoadContext>>,
}
pub(crate) struct RoutePlan<Sel = DefaultWorkerSelector>
where
Sel: WorkerSelector<ModelRuntimeConfig> + Send + 'static,
{
signals: RoutePlanSignals,
selection: WorkerSelection,
cleanup: KvRequestCleanup<Sel>,
affinity: Option<AffinityAcquire>,
budget: CleanupBudget,
}
pub(crate) struct RoutePreview {
request_id: String,
phase: RequestPhase,
signals: RoutePlanSignals,
budget: CleanupBudget,
}
#[derive(Clone, Copy, Debug)]
pub(crate) struct RoutePlanSignals {
pub(crate) worker: WorkerWithDpRank,
pub(crate) overlap_blocks: u32,
pub(crate) cached_tokens: usize,
pub(crate) potential_decode_blocks: u64,
pub(crate) total_kv_blocks: Option<u64>,
}
impl RoutePreview {
pub(crate) fn signals(&self) -> RoutePlanSignals {
self.signals
}
#[cfg(test)]
pub(crate) fn cleanup_budget_remaining(&self) -> std::time::Duration {
self.budget.remaining()
}
}
impl RoutePlanSignals {
pub(crate) fn decode_load_exceeds(self, threshold: f64) -> Option<bool> {
let total_kv_blocks = self.total_kv_blocks?;
Some(self.potential_decode_blocks as f64 > threshold * total_kv_blocks as f64)
}
}
impl<Sel> RoutePlan<Sel>
where
Sel: WorkerSelector<ModelRuntimeConfig> + Send + 'static,
{
pub(crate) fn signals(&self) -> RoutePlanSignals {
self.signals
}
#[cfg(test)]
pub(crate) fn cleanup_budget_remaining(&self) -> std::time::Duration {
self.budget.remaining()
}
#[cfg(test)]
pub(crate) async fn abort(self) {
self.cleanup.finish().await;
}
}
pub type KvPushRouter<Sel = DefaultWorkerSelector> = RoutingHost<Sel>;
impl<Sel> RoutingHost<Sel>
where
Sel: WorkerSelector<ModelRuntimeConfig> + Send + 'static,
{
pub fn new(
inner: PushRouter<PreprocessedRequest, Annotated<LLMEngineOutput>>,
kv_router: Arc<KvRouter<Sel>>,
session_affinity_ttl: Option<Duration>,
) -> Result<Self, Error> {
let affinity = session_affinity_ttl
.map(AffinityCoordinator::new)
.transpose()?;
Ok(Self::new_with_coordinator(
inner,
kv_router,
affinity,
SessionAffinityMode::Hard,
))
}
pub fn new_with_load_context(
inner: PushRouter<PreprocessedRequest, Annotated<LLMEngineOutput>>,
kv_router: Arc<KvRouter<Sel>>,
load_context: Arc<crate::kv_router::RoutingLoadContext>,
session_affinity_ttl: Option<Duration>,
session_affinity_mode: SessionAffinityMode,
) -> Result<Self, Error> {
let affinity = session_affinity_ttl
.map(AffinityCoordinator::new)
.transpose()?;
Ok(Self::new_with_load_context_and_coordinator(
inner,
kv_router,
load_context,
affinity,
session_affinity_mode,
))
}
pub(crate) fn new_with_coordinator(
inner: PushRouter<PreprocessedRequest, Annotated<LLMEngineOutput>>,
kv_router: Arc<KvRouter<Sel>>,
affinity: Option<AffinityCoordinator>,
session_affinity_mode: SessionAffinityMode,
) -> Self {
Self::new_with_optional_load_context_and_coordinator(
inner,
kv_router,
None,
affinity,
session_affinity_mode,
)
}
pub(crate) fn new_with_load_context_and_coordinator(
inner: PushRouter<PreprocessedRequest, Annotated<LLMEngineOutput>>,
kv_router: Arc<KvRouter<Sel>>,
load_context: Arc<crate::kv_router::RoutingLoadContext>,
affinity: Option<AffinityCoordinator>,
session_affinity_mode: SessionAffinityMode,
) -> Self {
Self::new_with_optional_load_context_and_coordinator(
inner,
kv_router,
Some(load_context),
affinity,
session_affinity_mode,
)
}
fn new_with_optional_load_context_and_coordinator(
inner: PushRouter<PreprocessedRequest, Annotated<LLMEngineOutput>>,
kv_router: Arc<KvRouter<Sel>>,
load_context: Option<Arc<crate::kv_router::RoutingLoadContext>>,
affinity: Option<AffinityCoordinator>,
session_affinity_mode: SessionAffinityMode,
) -> Self {
let request_metrics =
RouterRequestMetrics::from_component(kv_router.client().endpoint.component());
RoutingHost {
inner,
policy: RoutingPolicy::Kv(kv_router),
request_metrics,
affinity,
session_affinity_mode,
hosted_occupancy: None,
lora: None,
routing_context: load_context,
}
}
#[cfg(test)]
pub(crate) fn new_builtin(
inner: PushRouter<PreprocessedRequest, Annotated<LLMEngineOutput>>,
load_context: Arc<crate::kv_router::RoutingLoadContext>,
) -> Result<Self, Error> {
Self::new_builtin_with_capabilities(
inner,
load_context,
None,
SessionAffinityMode::Hard,
None,
)
}
pub(crate) fn new_builtin_with_coordinator(
inner: PushRouter<PreprocessedRequest, Annotated<LLMEngineOutput>>,
load_context: Arc<crate::kv_router::RoutingLoadContext>,
affinity: Option<AffinityCoordinator>,
session_affinity_mode: SessionAffinityMode,
) -> Result<Self, Error> {
Self::new_builtin_with_capabilities(
inner,
load_context,
affinity,
session_affinity_mode,
None,
)
}
pub(crate) fn new_builtin_with_capabilities(
inner: PushRouter<PreprocessedRequest, Annotated<LLMEngineOutput>>,
load_context: Arc<crate::kv_router::RoutingLoadContext>,
affinity: Option<AffinityCoordinator>,
session_affinity_mode: SessionAffinityMode,
lora: Option<(Arc<LoraFilter>, Arc<LoadEstimator>)>,
) -> Result<Self, Error> {
if affinity.is_some() && lora.is_some() {
anyhow::bail!("session affinity and LoRA filtering cannot both be enabled");
}
let policy = match inner.router_mode() {
RouterMode::Direct => RoutingPolicy::Direct,
RouterMode::DeviceAwareWeighted => RoutingPolicy::DeviceAwareWeighted,
mode => {
RoutingPolicy::Builtin(BuiltinWorkerSelector::new(mode).ok_or_else(|| {
anyhow::anyhow!("{mode:?} routing is not a first-party policy")
})?)
}
};
let required_worker_inputs = match &policy {
RoutingPolicy::Builtin(selector) => selector.required_worker_inputs(),
RoutingPolicy::DeviceAwareWeighted => WorkerInputs::OCCUPANCY,
RoutingPolicy::Direct => WorkerInputs::NONE,
RoutingPolicy::Kv(_) => unreachable!(),
};
let hosted_occupancy = matches!(&policy, RoutingPolicy::Builtin(_))
.then_some(required_worker_inputs.contains(WorkerInputs::OCCUPANCY))
.unwrap_or(false)
.then(|| HostedOccupancy::new(&inner))
.transpose()?;
if lora.is_some()
&& !matches!(
inner.router_mode(),
RouterMode::RoundRobin | RouterMode::Random
)
{
anyhow::bail!(
"LoRA filtering is unsupported with {:?} routing",
inner.router_mode()
);
}
let lora_selector = lora.as_ref().map(|_| {
BuiltinWorkerSelector::new(inner.router_mode())
.expect("LoRA routing mode was validated above")
});
let request_metrics =
RouterRequestMetrics::from_component(inner.client.endpoint.component());
Ok(Self {
inner,
policy,
request_metrics,
affinity,
session_affinity_mode,
hosted_occupancy,
lora: lora
.zip(lora_selector)
.map(|((filter, load_estimator), selector)| LoraRouting {
filter,
load_estimator,
selector,
}),
routing_context: Some(load_context),
})
}
pub fn required_worker_inputs(&self) -> WorkerInputs {
match &self.policy {
RoutingPolicy::Kv(chooser) => chooser.required_worker_inputs(),
RoutingPolicy::Builtin(selector) => selector.required_worker_inputs(),
RoutingPolicy::Direct => WorkerInputs::NONE,
RoutingPolicy::DeviceAwareWeighted => WorkerInputs::OCCUPANCY,
}
}
#[cfg(test)]
pub(crate) fn occupancy_for_test(&self, worker_id: u64) -> u64 {
self.inner.occupancy_for_test(worker_id)
}
pub fn kv_router(&self) -> &Arc<KvRouter<Sel>> {
self.kv_router_if_enabled()
.expect("routing host has no KV capability")
}
pub(crate) fn kv_router_if_enabled(&self) -> Option<&Arc<KvRouter<Sel>>> {
match &self.policy {
RoutingPolicy::Kv(chooser) => Some(chooser),
RoutingPolicy::Builtin(_)
| RoutingPolicy::Direct
| RoutingPolicy::DeviceAwareWeighted => None,
}
}
pub(crate) fn peek_next_worker(&self) -> Option<u64> {
match &self.policy {
RoutingPolicy::Builtin(selector) => match &self.hosted_occupancy {
Some(occupancy) => occupancy.peek(&self.inner, selector),
None => self
.inner
.with_selectable_worker_ids(|ids| {
selector.peek_worker(
dynamo_kv_router::selector::WorkerSelectionInput::hosted(ids, None),
)
})
.ok()
.and_then(Result::ok),
},
RoutingPolicy::DeviceAwareWeighted => self.inner.peek_next_worker(),
RoutingPolicy::Direct => None,
RoutingPolicy::Kv(_) => None,
}
}
fn validate_explicit_worker(
&self,
request: &PreprocessedRequest,
phase: RequestPhase,
) -> Result<(), Error> {
let Some(routing) = request.routing.as_ref() else {
return Ok(());
};
let target = match phase {
RequestPhase::Prefill => routing
.prefill_worker_id
.map(|id| (id, "prefill_worker_id")),
RequestPhase::Decode | RequestPhase::Aggregated => {
routing.decode_worker_id.map(|id| (id, "decode_worker_id"))
}
}
.or_else(|| {
routing
.backend_instance_id
.map(|id| (id, "backend_instance_id"))
});
if let Some((worker_id, field)) = target
&& !self.inner.client.is_instance_discovered(worker_id)
{
return Err(invalid_argument(format!(
"nvext.{field}={worker_id} does not identify a known worker"
)));
}
Ok(())
}
fn affinity_target_is_valid(&self, target: AffinityTarget) -> bool {
if !self.inner.client.is_instance_discovered(target.worker_id) {
return false;
}
let Some(kv_router) = self.kv_router_if_enabled() else {
return true;
};
let workers = kv_router.workers_with_configs.borrow();
let Some(config) = workers.get(&target.worker_id) else {
return true;
};
let Some(dp_rank) = target.dp_rank else {
return true;
};
let start = config.data_parallel_start_rank();
let end = start.saturating_add(config.data_parallel_size());
(start..end).contains(&dp_rank)
}
#[allow(clippy::too_many_arguments)]
async fn acquire_affinity_slot(
&self,
affinity: &AffinityCoordinator,
session_id: &SessionAffinityId,
requested_target: Option<AffinityTarget>,
context: &dyn AsyncEngineContext,
phase: RequestPhase,
staged_kv: StagedKv,
budget: &CleanupBudget,
) -> Result<AffinityAcquire, Error> {
match DispatchCancellation::for_request(phase, staged_kv) {
DispatchCancellation::CancelWhenStopped => {
affinity
.acquire_with_context(session_id, requested_target, context)
.await
}
DispatchCancellation::DispatchWhenStopped => await_with_cleanup_policy(
context,
phase,
staged_kv,
"affinity.acquire",
budget,
affinity.acquire(session_id, requested_target),
)
.await
.and_then(|result| result),
}
}
async fn select_with_session_affinity<T, Select, SelectionFuture>(
&self,
request: &SingleIn<PreprocessedRequest>,
phase: RequestPhase,
is_query_only: bool,
budget: &CleanupBudget,
mut select: Select,
) -> Result<(T, Option<AffinityAcquire>), Error>
where
Select: FnMut(Option<AffinityTarget>) -> SelectionFuture,
SelectionFuture: Future<Output = Result<T, Error>>,
{
let staged_kv = StagedKv::for_request(request.content());
let Some(affinity) = self.affinity.as_ref() else {
return Ok((select(None).await?, None));
};
let Some(session_id) = affinity_id(request)? else {
return Ok((select(None).await?, None));
};
let explicit = explicit_target(request.content(), phase)?;
if is_query_only {
let target = affinity.query_target(&session_id, explicit)?;
return Ok((select(target).await?, None));
}
let request_context = request.context();
let operation = self
.acquire_affinity_slot(
affinity,
&session_id,
explicit,
request_context.as_ref(),
phase,
staged_kv,
budget,
)
.await?;
let target = operation.target();
match select(target).await {
Ok(selection) => Ok((selection, Some(operation))),
Err(error) if is_cancelled(&error) => Err(error),
Err(_error)
if self.session_affinity_mode == SessionAffinityMode::Hard
&& explicit.is_none()
&& target.is_some_and(|target| !self.affinity_target_is_valid(target)) =>
{
operation.invalidate();
let retry = self
.acquire_affinity_slot(
affinity,
&session_id,
None,
request_context.as_ref(),
phase,
staged_kv,
budget,
)
.await?;
let selection = select(retry.target()).await?;
Ok((selection, Some(retry)))
}
Err(error) => Err(error),
}
}
pub(crate) async fn select_and_dispatch_prefill<M, F>(
&self,
request: SingleIn<PreprocessedRequest>,
prepare: F,
) -> Result<(M, ManyOut<Annotated<LLMEngineOutput>>), Error>
where
F: FnOnce(&mut PreprocessedRequest, AffinityTarget) -> Result<M, Error>,
{
match &self.policy {
RoutingPolicy::Kv(_) => self.select_and_dispatch_kv_prefill(request, prepare).await,
RoutingPolicy::Builtin(_)
| RoutingPolicy::Direct
| RoutingPolicy::DeviceAwareWeighted => {
self.select_and_dispatch_builtin(request, RequestPhase::Prefill, prepare)
.await
}
}
}
}
#[async_trait]
impl<Sel> AsyncEngine<SingleIn<PreprocessedRequest>, ManyOut<Annotated<LLMEngineOutput>>, Error>
for RoutingHost<Sel>
where
Sel: WorkerSelector<ModelRuntimeConfig> + Send + 'static,
{
async fn generate(
&self,
request: SingleIn<PreprocessedRequest>,
) -> Result<ManyOut<Annotated<LLMEngineOutput>>, Error> {
let budget = CleanupBudget::default();
if !matches!(&self.policy, RoutingPolicy::Kv(_)) {
let phase = request
.tracker
.as_ref()
.map(|tracker| tracker.phase())
.unwrap_or(RequestPhase::Aggregated);
return self
.select_and_dispatch_builtin(request, phase, |_, _| Ok(()))
.await
.map(|(_, stream)| stream);
}
let is_query_only = request.get_annotation_value("query_instance_id").is_some();
let phase = request
.tracker
.as_ref()
.map(|tracker| tracker.phase())
.unwrap_or(RequestPhase::Aggregated);
let phase_label = phase.to_string();
let route_guard = StageGuard::new(STAGE_ROUTE, &phase_label);
let (mut selection, mut operation) = self
.select_with_affinity(&request, phase, is_query_only, &budget)
.await?;
if is_query_only {
let routing_parts = RoutingRequestParts::new(&request);
if let Some(ref tracker) = request.tracker {
let isl_blocks = routing_parts
.token_ids
.len()
.div_ceil(self.kv_router().block_size() as usize);
tracker.record_kv_hit(selection.effective_overlap_blocks, isl_blocks);
tracker.record_isl(routing_parts.token_ids.len(), Some(selection.cached_tokens));
tracker.record_worker(
selection.worker.worker_id,
Some(selection.worker.dp_rank),
self.kv_router().worker_type(),
);
tracker.record_router_queue_depth(self.kv_router().pending_count());
}
self.request_metrics
.input_sequence_tokens
.observe(request.token_ids.len() as f64);
let stream_context = request.context().clone();
let worker_id_info = request
.tracker
.as_ref()
.and_then(|tracker| tracker.get_worker_info());
tracing::trace!(
?phase,
worker_id = selection.worker.worker_id,
?worker_id_info,
"Returning worker selection (query-only mode)"
);
let output = LLMEngineOutput {
routing_data: Some(RoutingData {
worker_id: worker_id_info,
token_ids: Some(request.token_ids.clone()),
..Default::default()
}),
..Default::default()
};
let response = Annotated::from_data(output);
let stream = stream::iter(vec![response]);
return Ok(ResponseStream::new(Box::pin(stream), stream_context));
}
let guard = match self
.track_selection(&request, &mut selection, phase, false, &budget)
.await
{
Ok(guard) => guard,
Err(error) => return Err(error),
};
drop(route_guard);
let selected_target = route_target(selection.worker);
let stream = match self
.dispatch_selection(request, selection, guard, &budget)
.await
{
Ok(stream) => stream,
Err(error) => {
if self.session_affinity_mode == SessionAffinityMode::Hard
&& !self.affinity_target_is_valid(selected_target)
&& let Some(operation) = operation.take()
{
operation.invalidate();
}
return Err(error);
}
};
match operation {
Some(operation) => {
operation.into_stream(selected_target, stream, self.session_affinity_mode)
}
None => Ok(stream),
}
}
}
enum ResponseItemOutcome {
Healthy,
DrainableTerminal,
Failed,
}
fn classify_response_item(item: &Annotated<LLMEngineOutput>) -> ResponseItemOutcome {
if item.error.is_some() || item.event.as_deref() == Some("error") {
return ResponseItemOutcome::Failed;
}
let terminal = item
.data
.as_ref()
.and_then(|data| data.finish_reason.as_ref())
.is_some_and(|reason| matches!(reason, FinishReason::Error(_) | FinishReason::Cancelled));
if terminal {
ResponseItemOutcome::DrainableTerminal
} else {
ResponseItemOutcome::Healthy
}
}
#[cfg(test)]
mod tests;