use std::collections::HashMap;
use std::fmt;
use std::future::Future;
use std::pin::Pin;
use serde_json::{json, Map, Value};
use crate::error::{AreevError, Result};
use crate::http::HttpClient;
use crate::types::{
ChatResumeRequest, ChatToolOutput, HarnessChatRequest, HarnessChatResponse, HarnessChatStatus,
};
pub struct Harness<'a> {
http: &'a HttpClient,
memory_id: String,
}
impl<'a> Harness<'a> {
pub(crate) fn new(http: &'a HttpClient, memory_id: String) -> Self {
Self { http, memory_id }
}
pub fn create(&self, name: &str, slug: &str) -> CreateHarnessBuilder<'_, 'a> {
CreateHarnessBuilder::new(self, name, slug)
}
pub async fn list(&self) -> Result<Value> {
let path = format!("/memories/{}/harnesses", self.memory_id);
self.http._get(&path, None).await
}
pub async fn get(&self, slug: &str) -> Result<Value> {
let path = format!("/memories/{}/harnesses/{}", self.memory_id, slug);
self.http._get(&path, None).await
}
pub async fn update(&self, slug: &str, if_match: &str, fields: Value) -> Result<Value> {
let mut body = match fields {
Value::Object(m) => m,
Value::Null => Map::new(),
other => {
return Err(AreevError::Other {
http_status: 0,
code: Some("CFG-E004".into()),
message: format!("update fields must be a JSON object, got {other}"),
body: Value::Null,
request_id: None,
});
}
};
body.insert("if_match".to_string(), Value::String(if_match.to_string()));
default_provider_type(&mut body);
let path = format!("/memories/{}/harnesses/{}", self.memory_id, slug);
self.http._patch(&path, Some(&Value::Object(body))).await
}
pub async fn delete(&self, slug: &str) -> Result<()> {
let path = format!("/memories/{}/harnesses/{}", self.memory_id, slug);
self.http._delete_no_retry(&path).await.map(|_| ())
}
pub async fn templates(&self) -> Result<Value> {
let path = format!("/memories/{}/harness-templates", self.memory_id);
self.http._get(&path, None).await
}
pub fn chat<'b>(&'b self, slug: &str, message: &str) -> ChatBuilder<'b, 'a> {
ChatBuilder::new(self, slug, message)
}
pub async fn chat_resume(
&self,
slug: &str,
session_id: &str,
tool_outputs: Vec<Value>,
) -> Result<Value> {
let body = json!({
"session_id": session_id,
"tool_outputs": tool_outputs,
});
let path = format!(
"/memories/{}/harnesses/{}/chat/resume",
self.memory_id, slug
);
self.http._post(&path, Some(&body)).await
}
pub async fn chat_cancel(&self, slug: &str, session_id: &str) -> Result<()> {
let path = format!(
"/memories/{}/harnesses/{}/chat/sessions/{}",
self.memory_id, slug, session_id
);
match self.http._delete(&path).await {
Ok(_) => Ok(()),
Err(e) => {
if e.http_status() == 404 {
Ok(())
} else {
Err(e)
}
}
}
}
pub fn chat_session(&self, slug: &str) -> ChatSessionBuilder<'_, 'a> {
ChatSessionBuilder::new(self, slug)
}
pub async fn chat_interactive<E>(
&self,
slug: &str,
conversation_id: &str,
user_message: &str,
executors: &ChatExecutors<E>,
) -> Result<HarnessChatResponse>
where
E: Fn(String, Value) -> Pin<Box<dyn Future<Output = Result<Value>> + Send>> + Send + Sync,
{
let conversation_id = if conversation_id.is_empty() {
gen_conversation_id()
} else {
conversation_id.to_string()
};
let req = HarnessChatRequest {
conversation_id,
user_message: user_message.to_string(),
model: None,
provider: None,
};
let mut resp = self.http.harness_chat(slug, &req).await?;
loop {
if matches!(resp.status, HarnessChatStatus::Completed) {
return Ok(resp);
}
let session_id = resp.session_id.clone().ok_or_else(|| AreevError::Other {
http_status: 500,
code: None,
message: "requires_action without session_id".into(),
body: Value::Null,
request_id: None,
})?;
let mut tool_outputs = Vec::with_capacity(resp.pending_tool_calls.len());
for pc in &resp.pending_tool_calls {
let executor = executors
.get(&pc.tool_name)
.ok_or_else(|| AreevError::Other {
http_status: 500,
code: None,
message: format!(
"no executor registered for tool '{}' (tool_call_id={})",
pc.tool_name, pc.tool_call_id
),
body: Value::Null,
request_id: None,
})?;
let args: Value = serde_json::from_str(&pc.arguments).unwrap_or(Value::Null);
let output = executor(pc.tool_name.clone(), args).await?;
tool_outputs.push(ChatToolOutput {
tool_call_id: pc.tool_call_id.clone(),
output,
is_error: None,
});
}
let resume_req = ChatResumeRequest {
session_id,
tool_outputs,
tenant_timestamp_ms: None,
};
resp = self.http.harness_chat_resume(slug, &resume_req).await?;
}
}
pub async fn provision(&self, slug: &str, template_id: &str) -> Result<Value> {
let body = json!({ "template_id": template_id });
let path = format!("/memories/{}/harnesses/{}/provision", self.memory_id, slug);
self.http._post(&path, Some(&body)).await
}
pub async fn list_conversations(&self, slug: &str, filters: Option<&Value>) -> Result<Value> {
let path = format!(
"/memories/{}/harnesses/{}/conversations",
self.memory_id, slug
);
self.http._get(&path, filters).await
}
pub async fn delete_conversation(&self, slug: &str, conversation_id: &str) -> Result<()> {
let path = format!(
"/memories/{}/harnesses/{}/conversations/{}",
self.memory_id, slug, conversation_id
);
self.http._delete(&path).await.map(|_| ())
}
pub async fn rotate_callback_secret(
&self,
slug: &str,
if_match: Option<&str>,
) -> Result<RotatedCallbackSecret> {
let path = format!(
"/memories/{}/harnesses/{}/tool_callback_secret/rotate",
self.memory_id, slug
);
let resp = self
.http
._post_if_match_no_retry(&path, Some(&json!({})), if_match)
.await?;
let get = |keys: &[&str]| -> Option<String> {
let map = resp.as_object()?;
for k in keys {
if let Some(s) = map.get(*k).and_then(|v| v.as_str()) {
if !s.is_empty() {
return Some(s.to_string());
}
}
}
None
};
Ok(RotatedCallbackSecret {
secret: get(&["secret", "callback_secret"]).unwrap_or_default(),
callback_secret_fp: get(&["callback_secret_fp", "fingerprint"]),
rotated_at_ms: resp.get("rotated_at_ms").and_then(|v| v.as_i64()),
})
}
}
#[derive(Clone)]
pub struct RotatedCallbackSecret {
pub secret: String,
pub callback_secret_fp: Option<String>,
pub rotated_at_ms: Option<i64>,
}
impl fmt::Debug for RotatedCallbackSecret {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("RotatedCallbackSecret")
.field("secret", &"[REDACTED]")
.field("callback_secret_fp", &self.callback_secret_fp)
.field("rotated_at_ms", &self.rotated_at_ms)
.finish()
}
}
pub struct ChatBuilder<'b, 'a> {
harness: &'b Harness<'a>,
slug: String,
body: Map<String, Value>,
}
impl<'b, 'a> ChatBuilder<'b, 'a> {
fn new(harness: &'b Harness<'a>, slug: &str, message: &str) -> Self {
let mut body = Map::new();
body.insert(
"user_message".to_string(),
Value::String(message.to_string()),
);
Self {
harness,
slug: slug.to_string(),
body,
}
}
pub fn conversation_id(mut self, conversation_id: &str) -> Self {
self.body.insert(
"conversation_id".to_string(),
Value::String(conversation_id.to_string()),
);
self
}
pub fn provider(mut self, provider: &str) -> Self {
self.body
.insert("provider".to_string(), Value::String(provider.to_string()));
self
}
pub fn model(mut self, model: &str) -> Self {
self.body
.insert("model".to_string(), Value::String(model.to_string()));
self
}
pub async fn send(mut self) -> Result<Value> {
self.body
.entry("conversation_id".to_string())
.or_insert_with(|| Value::String(gen_conversation_id()));
let path = format!(
"/memories/{}/harnesses/{}/chat",
self.harness.memory_id, self.slug
);
self.harness
.http
._post(&path, Some(&Value::Object(self.body)))
.await
}
}
pub struct CreateHarnessBuilder<'b, 'a> {
harness: &'b Harness<'a>,
body: Map<String, Value>,
}
impl<'b, 'a> CreateHarnessBuilder<'b, 'a> {
fn new(harness: &'b Harness<'a>, name: &str, slug: &str) -> Self {
let mut body = Map::new();
body.insert("name".to_string(), Value::String(name.to_string()));
body.insert("slug".to_string(), Value::String(slug.to_string()));
Self { harness, body }
}
fn set(mut self, key: &str, value: &str) -> Self {
self.body
.insert(key.to_string(), Value::String(value.to_string()));
self
}
pub fn description(self, description: &str) -> Self {
self.set("description", description)
}
pub fn template_id(self, template_id: &str) -> Self {
self.set("template_id", template_id)
}
pub fn assemble_query(self, assemble_query: &str) -> Self {
self.set("assemble_query", assemble_query)
}
pub fn avatar(self, avatar: &str) -> Self {
self.set("avatar", avatar)
}
pub fn accent_color(self, accent_color: &str) -> Self {
self.set("accent_color", accent_color)
}
pub fn welcome_message(self, welcome_message: &str) -> Self {
self.set("welcome_message", welcome_message)
}
pub fn llm_config(mut self, provider_id: &str, model: &str) -> Self {
self.body.insert(
"llm_config".to_string(),
json!({
"provider_id": provider_id,
"model": model,
}),
);
self
}
pub fn provider_type(mut self, provider_type: &str) -> Self {
if let Some(Value::Object(cfg)) = self.body.get_mut("llm_config") {
cfg.insert(
"provider_type".to_string(),
Value::String(provider_type.to_string()),
);
}
self
}
pub async fn send(mut self) -> Result<Value> {
default_provider_type(&mut self.body);
let path = format!("/memories/{}/harnesses", self.harness.memory_id);
self.harness
.http
._post(&path, Some(&Value::Object(self.body)))
.await
}
}
fn default_provider_type(body: &mut Map<String, Value>) {
if let Some(Value::Object(cfg)) = body.get_mut("llm_config") {
let has_provider_id = cfg
.get("provider_id")
.and_then(Value::as_str)
.is_some_and(|s| !s.is_empty());
let has_provider_type = cfg
.get("provider_type")
.and_then(Value::as_str)
.is_some_and(|s| !s.is_empty());
if has_provider_id && !has_provider_type {
let pid = cfg["provider_id"].clone();
cfg.insert("provider_type".to_string(), pid);
}
}
}
pub struct ChatSessionBuilder<'b, 'a, E = ChatExecutor> {
harness: &'b Harness<'a>,
slug: String,
conversation_id: Option<String>,
message: Option<String>,
executors: Option<&'b ChatExecutors<E>>,
}
impl<'b, 'a, E> ChatSessionBuilder<'b, 'a, E> {
fn new(harness: &'b Harness<'a>, slug: &str) -> Self {
Self {
harness,
slug: slug.to_string(),
conversation_id: None,
message: None,
executors: None,
}
}
pub fn conversation_id(mut self, conversation_id: impl Into<String>) -> Self {
self.conversation_id = Some(conversation_id.into());
self
}
pub fn message(mut self, message: impl Into<String>) -> Self {
self.message = Some(message.into());
self
}
pub fn executors(mut self, executors: &'b ChatExecutors<E>) -> Self {
self.executors = Some(executors);
self
}
pub async fn send(self) -> Result<HarnessChatResponse>
where
E: Fn(String, Value) -> Pin<Box<dyn Future<Output = Result<Value>> + Send>> + Send + Sync,
{
let message = self.message.ok_or_else(|| AreevError::Validation {
code: Some("SDK-E003".to_string()),
http_status: 0,
message: "chat_session: .message(...) is required".to_string(),
body: Value::Null,
request_id: None,
})?;
let executors = self.executors.ok_or_else(|| AreevError::Validation {
code: Some("SDK-E003".to_string()),
http_status: 0,
message: "chat_session: .executors(...) is required".to_string(),
body: Value::Null,
request_id: None,
})?;
let conversation_id = self.conversation_id.unwrap_or_default();
self.harness
.chat_interactive(&self.slug, &conversation_id, &message, executors)
.await
}
}
fn gen_conversation_id() -> String {
format!("conv-{}", uuid::Uuid::new_v4())
}
pub type ChatExecutor =
Box<dyn Fn(String, Value) -> Pin<Box<dyn Future<Output = Result<Value>> + Send>> + Send + Sync>;
pub struct ChatExecutors<E = ChatExecutor> {
inner: HashMap<String, E>,
}
impl<E> Default for ChatExecutors<E> {
fn default() -> Self {
Self {
inner: HashMap::new(),
}
}
}
impl<E> ChatExecutors<E> {
pub fn new() -> Self {
Self::default()
}
pub fn insert(&mut self, name: impl Into<String>, exec: E) -> &mut Self {
self.inner.insert(name.into(), exec);
self
}
pub fn get(&self, name: &str) -> Option<&E> {
self.inner.get(name)
}
pub fn len(&self) -> usize {
self.inner.len()
}
pub fn is_empty(&self) -> bool {
self.inner.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::default_provider_type;
use serde_json::{json, Map, Value};
fn map(v: Value) -> Map<String, Value> {
v.as_object().cloned().unwrap()
}
#[test]
fn defaults_provider_type_to_provider_id_when_absent() {
let mut body = map(json!({
"llm_config": {"provider_id": "openai", "model": "gpt-4o-mini"}
}));
default_provider_type(&mut body);
assert_eq!(body["llm_config"]["provider_type"], "openai");
}
#[test]
fn keeps_explicit_provider_type() {
let mut body = map(json!({
"llm_config": {"provider_id": "openrouter", "provider_type": "openai", "model": "x"}
}));
default_provider_type(&mut body);
assert_eq!(body["llm_config"]["provider_type"], "openai");
}
#[test]
fn noop_without_llm_config_or_provider_id() {
let mut body = map(json!({"name": "Bot"}));
default_provider_type(&mut body);
assert!(!body.contains_key("llm_config"));
let mut body2 = map(json!({"llm_config": {"model": "x"}}));
default_provider_type(&mut body2);
assert!(body2["llm_config"].get("provider_type").is_none());
}
}