use std::sync::atomic::{AtomicUsize, Ordering};
use funera_orchestrate::middleware::*;
use funera_orchestrate::middleware_bundle::MiddlewareBundle;
use funera_orchestrate::{Agent, AgentEvent, AgentRuntime, DeepSeekProvider};
struct EventLogger;
impl InspectorMiddleware<AgentEvent> for EventLogger {
fn name(&self) -> &str {
"event_logger"
}
fn inspect(&self, event: &AgentEvent) -> Result<(), InspectorError> {
match event {
AgentEvent::Text(t) => eprintln!("[log] token: {t}"),
AgentEvent::Reasoning(r) => eprintln!("[log] reasoning: {r}"),
AgentEvent::ToolCallRequest { name, args, .. } => {
eprintln!("[log] tool_call: {name}({args})")
}
AgentEvent::ToolCallResult { name, result, .. } => {
let status = match result {
Ok(_) => "ok",
Err(_) => "err",
};
eprintln!("[log] tool_result: {name} ({status})");
}
AgentEvent::TurnStart => eprintln!("[log] --- turn start ---"),
AgentEvent::TurnEnd { .. } => eprintln!("[log] --- turn end ---"),
AgentEvent::Done => eprintln!("[log] done"),
AgentEvent::Error(e) => eprintln!("[log] error: {e}"),
AgentEvent::ToolApprovalRequired {
tool_name, reason, ..
} => {
eprintln!("[log] approval required: {tool_name} — {reason}")
}
}
Ok(())
}
}
struct TurnCounter {
count: AtomicUsize,
}
impl InspectorMiddleware<AgentEvent> for TurnCounter {
fn name(&self) -> &str {
"turn_counter"
}
fn inspect(&self, event: &AgentEvent) -> Result<(), InspectorError> {
if matches!(event, AgentEvent::TurnStart) {
let n = self.count.fetch_add(1, Ordering::Relaxed) + 1;
eprintln!("[counter] turn #{n}");
}
Ok(())
}
}
struct Censor {
words: Vec<&'static str>,
}
impl MutatorMiddleware<AgentEvent> for Censor {
fn name(&self) -> &str {
"censor"
}
fn process(&self, event: AgentEvent) -> MutatorAction<AgentEvent> {
match event {
AgentEvent::Text(t) => {
let orig = t.clone();
let mut result = t;
for w in &self.words {
result = result.replace(w, "***");
}
if result != orig {
eprintln!("[censor] modified token (sensitive word detected)");
MutatorAction::Modify(AgentEvent::Text(result))
} else {
MutatorAction::Pass
}
}
_e => MutatorAction::Pass,
}
}
}
struct BlockTool {
tool_name: String,
}
impl MutatorMiddleware<AgentEvent> for BlockTool {
fn name(&self) -> &str {
"block_tool"
}
fn process(&self, event: AgentEvent) -> MutatorAction<AgentEvent> {
if let AgentEvent::ToolCallRequest { name, .. } = &event
&& self.tool_name.eq_ignore_ascii_case(name)
{
eprintln!("[block_tool] blocked tool call: {name}");
return MutatorAction::Block {
reason: format!("tool '{}' is blocked", name),
};
}
MutatorAction::Pass
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let bundle = MiddlewareBundle::from_chain(
MiddlewareChain::<AgentEvent>::new()
.with_inspectors((
EventLogger,
TurnCounter {
count: AtomicUsize::new(0),
},
))
.with_mutator(Censor {
words: vec!["secret", "password"],
})
.with_mutator(BlockTool {
tool_name: "shell".into(),
}),
);
let runtime = AgentRuntime::<DeepSeekProvider>::builder()
.api_key(std::env::var("OPENAI_API_KEY")?)
.with_middleware_bundle(bundle)
.build()?;
let agent = Agent::builder()
.system_prompt("You are a helpful assistant. You have access to shell and other tools.")
.build();
let (_runtime, response) = agent
.send("Say hello! Also, my password is 'my_secret_123'.", runtime)
.await?
.await?;
println!("Response: {}", response.content);
Ok(())
}