#![allow(missing_docs)]
use std::collections::HashMap;
use serde::{Deserialize, Serialize};
pub type AnswersByToolCall = HashMap<String, HashMap<String, String>>;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Workspace {
pub id: String,
pub name: String,
#[serde(default)]
pub created_at: i64,
#[serde(default)]
pub updated_at: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkspacesResponse {
pub workspaces: Vec<Workspace>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Agent {
pub uid: String,
pub name: String,
#[serde(default)]
pub description: String,
#[serde(default)]
pub mode: String,
#[serde(default)]
pub icon: String,
#[serde(default)]
pub is_published: bool,
#[serde(default)]
pub published_at: i64,
#[serde(default)]
pub created_at: i64,
#[serde(default)]
pub updated_at: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentsResponse {
pub agents: Vec<Agent>,
#[serde(default)]
pub total: i32,
}
#[derive(Debug, Serialize, Default, Clone)]
pub struct GetAgentsOptions {
#[serde(skip_serializing_if = "Option::is_none")]
page: Option<i32>,
#[serde(skip_serializing_if = "Option::is_none")]
limit: Option<i32>,
#[serde(skip_serializing_if = "Option::is_none")]
name: Option<String>,
}
impl GetAgentsOptions {
#[inline]
pub fn new() -> Self {
Default::default()
}
#[inline]
#[must_use]
pub fn page(self, page: i32) -> Self {
Self {
page: Some(page),
..self
}
}
#[inline]
#[must_use]
pub fn limit(self, limit: i32) -> Self {
Self {
limit: Some(limit),
..self
}
}
#[inline]
#[must_use]
pub fn name(self, name: impl Into<String>) -> Self {
Self {
name: Some(name.into()),
..self
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ConversationStatus {
Succeeded,
Interrupted,
Failed,
Stopped,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Reference {
#[serde(default)]
pub index: i32,
#[serde(default)]
pub title: String,
#[serde(default)]
pub url: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Question {
pub question: String,
#[serde(default)]
pub options: Vec<QuestionOption>,
#[serde(default)]
pub multi_select: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QuestionOption {
#[serde(default)]
pub description: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Interrupt {
pub node_id: String,
pub tool_call_id: String,
#[serde(default)]
pub questions: Vec<Question>,
#[serde(default)]
pub message_id: i64,
#[serde(default)]
pub chat_id: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentError {
#[serde(default)]
pub code: i32,
#[serde(default)]
pub message: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ConversationResponse {
pub chat_uid: String,
#[serde(deserialize_with = "crate::serde_utils::deserialize_string_or_int_as_string")]
pub message_id: String,
pub status: ConversationStatus,
#[serde(default)]
pub answer: String,
#[serde(default)]
pub references: Option<Vec<Reference>>,
#[serde(default)]
pub elapsed_time: f64,
#[serde(default)]
pub interrupt: Option<Interrupt>,
#[serde(default)]
pub error: Option<AgentError>,
}
impl ConversationResponse {
pub(crate) fn from_stream_parts(
started: Option<(String, String)>,
payload: WorkflowFinishedPayload,
) -> Self {
let (chat_uid, message_id) = started.unwrap_or_default();
let error = (payload.status == ConversationStatus::Failed).then_some(AgentError {
code: payload.error_code,
message: payload.error_message,
});
Self {
chat_uid,
message_id,
status: payload.status,
answer: payload.outputs.answer.unwrap_or_default(),
references: payload.outputs.references,
elapsed_time: payload.elapsed_time,
interrupt: None,
error,
}
}
pub(crate) fn from_stream_interrupt(
started: Option<(String, String)>,
interrupt: Interrupt,
) -> Self {
let (chat_uid, message_id) = started.unwrap_or_default();
Self {
chat_uid,
message_id,
status: ConversationStatus::Interrupted,
answer: String::new(),
references: None,
elapsed_time: 0.0,
interrupt: Some(interrupt),
error: None,
}
}
}
#[derive(Debug, Clone, Deserialize)]
pub struct ChatStartedPayload {
pub chat_uid: String,
#[serde(deserialize_with = "crate::serde_utils::deserialize_string_or_int_as_string")]
pub message_id: String,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct MessagePayload {
#[serde(default)]
pub text: String,
#[serde(default, rename = "type")]
pub message_type: String,
#[serde(default)]
pub key: String,
#[serde(default)]
pub started_at: i64,
#[serde(default)]
pub stage: String,
#[serde(default)]
pub stage_title: String,
#[serde(default)]
pub stage_finished_title: String,
#[serde(default)]
pub outputs: Option<serde_json::Value>,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct WorkflowOutputs {
#[serde(default)]
pub answer: Option<String>,
#[serde(default)]
pub references: Option<Vec<Reference>>,
}
#[derive(Debug, Clone, Deserialize)]
pub struct WorkflowFinishedPayload {
pub status: ConversationStatus,
#[serde(default)]
pub elapsed_time: f64,
#[serde(default)]
pub outputs: WorkflowOutputs,
#[serde(default)]
pub error: String,
#[serde(default)]
pub error_code: i32,
#[serde(default)]
pub error_message: String,
#[serde(default)]
pub error_args: Option<serde_json::Value>,
#[serde(default)]
pub process_data: Vec<serde_json::Value>,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct WorkflowStartedInputs {
#[serde(default)]
pub chat_id: i64,
#[serde(default)]
pub chat_uid: String,
#[serde(
default,
deserialize_with = "crate::serde_utils::deserialize_string_or_int_as_string"
)]
pub message_id: String,
#[serde(default)]
pub query: String,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct WorkflowStartedPayload {
#[serde(default)]
pub hit_cache: bool,
#[serde(default)]
pub inputs: WorkflowStartedInputs,
#[serde(default)]
pub started_at: i64,
#[serde(default)]
pub workflow_id: i64,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ChatFinishedPayload {
#[serde(default)]
pub chat_id: i64,
#[serde(default)]
pub chat_uid: String,
#[serde(
default,
deserialize_with = "crate::serde_utils::deserialize_string_or_int_as_string"
)]
pub message_id: String,
#[serde(default)]
pub error: String,
#[serde(default)]
pub error_message: String,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ChatTitleUpdatedPayload {
#[serde(default)]
pub chat_id: i64,
#[serde(default)]
pub chat_uid: String,
#[serde(default)]
pub source: String,
#[serde(default)]
pub title: String,
#[serde(default)]
pub updated_at: i64,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ThinkingStartedPayload {
#[serde(default)]
pub started_at: i64,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ThinkingFinishedPayload {
#[serde(default)]
pub finished_at: i64,
#[serde(default)]
pub elapsed_time: i32,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct NodeToolUseStartedPayload {
#[serde(default)]
pub tool_use_id: String,
#[serde(default)]
pub tool_name: String,
#[serde(default)]
pub tool_func_name: String,
#[serde(default)]
pub tool_args: String,
#[serde(default)]
pub tips: String,
#[serde(default)]
pub tip_chips: Vec<String>,
#[serde(default)]
pub iteration: i32,
#[serde(default)]
pub started_at: i64,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct NodeToolUseOutputs {
#[serde(default)]
pub references: Option<Vec<Reference>>,
#[serde(default)]
pub reference_domains: Option<Vec<String>>,
#[serde(default)]
pub query: Option<String>,
#[serde(default)]
pub text: Option<String>,
#[serde(default)]
pub tool_args: Option<serde_json::Value>,
#[serde(default)]
pub data: Option<serde_json::Value>,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct NodeToolUseFinishedPayload {
#[serde(default)]
pub tool_use_id: String,
#[serde(default)]
pub status: String,
#[serde(default)]
pub error: String,
#[serde(default)]
pub elapsed_time: f64,
#[serde(default)]
pub started_at: i64,
#[serde(default)]
pub tool_name: String,
#[serde(default)]
pub tool_func_name: String,
#[serde(default)]
pub tool_args: String,
#[serde(default)]
pub tool_type: String,
#[serde(default)]
pub tips: String,
#[serde(default)]
pub tip_chips: Vec<String>,
#[serde(default)]
pub iteration: i32,
#[serde(default)]
pub is_thinking: bool,
#[serde(default)]
pub outputs: NodeToolUseOutputs,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct SubagentStartedPayload {
#[serde(default)]
pub node_id: String,
#[serde(default)]
pub tool_use_id: String,
#[serde(default)]
pub started_at: i64,
#[serde(default)]
pub goal: String,
#[serde(default)]
pub prompt: String,
#[serde(default)]
pub subagent_id: String,
#[serde(default)]
pub tools: Vec<serde_json::Value>,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct SubagentProgressPayload {
#[serde(default)]
pub node_id: String,
#[serde(default)]
pub parent_tool_call_id: String,
#[serde(default)]
pub subagent_tool_name: String,
#[serde(default)]
pub subagent_tool_args: String,
#[serde(default)]
pub subagent_status: String,
#[serde(default)]
pub subagent_duration_ms: i64,
#[serde(default)]
pub subagent_iteration: i32,
#[serde(default)]
pub started_at: i64,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct SubagentOutputs {
#[serde(default)]
pub goal: Option<String>,
#[serde(default)]
pub result: Option<String>,
#[serde(default)]
pub subagent_tools: Option<Vec<serde_json::Value>>,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct SubagentFinishedPayload {
#[serde(default)]
pub node_id: String,
#[serde(default)]
pub tool_use_id: String,
#[serde(default)]
pub status: String,
#[serde(default)]
pub started_at: i64,
#[serde(default)]
pub elapsed_time: f64,
#[serde(default)]
pub error: String,
#[serde(default)]
pub outputs: SubagentOutputs,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct AgentToolStartedPayload {
#[serde(default)]
pub node_id: String,
#[serde(default)]
pub tool_use_id: String,
#[serde(default)]
pub agent_tool_name: String,
#[serde(default)]
pub title: String,
#[serde(default)]
pub started_at: i64,
#[serde(default)]
pub tool_args: String,
#[serde(default)]
pub tool_name: String,
#[serde(default)]
pub tips: String,
#[serde(default)]
pub tip_chips: Vec<String>,
#[serde(default)]
pub is_thinking: bool,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct AgentToolProgressPayload {
#[serde(default)]
pub node_id: String,
#[serde(default)]
pub parent_tool_call_id: String,
#[serde(default)]
pub agent_tool_name: String,
#[serde(default)]
pub inner_tool_name: String,
#[serde(default)]
pub inner_tool_args: String,
#[serde(default)]
pub status: String,
#[serde(default)]
pub duration_ms: i64,
#[serde(default)]
pub started_at: i64,
#[serde(default)]
pub is_thinking: bool,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct AgentToolFinishedPayload {
#[serde(default)]
pub node_id: String,
#[serde(default)]
pub tool_use_id: String,
#[serde(default)]
pub agent_tool_name: String,
#[serde(default)]
pub status: String,
#[serde(default)]
pub started_at: i64,
#[serde(default)]
pub elapsed_time: f64,
#[serde(default)]
pub error: String,
#[serde(default)]
pub tool_args: String,
#[serde(default)]
pub outputs: Option<serde_json::Value>,
#[serde(default)]
pub tool_type: String,
#[serde(default)]
pub tips: String,
#[serde(default)]
pub tip_chips: Vec<String>,
#[serde(default)]
pub is_thinking: bool,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct QueryMaskedPayload {
#[serde(default)]
pub raw_query: String,
#[serde(default)]
pub masked_query: String,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct PlanChangedPayload {
#[serde(default)]
pub node_id: String,
#[serde(default)]
pub started_at: i64,
#[serde(default)]
pub outputs: Option<serde_json::Value>,
#[serde(default)]
pub tool_name: String,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ContextCompressStartedPayload {
#[serde(default)]
pub started_at: String,
#[serde(default)]
pub inputs: Option<serde_json::Value>,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ContextCompressFinishedPayload {
#[serde(default)]
pub created_at: String,
#[serde(default)]
pub inputs: Option<serde_json::Value>,
#[serde(default)]
pub outputs: Option<serde_json::Value>,
}
#[derive(Debug, Clone)]
pub enum ConversationStreamEvent {
ChatStarted(ChatStartedPayload),
WorkflowStarted(WorkflowStartedPayload),
Message(MessagePayload),
Ping,
ThinkingStarted(ThinkingStartedPayload),
ThinkingFinished(ThinkingFinishedPayload),
NodeToolUseStarted(NodeToolUseStartedPayload),
NodeToolUseFinished(NodeToolUseFinishedPayload),
SubagentStarted(SubagentStartedPayload),
SubagentProgress(SubagentProgressPayload),
SubagentFinished(SubagentFinishedPayload),
AgentToolStarted(AgentToolStartedPayload),
AgentToolProgress(AgentToolProgressPayload),
AgentToolFinished(AgentToolFinishedPayload),
HumanInteractionRequired(ConversationResponse),
QueryMasked(QueryMaskedPayload),
PlanChanged(PlanChangedPayload),
ContextCompressStarted(ContextCompressStartedPayload),
ContextCompressFinished(ContextCompressFinishedPayload),
ChatFinished(ChatFinishedPayload),
WorkflowFinished(ConversationResponse),
ChatTitleUpdated(ChatTitleUpdatedPayload),
Other {
event: String,
data: serde_json::Value,
},
}
#[cfg(test)]
mod tests {
use super::*;
const SUCCEEDED_JSON: &str = r#"{
"chat_uid": "ct_9f2c1a5b",
"message_id": "42",
"status": "succeeded",
"answer": "Tesla (TSLA.US) recently...",
"references": [
{ "index": 1, "title": "...", "url": "..." }
],
"elapsed_time": 3.21
}"#;
const INTERRUPTED_JSON: &str = r#"{
"chat_uid": "ct_9f2c1a5b",
"message_id": "43",
"status": "interrupted",
"answer": "",
"references": null,
"elapsed_time": 1.05,
"interrupt": {
"node_id": "n_ask_human",
"tool_call_id": "call_abc123",
"questions": [
{
"question": "Which time range would you like to check?",
"options": [
{ "description": "Past week" },
{ "description": "Past month" }
],
"multi_select": false
}
],
"message_id": 43,
"chat_id": 1001
}
}"#;
#[test]
fn deserialize_succeeded_conversation_response() {
let resp: ConversationResponse = serde_json::from_str(SUCCEEDED_JSON).unwrap();
assert_eq!(resp.chat_uid, "ct_9f2c1a5b");
assert_eq!(resp.message_id, "42");
assert_eq!(resp.status, ConversationStatus::Succeeded);
assert_eq!(resp.answer, "Tesla (TSLA.US) recently...");
assert_eq!(resp.references.as_ref().unwrap().len(), 1);
assert_eq!(resp.references.as_ref().unwrap()[0].index, 1);
assert!((resp.elapsed_time - 3.21).abs() < f64::EPSILON);
assert!(resp.interrupt.is_none());
assert!(resp.error.is_none());
}
#[test]
fn deserialize_interrupted_conversation_response() {
let resp: ConversationResponse = serde_json::from_str(INTERRUPTED_JSON).unwrap();
assert_eq!(resp.status, ConversationStatus::Interrupted);
let interrupt = resp.interrupt.expect("interrupt");
assert_eq!(interrupt.node_id, "n_ask_human");
assert_eq!(interrupt.tool_call_id, "call_abc123");
assert_eq!(interrupt.message_id, 43);
assert_eq!(interrupt.chat_id, 1001);
assert_eq!(interrupt.questions.len(), 1);
assert_eq!(interrupt.questions[0].options.len(), 2);
assert!(!interrupt.questions[0].multi_select);
}
#[test]
fn deserialize_chat_started_payload_with_numeric_message_id() {
let json = r#"{"chat_uid":"ct_9f2c1a5b","message_id":42}"#;
let payload: ChatStartedPayload = serde_json::from_str(json).unwrap();
assert_eq!(payload.chat_uid, "ct_9f2c1a5b");
assert_eq!(payload.message_id, "42");
}
#[test]
fn deserialize_message_payload() {
let json = r#"{"text":"Tesla"}"#;
let payload: MessagePayload = serde_json::from_str(json).unwrap();
assert_eq!(payload.text, "Tesla");
}
#[test]
fn deserialize_message_payload_with_full_fields() {
let json =
r#"{"text":"Tesla","type":"answer","key":"n_llm_1:answer","started_at":1752048000}"#;
let payload: MessagePayload = serde_json::from_str(json).unwrap();
assert_eq!(payload.text, "Tesla");
assert_eq!(payload.message_type, "answer");
assert_eq!(payload.key, "n_llm_1:answer");
assert_eq!(payload.started_at, 1752048000);
}
#[test]
fn deserialize_workflow_finished_payload() {
let json = r#"{"status":"succeeded","elapsed_time":3.21,"outputs":{"answer":"Tesla (TSLA.US) recently..."}}"#;
let payload: WorkflowFinishedPayload = serde_json::from_str(json).unwrap();
assert_eq!(payload.status, ConversationStatus::Succeeded);
assert!((payload.elapsed_time - 3.21).abs() < f64::EPSILON);
assert_eq!(
payload.outputs.answer.as_deref(),
Some("Tesla (TSLA.US) recently...")
);
let resp = ConversationResponse::from_stream_parts(
Some(("ct_9f2c1a5b".to_string(), "42".to_string())),
payload,
);
assert_eq!(resp.chat_uid, "ct_9f2c1a5b");
assert_eq!(resp.message_id, "42");
assert_eq!(resp.answer, "Tesla (TSLA.US) recently...");
assert!(resp.interrupt.is_none());
assert!(resp.error.is_none());
}
#[test]
fn deserialize_workflow_finished_payload_with_failure() {
let json = r#"{"status":"failed","elapsed_time":0.8,"error":"upstream timeout","error_code":500,"error_message":"Something went wrong, please try again"}"#;
let payload: WorkflowFinishedPayload = serde_json::from_str(json).unwrap();
assert_eq!(payload.status, ConversationStatus::Failed);
assert_eq!(payload.error, "upstream timeout");
assert_eq!(payload.error_code, 500);
assert_eq!(
payload.error_message,
"Something went wrong, please try again"
);
let resp = ConversationResponse::from_stream_parts(None, payload);
assert_eq!(resp.status, ConversationStatus::Failed);
let error = resp.error.expect("error");
assert_eq!(error.code, 500);
assert_eq!(error.message, "Something went wrong, please try again");
}
#[test]
fn conversation_response_from_stream_interrupt() {
let json = r#"{"node_id":"n_ask_human","tool_call_id":"call_abc123","questions":[{"question":"Which time range would you like to check?","options":[{"description":"Past week"},{"description":"Past month"}],"multi_select":false}],"message_id":43,"chat_id":1001}"#;
let interrupt: Interrupt = serde_json::from_str(json).unwrap();
let resp = ConversationResponse::from_stream_interrupt(
Some(("ct_9f2c1a5b".to_string(), "43".to_string())),
interrupt,
);
assert_eq!(resp.chat_uid, "ct_9f2c1a5b");
assert_eq!(resp.message_id, "43");
assert_eq!(resp.status, ConversationStatus::Interrupted);
let interrupt = resp.interrupt.expect("interrupt");
assert_eq!(interrupt.node_id, "n_ask_human");
assert_eq!(interrupt.questions.len(), 1);
}
#[test]
fn deserialize_node_tool_use_finished_payload() {
let json = r#"{"tool_use_id":"call_abc123","status":"succeeded","elapsed_time":1.42,"tool_name":"Web Search","tool_func_name":"web_search","tool_args":"{\"query\":\"TSLA stock news\"}","tool_type":"builtin","tips":"Searched the web","iteration":1,"is_thinking":true,"outputs":{"query":"TSLA stock news","references":[{"index":1,"title":"...","url":"..."}]}}"#;
let payload: NodeToolUseFinishedPayload = serde_json::from_str(json).unwrap();
assert_eq!(payload.tool_use_id, "call_abc123");
assert_eq!(payload.status, "succeeded");
assert_eq!(payload.tool_func_name, "web_search");
assert!(payload.is_thinking);
assert_eq!(payload.outputs.query.as_deref(), Some("TSLA stock news"));
assert_eq!(payload.outputs.references.as_ref().unwrap().len(), 1);
}
#[test]
fn deserialize_plan_changed_payload_picks_up_sibling_tool_name() {
let mut payload: PlanChangedPayload =
serde_json::from_str(r#"{"node_id":"n_plan","started_at":1752048000}"#).unwrap();
payload.tool_name = "planner".to_string();
assert_eq!(payload.node_id, "n_plan");
assert_eq!(payload.tool_name, "planner");
}
#[test]
fn deserialize_workspaces_response() {
let json = r#"{
"workspaces": [
{ "id": "1001", "name": "My Workspace", "created_at": 1742000000, "updated_at": 1742001000 }
]
}"#;
let resp: WorkspacesResponse = serde_json::from_str(json).unwrap();
assert_eq!(resp.workspaces.len(), 1);
assert_eq!(resp.workspaces[0].id, "1001");
}
#[test]
fn deserialize_agents_response() {
let json = r#"{
"agents": [
{
"uid": "ag_7d3f9b2c",
"name": "US Stock Analyst",
"description": "Answers US stock questions with market and fundamental data",
"mode": "chat",
"icon": "https://cdn.longbridge.com/icons/agent.png",
"is_published": true,
"published_at": 1742000000,
"created_at": 1741000000,
"updated_at": 1742001000
}
],
"total": 12
}"#;
let resp: AgentsResponse = serde_json::from_str(json).unwrap();
assert_eq!(resp.total, 12);
assert_eq!(resp.agents[0].uid, "ag_7d3f9b2c");
assert!(resp.agents[0].is_published);
}
}