use std::collections::BTreeMap;
use serde_json::{Map, Value};
use crate::canonical::{CanonicalRequest, Content};
use crate::ingress::anthropic_messages::AnthAcc;
use crate::store::{content_key, Clock};
pub struct IngressState {
pub(crate) stream: bool,
pub(crate) include_usage: bool,
pub(crate) created: u64,
pub(crate) fallback_model: String,
pub(crate) id: Option<String>,
pub(crate) model: Option<String>,
pub(crate) pending: Vec<String>,
pub(crate) adaptations: Vec<String>,
pub(crate) slots: BTreeMap<u32, Slot>,
pub(crate) text: String,
pub(crate) refusal: String,
pub(crate) tools: Vec<ToolAcc>,
pub(crate) finish: Option<String>,
pub(crate) usage: Option<Value>,
pub(crate) error: Option<Value>,
pub(crate) status: Option<u16>,
pub(crate) blocks: Vec<Content>,
stash: Vec<(String, Vec<u8>)>,
pub(crate) anth: AnthAcc,
}
pub(crate) enum Slot {
Text,
Tool(usize),
Thinking(ThinkAcc),
Skip,
}
pub(crate) struct ToolAcc {
pub(crate) id: String,
pub(crate) name: String,
pub(crate) args: String,
pub(crate) signature: Option<String>,
}
#[derive(Default)]
pub(crate) struct ThinkAcc {
pub(crate) text: String,
pub(crate) id: Option<String>,
pub(crate) signature: Option<String>,
pub(crate) encrypted: Option<String>,
}
impl IngressState {
pub fn for_request(
req: &CanonicalRequest,
adaptations: Vec<String>,
clock: &dyn Clock,
) -> IngressState {
let include_usage = req
.extra
.get("stream_options")
.is_some_and(|o| o["include_usage"] == Value::Bool(true));
IngressState {
stream: req.stream == Some(true),
include_usage,
created: clock.now(),
fallback_model: req.model.clone(),
id: None,
model: None,
pending: adaptations.clone(),
adaptations,
slots: BTreeMap::new(),
text: String::new(),
refusal: String::new(),
tools: Vec::new(),
finish: None,
usage: None,
error: None,
status: None,
blocks: Vec::new(),
stash: Vec::new(),
anth: AnthAcc::default(),
}
}
pub fn status(&self) -> u16 {
self.status.unwrap_or(200)
}
pub fn take_stash(&mut self) -> Vec<(String, Vec<u8>)> {
std::mem::take(&mut self.stash)
}
pub(crate) fn wire_id(&self) -> String {
self.id
.clone()
.unwrap_or_else(|| format!("chatcmpl-brazen-{}", self.created))
}
pub(crate) fn wire_model(&self) -> &str {
self.model.as_deref().unwrap_or(&self.fallback_model)
}
pub(crate) fn close(&mut self, index: u32) {
match self.slots.remove(&index) {
Some(Slot::Thinking(t))
if t.id.is_some() || t.signature.is_some() || t.encrypted.is_some() =>
{
self.blocks.push(Content::Thinking {
text: t.text,
signature: t.signature,
id: t.id,
encrypted_content: t.encrypted,
});
}
Some(Slot::Tool(i)) if self.tools[i].signature.is_some() => {
let t = &self.tools[i];
self.blocks.push(Content::ToolUse {
id: t.id.clone(),
name: t.name.clone(),
input: parse_args(&t.args),
signature: t.signature.clone(),
});
}
_ => {}
}
}
pub(crate) fn finish_stash(&mut self) {
if self.blocks.is_empty() {
return;
}
let payload = serde_json::to_vec(&self.blocks).unwrap_or_default();
if self.tools.is_empty() {
self.stash.push((content_key(&self.text), payload));
} else {
for t in &self.tools {
self.stash.push((t.id.clone(), payload.clone()));
}
}
}
}
fn parse_args(args: &str) -> Value {
if args.is_empty() {
return Value::Object(Map::new());
}
serde_json::from_str(args).unwrap_or(Value::Null)
}