Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions .cargo/mutants.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
# cargo-mutants configuration.
#
# Options docs: https://github.com/sourcefrog/cargo-mutants

exclude_re = [
# Equivalent mutants in EditTool::apply_ops (funera_builtin_tools/src/edit.rs:163).
# Verified anchors always satisfy `idx <= lines.len() - 1`, so
# `idx >= lines.len() - 1` and its `+ 1` / `/ 1` alternatives compute the
# same insert position in every reachable state — no black-box test can
# distinguish them.
"replace - with \\+ in EditTool::apply_ops",
"replace - with / in EditTool::apply_ops",
# Equivalent in EditTool::apply_ops (edit.rs:157): when the range bounds
# check is reached, both operands are statically false — the end anchor was
# verified within the file (end_idx < len) and the start guard at line 135
# guarantees end_idx >= idx. `false || false` and `false && false` are
# indistinguishable.
"replace \\|\\| with && in EditTool::apply_ops",
# Equivalent on Windows: `ExitStatus::code()` is always `Some` there, so the
# `unwrap_or(-1)` default in ShellTool::format_output is unreachable and the
# `- 1` replacement cannot be observed. (On Unix it can be exercised by a
# signal-death test; remove this pattern to re-enable it on such CI.)
"delete - in format_output",
]
29 changes: 29 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,35 @@ and this project adheres to [Semantic Versioning](https://semver.org/).

### Added

- Cancellation — every entry point (`fire` → `FireHandle`, `fire_stream`,
`send`, `send_stream`) carries a per-call `CancellationToken`; `cancel()`
or dropping the handle immediately stops receiving output, abandons
in-flight tool executions (workers survive), and notifies middleware via
`AgentEvent::Cancelled` (never written to session history). `fire` now
returns a `FireHandle` — `await` it for the `ChatResponse`.
- `serde` feature on `funera-core` — gated `Serialize`/`Deserialize` derives
for the chat message types (`FuneraMessage`, `MsgVariant`, `Role`,
`TextMessage`, `ToolRequestMessage`, `ToolResponseMessage`).
- `AgentRuntime::session_id` — stable conversation id reused by `send` /
`send_stream`; `fire` / `fire_stream` keep a fresh one-shot id (fork
semantics).
- Token usage tracking — `TokenUsage` + `TokenEvent::Usage`; requests ask for
`stream_options.include_usage`; `ChatResponse.usage` carries the last turn's
usage. Cost computation is left to callers.
- Parallel tool execution — the tool bus fans out to `max_concurrent_tools`
(default 4) workers via `AgentRuntimeBuilder::max_concurrent_tools`, so
multiple tool calls in one turn run concurrently.
- Reasoning levels — `ReasoningLevel`
(`Off`/`Minimal`/`Low`/`Medium`/`High`/`XHigh`/`Max`, default `Medium`) with
hot-reload via `AgentRuntime::set_reasoning_level`; providers interpret it
(DeepSeek mirrors the official harness: `thinking` + `reasoning_effort` for
`high`/`max`; OpenAI maps to `reasoning_effort`, clamping `xhigh`/`max` to
`high`).
- `Tool::is_shell_tool` marker (default `false`) — tools that execute shell
commands opt in explicitly, and `ShellPolicy` scrutiny is now driven by the
marker instead of a hard-coded tool-name list
(`shell`/`bash`/`sh`/`cmd`/`powershell`). `ToolPolicy::check_shell_command`
now takes the tool itself.
- `FuneraEnv::effect` / `dispose` — reversible effects with LIFO teardown: every
registration can be paired with its inverse, and disposal runs all inverses in
reverse registration order (idempotent and panic-isolated). `EnvActor` calls
Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

30 changes: 23 additions & 7 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
.system_prompt("You are a helpful assistant.")
.build();

let resp = agent.fire("Hello!", &runtime).await?;
let resp = agent.fire("Hello!", &runtime).await?.await?.await?;
println!("{}", resp.content);
Ok(())
}
Expand Down Expand Up @@ -179,7 +179,8 @@ sequenceDiagram
- **Skill system** — load prompt templates from YAML-frontmatter Markdown files
- **Middleware pipeline** — intercept agent events with inspectors (read-only, parallel) and mutators (pass/modify/block, sequential)
- **Security layer** — tool/shell policies, path allowlisting, audit logging, secure API key storage
- **Type-state session** — compile-time enforcement of session ownership (`Idle` / `Acquired`)
- **Type-state session** — compile-time enforcement of session ownership (Idle / Acquired)
- **Cancellation** — every call returns a cancellable handle ( ire → FireHandle, ire_stream / send / send_stream); cancel() or dropping the handle immediately stops receiving output, abandons in-flight tool executions, and notifies middleware via AgentEvent::Cancelled (never written to session history).

## Examples

Expand All @@ -192,6 +193,7 @@ sequenceDiagram
| `middleware` | Inspector/Mutator middleware pipeline |
| `reversible_effects` | LIFO teardown of registered effects (no LLM) |
| `tool_policy` / `security` / `sandbox` | Security policies, audit, and sandboxing |
| `cancel` | Cancellation: explicit `cancel()`, drop-to-cancel, timeout via `select!` |

## Reversible effects

Expand All @@ -217,9 +219,11 @@ env.effect(|| {
disposer is caught and logged; the rest still run).
- `EnvActor` calls `dispose()` automatically once the runtime is dropped, so anything registered
against the env is reverted — no memory or service leaks.
- For tool registrations, the safe inverse is `remove_tool_if_same`: it removes a tool only if
the registered entry is the *same* `Arc`, so a stale teardown never deletes a replacement
tool that reuses the same name.
- Tool registrations are paired with their exact inverse internally
(`remove_tool_if_same` removes a tool only if the registered entry is the
*same* `Arc`, so a stale teardown never deletes a replacement that reuses
the same name); external extensions register tools through the runtime API
and do not reach into the tool registry themselves.

Run the demo:

Expand Down Expand Up @@ -277,7 +281,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
.system_prompt("You are a helpful assistant.")
.build();

let resp = agent.fire("Hello!", &runtime).await?;
let resp = agent.fire("Hello!", &runtime).await?.await?.await?;
println!("{}", resp.content);
Ok(())
}
Expand Down Expand Up @@ -387,6 +391,18 @@ let runtime = AgentRuntime::<DeepSeekProvider>::builder()
.build()?;
```

If your tool executes shell commands, opt in explicitly via `is_shell_tool`
so the security layer's `ShellPolicy` applies — based on the marker, not on
the tool's registered name:

```rust
impl Tool for Shell {
fn name(&self) -> &str { "shell" }
// ...
fn is_shell_tool(&self) -> bool { true }
}
```

### Security configuration

Requires the `security` feature (and optionally `funera-builtin-tools`, `sandbox`):
Expand All @@ -410,7 +426,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
.build()?;

let agent = Agent::builder().build();
let resp = agent.fire("List the git log.", &runtime).await?;
let resp = agent.fire("List the git log.", &runtime).await?.await?.await?;
println!("{}", resp.content);
Ok(())
}
Expand Down
17 changes: 17 additions & 0 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,3 +35,20 @@ env is torn down:
This guarantees that anything registered against the env — tools, listeners, connections — is
reverted on teardown, preventing memory and service leaks. See
[Reversible Effects](concepts/effects.md) for the pattern.

## Cancellation

Every agent call returns a cancellable handle — `FireHandle` (one-shot `fire`),
`FireStreamHandle`, `SendHandle`, `SendStreamHandle`. `cancel()` — or simply
dropping the handle — cancels the call's `CancellationToken`. The `ReActLoop`
exits cooperatively at its next blocking point (stream consumption, the initial
provider request, in-flight tool execution) and:

- stops receiving output immediately;
- abandons in-flight tool executions (the shared `ToolExecutor` workers survive);
- emits `AgentEvent::Cancelled` through the middleware chain and to event
subscribers — **never** into session history — so middleware holding external
services can clean up.

Awaiting a cancelled handle yields `OrchestrateError::Cancelled`; `recv()` on a
streaming handle returns `None` once cancelled.
31 changes: 26 additions & 5 deletions docs/concepts/effects.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,13 +28,33 @@ env.effect(|| {
sender is gone — i.e. when the runtime is dropped. You can also call `env.dispose()` manually
to tear down an env you created directly.

## Leak-safe tool registration
## Tool registration is an internal pairing — not an extension API

The safe inverse of `add_tool` is `remove_tool_if_same`: it removes the tool only if the
registered entry is the *same* `Arc` value, so a stale teardown can never delete a replacement
tool that reuses the same name:
Inside the framework, a tool registration is paired with its exact inverse so
teardown is leak-safe: the safe inverse of `add_tool` is `remove_tool_if_same`,
which removes a tool only if the registered entry is the *same* `Arc` value. A
stale teardown can therefore never delete a replacement tool that reuses the
same name. `EnvActor` performs this pairing for the registrations it owns, and
`dispose()` runs the paired inverses in LIFO order.

This registry-level pairing is an **internal implementation detail**, not a
public extension pattern:

- `FuneraEnv::tool_registry` and its mutation methods (`add_tool`,
`remove_tool_if_same`) are crate-private; external code cannot reach them.
- External extensions should not couple their effects to tool-registry state.
If you need tools, register them through the public runtime API —
`with_tool_instance` / `add_tool` / `remove_tool` on `AgentRuntime` — and
let the framework own the registration lifecycle (including disposal).
- Reserve `FuneraEnv::effect` for effects you create yourself (connections,
listeners, resources), whose inverses you control.

For reference, the internal pattern — how the framework's own env pairs a
registration with its inverse (mirrored in the crate's tests):

```rust,ignore
// internal: `tool_registry` and `add_tool` are crate-private — this is how
// the framework pairs its own registrations, not a public extension API.
let tool: Arc<dyn Tool> = Arc::new(MyTool);
env.add_tool(Arc::clone(&tool)).await;

Expand All @@ -50,7 +70,8 @@ env.effect(move || {
});
```

When the env is disposed, exactly this tool is removed — never a later replacement.
When the env is disposed, exactly this tool is removed — never a later
replacement.

## Run the example

Expand Down
2 changes: 1 addition & 1 deletion docs/examples/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,4 +17,4 @@ LLM examples require an API key:

- `minimal`, `multi_turn`, `streaming`, `custom_tool`, `middleware`, `skills`,
`builtin_tools`, `session_reset`, `raw_events`, `multi_runtime`, `sandbox`,
`security`, `streaming_with_tools`
`security`, `streaming_with_tools`, `cancel`
4 changes: 2 additions & 2 deletions docs/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
.system_prompt("You are a helpful assistant.")
.build();

let resp = agent.fire("Hello!", &runtime).await?;
let resp = agent.fire("Hello!", &runtime).await?.await?.await?;
println!("{}", resp.content);
Ok(())
}
Expand All @@ -39,5 +39,5 @@ Examples that do not require an LLM API key:

```bash
cargo run -p funera-orchestrate --example reversible_effects
cargo run -p funera-orchestrate --example tool_policy
cargo run -p funera-orchestrate --example tool_policy --features security
```
1 change: 1 addition & 0 deletions funera-orchestrate/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ thiserror = "2.0.18"
anyhow = "1.0.103"
futures = "0.3.32"
tracing = "0.1"
tokio-util = "0.7.18"

[features]
default = ["deepseek", "tool"]
Expand Down
1 change: 1 addition & 0 deletions funera-orchestrate/examples/builtin_tools.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
// Ask the agent to read Cargo.toml (it will use the Read tool)
let resp = agent
.fire("Read Cargo.toml and tell me the dependencies.", &runtime)
.await?
.await?;
println!("{}", resp.content);

Expand Down
123 changes: 123 additions & 0 deletions funera-orchestrate/examples/cancel.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
//! Cancellation — interrupt an agent call mid-flight.
//!
//! ```bash
//! cargo run -p funera-orchestrate --example cancel
//! # with the middleware demo:
//! cargo run -p funera-orchestrate --example cancel --features middleware
//! ```
//!
//! Every funera entry point returns a cancellable handle:
//!
//! - [`Agent::fire`] / [`Agent::fire_stream`] → one-shot call
//! - [`Agent::send`] / [`Agent::send_stream`] → multi-turn call
//!
//! [`cancel()`](FireStreamHandle::cancel) immediately stops receiving output,
//! abandons in-flight tool executions, and emits [`AgentEvent::Cancelled`] to
//! subscribers — and through the middleware chain when the `middleware`
//! feature is enabled — so middleware holding external services can clean up.
//! **Dropping the handle cancels automatically**: no detached background work
//! is left running.
//!
//! Requires an API key (set `OPENAI_API_KEY`, or call `.api_key()`).

use std::time::Duration;

use funera_orchestrate::{Agent, AgentEvent, AgentRuntime, DeepSeekProvider};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
#[cfg_attr(not(feature = "middleware"), allow(unused_mut))]
let mut builder = AgentRuntime::<DeepSeekProvider>::builder()
.api_key(std::env::var("OPENAI_API_KEY")?)
.model(std::env::var("OPENAI_MODEL").unwrap_or_else(|_| "deepseek-v4-flash".into()));

// A tiny middleware inspector that logs AgentEvent::Cancelled, so the
// "cancellation reaches middleware" contract is visible. Requires the
// `middleware` feature.
#[cfg(feature = "middleware")]
{
use funera_orchestrate::middleware::{InspectorError, InspectorMiddleware};
use funera_orchestrate::middleware_bundle::MiddlewareBundle;

struct CancelLogger;
impl InspectorMiddleware<AgentEvent> for CancelLogger {
fn name(&self) -> &str {
"cancel_logger"
}
fn inspect(&self, event: &AgentEvent) -> Result<(), InspectorError> {
if matches!(event, AgentEvent::Cancelled) {
eprintln!(
"[middleware] saw AgentEvent::Cancelled — clean up external services here"
);
}
Ok(())
}
}

let bundle = MiddlewareBundle::from_chain(
funera_orchestrate::middleware::MiddlewareChain::<AgentEvent>::new()
.with_inspector(CancelLogger),
);
builder = builder.with_middleware_bundle(bundle);
}

let runtime = builder.build()?;
let agent = Agent::builder()
.system_prompt("You are a helpful assistant.")
.build();

// ── 1. Explicit cancel() on a streaming call ─────────────────
let mut stream = agent
.fire_stream("Write a very long essay about Rust.", &runtime)
.await?;

// Subscribe before cancelling so we can observe AgentEvent::Cancelled.
let mut events = agent.subscribe_events();
tokio::time::sleep(Duration::from_millis(200)).await;
stream.cancel();

// recv() returns None once cancelled.
assert!(
stream.recv().await.is_none(),
"recv should end after cancel"
);

// The cancellation event reaches subscribers (and middleware, if enabled).
loop {
match events.recv().await {
Ok(AgentEvent::Cancelled) => {
println!("[1] explicit cancel: subscriber received AgentEvent::Cancelled");
break;
}
Ok(_) => {}
Err(_) => break,
}
}

// Awaiting a cancelled handle yields OrchestrateError::Cancelled.
match stream.await {
Err(funera_orchestrate::OrchestrateError::Cancelled) => {
println!("[1] explicit cancel: wait() reported OrchestrateError::Cancelled");
}
other => println!("[1] unexpected result: {other:?}"),
}

// ── 2. Dropping the handle auto-cancels ─────────────────────
let handle = agent.fire("Explain the borrow checker.", &runtime).await?;
// Drop without awaiting: the background call is cancelled, nothing leaks.
drop(handle);
println!("[2] dropped handle: background call auto-cancelled");

// ── 3. Timeout via tokio::select! ───────────────────────────
let handle = agent.fire("Write a 10-page novel.", &runtime).await?;
let resp = tokio::select! {
r = handle => r?,
_ = tokio::time::sleep(Duration::from_secs(3)) => {
println!("[3] timeout: the fire future was dropped, call cancelled");
return Ok(());
}
};
println!("[3] completed before timeout: {} chars", resp.content.len());

Ok(())
}
1 change: 1 addition & 0 deletions funera-orchestrate/examples/middleware.rs
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ impl InspectorMiddleware<AgentEvent> for EventLogger {
} => {
eprintln!("[log] approval required: {tool_name} — {reason}")
}
AgentEvent::Cancelled => eprintln!("[log] call cancelled"),
}
Ok(())
}
Expand Down
1 change: 1 addition & 0 deletions funera-orchestrate/examples/minimal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
// fire() shares the runtime (&) — no session state is mutated
let resp = agent
.fire("Tell me about Rust programming.", &runtime)
.await?
.await?;

println!("=== Response ===");
Expand Down
Loading
Loading