openkind_engine/
dispatch.rs1use openkind_core::{Answer, ResponseContract, SystemRequest, SystemResponse};
4use tracing::instrument;
5
6use crate::error::{EngineError, EngineResult};
7use crate::registry::EngineRegistry;
8
9#[instrument(skip(req, registry), fields(model = %req.model, n_questions = req.questions.len()))]
12pub async fn dispatch(
13 req: SystemRequest,
14 registry: &EngineRegistry,
15) -> EngineResult<SystemResponse> {
16 metrics::counter!("openkind_requests_total").increment(1);
17
18 let start = std::time::Instant::now();
19 let requested_model = req.model.clone();
20 let engine = registry
21 .get(&req.model)
22 .ok_or_else(|| EngineError::UnknownModel(req.model.clone()))?;
23
24 let contract = ResponseContract::from_request(&req)?;
25 let input_tokens = engine.estimate_input_tokens(&req);
26
27 let mut resp = engine.evaluate(req).await?;
28
29 if resp.model != requested_model {
30 return Err(EngineError::Backend {
31 backend: engine.backend_id().to_string(),
32 message: format!(
33 "backend returned model `{}`, expected registered alias `{requested_model}`",
34 resp.model
35 ),
36 });
37 }
38
39 if let Err(validation) = contract.validate(&resp) {
43 return Err(EngineError::Backend {
44 backend: engine.backend_id().to_string(),
45 message: format!("backend returned an invalid response: {validation}"),
46 });
47 }
48
49 if resp.usage.input_tokens == 0 {
52 resp.usage.input_tokens = input_tokens;
53 }
54 if resp.usage.output_tokens == 0 {
55 resp.usage.output_tokens = estimate_output_tokens(&resp);
56 }
57
58 let elapsed_ms = start.elapsed().as_millis() as f64;
59 metrics::histogram!("openkind_request_duration_ms").record(elapsed_ms);
60 metrics::counter!("openkind_responses_total").increment(1);
61
62 Ok(resp)
63}
64
65fn estimate_output_tokens(resp: &SystemResponse) -> u32 {
66 let sum: usize = resp
69 .answers
70 .values()
71 .map(|a| match a {
72 Answer::Noul(_) => 1,
73 Answer::Choice(_) => 1,
74 Answer::Score(_) => 4,
75 })
76 .sum();
77 u32::try_from(sum).unwrap_or(u32::MAX)
78}