#![deny(unsafe_code)]
#![warn(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
pub mod billing;
pub mod component_source;
pub mod config;
pub mod config_provider;
pub mod cost;
pub mod dispatch_ledger;
pub mod dw;
pub mod error;
pub mod graph;
pub mod guardrail;
pub mod guardrail_provider;
pub mod http_provider;
pub mod knowledge;
pub mod layered_provider;
pub mod llm;
pub mod llm_credential;
pub mod llm_extension;
#[cfg(feature = "greentic-llm-backend")]
pub mod llm_greentic;
pub mod llm_openai;
pub mod long_term;
pub mod r#loop;
pub mod manifest_provider;
pub mod manifest_tools;
pub mod mcp_local;
pub mod mcp_secrets;
pub mod mcp_source;
pub mod mcp_store_pull;
pub mod memory;
pub mod short_term;
pub mod state;
pub mod state_redis;
pub mod telemetry;
pub mod tenant;
pub mod tools;
#[cfg(feature = "test-mock")]
pub mod mock;
#[cfg(any(test, feature = "test-mock"))]
pub mod test_support;
#[cfg(feature = "serve")]
pub mod serve;
pub use component_source::{
ComponentInvoker, ComponentOperation, ComponentToolCatalog, ComponentToolEntry,
ComponentToolSource,
};
pub use config::{
AgentConfig, AgentLimits, LlmProviderRef, MemoryProviderRef, MemorySettings, ToolRef,
};
pub use config_provider::{CachingConfigProvider, ConfigProvider, InMemoryConfigProvider};
#[cfg(feature = "test-mock")]
pub use cost::MockTokenMeter;
pub use cost::{RedisTokenMeter, TokenMeter};
pub use dispatch_ledger::{DispatchLedger, NoopDispatchLedger, RedisDispatchLedger};
pub use error::{AgentError, ConfigError, LlmError, MemoryError, StateError, TerminationReason};
pub use graph::http_provider::{CachingGraphProvider, HttpGraphProvider};
pub use http_provider::HttpConfigProvider;
pub use layered_provider::LayeredConfigProvider;
pub use llm::{LlmBackend, LlmRequest, LlmResponse, RetryingLlmBackend};
pub use llm_extension::{
BridgeCredential, ExtensionLlmBackend, LlmExtensionInvoker, RuntimeInvoker,
};
#[cfg(feature = "greentic-llm-backend")]
pub use llm_greentic::GreenticLlmBackend;
pub use llm_openai::{OpenAiLlmBackend, encode_tool_name, split_tool_name};
pub use long_term::{
EpisodeIngest, EpisodeSource, IngestOutcome, LongTermMemory, LongTermMemoryError, RecallQuery,
RecalledFact,
};
pub use manifest_provider::ManifestToolOverlayProvider;
pub use mcp_source::{
MCP_ROLE_AGENTIC_WORKER, MCP_ROLE_FLOW_EDITOR, McpRoute, McpToolCatalog, McpToolEntry,
McpToolSource, dispatch_route,
};
pub use memory::{InMemoryMemoryProvider, MemoryProvider, MemoryQuery, MemoryRecord};
pub use state::{AgentStateStore, ChatMessage, ConversationState, SessionLock};
pub use state_redis::RedisAgentStateStore;
pub use telemetry::{OtelTelemetry, StepTelemetryCtx, Telemetry};
pub use tenant::TenantContext;
pub use tools::{RedisToolLedger, ToolLedger};
use std::sync::Arc;
pub trait StepObserver: Send + Sync {
fn wants_streaming(&self) -> bool {
false
}
fn on_token_delta(&self, _chunk: &str) {}
fn on_tool_call(&self, _name: &str, _call_id: &str) {}
fn on_tool_result(&self, _name: &str, _call_id: &str, _result: &serde_json::Value) {}
}
pub struct NoopStepObserver;
impl StepObserver for NoopStepObserver {}
pub struct AgentRuntime {
pub(crate) config_provider: Arc<dyn ConfigProvider>,
pub(crate) state_store: Arc<dyn AgentStateStore>,
pub(crate) ext_runtime: Arc<greentic_ext_runtime::ExtensionRuntime>,
pub(crate) llm: Arc<dyn LlmBackend>,
pub(crate) telemetry: Arc<dyn Telemetry>,
pub(crate) token_meter: Arc<dyn TokenMeter>,
pub(crate) billing_meter: Arc<dyn crate::billing::BillingMeter>,
pub(crate) ledger: Arc<dyn ToolLedger>,
pub(crate) mcp: Option<Arc<crate::mcp_source::McpToolSource>>,
pub(crate) guardrail_policy: Arc<dyn crate::guardrail::GuardrailPolicy>,
pub(crate) guardrail_evaluator: Arc<dyn crate::guardrail::GuardrailEvaluator>,
pub(crate) components: Option<Arc<crate::component_source::ComponentToolSource>>,
pub(crate) long_term_memory: Option<Arc<dyn long_term::LongTermMemory>>,
pub(crate) knowledge: Option<Arc<dyn knowledge::Knowledge>>,
pub(crate) short_term_memory: Option<Arc<dyn crate::memory::MemoryProvider>>,
}
impl AgentRuntime {
#[allow(clippy::too_many_arguments)]
pub fn new(
config_provider: Arc<dyn ConfigProvider>,
state_store: Arc<dyn AgentStateStore>,
ext_runtime: Arc<greentic_ext_runtime::ExtensionRuntime>,
llm: Arc<dyn LlmBackend>,
telemetry: Arc<dyn Telemetry>,
token_meter: Arc<dyn TokenMeter>,
ledger: Arc<dyn ToolLedger>,
mcp: Option<Arc<crate::mcp_source::McpToolSource>>,
) -> Self {
Self {
config_provider,
state_store,
ext_runtime,
llm,
telemetry,
token_meter,
billing_meter: Arc::new(crate::billing::NoopBillingMeter),
ledger,
mcp,
guardrail_policy: Arc::new(crate::guardrail::NoMandatoryGuardrails),
guardrail_evaluator: Arc::new(crate::guardrail::AcceptAllEvaluator),
components: None,
long_term_memory: None,
knowledge: None,
short_term_memory: None,
}
}
pub fn with_guardrails(
mut self,
policy: Arc<dyn crate::guardrail::GuardrailPolicy>,
evaluator: Arc<dyn crate::guardrail::GuardrailEvaluator>,
) -> Self {
self.guardrail_policy = policy;
self.guardrail_evaluator = evaluator;
self
}
#[must_use]
pub fn with_billing_meter(mut self, meter: Arc<dyn crate::billing::BillingMeter>) -> Self {
self.billing_meter = meter;
self
}
#[must_use]
pub fn with_component_source(
mut self,
components: Option<Arc<crate::component_source::ComponentToolSource>>,
) -> Self {
self.components = components;
self
}
#[must_use]
pub fn with_long_term_memory(mut self, memory: Arc<dyn long_term::LongTermMemory>) -> Self {
self.long_term_memory = Some(memory);
self
}
#[must_use]
pub fn with_short_term_memory(
mut self,
memory: Arc<dyn crate::memory::MemoryProvider>,
) -> Self {
self.short_term_memory = Some(memory);
self
}
pub async fn remember_episode(
&self,
tenant: &TenantContext,
episode: long_term::EpisodeIngest,
) -> Result<long_term::IngestOutcome, long_term::LongTermMemoryError> {
let memory = self.long_term_memory.as_ref().ok_or_else(|| {
long_term::LongTermMemoryError::NotConfigured("long-term memory not wired".into())
})?;
let ctx = long_term::to_types_tenant(tenant)?;
memory.ingest_episode(&ctx, episode).await
}
pub async fn recall_long_term(
&self,
tenant: &TenantContext,
query: long_term::RecallQuery,
) -> Result<Vec<long_term::RecalledFact>, long_term::LongTermMemoryError> {
let memory = self.long_term_memory.as_ref().ok_or_else(|| {
long_term::LongTermMemoryError::NotConfigured("long-term memory not wired".into())
})?;
let ctx = long_term::to_types_tenant(tenant)?;
memory.recall(&ctx, query).await
}
#[must_use]
pub fn with_knowledge(mut self, knowledge: Arc<dyn knowledge::Knowledge>) -> Self {
self.knowledge = Some(knowledge);
self
}
pub async fn search_knowledge(
&self,
tenant: &TenantContext,
query: knowledge::KnowledgeQuery,
) -> knowledge::KnowledgeResult<Vec<knowledge::RetrievedChunk>> {
let kb = self
.knowledge
.as_ref()
.ok_or(knowledge::KnowledgeError::NotConfigured)?;
let ctx = knowledge::to_types_tenant(tenant)?;
kb.search(&ctx, query).await
}
pub async fn step(
&self,
tenant: TenantContext,
session_id: &str,
agent_id: &str,
message: AgentInput,
) -> Result<AgentOutput, AgentError> {
self.step_with_observer(
tenant,
session_id,
agent_id,
message,
Arc::new(NoopStepObserver),
)
.await
}
pub async fn step_with_observer(
&self,
tenant: TenantContext,
session_id: &str,
agent_id: &str,
message: AgentInput,
observer: Arc<dyn StepObserver>,
) -> Result<AgentOutput, AgentError> {
r#loop::run_step(self, tenant, session_id, agent_id, message, observer).await
}
}
#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
pub struct AgentInput {
pub text: String,
}
#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
pub struct AgentOutput {
pub reply: String,
pub trail: Vec<AgentStep>,
pub terminated_by: TerminationReason,
}
#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum AgentStep {
ToolCall {
name: String,
call_id: String,
result: serde_json::Value,
},
ToolCallReused {
name: String,
call_id: String,
},
ToolCallBlocked {
name: String,
reason: String,
},
Reply {
text: String,
},
}