use std::collections::BTreeMap;
use openai_frontend::{
OpenAiError, OpenAiResult, SystemOneAnswer, SystemOneJson, SystemOneQuestion, SystemOneRequest,
SystemOneResponse, SystemOneUsage,
};
use sha2::{Digest, Sha256};
use skippy_runtime::{
DecisionError, DecisionModel, DecisionOutput, DecisionQuestion, DecisionQuestionKind,
DecisionRequest, DecisionValue,
};
use crate::frontend::{OpenAiBackendMode, StageOpenAiBackend};
mod laya;
pub use laya::LayaSystemOneBackend;
const MAX_CHOICES: usize = 26;
const MAX_SCORE_LEVELS: usize = 10;
impl StageOpenAiBackend {
pub(super) fn run_system_one(
&self,
request: SystemOneRequest,
) -> OpenAiResult<SystemOneResponse> {
self.validate_system_one_request(&request)?;
let decision = decision_request(&request)?;
let output = self
.iteration_scheduler
.execute_runtime("system-one-read", move |runtime| {
runtime.model.decide(&decision).map_err(decision_error)
})?;
response(request, output)
}
fn validate_system_one_request(&self, request: &SystemOneRequest) -> OpenAiResult<()> {
validate_request_fields(request, &self.model_id)?;
match &self.mode {
OpenAiBackendMode::LocalRuntime => Ok(()),
OpenAiBackendMode::EmbeddedStageZero { config, .. } if config.downstream.is_none() => {
Ok(())
}
OpenAiBackendMode::EmbeddedStageZero { .. } => Err(OpenAiError::unsupported(
"System One reads currently require a complete model on one Skippy worker",
)),
}
}
}
fn run_on_model(
model: &dyn DecisionModel,
model_id: &str,
request: SystemOneRequest,
) -> OpenAiResult<SystemOneResponse> {
validate_request_fields(&request, model_id)?;
let decision = decision_request(&request)?;
let output = model.decide(&decision).map_err(decision_error)?;
response(request, output)
}
fn validate_request_fields(request: &SystemOneRequest, model_id: &str) -> OpenAiResult<()> {
const ALIASES: &[&str] = &["openjev-latest", "openjev-0.1", "jev-latest", "jev-preview"];
if request.model != model_id && !ALIASES.contains(&request.model.as_str()) {
return Err(OpenAiError::invalid_request(format!(
"model {:?} is not loaded; use {:?} or openjev-latest",
request.model, model_id
)));
}
if request.questions.is_empty() {
return Err(OpenAiError::invalid_request(
"System One needs at least one question",
));
}
if request
.images
.as_ref()
.is_some_and(|images| !images.is_empty())
|| request.steps.is_some_and(|steps| steps != 1)
|| request.samples.is_some_and(|samples| samples != 1)
|| request.think.is_some_and(|think| think != 0)
|| request.sequential.unwrap_or(false)
{
return Err(OpenAiError::unsupported(
"this PoC supports one text-only System One read; images, multiple steps/samples, thinking, and sequential reads are not yet supported",
));
}
Ok(())
}
fn decision_request(request: &SystemOneRequest) -> OpenAiResult<DecisionRequest> {
let canonical = serde_json::to_vec(&(&request.state, &request.questions)).map_err(|error| {
OpenAiError::invalid_request(format!("serialize System One request: {error}"))
})?;
let questions = request
.questions
.iter()
.map(|(key, question)| decision_question(key, question))
.collect::<OpenAiResult<Vec<_>>>()?;
Ok(DecisionRequest {
state: decision_value(&request.state),
questions,
seed: Sha256::digest(canonical).into(),
})
}
fn decision_question(key: &str, question: &SystemOneQuestion) -> OpenAiResult<DecisionQuestion> {
let (instructions, kind) = match question {
SystemOneQuestion::Noul {
instructions,
criteria,
} => (
instructions,
DecisionQuestionKind::Noul {
when_true: criteria
.as_ref()
.and_then(|criteria| criteria.r#true.as_ref())
.map(decision_value),
when_false: criteria
.as_ref()
.and_then(|criteria| criteria.r#false.as_ref())
.map(decision_value),
},
),
SystemOneQuestion::Choice {
instructions,
criteria,
} => {
if !(2..=MAX_CHOICES).contains(&criteria.len()) {
return Err(OpenAiError::invalid_request(format!(
"question {key:?}: choice criteria must contain 2 to {MAX_CHOICES} options"
)));
}
(
instructions,
DecisionQuestionKind::Choice {
options: criteria
.iter()
.map(|(name, description)| (name.to_string(), decision_value(description)))
.collect(),
},
)
}
SystemOneQuestion::Score {
instructions,
criteria,
} => {
if !(2..=MAX_SCORE_LEVELS).contains(&criteria.len()) {
return Err(OpenAiError::invalid_request(format!(
"question {key:?}: score criteria must contain 2 to {MAX_SCORE_LEVELS} levels"
)));
}
(
instructions,
DecisionQuestionKind::Score {
levels: criteria.iter().map(decision_value).collect(),
},
)
}
};
Ok(DecisionQuestion {
key: key.to_string(),
instructions: instructions.as_ref().map(decision_value),
kind,
})
}
fn decision_value(value: &SystemOneJson) -> DecisionValue {
match value {
SystemOneJson::Null => DecisionValue::Null,
SystemOneJson::Bool(value) => DecisionValue::Bool(*value),
SystemOneJson::Number(value) => DecisionValue::Number(value.clone()),
SystemOneJson::String(value) => DecisionValue::String(value.clone()),
SystemOneJson::Array(values) => {
DecisionValue::Array(values.iter().map(decision_value).collect())
}
SystemOneJson::Object(object) => DecisionValue::Object(
object
.iter()
.map(|(key, value)| (key.to_string(), decision_value(value)))
.collect(),
),
}
}
fn decision_error(error: DecisionError) -> OpenAiError {
match error {
DecisionError::InvalidRequest(message) => OpenAiError::invalid_request(message),
DecisionError::Unsupported(message) => OpenAiError::unsupported(message),
DecisionError::Backend(error) => OpenAiError::backend(format!("{error:#}")),
}
}
fn response(request: SystemOneRequest, output: DecisionOutput) -> OpenAiResult<SystemOneResponse> {
Ok(SystemOneResponse {
answers: answers(&request.questions, &output.probabilities)?,
model: request.model,
usage: SystemOneUsage {
input_tokens: u32::try_from(output.input_tokens).unwrap_or(u32::MAX),
output_tokens: 0,
},
})
}
fn confidence(probabilities: &[f32]) -> f32 {
let entropy = -probabilities
.iter()
.filter(|probability| **probability > 0.0)
.map(|probability| probability * probability.ln())
.sum::<f32>();
(1.0 - entropy / (probabilities.len() as f32).ln()).clamp(0.0, 1.0)
}
fn answers(
questions: &BTreeMap<String, SystemOneQuestion>,
probabilities: &[Vec<f32>],
) -> OpenAiResult<BTreeMap<String, SystemOneAnswer>> {
if questions.len() != probabilities.len() {
return Err(OpenAiError::backend(
"System One output did not match the question count",
));
}
questions
.iter()
.zip(probabilities)
.map(|((key, question), probabilities)| {
let answer = match question {
SystemOneQuestion::Noul { .. } => SystemOneAnswer::Noul {
noul: probabilities[0],
},
SystemOneQuestion::Choice { criteria, .. } => {
let names = criteria.iter().map(|(name, _)| name).collect::<Vec<_>>();
let selected = probabilities
.iter()
.enumerate()
.max_by(|left, right| left.1.total_cmp(right.1))
.map(|(index, _)| index)
.unwrap_or(0);
SystemOneAnswer::Choice {
choice: names[selected].to_string(),
probabilities: names
.iter()
.zip(probabilities)
.map(|(name, probability)| (name.to_string(), *probability))
.collect(),
confidence: confidence(probabilities),
}
}
SystemOneQuestion::Score { criteria, .. } => SystemOneAnswer::Score {
score: probabilities
.iter()
.enumerate()
.map(|(index, probability)| index as f32 * probability)
.sum(),
legend: criteria
.iter()
.enumerate()
.map(|(index, level)| (index.to_string(), level.to_value()))
.collect(),
probabilities: probabilities
.iter()
.enumerate()
.map(|(index, probability)| (index.to_string(), *probability))
.collect(),
confidence: confidence(probabilities),
},
};
Ok((key.clone(), answer))
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
fn request(json: &str) -> SystemOneRequest {
serde_json::from_str(json).expect("request")
}
#[test]
fn confidence_is_zero_for_uniform_distribution() {
assert!(confidence(&[0.5, 0.5]).abs() < f32::EPSILON);
}
#[test]
fn confidence_is_one_for_certain_distribution() {
assert!((confidence(&[1.0, 0.0]) - 1.0).abs() < f32::EPSILON);
}
#[test]
fn decision_request_keeps_request_order_for_options_and_state() {
let request = request(
r#"{"model":"m","state":{"zeta":1,"alpha":2},"questions":{"team":{"type":"choice","criteria":{"support":"","billing":""}}}}"#,
);
let decision = decision_request(&request).unwrap();
let DecisionValue::Object(state) = &decision.state else {
panic!("object state")
};
assert_eq!(state[0].0, "zeta");
let DecisionQuestionKind::Choice { options } = &decision.questions[0].kind else {
panic!("choice")
};
assert_eq!(options[0].0, "support");
}
const SEED_REQUEST: &str = r#"{"model":"m","state":{"zeta":{"b":[1,2.5,"x"],"a":null},"alpha":"hi"},"questions":{"team":{"type":"choice","instructions":{"y":1,"x":2},"criteria":{"support":"faults","billing":{"z":true,"a":false}}},"billing":{"type":"noul","instructions":"Billing?","criteria":{"true":"a charge","false":{"k":2,"j":1}}},"urgency":{"type":"score","criteria":["low",{"m":1,"l":0},"high"]}}}"#;
#[test]
fn the_seed_is_the_digest_diffusiongemma_has_always_used() {
let seed = decision_request(&request(SEED_REQUEST)).unwrap().seed;
let hex = seed
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();
assert_eq!(
hex,
"3e7216c3eca8bf70ed5ac23016188815058b30d8155e1adcb745e70a5c8d98dd"
);
}
#[test]
fn jev_bounds_are_enforced_before_any_backend() {
let one_choice = request(
r#"{"model":"m","state":"s","questions":{"q":{"type":"choice","criteria":{"only":""}}}}"#,
);
assert!(decision_request(&one_choice).is_err());
}
#[test]
fn answers_map_request_order_probabilities() {
let request = request(
r#"{"model":"m","state":"s","questions":{"billing":{"type":"noul"},"team":{"type":"choice","criteria":{"support":"","billing":""}}}}"#,
);
let answers = answers(&request.questions, &[vec![0.8, 0.2], vec![0.3, 0.7]]).unwrap();
assert_eq!(answers["billing"], SystemOneAnswer::Noul { noul: 0.8 });
let SystemOneAnswer::Choice {
choice,
probabilities,
..
} = &answers["team"]
else {
panic!("choice answer")
};
assert_eq!(choice, "billing");
assert_eq!(probabilities["support"], 0.3);
}
}