use std::{collections::HashMap, sync::Arc};
use connectrpc::{ConnectError, ErrorCode};
use futures::channel::mpsc;
use polyc_agent::{
DelegateRecord, DispatchRecorder, HandoffRequest, PendingApproval, RunTurnOptions,
ToolExecutor, TurnResult, TurnStreamEvent, UnattendedDenial, llm_stop_to_wire_i32,
retry::Clock, run_turn_with, wire_to_llm,
};
use polyc_llm::{
CacheHint, DynProvider, LlmError, LlmErrorKind, LlmProvider, Message as LlmMessage, Usage,
into_dyn, turn::StubProvider,
};
use polyc_proto::humanize_tool_name;
use polyc_proto::proto::polychrome::agent::v1::Message as WireMessage;
use polyc_proto::proto::polychrome::harness::v1::{
ApprovalResponse as WireApprovalResponse, DelegateDescriptor as WireDelegateDescriptor,
DelegateFact as WireDelegateFact, ExecutionLabel as WireExecutionLabel,
HandoffRequest as WireHandoffRequest, HarnessMessage, PendingApproval as WirePendingApproval,
PendingQuestion as WirePendingQuestion, QuestionAnswer as WireQuestionAnswer,
QuestionOptionWire as WireQuestionOptionWire, StepOutcomeProposal as WireStepOutcomeProposal,
StopReason as WireStopReason, TurnBatch as WireTurnBatch, TurnFailure as WireTurnFailure,
TurnFailureKind as WireTurnFailureKind, TurnInput as WireTurnInput,
UnattendedDenialFact as WireUnattendedDenialFact, Usage as WireUsage,
harness_message::Frame as HarnessFrame,
};
pub mod broker;
pub mod execution;
pub mod model;
use execution::ExecutionLabel;
#[must_use]
pub fn frame(label: &ExecutionLabel, payload: impl Into<HarnessFrame>) -> HarnessMessage {
HarnessMessage {
frame: Some(payload.into()),
label: buffa::MessageField::some(WireExecutionLabel::from(label)),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
#[derive(Clone)]
pub struct RegisteredProvider {
pub provider: Arc<DynProvider>,
pub default_model: String,
}
pub struct Resolved {
pub provider: Arc<DynProvider>,
pub provider_name: String,
pub model: String,
}
#[derive(Clone)]
pub struct ProviderSet {
providers: Arc<HashMap<String, RegisteredProvider>>,
default_provider: String,
}
impl ProviderSet {
#[must_use]
pub fn new(
mut providers: HashMap<String, RegisteredProvider>,
default_provider: impl Into<String>,
) -> Self {
providers
.entry("stub".to_owned())
.or_insert_with(|| RegisteredProvider {
provider: into_dyn(StubProvider),
default_model: "stub".to_owned(),
});
let requested_default = default_provider.into();
let default_provider = if providers.contains_key(&requested_default) {
requested_default
} else {
tracing::error!(
requested = %requested_default,
"configured default provider is not registered; falling back to stub"
);
"stub".to_owned()
};
Self {
providers: Arc::new(providers),
default_provider,
}
}
#[must_use]
pub fn proxy_only(
provider_name: impl Into<String>,
provider: Arc<DynProvider>,
default_model: impl Into<String>,
) -> Self {
let provider_name = provider_name.into();
let mut providers = HashMap::new();
providers.insert(
provider_name.clone(),
RegisteredProvider {
provider,
default_model: default_model.into(),
},
);
Self {
providers: Arc::new(providers),
default_provider: provider_name,
}
}
#[must_use]
pub fn backend(&self, name: &str) -> Option<(Arc<DynProvider>, String)> {
self.providers
.get(name)
.map(|reg| (Arc::clone(®.provider), reg.default_model.clone()))
}
#[must_use]
pub fn resolve(&self, provider: &str, model: &str) -> Resolved {
let requested = if provider.is_empty() {
self.default_provider.as_str()
} else {
provider
};
let (name, reg) = self.providers.get(requested).map_or_else(
|| {
tracing::warn!(
requested = %requested,
default = %self.default_provider,
registered = ?self.providers.keys().collect::<Vec<_>>(),
"unknown provider selector; falling back to the default provider"
);
let reg = self
.providers
.get(self.default_provider.as_str())
.expect("default_provider is a registered key (ProviderSet::new invariant)");
(self.default_provider.as_str(), reg)
},
|reg| (requested, reg),
);
let model = if model.is_empty() {
reg.default_model.clone()
} else {
model.to_owned()
};
Resolved {
provider: reg.provider.clone(),
provider_name: name.to_owned(),
model,
}
}
}
impl Default for ProviderSet {
fn default() -> Self {
Self::new(HashMap::new(), "stub")
}
}
#[derive(Debug, Default)]
pub struct VerifiedApprovals {
decisions: Vec<polyc_agent::ApprovalDecision>,
session_approved: HashMap<String, polyc_capability::CapabilitySet>,
}
impl VerifiedApprovals {
#[must_use]
pub fn approved_count(&self) -> usize {
self.decisions
.iter()
.filter(|decision| decision.approved)
.count()
}
#[must_use]
pub fn denied_count(&self) -> usize {
self.decisions
.iter()
.filter(|decision| !decision.approved)
.count()
}
#[must_use]
pub fn session_approved_count(&self) -> usize {
self.session_approved.len()
}
}
#[must_use]
pub fn verify_approval_responses(
responses: &[WireApprovalResponse],
current_caller: &str,
current_sandbox_mode: &str,
) -> VerifiedApprovals {
let mut out = VerifiedApprovals::default();
for resp in responses {
if !polyc_crypto::approval::verify_wire_response(
&resp.request_id,
&resp.tool_name,
&resp.args_json,
&resp.modified_args_json,
resp.approved,
resp.approved_for_session,
&resp.covered_capabilities,
&resp.caller,
&resp.approver,
&resp.sandbox_mode,
&resp.reason,
&resp.injected_context,
&resp.conversation_id,
&resp.nonce,
&resp.turn_id,
resp.routine_grant,
&resp.tool_descriptor_hash,
&resp.grant_scope,
&resp.signer_pk_hex,
&resp.signature_hex,
) {
tracing::warn!(
request_id = %resp.request_id,
"rejected approval_response: signature failed to verify"
);
continue;
}
if resp.approved {
if polyc_crypto::approval::is_session_grant_for(
resp.approved,
resp.approved_for_session,
&resp.caller,
current_caller,
) && resp.sandbox_mode == current_sandbox_mode
{
let (covered, unknown) = polyc_capability::CapabilitySet::from_names(
resp.covered_capabilities.iter().map(String::as_str),
);
if !unknown.is_empty() {
tracing::warn!(
tool = %resp.tool_name,
?unknown,
"session grant carries unknown covered-capability names; ignoring them"
);
}
let entry = out
.session_approved
.entry(resp.tool_name.clone())
.or_default();
*entry = entry.union(covered);
}
if !resp.already_executed {
out.decisions
.push(polyc_agent::ApprovalDecision::from(resp));
}
} else {
tracing::info!(
request_id = %resp.request_id,
"approval_response verified but denied"
);
out.decisions
.push(polyc_agent::ApprovalDecision::from(resp));
}
}
out
}
#[must_use]
pub fn verify_routine_tool_grants(
grants: &[WireApprovalResponse],
current_caller: &str,
current_sandbox_mode: &str,
expected_fire_conversation: &str,
) -> polyc_agent::RoutineGrantSet {
let mut out = polyc_agent::RoutineGrantSet::default();
for resp in grants {
if !resp.routine_grant {
continue;
}
if !polyc_crypto::approval::verify_wire_response(
&resp.request_id,
&resp.tool_name,
&resp.args_json,
&resp.modified_args_json,
resp.approved,
resp.approved_for_session,
&resp.covered_capabilities,
&resp.caller,
&resp.approver,
&resp.sandbox_mode,
&resp.reason,
&resp.injected_context,
&resp.conversation_id,
&resp.nonce,
&resp.turn_id,
resp.routine_grant,
&resp.tool_descriptor_hash,
&resp.grant_scope,
&resp.signer_pk_hex,
&resp.signature_hex,
) {
tracing::warn!(
tool = %resp.tool_name,
"rejected routine_tool_grant: signature failed to verify"
);
continue;
}
if resp.conversation_id != expected_fire_conversation {
tracing::warn!(
tool = %resp.tool_name,
"rejected routine_tool_grant: not bound to this turn's fire conversation"
);
continue;
}
if !polyc_crypto::approval::is_session_grant_for(
resp.approved,
resp.approved_for_session,
&resp.caller,
current_caller,
) || resp.sandbox_mode != current_sandbox_mode
{
tracing::warn!(
tool = %resp.tool_name,
"rejected routine_tool_grant: not bound to this turn's caller or sandbox mode"
);
continue;
}
let (covered, unknown) = polyc_capability::CapabilitySet::from_names(
resp.covered_capabilities.iter().map(String::as_str),
);
if !unknown.is_empty() {
tracing::warn!(
tool = %resp.tool_name,
?unknown,
"routine_tool_grant carries unknown covered-capability names; ignoring them"
);
}
match resp.grant_scope.as_str() {
"tool" => {
out.per_tool.insert(
resp.tool_name.clone(),
polyc_agent::routine_grant::PerToolGrant {
covered,
descriptor_hash: resp.tool_descriptor_hash.clone(),
},
);
}
"blanket_below_high" | "blanket_all" => {
out.blanket = Some(polyc_agent::routine_grant::BlanketGrant {
include_high: resp.grant_scope == "blanket_all",
});
}
_ => {
tracing::warn!(
tool = %resp.tool_name,
grant_scope = %resp.grant_scope,
"rejected routine_tool_grant: unrecognized grant_scope"
);
}
}
}
out
}
#[must_use]
pub fn verify_question_answers(
answers: &[WireQuestionAnswer],
) -> Vec<polyc_agent::question::VerifiedAnswer> {
use polyc_agent::question::AnswerState;
use std::str::FromStr as _;
let mut out = Vec::new();
for a in answers {
let Ok(state) = AnswerState::from_str(&a.state) else {
tracing::warn!(
call_id = %a.call_id,
index = a.index,
state = %a.state,
"rejected question_response: unrecognized state"
);
continue;
};
let selected_index = (!matches!(state, AnswerState::Declined)).then_some(a.selected_index);
if !polyc_crypto::question::verify_wire_answer(
&a.turn_id,
&a.call_id,
a.index,
&a.question_args_json,
&a.state,
selected_index,
&a.selected_label,
&a.answered_by,
&a.conversation_id,
&a.nonce,
&a.signer_pk_hex,
&a.signature_hex,
) {
tracing::warn!(
call_id = %a.call_id,
index = a.index,
"rejected question_response: signature failed to verify"
);
continue;
}
out.push(polyc_agent::question::VerifiedAnswer {
turn_id: a.turn_id.clone(),
call_id: a.call_id.clone(),
index: a.index,
state,
selected_index,
selected_label: a.selected_label.clone(),
answered_by: a.answered_by.clone(),
});
}
out
}
#[must_use]
pub fn current_sandbox_mode() -> String {
polyc_tools::current_sandbox_mode()
}
#[must_use]
pub fn wire_title(agent_title: &str, name: &str) -> String {
if agent_title.is_empty() {
humanize_tool_name(name)
} else {
agent_title.to_owned()
}
}
#[must_use]
pub fn connect_error_from_turn(kind: LlmErrorKind, msg: String) -> ConnectError {
match kind {
LlmErrorKind::RateLimit => ConnectError::resource_exhausted(msg),
LlmErrorKind::Timeout => ConnectError::deadline_exceeded(msg),
LlmErrorKind::Unavailable => ConnectError::unavailable(msg),
LlmErrorKind::Auth => ConnectError::new(ErrorCode::Unauthenticated, msg),
LlmErrorKind::BadRequest => ConnectError::invalid_argument(msg),
LlmErrorKind::Ambiguous => ConnectError::new(ErrorCode::Unknown, msg),
LlmErrorKind::Other => ConnectError::internal(msg),
}
}
#[must_use]
pub fn classify_turn_error(kind: LlmErrorKind, message: String) -> CapturedTurnError {
let kind = match kind {
LlmErrorKind::RateLimit => WireTurnFailureKind::RateLimit,
LlmErrorKind::Timeout => WireTurnFailureKind::Timeout,
LlmErrorKind::Unavailable => WireTurnFailureKind::Unavailable,
LlmErrorKind::Auth => WireTurnFailureKind::Auth,
LlmErrorKind::BadRequest => WireTurnFailureKind::BadRequest,
LlmErrorKind::Other => WireTurnFailureKind::Other,
LlmErrorKind::Ambiguous => return CapturedTurnError::Ambiguous,
};
CapturedTurnError::Failed(WireTurnFailure {
kind: kind.into(),
message,
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
}
#[derive(Debug, Clone)]
pub enum CapturedTurnError {
Failed(WireTurnFailure),
Ambiguous,
}
impl CapturedTurnError {
#[must_use]
pub fn connect_error(&self) -> ConnectError {
match self {
Self::Ambiguous => ConnectError::new(
ErrorCode::Unknown,
"the model attempt's outcome is not known, so this turn records none",
),
Self::Failed(failure) => {
let message = failure.message.clone();
match failure.kind.as_known() {
Some(WireTurnFailureKind::RateLimit) => {
ConnectError::resource_exhausted(message)
}
Some(WireTurnFailureKind::Timeout) => ConnectError::deadline_exceeded(message),
Some(WireTurnFailureKind::Unavailable) => ConnectError::unavailable(message),
Some(WireTurnFailureKind::Auth) => {
ConnectError::new(ErrorCode::Unauthenticated, message)
}
Some(WireTurnFailureKind::BadRequest) => {
ConnectError::invalid_argument(message)
}
_ => ConnectError::internal(message),
}
}
}
}
}
pub async fn run_turn_captured(
provider: &(impl LlmProvider + ?Sized),
tools: &(impl ToolExecutor + ?Sized),
model: &str,
input: Vec<LlmMessage>,
options: RunTurnOptions,
) -> Result<TurnResult, CapturedTurnError> {
let turn = run_turn_with(provider, tools, model, input, options)
.await
.map_err(|e| classify_turn_error(e.kind(), e.to_string()))?;
if let Some(failure) = turn.mid_stream_failure {
return Err(classify_turn_error(failure.kind, failure.message));
}
Ok(turn)
}
pub async fn run_turn_captured_under_grant(
provider: &(impl LlmProvider + ?Sized),
tools: &(impl ToolExecutor + ?Sized),
model: &str,
input: Vec<LlmMessage>,
options: RunTurnOptions,
granted: polyc_capability::CapabilitySet,
) -> Result<TurnResult, CapturedTurnError> {
polyc_agent::with_execution_capabilities(
granted,
run_turn_captured(provider, tools, model, input, options),
)
.await
}
#[allow(clippy::struct_excessive_bools)]
#[derive(Clone, Copy, Debug, Default)]
pub struct TurnToggles {
pub escalate_sandbox_denials: bool,
pub untrusted_context: bool,
pub unattended: bool,
pub escape_hatch: bool,
pub fire_dispatch: bool,
pub deployment_approve_all_dangerous: bool,
}
impl TurnToggles {
#[must_use]
pub fn from_wire(turn_input: &WireTurnInput) -> Self {
Self {
escalate_sandbox_denials: turn_input.escalate_sandbox_denials,
untrusted_context: turn_input.untrusted_context,
unattended: turn_input.unattended,
escape_hatch: turn_input
.retrieval
.as_option()
.is_some_and(|cfg| cfg.escape_hatch),
fire_dispatch: turn_input.fire_dispatch,
deployment_approve_all_dangerous: turn_input.deployment_approve_all_dangerous,
}
}
}
#[must_use]
pub fn resolve_delegate_descriptors(
providers: &ProviderSet,
tools: &(impl ToolExecutor + ?Sized),
wire: &[WireDelegateDescriptor],
) -> Vec<polyc_agent::DelegateDescriptor> {
let all_specs = tools.specs();
wire.iter()
.map(|d| {
let resolved = providers.resolve(&d.provider, &d.model);
let tool_specs: Vec<_> = all_specs
.iter()
.filter(|s| {
d.builtin_tools.iter().any(|name| name == &s.name)
|| d.tools_enabled.iter().any(|label| {
s.name
.strip_prefix(label.as_str())
.and_then(|rest| rest.strip_prefix("__"))
.is_some()
})
})
.cloned()
.collect();
polyc_agent::DelegateDescriptor {
agent_id: d.agent_id.clone(),
instructions: (!d.instructions.is_empty()).then(|| d.instructions.clone()),
provider: resolved.provider,
provider_name: resolved.provider_name,
model: resolved.model,
tool_specs,
max_steps: usize::try_from(d.max_steps).unwrap_or(usize::MAX),
native_search_allowed: d
.builtin_tools
.iter()
.any(|n| n == polyc_tools::web::NATIVE_SEARCH_GROUNDING),
share_in: polyc_agent::delegate::ShareInCeiling {
allow: d.share_in_allow.clone(),
max_files: usize::try_from(d.share_in_max_files).unwrap_or(usize::MAX),
max_bytes: d.share_in_max_bytes,
},
}
})
.collect()
}
#[must_use]
#[allow(clippy::too_many_arguments)] #[allow(clippy::implicit_hasher)] pub fn run_turn_options(
approvals: VerifiedApprovals,
routine_tool_grants: polyc_agent::RoutineGrantSet,
toggles: TurnToggles,
prompt_cache_key: String,
step_budget: u32,
stream_tx: Option<mpsc::Sender<TurnStreamEvent>>,
dispatch_recorder: Option<Arc<dyn DispatchRecorder>>,
clock: Option<Arc<dyn Clock + Send + Sync>>,
native_search_allowed: bool,
question_answers: Vec<polyc_agent::question::VerifiedAnswer>,
turn_start_unix_ms: u64,
) -> RunTurnOptions {
let routine_tool_grants = if toggles.fire_dispatch && toggles.deployment_approve_all_dangerous {
polyc_agent::RoutineGrantSet {
blanket: Some(polyc_agent::routine_grant::BlanketGrant { include_high: true }),
..routine_tool_grants
}
} else {
routine_tool_grants
};
RunTurnOptions {
approval_decisions: approvals.decisions,
session_approved_tools: approvals.session_approved,
routine_tool_grants,
stream_tx,
max_steps: (step_budget > 0).then(|| {
let cap = usize::try_from(step_budget).unwrap_or(usize::MAX);
cap.min(polyc_agent::resolve_default_max_steps())
}),
native_search_allowed,
escalate_sandbox_denials: toggles.escalate_sandbox_denials,
untrusted_context_seed: toggles.untrusted_context,
unattended: toggles.unattended,
escape_hatch: toggles.escape_hatch,
fire_dispatch: toggles.fire_dispatch,
dispatch_recorder,
cache_hint: CacheHint::from_key(prompt_cache_key),
clock,
delegate_descriptors: Vec::new(),
delegate_max_fanout: None,
delegate_turn_budget: None,
is_delegated_worker: false,
question_answers,
turn_start_unix_ms: (turn_start_unix_ms != 0).then_some(turn_start_unix_ms),
}
}
fn wire_pending_approval(p: PendingApproval) -> WirePendingApproval {
let title = wire_title(&p.title, &p.name);
WirePendingApproval {
turn_id: p.occurrence_turn_id.unwrap_or_default(),
id: p.id,
tool_name: p.name,
args_json: p.args_json,
title,
sandbox_mode: current_sandbox_mode(),
reason: p.reason,
missing_capabilities: p.missing_capabilities,
tool_descriptor_hash: p.tool_descriptor_hash,
required_capabilities: p.required_capabilities,
fire_dispatch: p.fire_dispatch,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn wire_pending_question(p: polyc_agent::question::PendingQuestion) -> WirePendingQuestion {
WirePendingQuestion {
turn_id: p.occurrence_turn_id.unwrap_or_default(),
call_id: p.call_id,
index: p.index,
header: p.item.header,
question: p.item.question,
options: p
.item
.options
.into_iter()
.map(|o| WireQuestionOptionWire {
label: o.label,
description: o.description,
recommended: o.recommended,
__buffa_unknown_fields: buffa::UnknownFields::default(),
})
.collect(),
args_json: p.args_json,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn wire_unattended_denial_fact(d: UnattendedDenial) -> WireUnattendedDenialFact {
WireUnattendedDenialFact {
tool: d.tool,
args_json: d.args_json,
missing_capabilities: d.missing_capabilities,
reason: d.reason,
tool_descriptor_hash: d.tool_descriptor_hash,
required_capabilities: d.required_capabilities,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn wire_delegate_fact(d: DelegateRecord) -> WireDelegateFact {
WireDelegateFact {
sub_agent_id: d.sub_agent_id,
target_agent_id: d.target_agent_id,
task: d.task,
resolved_provider: d.resolved_provider,
resolved_model: d.resolved_model,
input_tokens: d.usage.input_tokens,
output_tokens: d.usage.output_tokens,
succeeded: d.succeeded,
error: d.error,
first_party: d.first_party,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn wire_handoff_request(h: HandoffRequest) -> WireHandoffRequest {
WireHandoffRequest {
child_agent_id: h.child_agent_id,
reason: h.reason,
max_carry: u32::try_from(h.max_carry).unwrap_or(u32::MAX),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn wire_usage(u: Usage) -> WireUsage {
WireUsage {
input_tokens: u.input_tokens,
output_tokens: u.output_tokens,
cache_read_input_tokens: u.cache_read_input_tokens,
cache_creation_input_tokens: u.cache_creation_input_tokens,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
#[must_use]
pub fn build_step_outcome_proposal(
label: &ExecutionLabel,
messages: Vec<WireMessage>,
) -> HarnessMessage {
frame(
label,
WireStepOutcomeProposal {
messages,
__buffa_unknown_fields: buffa::UnknownFields::default(),
},
)
}
#[must_use]
pub fn build_final_batch(label: &ExecutionLabel, turn: TurnResult) -> HarnessMessage {
tracing::info!(
output_messages = turn.messages.len(),
input_tokens = turn.usage.input_tokens,
output_tokens = turn.usage.output_tokens,
cache_read_input_tokens = turn.usage.cache_read_input_tokens,
cache_creation_input_tokens = turn.usage.cache_creation_input_tokens,
pending_approvals = turn.pending_approvals.len(),
handoff_requested = turn.handoff.is_some(),
stop = ?turn.stop,
"harness turn complete"
);
let wire_pending: Vec<WirePendingApproval> = turn
.pending_approvals
.into_iter()
.map(wire_pending_approval)
.collect();
let wire_pending_questions: Vec<WirePendingQuestion> = turn
.pending_questions
.into_iter()
.map(wire_pending_question)
.collect();
let wire_unattended_denials: Vec<WireUnattendedDenialFact> = turn
.unattended_denials
.into_iter()
.map(wire_unattended_denial_fact)
.collect();
let wire_delegate_records: Vec<WireDelegateFact> = turn
.delegate_records
.into_iter()
.map(wire_delegate_fact)
.collect();
let wire_routine_grant_drift = turn.routine_grant_drift;
let wire_handoff = turn.handoff.map(wire_handoff_request);
let stop_reason_i32 = turn.stop.map_or(
WireStopReason::STOP_REASON_UNSPECIFIED as i32,
llm_stop_to_wire_i32,
);
frame(
label,
WireTurnBatch {
messages: Vec::new(),
usage: buffa::MessageField::some(wire_usage(turn.usage)),
pending_approvals: wire_pending,
pending_questions: wire_pending_questions,
handoff: wire_handoff.map_or_else(buffa::MessageField::none, buffa::MessageField::some),
stop_reason: buffa::EnumValue::from(stop_reason_i32),
unattended_denials: wire_unattended_denials,
delegate_records: wire_delegate_records,
unattended_denial_aborted: turn.fire_stopped,
routine_grant_drift: wire_routine_grant_drift,
__buffa_unknown_fields: buffa::UnknownFields::default(),
},
)
}
#[must_use]
pub fn decode_turn_messages(messages: &[WireMessage]) -> Vec<LlmMessage> {
let mut input: Vec<LlmMessage> = Vec::new();
for msg in messages {
let m = wire_to_llm(msg);
if !m.content.is_empty() {
input.push(m);
}
}
input
}
pub async fn run_wire_turn(
providers: &ProviderSet,
request: &HarnessMessage,
expected_audience: &execution::ExecutionAudience,
tools: &(impl ToolExecutor + ?Sized),
clock: Option<Arc<dyn Clock + Send + Sync>>,
) -> Result<Vec<HarnessMessage>, ConnectError> {
let Some(HarnessFrame::Input(turn_input)) = request.frame.as_ref() else {
return Err(ConnectError::invalid_argument(
"a buffered Execution turn must open with an input frame",
));
};
let grant = execution::required_grant(turn_input.execution.as_option())
.map_err(|error| execution::connect_error_from_execution(&error))?;
let mut session =
execution::ExecutionSession::admit(&grant, expected_audience, execution::now())
.map_err(|error| execution::connect_error_from_execution(&error))?;
session
.check_opening_envelope(
&execution::required_label(request.label.as_option())
.map_err(|error| execution::connect_error_from_execution(&error))?,
)
.map_err(|error| execution::connect_error_from_execution(&error))?;
let input = decode_turn_messages(&turn_input.messages);
let approvals = verify_approval_responses(
&turn_input.approval_responses,
&turn_input.caller,
¤t_sandbox_mode(),
);
let routine_tool_grants = verify_routine_tool_grants(
&turn_input.routine_tool_grants,
&turn_input.caller,
¤t_sandbox_mode(),
grant.conversation().as_str(),
);
let resolved = providers.resolve(&turn_input.provider, &turn_input.model);
let native_search_allowed = !turn_input.scope_builtin_tools
|| turn_input
.builtin_tools
.iter()
.any(|n| n == polyc_tools::web::NATIVE_SEARCH_GROUNDING);
let question_answers = verify_question_answers(&turn_input.question_answers);
let mut options = run_turn_options(
approvals,
routine_tool_grants,
TurnToggles::from_wire(turn_input),
turn_input.prompt_cache_key.clone(),
turn_input.step_budget,
None,
None,
clock,
native_search_allowed,
question_answers,
turn_input.turn_start_unix_ms,
);
options.delegate_descriptors =
resolve_delegate_descriptors(providers, tools, &turn_input.delegate_descriptors);
options.delegate_max_fanout = turn_input.delegate_fanout_cap;
options.delegate_turn_budget = turn_input.delegate_turn_call_budget;
let mut turn = run_turn_captured_under_grant(
resolved.provider.as_ref(),
tools,
&resolved.model,
input,
options,
grant.capabilities().granted(),
)
.await
.map_err(|captured| captured.connect_error())?;
let mut frames = Vec::with_capacity(2);
let bodies = std::mem::take(&mut turn.messages);
if !bodies.is_empty() {
let proposal_label = session
.stamp_proposal()
.map_err(|error| execution::connect_error_from_execution(&error))?;
frames.push(build_step_outcome_proposal(&proposal_label, bodies));
}
let batch_label = session
.stamp_proposal()
.map_err(|error| execution::connect_error_from_execution(&error))?;
frames.push(build_final_batch(&batch_label, turn));
Ok(frames)
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs)]
#[test]
fn a_proxy_only_set_resolves_every_delegate_to_the_one_brokered_backend() {
let set = super::ProviderSet::proxy_only(
"brokered",
polyc_llm::into_dyn(polyc_llm::turn::StubProvider),
"approved-model",
);
let named = set.resolve("brokered", "");
assert_eq!(named.provider_name, "brokered");
assert_eq!(named.model, "approved-model");
assert_eq!(set.resolve("", "").provider_name, "brokered");
assert_eq!(
set.resolve("stub", "").provider_name,
"brokered",
"a delegate must not be able to name its way onto a stub backend"
);
assert!(
set.backend("stub").is_none(),
"no stub backend may be registered in a proxy-only set"
);
let general = super::ProviderSet::new(std::collections::HashMap::new(), "brokered");
assert!(general.backend("stub").is_some());
assert_eq!(general.resolve("stub", "").provider_name, "stub");
}
use super::{
CapturedTurnError, ProviderSet, RegisteredProvider, TurnToggles, VerifiedApprovals,
WireTurnFailureKind, classify_turn_error, connect_error_from_turn, run_turn_options,
wire_title,
};
use connectrpc::ErrorCode;
use polyc_llm::{LlmErrorKind, into_dyn, turn::StubProvider};
use std::collections::HashMap;
fn reg(default_model: &str) -> RegisteredProvider {
RegisteredProvider {
provider: into_dyn(StubProvider),
default_model: default_model.to_owned(),
}
}
fn provider_set(entries: &[(&str, &str)], default_provider: &str) -> ProviderSet {
let map: HashMap<String, RegisteredProvider> = entries
.iter()
.map(|(name, model)| ((*name).to_owned(), reg(model)))
.collect();
ProviderSet::new(map, default_provider)
}
#[test]
fn resolve_uses_named_provider_and_its_default_model() {
let set = provider_set(&[("vertex", "gemini"), ("openai", "llama3.2")], "vertex");
let r = set.resolve("openai", "");
assert_eq!(r.provider_name, "openai");
assert_eq!(r.model, "llama3.2");
let r = set.resolve("openai", "mixtral");
assert_eq!(r.model, "mixtral");
}
#[test]
fn resolve_empty_provider_uses_default() {
let set = provider_set(&[("openai", "llama3.2")], "openai");
let r = set.resolve("", "");
assert_eq!(r.provider_name, "openai");
assert_eq!(r.model, "llama3.2");
}
#[test]
fn resolve_unknown_provider_falls_back_to_default() {
let set = provider_set(&[("openai", "llama3.2")], "openai");
let r = set.resolve("openia", "");
assert_eq!(r.provider_name, "openai");
assert_eq!(r.model, "llama3.2");
}
#[test]
fn turn_error_kind_maps_to_connect_code() {
let cases = [
(LlmErrorKind::RateLimit, ErrorCode::ResourceExhausted),
(LlmErrorKind::Timeout, ErrorCode::DeadlineExceeded),
(LlmErrorKind::Unavailable, ErrorCode::Unavailable),
(LlmErrorKind::Auth, ErrorCode::Unauthenticated),
(LlmErrorKind::BadRequest, ErrorCode::InvalidArgument),
(LlmErrorKind::Other, ErrorCode::Internal),
];
for (kind, expected) in cases {
assert_eq!(
connect_error_from_turn(kind, "boom".to_owned()).code,
expected
);
}
}
#[test]
fn resolve_delegate_descriptors_derives_native_search_allowed_from_builtin_tools() {
use super::WireDelegateDescriptor;
let providers = provider_set(&[("vertex", "gemini-x")], "vertex");
let tools = polyc_tools::ToolRegistry::scoped(std::iter::empty());
let grounded = WireDelegateDescriptor {
agent_id: "researcher".to_owned(),
provider: "vertex".to_owned(),
model: "gemini-x".to_owned(),
builtin_tools: vec![
"web_fetch".to_owned(),
polyc_tools::web::NATIVE_SEARCH_GROUNDING.to_owned(),
],
max_steps: 4,
..WireDelegateDescriptor::default()
};
let ungrounded = WireDelegateDescriptor {
agent_id: "researcher".to_owned(),
provider: "vertex".to_owned(),
model: "gemini-x".to_owned(),
builtin_tools: vec!["web_fetch".to_owned()],
max_steps: 4,
..WireDelegateDescriptor::default()
};
let out = super::resolve_delegate_descriptors(&providers, &tools, &[grounded, ungrounded]);
assert_eq!(out.len(), 2);
assert!(
out[0].native_search_allowed,
"a descriptor naming web_search_grounding in builtin_tools must resolve to \
native_search_allowed: true"
);
assert!(
!out[1].native_search_allowed,
"a descriptor NOT naming web_search_grounding must resolve to \
native_search_allowed: false — never hardcoded true, or every worker would \
inherit a capability its own agent manifest never granted"
);
}
#[test]
fn every_definite_kind_classifies_to_a_durable_failure() {
let cases = [
(LlmErrorKind::RateLimit, WireTurnFailureKind::RateLimit),
(LlmErrorKind::Timeout, WireTurnFailureKind::Timeout),
(LlmErrorKind::Unavailable, WireTurnFailureKind::Unavailable),
(LlmErrorKind::Auth, WireTurnFailureKind::Auth),
(LlmErrorKind::BadRequest, WireTurnFailureKind::BadRequest),
(LlmErrorKind::Other, WireTurnFailureKind::Other),
];
for (kind, expected) in cases {
let CapturedTurnError::Failed(failure) = classify_turn_error(kind, "boom".to_owned())
else {
panic!("kind {kind:?} is definite and must carry a durable failure");
};
assert_eq!(
failure.kind, expected,
"kind {kind:?} should map to {expected:?}"
);
assert_eq!(failure.message, "boom");
}
}
#[test]
fn an_ambiguous_outcome_carries_no_durable_failure() {
let captured = classify_turn_error(LlmErrorKind::Ambiguous, "boom".to_owned());
assert!(
matches!(captured, CapturedTurnError::Ambiguous),
"an ambiguous kind may not become a persistable failure"
);
assert_eq!(captured.connect_error().code, ErrorCode::Unknown);
}
fn options_with_step_budget(step_budget: u32) -> polyc_agent::RunTurnOptions {
run_turn_options(
VerifiedApprovals::default(),
polyc_agent::RoutineGrantSet::default(),
TurnToggles::default(),
String::new(),
step_budget,
None,
None,
None,
true,
Vec::new(),
0,
)
}
#[test]
fn deployment_approve_all_dangerous_merges_a_blanket_all_grant_on_a_fire_dispatch() {
let opts = run_turn_options(
VerifiedApprovals::default(),
polyc_agent::RoutineGrantSet::default(),
TurnToggles {
fire_dispatch: true,
deployment_approve_all_dangerous: true,
..TurnToggles::default()
},
String::new(),
0,
None,
None,
None,
true,
Vec::new(),
0,
);
assert!(
opts.routine_tool_grants
.blanket
.is_some_and(|b| b.include_high),
"the deployment override must merge a blanket_all-shaped grant"
);
}
#[test]
fn deployment_approve_all_dangerous_is_inert_off_a_fire_dispatch() {
let opts = run_turn_options(
VerifiedApprovals::default(),
polyc_agent::RoutineGrantSet::default(),
TurnToggles {
fire_dispatch: false,
deployment_approve_all_dangerous: true,
..TurnToggles::default()
},
String::new(),
0,
None,
None,
None,
true,
Vec::new(),
0,
);
assert!(
opts.routine_tool_grants.blanket.is_none(),
"the override must never apply outside a fire dispatch"
);
}
#[test]
fn run_turn_options_step_budget_cap_below_base_lowers_max_steps() {
let base = polyc_agent::resolve_default_max_steps();
let cap = u32::try_from(base.saturating_sub(1).max(1)).unwrap_or(1);
let options = options_with_step_budget(cap);
assert_eq!(
options.max_steps,
Some(cap as usize),
"a cap below the baseline lowers max_steps to exactly the cap"
);
}
#[test]
fn run_turn_options_step_budget_cap_above_base_does_not_raise_max_steps() {
let base = polyc_agent::resolve_default_max_steps();
let cap = u32::try_from(base).unwrap_or(u32::MAX).saturating_add(100);
let options = options_with_step_budget(cap);
assert_eq!(
options.max_steps,
Some(base),
"a cap above the baseline never raises the budget past it"
);
}
#[test]
fn run_turn_options_zero_step_budget_leaves_max_steps_unset() {
let options = options_with_step_budget(0);
assert_eq!(
options.max_steps, None,
"unset ⇒ None, so the agent loop's own resolution applies"
);
}
#[test]
fn new_inserts_stub_and_coerces_missing_default() {
let set = provider_set(&[("openai", "llama3.2")], "ghost");
let r = set.resolve("", "");
assert_eq!(r.provider_name, "stub");
assert_eq!(r.model, "stub");
assert_eq!(set.resolve("stub", "").model, "stub");
}
#[test]
fn wire_title_prefers_curated_then_humanizes() {
assert_eq!(
wire_title("Pay for & fetch a web page", "paid_fetch"),
"Pay for & fetch a web page"
);
assert_eq!(wire_title("", "delete_file"), "Delete file");
assert_eq!(wire_title("", "send_message"), "Send message");
}
use super::{WireApprovalResponse, verify_approval_responses};
fn wire_response(
request_id: &str,
tool: &str,
args: &str,
approved_for_session: bool,
caller: &str,
sandbox_mode: &str,
) -> WireApprovalResponse {
use polyc_crypto::approval::{ApprovalSigner, response_payload};
let signer = ApprovalSigner::from_seed(9);
let (payload, _sig, _pk) = response_payload(
request_id,
tool,
args,
"",
true,
approved_for_session,
&[],
caller,
"",
sandbox_mode,
"ok",
"",
"conv-h",
"nonce-h",
"00000000-0000-0000-0000-000000000001",
&signer,
);
let d = polyc_crypto::approval::decode_response_full(&payload).expect("decode");
WireApprovalResponse {
request_id: d.request_id,
tool_name: d.tool_name,
args_json: d.args_json,
modified_args_json: d.modified_args_json,
injected_context: d.injected_context,
approved: d.approved,
approved_for_session: d.approved_for_session,
caller: d.caller,
approver: d.approver,
sandbox_mode: d.sandbox_mode,
reason: d.reason,
conversation_id: d.conversation_id,
nonce: d.nonce,
covered_capabilities: d.covered_capabilities,
signer_pk_hex: d.signer_pk_hex,
signature_hex: d.signature_hex,
turn_id: d.turn_id,
routine_grant: d.routine_grant,
tool_descriptor_hash: d.tool_descriptor_hash,
grant_scope: d.grant_scope,
already_executed: false,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
const M: &str = "workspace-write";
#[test]
fn executed_session_grant_keeps_tool_memory_but_not_the_one_shot() {
let mut resp = wire_response("call-1", "grep", "{}", true, "A", M);
resp.already_executed = true;
let v = verify_approval_responses(&[resp], "A", M);
assert!(
v.session_approved.contains_key("grep"),
"per-tool session memory survives execution"
);
assert!(
v.decisions.is_empty(),
"the spent one-shot tuple is NOT re-armed for execution"
);
}
#[test]
fn session_grant_applies_only_to_its_bound_caller() {
let responses = vec![wire_response("call-1", "grep", "{}", true, "A", M)];
let for_a = verify_approval_responses(&responses, "A", M);
assert!(for_a.session_approved.contains_key("grep"));
assert_eq!(for_a.decisions.len(), 1);
let for_b = verify_approval_responses(&responses, "B", M);
assert!(
for_b.session_approved.is_empty(),
"A's session grant must NOT auto-approve B"
);
assert_eq!(for_b.decisions.len(), 1);
}
#[test]
fn session_grant_is_bound_to_its_sandbox_mode() {
let responses = vec![wire_response("call-1", "grep", "{}", true, "A", M)];
assert!(
verify_approval_responses(&responses, "A", M)
.session_approved
.contains_key("grep")
);
let other = verify_approval_responses(&responses, "A", "danger-full-access");
assert!(
other.session_approved.is_empty(),
"a grant from one mode must not auto-approve under another"
);
assert_eq!(other.decisions.len(), 1);
}
#[test]
fn non_session_approval_yields_no_session_memory() {
let responses = vec![wire_response("call-1", "grep", "{}", false, "A", M)];
let v = verify_approval_responses(&responses, "A", M);
assert!(v.session_approved.is_empty());
assert_eq!(v.decisions.len(), 1);
}
#[test]
fn tampered_caller_fails_verification_and_is_dropped() {
let mut resp = wire_response("call-1", "grep", "{}", true, "A", M);
resp.caller = "ATTACKER".to_owned();
let v = verify_approval_responses(&[resp], "ATTACKER", M);
assert!(v.decisions.is_empty(), "a tampered response is dropped");
assert!(v.session_approved.is_empty());
}
use super::verify_routine_tool_grants;
fn wire_grant(
tool: &str,
caller: &str,
sandbox_mode: &str,
descriptor_hash: &str,
grant_scope: &str,
) -> WireApprovalResponse {
use polyc_crypto::approval::{ApprovalSigner, routine_grant_payload};
let signer = ApprovalSigner::from_seed(9);
let (payload, _sig, _pk) = routine_grant_payload(
"call-1",
tool,
"{}",
"",
true,
caller,
"",
sandbox_mode,
"owner approved during setup",
&["mutate-external".to_owned()],
"conv-h",
"nonce-h",
"00000000-0000-0000-0000-000000000001",
descriptor_hash,
grant_scope,
&signer,
);
let d = polyc_crypto::approval::decode_response_full(&payload).expect("decode");
WireApprovalResponse {
request_id: d.request_id,
tool_name: d.tool_name,
args_json: d.args_json,
modified_args_json: d.modified_args_json,
injected_context: d.injected_context,
approved: d.approved,
approved_for_session: d.approved_for_session,
caller: d.caller,
approver: d.approver,
sandbox_mode: d.sandbox_mode,
reason: d.reason,
conversation_id: d.conversation_id,
nonce: d.nonce,
covered_capabilities: d.covered_capabilities,
signer_pk_hex: d.signer_pk_hex,
signature_hex: d.signature_hex,
turn_id: d.turn_id,
routine_grant: d.routine_grant,
tool_descriptor_hash: d.tool_descriptor_hash,
grant_scope: d.grant_scope,
already_executed: false,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
#[test]
fn a_matching_per_tool_grant_verifies_and_buckets_by_tool_name() {
let grant = wire_grant("send_message", "A", M, "sha256:abc", "tool");
let out = verify_routine_tool_grants(&[grant], "A", M, "conv-h");
let entry = out
.per_tool
.get("send_message")
.expect("the grant buckets under its tool name");
assert_eq!(entry.descriptor_hash, "sha256:abc");
assert!(out.blanket.is_none());
}
#[test]
fn a_blanket_grant_verifies_and_buckets_as_blanket() {
let grant = wire_grant("", "A", M, "", "blanket_all");
let out = verify_routine_tool_grants(&[grant], "A", M, "conv-h");
assert!(out.per_tool.is_empty());
assert!(
out.blanket.is_some_and(|b| b.include_high),
"blanket_all sets include_high"
);
}
#[test]
fn a_grant_for_a_different_caller_is_dropped() {
let grant = wire_grant("send_message", "A", M, "sha256:abc", "tool");
let out = verify_routine_tool_grants(&[grant], "B", M, "conv-h");
assert!(
out.per_tool.is_empty(),
"A's grant must not apply to a turn whose caller is B"
);
}
#[test]
fn a_grant_under_a_different_sandbox_mode_is_dropped() {
let grant = wire_grant("send_message", "A", M, "sha256:abc", "tool");
let out = verify_routine_tool_grants(&[grant], "A", "danger-full-access", "conv-h");
assert!(
out.per_tool.is_empty(),
"a grant minted under one mode must not apply under another"
);
}
#[test]
fn a_tampered_grant_fails_verification_and_is_dropped() {
let mut grant = wire_grant("send_message", "A", M, "sha256:abc", "tool");
grant.tool_descriptor_hash = "sha256:evil".to_owned();
let out = verify_routine_tool_grants(&[grant], "A", M, "conv-h");
assert!(out.per_tool.is_empty(), "a tampered grant is dropped");
}
#[test]
fn a_response_with_no_routine_grant_marker_is_ignored() {
let resp = wire_response("call-1", "grep", "{}", true, "A", M);
let out = verify_routine_tool_grants(&[resp], "A", M, "conv-h");
assert!(out.per_tool.is_empty());
assert!(out.blanket.is_none());
}
#[test]
fn a_grant_for_a_different_fire_conversation_is_dropped() {
let grant = wire_grant("send_message", "A", M, "sha256:abc", "tool");
let out = verify_routine_tool_grants(&[grant], "A", M, "conv-b");
assert!(
out.per_tool.is_empty(),
"a grant signed for conversation A must not apply to conversation B"
);
}
#[test]
fn a_signed_unknown_grant_scope_yields_neither_per_tool_nor_blanket() {
let grant = wire_grant("send_message", "A", M, "sha256:abc", "tools");
let out = verify_routine_tool_grants(&[grant], "A", M, "conv-h");
assert!(
out.per_tool.is_empty(),
"an unrecognized scope must not bucket as a per-tool grant"
);
assert!(
out.blanket.is_none(),
"an unrecognized scope must not widen into a blanket grant"
);
}
use super::{WireQuestionAnswer, verify_question_answers};
const TEST_TURN: &str = "018f47f0-5f70-7cc5-98df-123456789abc";
fn wire_question_answer(
call_id: &str,
index: u32,
state: &str,
selected_index: u32,
selected_label: &str,
answered_by: &str,
signer: &polyc_crypto::approval::ApprovalSigner,
) -> WireQuestionAnswer {
let signed_selected_index =
(state != polyc_crypto::question::DECLINED_STATE).then_some(selected_index);
let (payload, sig, pk) = polyc_crypto::question::answer_payload(
TEST_TURN,
call_id,
index,
r#"{"questions":[]}"#,
state,
signed_selected_index,
selected_label,
answered_by,
"conv-q",
"nonce-q",
signer,
);
let verified =
polyc_crypto::question::verify_signed_answer(&payload).expect("payload verifies");
WireQuestionAnswer {
turn_id: TEST_TURN.to_owned(),
call_id: verified.call_id,
index: verified.index,
question_args_json: verified.question_args_json,
state: verified.state,
selected_index: verified.selected_index.unwrap_or_default(),
selected_label: verified.selected_label,
answered_by: verified.answered_by,
conversation_id: verified.conversation_id,
nonce: verified.nonce,
signer_pk_hex: polyc_crypto::hex::lower(&pk),
signature_hex: polyc_crypto::hex::lower(&sig),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
#[test]
fn verify_question_answers_accepts_a_well_formed_signed_answer() {
let signer = polyc_crypto::approval::ApprovalSigner::from_seed(1);
let wire = wire_question_answer(
"call-1",
0,
polyc_crypto::question::ANSWERED_STATE,
1,
"Production",
"slack:T1:U9",
&signer,
);
let out = verify_question_answers(&[wire]);
assert_eq!(out.len(), 1);
assert_eq!(out[0].call_id, "call-1");
assert_eq!(out[0].index, 0);
assert_eq!(out[0].state, polyc_agent::question::AnswerState::Answered);
assert_eq!(out[0].selected_index, Some(1));
assert_eq!(out[0].selected_label, "Production");
}
#[test]
fn verify_question_answers_drops_a_tampered_answer() {
let signer = polyc_crypto::approval::ApprovalSigner::from_seed(1);
let mut wire = wire_question_answer(
"call-1",
0,
polyc_crypto::question::ANSWERED_STATE,
1,
"Production",
"slack:T1:U9",
&signer,
);
wire.selected_index = 0;
wire.selected_label = "Staging".to_owned();
assert!(
verify_question_answers(&[wire]).is_empty(),
"a tampered answer must never verify"
);
}
#[test]
fn verify_question_answers_drops_an_answer_signed_by_a_different_key() {
let signer = polyc_crypto::approval::ApprovalSigner::from_seed(1);
let other = polyc_crypto::approval::ApprovalSigner::from_seed(2);
let mut wire = wire_question_answer(
"call-1",
0,
polyc_crypto::question::ANSWERED_STATE,
1,
"Production",
"slack:T1:U9",
&signer,
);
wire.signer_pk_hex = polyc_crypto::hex::lower(&other.public_key_bytes());
assert!(verify_question_answers(&[wire]).is_empty());
}
#[test]
fn verify_question_answers_drops_an_unrecognized_state() {
let signer = polyc_crypto::approval::ApprovalSigner::from_seed(1);
let wire = wire_question_answer(
"call-1",
0,
"maybe",
1,
"Production",
"slack:T1:U9",
&signer,
);
assert!(verify_question_answers(&[wire]).is_empty());
}
#[test]
fn verify_question_answers_declined_carries_no_selection() {
let signer = polyc_crypto::approval::ApprovalSigner::from_seed(1);
let wire = wire_question_answer(
"call-1",
0,
polyc_crypto::question::DECLINED_STATE,
0,
"",
"slack:T1:U9",
&signer,
);
let out = verify_question_answers(&[wire]);
assert_eq!(out.len(), 1);
assert_eq!(out[0].state, polyc_agent::question::AnswerState::Declined);
assert_eq!(out[0].selected_index, None);
}
}