use std::time::{Duration, Instant};
use chrono::Utc;
use serde_json::Value;
use tap::{Pipe, Tap, TapFallible};
use crate::error::GatewayError;
use crate::gateway::{Gateway, RequestCtx};
use crate::jev_native::JevNativeProvider;
use crate::ledger::UsageEntry;
use crate::routing::classify::vertex_triggers;
use crate::routing::jev_router::{
self, Answers, DecisionOutcome, EffortPolicy, RoutePlan, RoutingMode, Selection,
};
use crate::routing::request::ChatRequest;
use crate::routing::table::JevRoute;
struct Decision {
answers: Answers,
input_tokens: u64,
output_tokens: u64,
}
impl Gateway {
pub(crate) async fn plan_route(
&self,
req: &ChatRequest,
ctx: &RequestCtx,
request_id: &str,
) -> Result<RoutePlan, GatewayError> {
let legs = self
.routes
.legs(&req.model)
.ok_or_else(|| GatewayError::UnknownModel(req.model.clone()))?;
let jev_route = self.routes.jev_route(&req.model);
let mode = jev_router::resolve_mode(jev_route.is_some(), req.routing_strategy.as_deref())
.map_err(GatewayError::BadRequest)?;
self.guard_input(req)?;
match (mode, jev_route) {
(RoutingMode::Jev | RoutingMode::StaticOverride, Some(route)) => {
self.plan_tiered(req, ctx, request_id, route, mode).await
}
_ => Ok(RoutePlan::static_legs(legs, client_sets_effort(req))),
}
}
async fn plan_tiered(
&self,
req: &ChatRequest,
ctx: &RequestCtx,
request_id: &str,
route: &JevRoute,
mode: RoutingMode,
) -> Result<RoutePlan, GatewayError> {
reject_jev_block(req)?;
let tiers = jev_router::eligible_tiers(&route.tiers, vertex_triggers(req));
jev_router::nearest_serving(&tiers, 0).ok_or_else(|| no_eligible_legs(req))?;
let (selection, answers) = match mode {
RoutingMode::StaticOverride => (
Selection {
tier: route.default_index(),
outcome: DecisionOutcome::StaticOverride,
bump: false,
},
None,
),
_ => self.decide(req, ctx, request_id, route).await,
};
let start = jev_router::nearest_serving(&tiers, selection.tier)
.ok_or_else(|| no_eligible_legs(req))?;
let policy = match client_sets_effort(req) {
true => EffortPolicy::Client,
false => EffortPolicy::Tier {
bump: selection.bump,
},
};
let decided = route.tiers[selection.tier].name.as_str();
self.metrics
.routing_decision(&req.model, decided, selection.outcome.metric_label());
RoutePlan {
mode,
legs: jev_router::order_legs(&tiers, start, policy),
tier_names: route.tiers.iter().map(|t| t.name.clone()).collect(),
decided: Some(selection.tier),
outcome: Some(selection.outcome),
client_effort: matches!(policy, EffortPolicy::Client),
}
.tap(|plan| {
tracing::info!(
target: "synapse::routing",
route = %req.model,
routing.mode = mode.as_str(),
routing.tier_decided = decided,
routing.outcome = selection.outcome.metric_label(),
routing.effort = plan.planned_effort(),
routing.difficulty_score = ?answers.map(|a| a.difficulty),
routing.confidence = ?answers.map(|a| a.confidence),
routing.needs_reasoning = ?answers.and_then(|a| a.needs_reasoning),
"route planned"
)
})
.pipe(Ok)
}
async fn decide(
&self,
req: &ChatRequest,
ctx: &RequestCtx,
request_id: &str,
route: &JevRoute,
) -> (Selection, Option<Answers>) {
let result = match &self.jev_native {
None => Err(DecisionOutcome::Unavailable),
Some(provider) => {
let started = Instant::now();
let body = serde_json::json!({
"model": route.router.model,
"state": jev_router::build_state(req),
"questions": jev_router::build_questions(&route.tiers),
});
tokio::time::timeout(
Duration::from_millis(route.router.timeout_ms),
call_jev(provider, body, &req.model),
)
.await
.tap_err(|_| {
tracing::warn!(
target: "synapse::routing",
route = %req.model,
error.kind = "timeout",
timeout_ms = route.router.timeout_ms,
"jev decision timed out; routing to default_tier"
)
})
.unwrap_or(Err(DecisionOutcome::Timeout))
.tap(|_| {
self.metrics
.routing_decision_duration(&req.model, started.elapsed().as_secs_f64())
})
}
}
.tap_ok(|d| {
self.record_route_decision(ctx, &req.model, request_id, &route.router.model, d)
});
let answers = result.as_ref().ok().map(|d| d.answers);
(
jev_router::select(result.map(|d| d.answers), route),
answers,
)
}
fn record_route_decision(
&self,
ctx: &RequestCtx,
route: &str,
request_id: &str,
model: &str,
d: &Decision,
) {
let attr = self.attribution_of(ctx, route);
self.ledger.enqueue(UsageEntry {
ts: Utc::now(),
tenant: self.tenant_of(ctx).to_string(),
workspace: attr.workspace,
user: attr.user,
thread: attr.thread,
message: attr.message,
route: route.to_string(),
provider: "typesafe".into(),
model: model.to_string(),
lane: "jev".into(),
input_tokens: d.input_tokens,
output_tokens: d.output_tokens,
cost_usd: self
.pricing
.cost_usd("typesafe", model, d.input_tokens, d.output_tokens),
request_id: request_id.to_string(),
status: "ok".into(),
op: "route_decision".into(),
user_task_type: attr.user_task_type,
ai_task_type: attr.ai_task_type,
});
}
}
async fn call_jev(
provider: &JevNativeProvider,
body: Value,
route: &str,
) -> Result<Decision, DecisionOutcome> {
let warn = |kind: &str, status: Option<u16>| {
tracing::warn!(
target: "synapse::routing",
route,
error.kind = kind,
http.status = ?status,
"jev decision failed; routing to default_tier"
)
};
let resp = provider
.evaluate(body)
.await
.tap_err(|_| warn("transport", None))
.map_err(|_| DecisionOutcome::Error)?
.error_for_status()
.tap_err(|e| warn("http_status", e.status().map(|s| s.as_u16())))
.map_err(|_| DecisionOutcome::Error)?;
let status = resp.status().as_u16();
let value: Value = resp
.json()
.await
.tap_err(|_| warn("unreadable_body", Some(status)))
.map_err(|_| DecisionOutcome::Error)?;
jev_router::parse_answers(&value["answers"])
.map(|answers| Decision {
answers,
input_tokens: value["usage"]["input_tokens"].as_u64().unwrap_or(0),
output_tokens: value["usage"]["output_tokens"].as_u64().unwrap_or(0),
})
.ok_or(DecisionOutcome::Error)
.tap_err(|_| warn("unparseable_answers", Some(status)))
}
fn client_sets_effort(req: &ChatRequest) -> bool {
match vertex_triggers(req) {
true => req
.vertex
.as_ref()
.is_some_and(|v| v.thinking_config.is_some()),
false => crate::routing::executor::client_effort(req).is_some(),
}
}
fn reject_jev_block(req: &ChatRequest) -> Result<(), GatewayError> {
req.jev
.as_ref()
.is_some_and(|j| !j.questions.is_empty() || j.extract.is_some())
.pipe(|has_block| match has_block {
true => Err(GatewayError::BadRequest(
"the jev extension is not supported on jev-routed routes".into(),
)),
false => Ok(()),
})
}
fn no_eligible_legs(req: &ChatRequest) -> GatewayError {
GatewayError::BadRequest(format!(
"no legs of route '{}' can serve this request (native Vertex features need a vertex leg)",
req.model
))
}