use super::messages::StreamMsg;
use crate::agent::runtime::PromptOverrides;
use crate::approval_facts::ApprovalFacts;
use crate::config::runtime::RuntimeConfig;
use crate::grant_token::{TurnPrimary, grant_token};
use crate::interactive::session_universe::SessionUniverse;
use async_trait::async_trait;
use saya_agent::{
AgentEvent, AgentEventSink, AgentMode, ApprovalChoice, ApprovalDecider, ApprovalDecision,
ApprovalPolicy, ChatMessage, SessionPolicy, ToolDefinition,
};
use saya_store::SqliteStateStore;
use std::sync::Arc;
use tokio::sync::mpsc::UnboundedSender;
use tokio::sync::oneshot;
pub(crate) struct ChannelSink {
pub(crate) tx: UnboundedSender<StreamMsg>,
}
#[async_trait]
impl AgentEventSink for ChannelSink {
async fn emit(&self, event: AgentEvent) {
let _ = self.tx.send(StreamMsg::Event(event));
}
}
pub(crate) struct ChannelApproval {
tx: UnboundedSender<StreamMsg>,
policy: SessionPolicy,
primary: TurnPrimary,
facts: ApprovalFacts,
journal: Option<Arc<saya_store::SessionJournal>>,
}
impl ChannelApproval {
pub(crate) fn new(
tx: UnboundedSender<StreamMsg>,
policy: SessionPolicy,
primary: TurnPrimary,
facts: ApprovalFacts,
journal: Option<Arc<saya_store::SessionJournal>>,
) -> Self {
Self {
tx,
policy,
primary,
facts,
journal,
}
}
}
#[async_trait]
impl ApprovalDecider for ChannelApproval {
async fn approve(&self, tool: &ToolDefinition, arguments: &serde_json::Value) -> bool {
if crate::interactive::session_deny::denied_call_program(
&tool.name,
arguments,
&self.facts.denied_programs,
)
.is_some()
{
return false;
}
let primary = self.primary.get();
let grant = grant_token(&tool.name, arguments, primary.as_deref(), &self.facts);
match self.policy.resolve(&tool.effect, grant.as_deref()) {
ApprovalDecision::Allow => true,
ApprovalDecision::Deny { .. } => false,
ApprovalDecision::Ask => {
let (respond, answer) = oneshot::channel();
let detail = crate::approval_facts::call_facts(
&tool.name,
arguments,
grant.as_deref(),
&self.facts,
primary.as_deref(),
Some(self.policy.grants()),
);
if self
.tx
.send(StreamMsg::ApprovalRequest {
tool: tool.name.clone(),
detail,
grant,
respond,
})
.is_err()
{
return false;
}
match answer.await {
Ok(choice) => {
let (_, warning) = crate::interactive::session_grants::record_prompt_answer(
&self.policy,
&choice,
self.journal.as_deref(),
);
if let Some(warning) = warning {
let _ = self.tx.send(StreamMsg::Notice(warning));
}
!matches!(choice, ApprovalChoice::Deny)
}
Err(_) => false,
}
}
}
}
fn refusal_detail(
&self,
tool: &ToolDefinition,
arguments: &serde_json::Value,
) -> Option<String> {
crate::interactive::session_deny::denied_call_program(
&tool.name,
arguments,
&self.facts.denied_programs,
)
.map(|program| crate::interactive::session_deny::denied_refusal(&program))
}
}
pub(crate) struct StreamRequest {
pub(crate) runtime: Arc<RuntimeConfig>,
pub(crate) prompt: String,
pub(crate) approval: ApprovalPolicy,
pub(crate) policy: SessionPolicy,
pub(crate) overrides: PromptOverrides,
pub(crate) history: Vec<ChatMessage>,
pub(crate) state_db: SqliteStateStore,
pub(crate) last_sql: Option<String>,
pub(crate) session: Arc<SessionUniverse>,
pub(crate) journal: Option<Arc<saya_store::SessionJournal>>,
pub(crate) agent_mode: AgentMode,
}
pub(crate) fn approval_capabilities() -> (bool, bool) {
(false, true)
}