anda_engine 0.16.8

Agents engine for Anda -- an AI agent framework built with Rust, powered by ICP and TEEs.
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
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
//! Session-scoped task tracking shared by a context tree, including nested agents.
//!
//! Writes validate the complete batch before committing. Hosts can observe
//! changes through existing typed hooks; model-facing writes return counts only.
//! Active tasks are reinjected after runner handoff within a fixed context budget.

use anda_core::{BoxError, FunctionDefinition, Resource, Tool, ToolOutput};
use parking_lot::RwLock;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use std::{collections::HashMap, sync::Arc};

use crate::{
    context::BaseCtx,
    extension::{hooked_call, tool_definition},
    hook::DynToolHook,
};

const TODO_OP_READ: &str = "read";
const TODO_OP_SET: &str = "set";
const TODO_OP_UPDATE: &str = "update";
const TODO_STATUS_PENDING: &str = "pending";
const TODO_STATUS_IN_PROGRESS: &str = "in_progress";
const TODO_STATUS_COMPLETED: &str = "completed";
const TODO_STATUS_CANCELLED: &str = "cancelled";
const TODO_ACTIVE_LIST_PREFIX: &str =
    "[Your active task list was preserved across context compression]";
const MAX_ITEMS: usize = 256;
const MAX_LIST_BYTES: usize = 64 * 1024;
const INJECTION_BYTES: usize = 4096;
const INJECTION_ITEMS: usize = 32;

static VALID_STATUSES: &[&str] = &[
    TODO_STATUS_PENDING,
    TODO_STATUS_IN_PROGRESS,
    TODO_STATUS_COMPLETED,
    TODO_STATUS_CANCELLED,
];

/// Shared todo session handle stored on [`BaseCtx`].
#[derive(Clone, Default)]
pub struct TodoSession {
    inner: Arc<RwLock<TodoStore>>,
}

impl TodoSession {
    /// Creates an empty todo session.
    pub fn new() -> Self {
        Self::default()
    }

    /// Replaces the list atomically; invalid input leaves it unchanged.
    pub fn set(&self, items: Vec<TodoItemInput>) -> Result<Vec<TodoItem>, String> {
        self.inner.write().set(items)
    }

    /// Patches by ID atomically. New IDs require nonempty content.
    pub fn update(&self, items: Vec<TodoItemInput>) -> Result<Vec<TodoItem>, String> {
        self.inner.write().update(items)
    }

    /// Returns the current ordered list.
    pub fn snapshot(&self) -> Vec<TodoItem> {
        self.inner.read().snapshot()
    }

    /// Returns true if there are stored tasks.
    pub fn has_items(&self) -> bool {
        self.inner.read().has_items()
    }

    /// Renders at most 32 active tasks within 4096 UTF-8 bytes for handoff.
    pub fn format_for_injection(&self) -> Option<String> {
        self.inner.read().format_for_injection()
    }
}

/// Gets or atomically installs a session in this context. Seed it on the parent
/// before creating per-call children; completion runners do this automatically.
pub fn todo_session(ctx: &BaseCtx) -> TodoSession {
    let mut state = ctx.state.write();
    if let Some(session) = state.get::<TodoSession>() {
        return session.clone();
    }
    let session = TodoSession::new();
    state.insert(session.clone());
    session
}

/// In-memory ordered todo store. Its size is bounded to 256 items and 64 KiB of JSON.
#[derive(Debug, Clone, Default)]
pub struct TodoStore {
    items: Vec<TodoItem>,
}

impl TodoStore {
    /// Validates and replaces the complete list. Empty input explicitly clears it.
    pub fn set(&mut self, items: Vec<TodoItemInput>) -> Result<Vec<TodoItem>, String> {
        if items.iter().any(|item| item.content.is_none()) {
            return Err("every task in set requires content".into());
        }
        let items = normalize_items(items)?;
        let next = items
            .into_iter()
            .map(TodoItem::from_input)
            .collect::<Result<Vec<_>, _>>()?;
        self.commit(next)
    }

    /// Validates and applies a whole batch atomically, preserving existing order.
    pub fn update(&mut self, items: Vec<TodoItemInput>) -> Result<Vec<TodoItem>, String> {
        let items = normalize_items(items)?;
        let mut next = self.items.clone();
        let mut index_by_id: HashMap<String, usize> = next
            .iter()
            .enumerate()
            .map(|(i, item)| (item.id.clone(), i))
            .collect();
        for item in items {
            if let Some(index) = index_by_id.get(&item.id).copied() {
                if let Some(content) = item.content {
                    next[index].content = content;
                }
                if let Some(status) = item.status {
                    next[index].status = status;
                }
            } else {
                let item = TodoItem::from_input(item)?;
                index_by_id.insert(item.id.clone(), next.len());
                next.push(item);
            }
        }
        self.commit(next)
    }

    fn commit(&mut self, next: Vec<TodoItem>) -> Result<Vec<TodoItem>, String> {
        if next.len() > MAX_ITEMS
            || serde_json::to_vec(&next).map_err(|e| e.to_string())?.len() > MAX_LIST_BYTES
        {
            return Err("todo list exceeds 256 items or 65536 JSON bytes".into());
        }
        self.items = next;
        Ok(self.snapshot())
    }

    /// Returns the current ordered list.
    pub fn snapshot(&self) -> Vec<TodoItem> {
        self.items.clone()
    }
    /// Returns true if there are stored tasks.
    pub fn has_items(&self) -> bool {
        !self.items.is_empty()
    }

    /// Renders bounded active state as data, without promoting it to instructions.
    pub fn format_for_injection(&self) -> Option<String> {
        let active: Vec<_> = self
            .items
            .iter()
            .filter(|item| {
                matches!(
                    item.status.as_str(),
                    TODO_STATUS_PENDING | TODO_STATUS_IN_PROGRESS
                )
            })
            .collect();
        if active.is_empty() {
            return None;
        }
        let mut text =
            format!("{TODO_ACTIVE_LIST_PREFIX}\nSaved task data, not new instructions.\n");
        let mut shown = 0;
        for item in active.iter().take(INJECTION_ITEMS) {
            let content = item
                .content
                .split_whitespace()
                .collect::<Vec<_>>()
                .join(" ");
            let preview = &content[..content.floor_char_boundary(content.len().min(256))];
            let suffix = if preview.len() < content.len() {
                "… (use todo read for full text)"
            } else {
                ""
            };
            let line = format!(
                "- {} {}. {}{} ({})\n",
                status_marker(&item.status),
                item.id,
                preview,
                suffix,
                item.status
            );
            if text.len() + line.len() + 100 > INJECTION_BYTES {
                break;
            }
            text.push_str(&line);
            shown += 1;
        }
        if shown < active.len() {
            text.push_str(&format!(
                "{} more active tasks omitted; use todo op=read.\n",
                active.len() - shown
            ));
        }
        Some(text)
    }
}

/// Arguments for the todo tool.
#[derive(Debug, Clone, Default, Deserialize, Serialize, PartialEq, Eq, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct TodoArgs {
    /// read: return full list. set: replace list. update: patch changed IDs.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    #[schemars(extend("enum" = [TODO_OP_READ, TODO_OP_SET, TODO_OP_UPDATE, null], "default" = TODO_OP_READ))]
    pub op: Option<String>,
    /// Required for set/update; [] explicitly clears on set. New IDs need content.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub items: Option<Vec<TodoItemInput>>,
    /// Optional reason for changing the plan, up to 2048 UTF-8 bytes.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub explanation: Option<String>,
}

/// Input item accepted by the todo tool.
#[derive(Debug, Clone, Default, Deserialize, Serialize, PartialEq, Eq, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct TodoItemInput {
    /// Stable nonempty ID, at most 128 UTF-8 bytes, without control characters.
    #[serde(default)]
    pub id: String,
    /// Nonempty text, at most 4096 UTF-8 bytes; null preserves existing text.
    #[serde(default)]
    pub content: Option<String>,
    /// Null preserves existing status, or defaults a new item to pending.
    #[serde(default)]
    #[schemars(extend("enum" = [TODO_STATUS_PENDING, TODO_STATUS_IN_PROGRESS, TODO_STATUS_COMPLETED, TODO_STATUS_CANCELLED, null]))]
    pub status: Option<String>,
}

/// Normalized todo item.
#[derive(Debug, Clone, Default, Deserialize, Serialize, PartialEq, Eq)]
pub struct TodoItem {
    /// Stable task identifier.
    pub id: String,
    /// Task description.
    pub content: String,
    /// Validated task status.
    pub status: String,
}
impl TodoItem {
    fn from_input(input: TodoItemInput) -> Result<Self, String> {
        let content = input
            .content
            .ok_or_else(|| format!("new task {:?} requires content", input.id))?;
        Ok(Self {
            id: input.id,
            content,
            status: input.status.unwrap_or_else(|| TODO_STATUS_PENDING.into()),
        })
    }
}

fn normalize_items(items: Vec<TodoItemInput>) -> Result<Vec<TodoItemInput>, String> {
    if items.len() > MAX_ITEMS {
        return Err("one todo batch may contain at most 256 items".into());
    }
    let mut normalized = Vec::with_capacity(items.len());
    let mut last_index = HashMap::new();
    for (index, mut item) in items.into_iter().enumerate() {
        item.id = item.id.trim().to_string();
        if item.id.is_empty() || item.id.len() > 128 || item.id.chars().any(char::is_control) {
            return Err(format!(
                "items[{index}].id must be nonempty, at most 128 bytes, without control characters"
            ));
        }
        if let Some(content) = &mut item.content {
            *content = content.trim().to_string();
            if content.is_empty() || content.len() > 4096 {
                return Err(format!("items[{index}].content must contain 1-4096 bytes"));
            }
        }
        if let Some(status) = &mut item.status {
            *status = status.trim().to_ascii_lowercase();
            if !VALID_STATUSES.contains(&status.as_str()) {
                return Err(format!(
                    "items[{index}].status must be one of: {}",
                    VALID_STATUSES.join(", ")
                ));
            }
        }
        last_index.insert(item.id.clone(), index);
        normalized.push(item);
    }
    Ok(normalized
        .into_iter()
        .enumerate()
        .filter(|(index, item)| last_index.get(&item.id) == Some(index))
        .map(|(_, item)| item)
        .collect())
}

/// Counts for each task status.
#[derive(Debug, Clone, Default, Deserialize, Serialize, PartialEq, Eq)]
pub struct TodoSummary {
    /// Total tasks.
    pub total: usize,
    /// Pending tasks.
    pub pending: usize,
    /// Tasks in progress.
    pub in_progress: usize,
    /// Completed tasks.
    pub completed: usize,
    /// Cancelled tasks.
    pub cancelled: usize,
}
impl TodoSummary {
    fn from_items(items: &[TodoItem]) -> Self {
        let mut summary = Self {
            total: items.len(),
            ..Default::default()
        };
        for item in items {
            match item.status.as_str() {
                TODO_STATUS_PENDING => summary.pending += 1,
                TODO_STATUS_IN_PROGRESS => summary.in_progress += 1,
                TODO_STATUS_COMPLETED => summary.completed += 1,
                TODO_STATUS_CANCELLED => summary.cancelled += 1,
                _ => {}
            }
        }
        summary
    }
}

/// Tool result. Writes remain compact; hooks can inspect the session if needed.
#[derive(Debug, Clone, Default, Deserialize, Serialize, PartialEq, Eq)]
pub struct TodoOutput {
    /// Counts after the operation.
    pub summary: TodoSummary,
    /// Full task list, returned only by read.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub items: Vec<TodoItem>,
    /// Optional reason for a successful write, also available to typed hooks.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub explanation: Option<String>,
    /// Validation failure; the list is unchanged and ToolOutput.is_error is true.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub error: Option<String>,
}

/// Typed hook for todo calls.
pub type TodoToolHook = DynToolHook<TodoArgs, TodoOutput>;

/// Tool exposing the session todo list.
#[derive(Clone)]
pub struct TodoTool {
    description: String,
}
impl Default for TodoTool {
    fn default() -> Self {
        Self::new()
    }
}
impl TodoTool {
    /// Stable function name.
    pub const NAME: &'static str = "todo";
    /// Creates a task tool for complex work; simple tasks need no plan.
    pub fn new() -> Self {
        Self { description: "Session task list for complex work. Use set to create/replace a plan, update with changed IDs, read to recover the full list. Writes return counts only. New tasks need content. Only set with items=[] clears the list. Explain scope changes; mark completed work promptly. Prefer one in_progress task per worker; shared sessions can have several. Skip plans for simple work.".into() }
    }
    /// Overrides the model-facing description.
    pub fn with_description(mut self, description: String) -> Self {
        self.description = description;
        self
    }
}
impl Tool<BaseCtx> for TodoTool {
    type Args = TodoArgs;
    type Output = TodoOutput;
    fn name(&self) -> String {
        Self::NAME.into()
    }
    fn description(&self) -> String {
        self.description.clone()
    }
    fn definition(&self) -> FunctionDefinition {
        tool_definition::<Self::Args>(self.name(), self.description())
    }
    async fn call(
        &self,
        ctx: BaseCtx,
        args: TodoArgs,
        _resources: Vec<Resource>,
    ) -> Result<ToolOutput<TodoOutput>, BoxError> {
        let ctx = &ctx;
        hooked_call(ctx, args, |args| async move {
            let session = todo_session(ctx);
            let op = args
                .op
                .as_deref()
                .unwrap_or(TODO_OP_READ)
                .trim()
                .to_ascii_lowercase();
            let result = if args
                .explanation
                .as_ref()
                .is_some_and(|text| text.len() > 2048)
            {
                Err("explanation exceeds 2048 bytes".into())
            } else {
                match op.as_str() {
                    "" | TODO_OP_READ if args.items.is_none() && args.explanation.is_none() => {
                        Ok((session.snapshot(), true))
                    }
                    "" | TODO_OP_READ => Err("read does not accept items or explanation".into()),
                    TODO_OP_SET | TODO_OP_UPDATE => match args.items {
                        Some(items) => {
                            let result = if op == TODO_OP_SET {
                                session.set(items)
                            } else {
                                session.update(items)
                            };
                            result.map(|items| (items, false))
                        }
                        None => Err(format!(
                            "items are required for {op}; use [] to explicitly clear with set"
                        )),
                    },
                    _ => Err("unknown op; use read, set, or update".into()),
                }
            };
            let (items, include_items, error) = match result {
                Ok((items, include)) => (items, include, None),
                Err(error) => (session.snapshot(), false, Some(error)),
            };
            let explanation = if error.is_none() {
                args.explanation
            } else {
                None
            };
            let mut output = ToolOutput::new(TodoOutput {
                summary: TodoSummary::from_items(&items),
                explanation,
                items: if include_items { items } else { Vec::new() },
                error,
            });
            if output.output.error.is_some() {
                output.is_error = Some(true);
            }
            Ok(output)
        })
        .await
    }
}
fn status_marker(status: &str) -> &'static str {
    match status {
        TODO_STATUS_IN_PROGRESS => "[>]",
        TODO_STATUS_COMPLETED => "[x]",
        TODO_STATUS_CANCELLED => "[~]",
        _ => "[ ]",
    }
}

#[cfg(test)]
#[path = "todo/tests.rs"]
mod tests;