use openkind_core::{Answer, ResponseContract, SystemRequest, SystemResponse};
use tracing::instrument;
use crate::error::{EngineError, EngineResult};
use crate::registry::EngineRegistry;
#[instrument(skip(req, registry), fields(model = %req.model, n_questions = req.questions.len()))]
pub async fn dispatch(
req: SystemRequest,
registry: &EngineRegistry,
) -> EngineResult<SystemResponse> {
metrics::counter!("openkind_requests_total").increment(1);
let start = std::time::Instant::now();
let requested_model = req.model.clone();
let engine = registry
.get(&req.model)
.ok_or_else(|| EngineError::UnknownModel(req.model.clone()))?;
let contract = ResponseContract::from_request(&req)?;
let input_tokens = engine.estimate_input_tokens(&req);
let mut resp = engine.evaluate(req).await?;
if resp.model != requested_model {
return Err(EngineError::Backend {
backend: engine.backend_id().to_string(),
message: format!(
"backend returned model `{}`, expected registered alias `{requested_model}`",
resp.model
),
});
}
if let Err(validation) = contract.validate(&resp) {
return Err(EngineError::Backend {
backend: engine.backend_id().to_string(),
message: format!("backend returned an invalid response: {validation}"),
});
}
if resp.usage.input_tokens == 0 {
resp.usage.input_tokens = input_tokens;
}
if resp.usage.output_tokens == 0 {
resp.usage.output_tokens = estimate_output_tokens(&resp);
}
let elapsed_ms = start.elapsed().as_millis() as f64;
metrics::histogram!("openkind_request_duration_ms").record(elapsed_ms);
metrics::counter!("openkind_responses_total").increment(1);
Ok(resp)
}
fn estimate_output_tokens(resp: &SystemResponse) -> u32 {
let sum: usize = resp
.answers
.values()
.map(|a| match a {
Answer::Noul(_) => 1,
Answer::Choice(_) => 1,
Answer::Score(_) => 4,
})
.sum();
u32::try_from(sum).unwrap_or(u32::MAX)
}