1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
//! `Chat` resource — app-style assistant chat over a memory: threads,
//! messages, tool dispatch, and SSE streaming.
//!
//! Mirrors the Python SDK's `client.chat.*` surface.
//!
//! This is the app's general assistant chat (thread-based). For the
//! harness conversational loop (Flow-A, caller-supplied tools) use
//! [`crate::Areev::harness`].
use serde_json::{json, Map, Value};
use crate::error::Result;
use crate::http::{HttpClient, SseStream};
/// Assistant chat + thread management.
///
/// Access via [`crate::Areev::chat`].
pub struct Chat<'a> {
http: &'a HttpClient,
memory_id: String,
}
impl<'a> Chat<'a> {
/// Internal constructor — use [`crate::Areev::chat`].
pub(crate) fn new(http: &'a HttpClient, memory_id: String) -> Self {
Self { http, memory_id }
}
/// Start a [`ChatSendBuilder`] for one non-streaming assistant turn.
pub fn send<'b>(&'b self, messages: Vec<Value>) -> ChatSendBuilder<'b, 'a> {
ChatSendBuilder::new(self, messages)
}
/// Open a token-by-token SSE stream for one assistant turn.
///
/// Returns an [`SseStream`] the caller pumps with
/// [`SseStream::next`]. **Not auto-retried** — a stream is a single
/// live connection; restart on failure.
///
/// ```no_run
/// # #[tokio::main]
/// # async fn main() -> areev::Result<()> {
/// # use serde_json::json;
/// let areev = areev::Areev::from_env();
/// let mut stream = areev
/// .chat()
/// .stream(vec![json!({"role": "user", "content": "hi"})], None, None)
/// .await?;
/// while let Some(chunk) = stream.next().await? {
/// print!("{chunk}");
/// }
/// # Ok(())
/// # }
/// ```
pub async fn stream(
&self,
messages: Vec<Value>,
provider: Option<&str>,
model: Option<&str>,
) -> Result<SseStream> {
let mut body = Map::new();
body.insert("messages".into(), Value::Array(messages));
if let Some(p) = provider {
body.insert("provider".into(), Value::String(p.to_string()));
}
if let Some(m) = model {
body.insert("model".into(), Value::String(m.to_string()));
}
let path = format!("/memories/{}/chat/stream", self.memory_id);
self.http
._stream_sse(&path, Some(&Value::Object(body)))
.await
}
/// Dispatch a single app-chat tool call server-side.
pub async fn dispatch_tool(&self, tool: &str, args: Value) -> Result<Value> {
let path = format!("/memories/{}/chat/tools/{}", self.memory_id, tool);
self.http._post(&path, Some(&args)).await
}
// ── Threads ──────────────────────────────────────────────────────
/// Create a chat thread.
pub async fn create_thread(&self, title: Option<&str>) -> Result<Value> {
let mut body = Map::new();
if let Some(t) = title {
body.insert("title".into(), Value::String(t.to_string()));
}
let path = format!("/memories/{}/chat/threads", self.memory_id);
self.http._post(&path, Some(&Value::Object(body))).await
}
/// List chat threads. `filters` is added as query parameters.
pub async fn list_threads(&self, filters: Option<&Value>) -> Result<Value> {
let path = format!("/memories/{}/chat/threads", self.memory_id);
self.http._get(&path, filters).await
}
/// Get a chat thread by id.
pub async fn get_thread(&self, thread_id: &str) -> Result<Value> {
let path = format!("/memories/{}/chat/threads/{}", self.memory_id, thread_id);
self.http._get(&path, None).await
}
/// Rename a chat thread.
pub async fn rename_thread(&self, thread_id: &str, title: &str) -> Result<Value> {
let body = json!({ "title": title });
let path = format!("/memories/{}/chat/threads/{}", self.memory_id, thread_id);
self.http._patch(&path, Some(&body)).await
}
/// Delete a chat thread and its messages.
pub async fn delete_thread(&self, thread_id: &str) -> Result<()> {
let path = format!("/memories/{}/chat/threads/{}", self.memory_id, thread_id);
self.http._delete(&path).await.map(|_| ())
}
/// List the messages in a thread. `filters` is added as query
/// parameters.
pub async fn list_thread_messages(
&self,
thread_id: &str,
filters: Option<&Value>,
) -> Result<Value> {
let path = format!(
"/memories/{}/chat/threads/{}/messages",
self.memory_id, thread_id
);
self.http._get(&path, filters).await
}
}
/// Builder for [`Chat::send`].
pub struct ChatSendBuilder<'b, 'a> {
chat: &'b Chat<'a>,
body: Map<String, Value>,
}
impl<'b, 'a> ChatSendBuilder<'b, 'a> {
fn new(chat: &'b Chat<'a>, messages: Vec<Value>) -> Self {
let mut body = Map::new();
body.insert("messages".into(), Value::Array(messages));
Self { chat, body }
}
/// Attach the turn to an existing thread.
pub fn thread_id(mut self, thread_id: &str) -> Self {
self.body
.insert("thread_id".into(), Value::String(thread_id.to_string()));
self
}
/// Stable conversation id for grouping turns.
pub fn conversation_id(mut self, conversation_id: &str) -> Self {
self.body.insert(
"conversation_id".into(),
Value::String(conversation_id.to_string()),
);
self
}
/// Override the LLM provider for this turn.
pub fn provider(mut self, provider: &str) -> Self {
self.body
.insert("provider".into(), Value::String(provider.to_string()));
self
}
/// Override the LLM model for this turn.
pub fn model(mut self, model: &str) -> Self {
self.body
.insert("model".into(), Value::String(model.to_string()));
self
}
/// Enable / disable server-side tool dispatch for this turn.
pub fn tools_enabled(mut self, tools_enabled: bool) -> Self {
self.body
.insert("tools_enabled".into(), Value::Bool(tools_enabled));
self
}
/// Set an arbitrary additional body field — escape hatch.
pub fn extra(mut self, key: &str, value: Value) -> Self {
self.body.insert(key.to_string(), value);
self
}
/// Issue the chat request and return the assembled assistant reply.
///
/// `POST /memories/{id}/chat` is **SSE-only** (`text/event-stream`) —
/// there is no JSON response. This consumes the stream, concatenates
/// the per-chunk content deltas into the full assistant message, and
/// returns:
///
/// ```json
/// { "role": "assistant", "content": "<full text>", "usage": { … },
/// "model": "…", "conversation_id": "…" }
/// ```
///
/// The `usage` / `model` / `conversation_id` fields are populated from
/// the terminal event when the server emits one. For token-by-token
/// consumption use [`Chat::stream`] instead.
pub async fn send(self) -> Result<Value> {
let path = format!("/memories/{}/chat", self.chat.memory_id);
let mut stream = self
.chat
.http
._stream_sse(&path, Some(&Value::Object(self.body)))
.await?;
let mut content = String::new();
let mut tail = Map::new();
while let Some(payload) = stream.next().await? {
let Ok(value) = serde_json::from_str::<Value>(&payload) else {
// Non-JSON SSE payload — treat the bare text as a delta.
content.push_str(&payload);
continue;
};
if let Some(delta) = extract_delta(&value) {
content.push_str(&delta);
} else if let Value::Object(map) = value {
// Terminal / metadata frame (usage, model, conversation_id).
tail = map;
}
}
let mut out = Map::new();
out.insert("role".into(), Value::String("assistant".into()));
out.insert("content".into(), Value::String(content));
for key in ["usage", "model", "conversation_id", "id"] {
if let Some(v) = tail.remove(key) {
out.insert(key.into(), v);
}
}
Ok(Value::Object(out))
}
}
/// Pull the assistant content delta out of one SSE chunk, supporting both
/// the cell's native `{"delta": "…"}` shape and the OpenAI-style
/// `chat.completion.chunk` shape (`choices[0].delta.content`).
fn extract_delta(value: &Value) -> Option<String> {
if let Some(s) = value.get("delta").and_then(|v| v.as_str()) {
return Some(s.to_string());
}
value
.get("choices")
.and_then(|c| c.as_array())
.and_then(|arr| arr.first())
.and_then(|c| c.get("delta"))
.and_then(|d| d.get("content"))
.and_then(|v| v.as_str())
.map(|s| s.to_string())
}