diff --git a/crates/ironclaw_reborn_composition/src/factory.rs b/crates/ironclaw_reborn_composition/src/factory.rs index d407172f64b..bb63a42cab6 100644 --- a/crates/ironclaw_reborn_composition/src/factory.rs +++ b/crates/ironclaw_reborn_composition/src/factory.rs @@ -43,9 +43,9 @@ use ironclaw_filesystem::{ MountDescriptor, RootFilesystem, StorageClass, }; use ironclaw_filesystem::{LocalFilesystem, ScopedFilesystem}; -#[cfg(any(feature = "libsql", feature = "postgres"))] -use ironclaw_host_api::runtime_policy::EffectiveRuntimePolicy; -use ironclaw_host_api::runtime_policy::{FilesystemBackendKind, ProcessBackendKind, SecretMode}; +use ironclaw_host_api::runtime_policy::{ + EffectiveRuntimePolicy, FilesystemBackendKind, ProcessBackendKind, SecretMode, +}; use ironclaw_host_api::{ EffectKind, ExtensionId, HostPath, MountPermissions, MountView, PackageId, RuntimeHttpEgress, UserId, VirtualPath, @@ -396,6 +396,7 @@ pub struct RebornLocalDevApprovalTestParts { pub(crate) struct RebornLocalRuntimeServices { pub(crate) approval_requests: Arc, pub(crate) capability_leases: Arc, + pub(crate) runtime_policy: Option, // Used in approval_test_support (cfg(test) only); suppress the dead-code // lint on non-test builds where that module is not compiled in. #[cfg_attr(not(test), allow(dead_code))] @@ -569,6 +570,7 @@ struct RebornLocalDevStoreGraph { struct RebornLocalDevStoreGraphInput { filesystem: Arc, owner_user_id: UserId, + runtime_policy: Option, skill_filesystem: Arc>, workspace_filesystem: Arc>, workspace_mounts: MountView, @@ -787,6 +789,7 @@ async fn build_local_dev(input: RebornBuildInput) -> Result, capability_policy: Arc, ) -> Arc { + let (approval_policy, resolved_profile) = local_dev_approval_policy(runtime_policy); + let gate_effects = capability_policy.approval_gate_effects(); + let gate_policy: Arc = Arc::new( + RuntimeProfileApprovalGatePolicy::new(resolved_profile, gate_effects), + ); + profile_approval_authorizer(approval_policy, gate_policy) +} + +pub(crate) fn local_dev_effects_require_approval( + runtime_policy: Option<&EffectiveRuntimePolicy>, + capability_policy: &LocalDevCapabilityPolicy, + effects: &[EffectKind], +) -> bool { + let (approval_policy, resolved_profile) = local_dev_approval_policy(runtime_policy); + RuntimeProfileApprovalGatePolicy::new( + resolved_profile, + capability_policy.approval_gate_effects(), + ) + .effects_require_approval(approval_policy, effects) +} + +fn local_dev_approval_policy( + runtime_policy: Option<&EffectiveRuntimePolicy>, +) -> (ApprovalPolicy, RuntimeProfile) { let approval_policy = runtime_policy .map(|policy| policy.approval_policy) .unwrap_or(ApprovalPolicy::AskDestructive); let resolved_profile = runtime_policy .map(|policy| policy.resolved_profile) .unwrap_or(RuntimeProfile::LocalDev); - let gate_effects = capability_policy.approval_gate_effects(); - let gate_policy: Arc = Arc::new( - RuntimeProfileApprovalGatePolicy::new(resolved_profile, gate_effects), - ); - profile_approval_authorizer(approval_policy, gate_policy) + (approval_policy, resolved_profile) } diff --git a/crates/ironclaw_reborn_composition/src/local_dev_capability_policy.toml b/crates/ironclaw_reborn_composition/src/local_dev_capability_policy.toml index cf4a1af5481..f01e4ccac6b 100644 --- a/crates/ironclaw_reborn_composition/src/local_dev_capability_policy.toml +++ b/crates/ironclaw_reborn_composition/src/local_dev_capability_policy.toml @@ -220,3 +220,9 @@ capability = "builtin.trigger_remove" effects = ["dispatch_capability", "external_write"] mounts = "ambient" network = "default" + +[[grants]] +capability = "builtin.outbound_delivery_target_set" +effects = ["dispatch_capability", "external_write"] +mounts = "ambient" +network = "default" diff --git a/crates/ironclaw_reborn_composition/src/runtime.rs b/crates/ironclaw_reborn_composition/src/runtime.rs index 69fd5b25aa1..d2f23da5cae 100644 --- a/crates/ironclaw_reborn_composition/src/runtime.rs +++ b/crates/ironclaw_reborn_composition/src/runtime.rs @@ -2375,6 +2375,7 @@ pub async fn build_reborn_runtime( model_gateway, milestone_sink.clone(), skill_activation_source.clone(), + None, ) .ok_or(RebornRuntimeError::HostRuntimeUnavailable)?; ( diff --git a/crates/ironclaw_reborn_composition/src/runtime/local_dev.rs b/crates/ironclaw_reborn_composition/src/runtime/local_dev.rs index 11769a79106..727313d5942 100644 --- a/crates/ironclaw_reborn_composition/src/runtime/local_dev.rs +++ b/crates/ironclaw_reborn_composition/src/runtime/local_dev.rs @@ -6,9 +6,10 @@ use std::{ use chrono::Utc; use uuid::Uuid; +use ironclaw_authorization::CapabilityLeaseStore; use ironclaw_host_api::{ - CapabilityId, ExecutionContext, ExtensionId, InvocationId, MountView, ResourceScope, - RuntimeKind, TrustClass, UserId, + CapabilityId, EffectKind, ExecutionContext, ExtensionId, InvocationId, MountView, + ResourceScope, RuntimeKind, TrustClass, UserId, }; use ironclaw_host_runtime::{ CapabilitySurfacePolicy, HostRuntime, SurfaceKind, @@ -20,6 +21,8 @@ use ironclaw_loop_support::{ HostManagedModelResponse, HostManagedToolResultContent, LoopCapabilityInputResolver, LoopCapabilityPortFactory, LoopCapabilityResultWriter, loop_driver_execution_extension_id, }; +use ironclaw_product_workflow::OutboundPreferencesProductFacade; +use ironclaw_run_state::ApprovalRequestStore; use ironclaw_threads::{ AppendCapabilityDisplayPreviewRequest, CapabilityDisplayPreviewEnvelope, CapabilityDisplayPreviewEnvelopeInput, CapabilityDisplayPreviewStatus, SessionThreadService, @@ -34,6 +37,7 @@ use ironclaw_turns::{ }, }; +use crate::local_dev_authorization::local_dev_effects_require_approval; use crate::local_dev_capability_policy::LocalDevCapabilityPolicy; use crate::local_dev_mounts::scoped_skill_management_mount_view; use crate::{ @@ -43,6 +47,7 @@ use crate::{ }; pub(super) mod extension_surface; +mod outbound_delivery; mod refreshing_capability_port; #[cfg(test)] mod shell_tests; @@ -51,6 +56,10 @@ mod surface_disclosure; mod synthetic_capability; use extension_surface::{LocalDevExtensionSurface, LocalDevExtensionSurfaceSource}; +#[cfg(test)] +pub(crate) use outbound_delivery::{ + OUTBOUND_DELIVERY_TARGET_SET_CAPABILITY_ID, OUTBOUND_DELIVERY_TARGETS_LIST_CAPABILITY_ID, +}; use refreshing_capability_port::{ RefreshingLocalDevCapabilityPortConfig, create_refreshing_local_dev_capability_port, }; @@ -75,11 +84,19 @@ pub(super) fn capability_wiring( model_gateway: Arc, milestone_sink: Arc, skill_activation_source: Option>, + outbound_preferences_facade: Option>, ) -> Option { let runtime = services.host_runtime.clone()?; let local_runtime = services.local_runtime.as_ref()?; let workspace_mounts = local_runtime.workspace_mounts.clone(); let memory_mounts = local_runtime.memory_mounts.clone(); + let approval_requests: Arc = local_runtime.approval_requests.clone(); + let capability_leases: Arc = local_runtime.capability_leases.clone(); + let outbound_delivery_target_set_requires_approval = local_dev_effects_require_approval( + local_runtime.runtime_policy.as_ref(), + policy.as_ref(), + &[EffectKind::ExternalWrite], + ); let extension_surface_source = LocalDevExtensionSurfaceSource::new(local_runtime.extension_management.clone()); let display_previews = Arc::new(CapabilityDisplayPreviewStore::default()); @@ -102,6 +119,10 @@ pub(super) fn capability_wiring( result_writer: Arc::clone(&capability_result_writer), milestone_sink, skill_activation_source, + outbound_preferences_facade, + outbound_delivery_target_set_requires_approval, + approval_requests, + capability_leases, }); let model_gateway: Arc = Arc::new( LocalDevResultHydratingModelGateway::new(model_gateway, capability_io), @@ -128,6 +149,10 @@ struct LocalDevLoopCapabilityPortFactory { result_writer: Arc, milestone_sink: Arc, skill_activation_source: Option>, + outbound_preferences_facade: Option>, + outbound_delivery_target_set_requires_approval: bool, + approval_requests: Arc, + capability_leases: Arc, } #[async_trait::async_trait] @@ -154,6 +179,11 @@ impl LoopCapabilityPortFactory for LocalDevLoopCapabilityPortFactory { result_writer: Arc::clone(&self.result_writer), milestone_sink: Arc::clone(&self.milestone_sink), skill_activation_source: self.skill_activation_source.clone(), + outbound_preferences_facade: self.outbound_preferences_facade.clone(), + outbound_delivery_target_set_requires_approval: self + .outbound_delivery_target_set_requires_approval, + approval_requests: Arc::clone(&self.approval_requests), + capability_leases: Arc::clone(&self.capability_leases), }) .await } diff --git a/crates/ironclaw_reborn_composition/src/runtime/local_dev/outbound_delivery.rs b/crates/ironclaw_reborn_composition/src/runtime/local_dev/outbound_delivery.rs new file mode 100644 index 00000000000..286c059e6f4 --- /dev/null +++ b/crates/ironclaw_reborn_composition/src/runtime/local_dev/outbound_delivery.rs @@ -0,0 +1,704 @@ +use std::sync::Arc; + +use async_trait::async_trait; +use ironclaw_authorization::{CapabilityLeaseError, CapabilityLeaseStatus, CapabilityLeaseStore}; +use ironclaw_host_api::{ + Action, ApprovalRequest, ApprovalRequestId, CapabilityGrantId, CapabilityId, CorrelationId, + InvocationFingerprint, InvocationId, Principal, ResourceEstimate, ResourceScope, UserId, +}; +use ironclaw_loop_support::{CapabilityResultWrite, loop_driver_execution_extension_id}; +use ironclaw_product_workflow::{ + OutboundPreferencesProductFacade, RebornOutboundDeliveryTargetId, RebornServicesError, + RebornServicesErrorCode, RebornSetOutboundPreferencesRequest, WebUiAuthenticatedCaller, +}; +use ironclaw_run_state::{ApprovalRequestStore, ApprovalStatus, RunStateError}; +use ironclaw_turns::{ + LoopGateRef, + run_profile::{ + AgentLoopHostError, AgentLoopHostErrorKind, CapabilityApprovalResume, CapabilityInputRef, + CapabilityOutcome, CapabilityProgress, CapabilityResultMessage, CapabilityResumeToken, + ConcurrencyHint, LoopRunContext, + }, +}; + +use crate::runtime::local_dev::synthetic_capability::{ + LocalDevSyntheticCapability, LocalDevSyntheticCapabilityDescriptor, + LocalDevSyntheticCapabilityHandler, LocalDevSyntheticCapabilityInvocation, +}; + +pub(crate) const OUTBOUND_DELIVERY_TARGETS_LIST_CAPABILITY_ID: &str = + "builtin.outbound_delivery_targets_list"; +const OUTBOUND_DELIVERY_TARGETS_LIST_PROVIDER_TOOL_NAME: &str = + "builtin__outbound_delivery_targets_list"; +const OUTBOUND_DELIVERY_TARGETS_LIST_DESCRIPTION: &str = "List available outbound delivery targets for final replies and routine/trigger results, such as Slack DMs or Slack channels. Use before saying a delivery product is unavailable or asking the user to reconnect it."; + +pub(crate) const OUTBOUND_DELIVERY_TARGET_SET_CAPABILITY_ID: &str = + "builtin.outbound_delivery_target_set"; +const OUTBOUND_DELIVERY_TARGET_SET_PROVIDER_TOOL_NAME: &str = + "builtin__outbound_delivery_target_set"; +const OUTBOUND_DELIVERY_TARGET_SET_DESCRIPTION: &str = "Set the current user's final-reply outbound delivery target, such as a Slack DM or Slack channel. Approval may be required before the preference is changed."; + +pub(super) fn outbound_delivery_capabilities( + facade: Arc, + fallback_user_id: UserId, + approval_requests: Arc, + capability_leases: Arc, + target_set_requires_approval: bool, +) -> Result, AgentLoopHostError> { + Ok(vec![ + LocalDevSyntheticCapability::new( + LocalDevSyntheticCapabilityDescriptor::new( + OUTBOUND_DELIVERY_TARGETS_LIST_CAPABILITY_ID, + OUTBOUND_DELIVERY_TARGETS_LIST_PROVIDER_TOOL_NAME, + OUTBOUND_DELIVERY_TARGETS_LIST_DESCRIPTION, + ConcurrencyHint::SafeForParallel, + outbound_delivery_targets_list_input_schema(), + )?, + Arc::new(OutboundDeliveryTargetsListHandler { + facade: Arc::clone(&facade), + fallback_user_id: fallback_user_id.clone(), + }), + ), + LocalDevSyntheticCapability::new( + LocalDevSyntheticCapabilityDescriptor::new( + OUTBOUND_DELIVERY_TARGET_SET_CAPABILITY_ID, + OUTBOUND_DELIVERY_TARGET_SET_PROVIDER_TOOL_NAME, + OUTBOUND_DELIVERY_TARGET_SET_DESCRIPTION, + ConcurrencyHint::Exclusive, + outbound_delivery_target_set_input_schema(), + )?, + Arc::new(OutboundDeliveryTargetSetHandler { + facade, + fallback_user_id, + approval_requests, + capability_leases, + requires_approval: target_set_requires_approval, + }), + ), + ]) +} + +struct OutboundDeliveryTargetsListHandler { + facade: Arc, + fallback_user_id: UserId, +} + +#[async_trait] +impl LocalDevSyntheticCapabilityHandler for OutboundDeliveryTargetsListHandler { + fn validate_provider_arguments( + &self, + arguments: &serde_json::Value, + ) -> Result<(), AgentLoopHostError> { + parse_optional_channel(arguments).map(|_| ()) + } + + async fn invoke( + &self, + invocation: LocalDevSyntheticCapabilityInvocation, + ) -> Result { + let channel_filter = parse_optional_channel(&invocation.input)?; + let caller = caller_for_run(&invocation, &self.fallback_user_id); + let mut response = self + .facade + .list_outbound_delivery_targets(caller) + .await + .map_err(|error| outbound_delivery_host_error("list_targets", error))?; + if let Some(channel_filter) = channel_filter { + response.targets.retain(|option| { + option + .target + .channel + .as_str() + .eq_ignore_ascii_case(channel_filter.as_str()) + }); + } + let count = response.targets.len(); + let output = serde_json::to_value(response).map_err(|error| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::Internal, + format!("outbound delivery target list output serialization failed: {error}"), + ) + })?; + write_completed_result( + invocation, + output, + format!("found {count} delivery target(s)"), + ) + .await + } +} + +struct OutboundDeliveryTargetSetHandler { + facade: Arc, + fallback_user_id: UserId, + approval_requests: Arc, + capability_leases: Arc, + requires_approval: bool, +} + +struct ApprovedDispatchLease { + scope: ResourceScope, + lease_id: CapabilityGrantId, +} + +#[async_trait] +impl LocalDevSyntheticCapabilityHandler for OutboundDeliveryTargetSetHandler { + fn validate_provider_arguments( + &self, + arguments: &serde_json::Value, + ) -> Result<(), AgentLoopHostError> { + parse_target_id(arguments).map(|_| ()) + } + + async fn invoke( + &self, + invocation: LocalDevSyntheticCapabilityInvocation, + ) -> Result { + if invocation.request.auth_resume.is_some() { + return Err(AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + "outbound delivery target setter does not support auth resume", + )); + } + + let input = invocation_replay_input(&invocation).clone(); + let target_id = parse_target_id(&input)?; + let approved_lease = if self.requires_approval { + match invocation.request.approval_resume.clone() { + Some(resume) => Some( + self.verify_approved_resume(&invocation, &resume, &input) + .await?, + ), + None => return self.request_approval(&invocation, &input, &target_id).await, + } + } else { + if invocation.request.approval_resume.is_some() { + return Err(AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + "outbound delivery target approval resume is not expected", + )); + } + None + }; + + let target_summary = target_id.as_str().to_string(); + let caller = caller_for_run(&invocation, &self.fallback_user_id); + if let Some(approved_lease) = approved_lease { + self.capability_leases + .consume(&approved_lease.scope, approved_lease.lease_id) + .await + .map_err(|error| approval_lease_error("consume_approval_lease", error))?; + } + let response = self + .facade + .set_outbound_preferences( + caller, + RebornSetOutboundPreferencesRequest { + final_reply_target_id: Some(target_id), + }, + ) + .await + .map_err(|error| outbound_delivery_host_error("set_target", error))?; + let output = serde_json::to_value(response).map_err(|error| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::Internal, + format!("outbound delivery target set output serialization failed: {error}"), + ) + })?; + write_completed_result( + invocation, + output, + format!("set delivery target to {target_summary}"), + ) + .await + } +} + +impl OutboundDeliveryTargetSetHandler { + async fn request_approval( + &self, + invocation: &LocalDevSyntheticCapabilityInvocation, + input: &serde_json::Value, + target_id: &RebornOutboundDeliveryTargetId, + ) -> Result { + let capability_id = outbound_delivery_target_set_capability_id()?; + let approval_request_id = ApprovalRequestId::new(); + let correlation_id = CorrelationId::new(); + let invocation_id = InvocationId::new(); + let estimate = ResourceEstimate::default(); + let scope = resource_scope_for_run( + &invocation.run_context, + &self.fallback_user_id, + invocation_id, + ); + let fingerprint = approval_fingerprint(&scope, &capability_id, &estimate, input)?; + self.approval_requests + .save_pending( + scope, + ApprovalRequest { + id: approval_request_id, + correlation_id, + requested_by: Principal::Extension(loop_driver_execution_extension_id( + &invocation.run_context, + )?), + action: Box::new(Action::Dispatch { + capability: capability_id, + estimated_resources: estimate.clone(), + }), + invocation_fingerprint: Some(fingerprint), + reason: format!( + "Change final reply delivery target to `{}`", + target_id.as_str() + ), + reusable_scope: None, + }, + ) + .await + .map_err(|error| approval_store_error("save_pending_approval", error))?; + + Ok(CapabilityOutcome::ApprovalRequired { + gate_ref: approval_gate_ref(approval_request_id)?, + safe_summary: "changing the outbound delivery target requires approval".to_string(), + approval_resume: Some(CapabilityApprovalResume { + approval_request_id, + resume_token: resume_token_from_invocation_id(invocation_id)?, + correlation_id, + input_ref: invocation.request.input_ref.clone(), + input: input.clone(), + estimate, + }), + }) + } + + async fn verify_approved_resume( + &self, + invocation: &LocalDevSyntheticCapabilityInvocation, + resume: &CapabilityApprovalResume, + input: &serde_json::Value, + ) -> Result { + if resume.input != *input { + return Err(AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + "outbound delivery target approval resume input does not match", + )); + } + + let capability_id = outbound_delivery_target_set_capability_id()?; + let invocation_id = invocation_id_from_resume_token(&resume.resume_token)?; + let scope = resource_scope_for_run( + &invocation.run_context, + &self.fallback_user_id, + invocation_id, + ); + let fingerprint = approval_fingerprint(&scope, &capability_id, &resume.estimate, input)?; + let approval_record = self + .approval_requests + .get(&scope, resume.approval_request_id) + .await + .map_err(|error| approval_store_error("load_approval", error))? + .ok_or_else(|| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::Unauthorized, + "outbound delivery target approval is unavailable", + ) + })?; + if approval_record.status != ApprovalStatus::Approved { + return Err(AgentLoopHostError::new( + AgentLoopHostErrorKind::Unauthorized, + "outbound delivery target approval has not been granted", + )); + } + if approval_record.request.correlation_id != resume.correlation_id { + return Err(AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + "outbound delivery target approval correlation does not match", + )); + } + if approval_record.request.invocation_fingerprint.as_ref() != Some(&fingerprint) { + return Err(AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + "outbound delivery target approval fingerprint does not match", + )); + } + if !approval_request_matches_capability( + approval_record.request.action.as_ref(), + &capability_id, + ) { + return Err(AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + "outbound delivery target approval action does not match", + )); + } + + let lease = self + .capability_leases + .leases_for_scope(&scope) + .await + .into_iter() + .find(|lease| { + lease.status == CapabilityLeaseStatus::Active + && lease.grant.capability == capability_id + && lease.grant.grantee == approval_record.request.requested_by + && lease.invocation_fingerprint.as_ref() == Some(&fingerprint) + }) + .ok_or_else(|| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::Unauthorized, + "outbound delivery target approval lease is unavailable", + ) + })?; + self.capability_leases + .claim(&scope, lease.grant.id, &fingerprint) + .await + .map_err(|error| approval_lease_error("claim_approval_lease", error))?; + Ok(ApprovedDispatchLease { + scope, + lease_id: lease.grant.id, + }) + } +} + +async fn write_completed_result( + invocation: LocalDevSyntheticCapabilityInvocation, + output: serde_json::Value, + safe_summary: String, +) -> Result { + let (result_ref, byte_len) = invocation + .result_writer + .write_capability_result(CapabilityResultWrite { + run_context: &invocation.run_context, + input_ref: invocation_effective_input_ref(&invocation), + invocation_id: InvocationId::new(), + capability_id: &invocation.request.capability_id, + output, + display_preview: None, + }) + .await?; + Ok(CapabilityOutcome::Completed(CapabilityResultMessage { + result_ref, + safe_summary, + progress: CapabilityProgress::MadeProgress, + terminate_hint: false, + byte_len, + })) +} + +fn invocation_replay_input( + invocation: &LocalDevSyntheticCapabilityInvocation, +) -> &serde_json::Value { + invocation + .request + .approval_resume + .as_ref() + .map(|resume| &resume.input) + .unwrap_or(&invocation.input) +} + +fn invocation_effective_input_ref( + invocation: &LocalDevSyntheticCapabilityInvocation, +) -> &CapabilityInputRef { + invocation + .request + .approval_resume + .as_ref() + .map(|resume| &resume.input_ref) + .unwrap_or(&invocation.request.input_ref) +} + +fn caller_for_run( + invocation: &LocalDevSyntheticCapabilityInvocation, + fallback_user_id: &UserId, +) -> WebUiAuthenticatedCaller { + WebUiAuthenticatedCaller::new( + invocation.run_context.scope.tenant_id.clone(), + effective_user_id(&invocation.run_context, fallback_user_id), + invocation.run_context.scope.agent_id.clone(), + invocation.run_context.scope.project_id.clone(), + ) +} + +fn resource_scope_for_run( + run_context: &LoopRunContext, + fallback_user_id: &UserId, + invocation_id: InvocationId, +) -> ResourceScope { + let mut scope = run_context.scope.to_resource_scope(); + scope.user_id = effective_user_id(run_context, fallback_user_id); + scope.invocation_id = invocation_id; + scope +} + +fn effective_user_id(run_context: &LoopRunContext, fallback_user_id: &UserId) -> UserId { + run_context + .scope + .explicit_owner_user_id() + .cloned() + .or_else(|| { + run_context + .actor + .as_ref() + .map(|actor| actor.user_id.clone()) + }) + .unwrap_or_else(|| fallback_user_id.clone()) +} + +fn outbound_delivery_targets_list_input_schema() -> serde_json::Value { + serde_json::json!({ + "type": "object", + "properties": { + "channel": { + "type": "string", + "description": "Optional product/channel filter such as slack." + } + }, + "additionalProperties": false + }) +} + +fn outbound_delivery_target_set_input_schema() -> serde_json::Value { + serde_json::json!({ + "type": "object", + "properties": { + "target_id": { + "type": "string", + "description": "Target id returned by builtin__outbound_delivery_targets_list." + } + }, + "required": ["target_id"], + "additionalProperties": false + }) +} + +fn parse_optional_channel(input: &serde_json::Value) -> Result, AgentLoopHostError> { + let input = input_object(input, "outbound delivery target list", &["channel"])?; + match input.get("channel") { + None => Ok(None), + Some(value) => value + .as_str() + .map(str::trim) + .filter(|value| !value.is_empty()) + .map(|value| Some(value.to_string())) + .ok_or_else(|| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + "outbound delivery target list channel must be a non-empty string", + ) + }), + } +} + +fn parse_target_id( + input: &serde_json::Value, +) -> Result { + let input = input_object(input, "outbound delivery target set", &["target_id"])?; + let value = input + .get("target_id") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + "outbound delivery target set target_id must be a string", + ) + })?; + RebornOutboundDeliveryTargetId::new(value).map_err(|reason| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + format!("outbound delivery target set target_id is invalid: {reason}"), + ) + }) +} + +fn input_object<'a>( + input: &'a serde_json::Value, + capability_name: &'static str, + allowed_fields: &[&str], +) -> Result<&'a serde_json::Map, AgentLoopHostError> { + let object = input.as_object().ok_or_else(|| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + format!("{capability_name} input must be an object"), + ) + })?; + if let Some(field) = object + .keys() + .find(|field| !allowed_fields.contains(&field.as_str())) + { + return Err(AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + format!("{capability_name} input contains unsupported field `{field}`"), + )); + } + Ok(object) +} + +fn outbound_delivery_target_set_capability_id() -> Result { + CapabilityId::new(OUTBOUND_DELIVERY_TARGET_SET_CAPABILITY_ID).map_err(|error| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::Internal, + format!("outbound delivery target set capability id is invalid: {error}"), + ) + }) +} + +fn approval_fingerprint( + scope: &ResourceScope, + capability_id: &CapabilityId, + estimate: &ResourceEstimate, + input: &serde_json::Value, +) -> Result { + InvocationFingerprint::for_dispatch(scope, capability_id, estimate, input).map_err(|error| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::Internal, + format!("outbound delivery target approval fingerprint could not be computed: {error}"), + ) + }) +} + +fn approval_request_matches_capability(action: &Action, capability_id: &CapabilityId) -> bool { + matches!(action, Action::Dispatch { capability, .. } if capability == capability_id) +} + +fn approval_gate_ref(request_id: ApprovalRequestId) -> Result { + LoopGateRef::new(format!("gate:approval-{request_id}")).map_err(|error| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::Internal, + format!("outbound delivery target approval gate ref is invalid: {error}"), + ) + }) +} + +fn resume_token_from_invocation_id( + invocation_id: InvocationId, +) -> Result { + CapabilityResumeToken::new(invocation_id.to_string()).map_err(|reason| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::Internal, + format!("outbound delivery target resume token is invalid: {reason}"), + ) + }) +} + +fn invocation_id_from_resume_token( + resume_token: &CapabilityResumeToken, +) -> Result { + InvocationId::parse(resume_token.as_str()).map_err(|error| { + AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + format!("outbound delivery target approval resume token is invalid: {error}"), + ) + }) +} + +fn outbound_delivery_host_error( + operation: &'static str, + error: RebornServicesError, +) -> AgentLoopHostError { + let kind = match error.code { + RebornServicesErrorCode::InvalidRequest | RebornServicesErrorCode::NotFound => { + AgentLoopHostErrorKind::InvalidInvocation + } + RebornServicesErrorCode::Unauthenticated | RebornServicesErrorCode::Forbidden => { + AgentLoopHostErrorKind::Unauthorized + } + RebornServicesErrorCode::Conflict | RebornServicesErrorCode::RateLimited => { + AgentLoopHostErrorKind::Unavailable + } + RebornServicesErrorCode::Unavailable => AgentLoopHostErrorKind::Unavailable, + RebornServicesErrorCode::Internal => AgentLoopHostErrorKind::Internal, + }; + ironclaw_loop_support::raw_agent_loop_host_error( + "local_dev_outbound_delivery", + operation, + kind, + "outbound delivery target operation failed", + error, + ) +} + +fn approval_store_error(operation: &'static str, error: RunStateError) -> AgentLoopHostError { + ironclaw_loop_support::raw_agent_loop_host_error( + "local_dev_outbound_delivery", + operation, + AgentLoopHostErrorKind::Unavailable, + "outbound delivery approval state operation failed", + error, + ) +} + +fn approval_lease_error( + operation: &'static str, + error: CapabilityLeaseError, +) -> AgentLoopHostError { + let kind = match error { + CapabilityLeaseError::UnknownLease { .. } + | CapabilityLeaseError::ExpiredLease { .. } + | CapabilityLeaseError::ExhaustedLease { .. } + | CapabilityLeaseError::UnclaimedFingerprintLease { .. } + | CapabilityLeaseError::FingerprintMismatch { .. } + | CapabilityLeaseError::InactiveLease { .. } => AgentLoopHostErrorKind::Unauthorized, + CapabilityLeaseError::Persistence { .. } + | CapabilityLeaseError::VersionMismatch + | CapabilityLeaseError::CasExhausted => AgentLoopHostErrorKind::Unavailable, + }; + ironclaw_loop_support::raw_agent_loop_host_error( + "local_dev_outbound_delivery", + operation, + kind, + "outbound delivery approval lease operation failed", + error, + ) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parse_optional_channel_rejects_empty_channel() { + let error = parse_optional_channel(&serde_json::json!({"channel": " "})) + .expect_err("empty channel should fail"); + + assert_eq!(error.kind, AgentLoopHostErrorKind::InvalidInvocation); + } + + #[test] + fn parse_optional_channel_rejects_non_object_input() { + let error = parse_optional_channel(&serde_json::Value::Null) + .expect_err("non-object input should fail"); + + assert_eq!(error.kind, AgentLoopHostErrorKind::InvalidInvocation); + } + + #[test] + fn parse_optional_channel_rejects_unknown_fields() { + let error = parse_optional_channel(&serde_json::json!({"unexpected": "value"})) + .expect_err("unknown fields should fail"); + + assert_eq!(error.kind, AgentLoopHostErrorKind::InvalidInvocation); + } + + #[test] + fn parse_target_id_requires_target_id() { + let error = + parse_target_id(&serde_json::json!({})).expect_err("missing target id should fail"); + + assert_eq!(error.kind, AgentLoopHostErrorKind::InvalidInvocation); + } + + #[test] + fn parse_target_id_rejects_malformed_target_id() { + let error = parse_target_id(&serde_json::json!({"target_id": "bad\nid"})) + .expect_err("malformed target id should fail"); + + assert_eq!(error.kind, AgentLoopHostErrorKind::InvalidInvocation); + } + + #[test] + fn parse_target_id_rejects_unknown_fields() { + let error = + parse_target_id(&serde_json::json!({"target_id": "slack:test", "unexpected": "value"})) + .expect_err("unknown fields should fail"); + + assert_eq!(error.kind, AgentLoopHostErrorKind::InvalidInvocation); + } +} diff --git a/crates/ironclaw_reborn_composition/src/runtime/local_dev/refreshing_capability_port.rs b/crates/ironclaw_reborn_composition/src/runtime/local_dev/refreshing_capability_port.rs index 68da9a9a270..98a95f7e081 100644 --- a/crates/ironclaw_reborn_composition/src/runtime/local_dev/refreshing_capability_port.rs +++ b/crates/ironclaw_reborn_composition/src/runtime/local_dev/refreshing_capability_port.rs @@ -1,10 +1,13 @@ use std::sync::{Arc, Mutex as StdMutex}; +use ironclaw_authorization::CapabilityLeaseStore; use ironclaw_host_api::{MountView, UserId}; use ironclaw_host_runtime::HostRuntime; use ironclaw_loop_support::{ HostRuntimeLoopCapabilityPortFactory, LoopCapabilityInputResolver, LoopCapabilityResultWriter, }; +use ironclaw_product_workflow::OutboundPreferencesProductFacade; +use ironclaw_run_state::ApprovalRequestStore; use ironclaw_turns::run_profile::{ AgentLoopHostError, AgentLoopHostErrorKind, CapabilityBatchInvocation, CapabilityBatchOutcome, CapabilityCallCandidate, CapabilityInvocation, CapabilityOutcome, LoopCapabilityPort, @@ -16,6 +19,7 @@ use tokio::sync::Mutex as AsyncMutex; use crate::local_dev_capability_policy::LocalDevCapabilityPolicy; use crate::runtime::LocalDevSelectableSkillContextSource; use crate::runtime::local_dev::extension_surface::LocalDevExtensionSurfaceSource; +use crate::runtime::local_dev::outbound_delivery::outbound_delivery_capabilities; use crate::runtime::local_dev::skill_activation::skill_activation_capability; use crate::runtime::local_dev::surface_disclosure::wrap_local_dev_surface_disclosure; use crate::runtime::local_dev::synthetic_capability::wrap_local_dev_synthetic_capabilities; @@ -35,6 +39,10 @@ pub(super) struct RefreshingLocalDevCapabilityPortConfig { pub(super) result_writer: Arc, pub(super) milestone_sink: Arc, pub(super) skill_activation_source: Option>, + pub(super) outbound_preferences_facade: Option>, + pub(super) outbound_delivery_target_set_requires_approval: bool, + pub(super) approval_requests: Arc, + pub(super) capability_leases: Arc, } pub(super) async fn create_refreshing_local_dev_capability_port( @@ -53,6 +61,11 @@ pub(super) async fn create_refreshing_local_dev_capability_port( result_writer: config.result_writer, milestone_sink: config.milestone_sink, skill_activation_source: config.skill_activation_source, + outbound_preferences_facade: config.outbound_preferences_facade, + outbound_delivery_target_set_requires_approval: config + .outbound_delivery_target_set_requires_approval, + approval_requests: config.approval_requests, + capability_leases: config.capability_leases, current: StdMutex::new(None), refresh_lock: AsyncMutex::new(()), }); @@ -76,6 +89,10 @@ struct RefreshingLocalDevCapabilityPort { result_writer: Arc, milestone_sink: Arc, skill_activation_source: Option>, + outbound_preferences_facade: Option>, + outbound_delivery_target_set_requires_approval: bool, + approval_requests: Arc, + capability_leases: Arc, current: StdMutex>>, refresh_lock: AsyncMutex<()>, } @@ -113,7 +130,7 @@ impl RefreshingLocalDevCapabilityPort { .with_capability_execution_mount(capability_id.clone(), self.memory_mounts.clone()); } let port = factory.for_run_context(self.run_context.clone()); - let synthetic_capabilities = match &self.skill_activation_source { + let mut synthetic_capabilities = match &self.skill_activation_source { Some(skill_activation_source) => { vec![skill_activation_capability(Arc::clone( skill_activation_source, @@ -121,6 +138,15 @@ impl RefreshingLocalDevCapabilityPort { } None => Vec::new(), }; + if let Some(outbound_preferences_facade) = &self.outbound_preferences_facade { + synthetic_capabilities.extend(outbound_delivery_capabilities( + Arc::clone(outbound_preferences_facade), + self.fallback_user_id.clone(), + Arc::clone(&self.approval_requests), + Arc::clone(&self.capability_leases), + self.outbound_delivery_target_set_requires_approval, + )?); + } let port = wrap_local_dev_synthetic_capabilities( port, synthetic_capabilities, diff --git a/crates/ironclaw_reborn_composition/src/runtime/local_dev/shell_tests.rs b/crates/ironclaw_reborn_composition/src/runtime/local_dev/shell_tests.rs index cd0c3b2302f..a839fb354d0 100644 --- a/crates/ironclaw_reborn_composition/src/runtime/local_dev/shell_tests.rs +++ b/crates/ironclaw_reborn_composition/src/runtime/local_dev/shell_tests.rs @@ -74,18 +74,12 @@ async fn local_dev_yolo_shell_translates_workspace_workdir_without_scoped_mounts .await .expect("local-dev services build"); let runtime = services.host_runtime.clone().expect("host runtime"); - let workspace_mounts = services + let local_runtime = services .local_runtime .as_ref() - .expect("local runtime substrate") - .workspace_mounts - .clone(); - let memory_mounts = services - .local_runtime - .as_ref() - .expect("local runtime substrate") - .memory_mounts - .clone(); + .expect("local runtime substrate"); // safety: test-only assertion in #[cfg(test)] module. + let workspace_mounts = local_runtime.workspace_mounts.clone(); + let memory_mounts = local_runtime.memory_mounts.clone(); let policy = Arc::new( crate::local_dev_capability_policy::local_dev_capability_policy().expect("policy parses"), ); @@ -103,6 +97,10 @@ async fn local_dev_yolo_shell_translates_workspace_workdir_without_scoped_mounts result_writer, milestone_sink: Arc::new(InMemoryLoopHostMilestoneSink::default()), skill_activation_source: None, + outbound_preferences_facade: None, + outbound_delivery_target_set_requires_approval: false, + approval_requests: local_runtime.approval_requests.clone(), + capability_leases: local_runtime.capability_leases.clone(), }; let run_context = run_context("shell-workdir").await; let port = factory diff --git a/crates/ironclaw_reborn_composition/src/runtime/local_dev/synthetic_capability.rs b/crates/ironclaw_reborn_composition/src/runtime/local_dev/synthetic_capability.rs index 875c99839a9..65bb4945758 100644 --- a/crates/ironclaw_reborn_composition/src/runtime/local_dev/synthetic_capability.rs +++ b/crates/ironclaw_reborn_composition/src/runtime/local_dev/synthetic_capability.rs @@ -340,10 +340,21 @@ impl LoopCapabilityPort for LocalDevSyntheticCapabilityPort { "synthetic capability call cites a stale capability surface", )); } - let input = self - .input_resolver - .resolve_capability_input(&self.run_context, &request.input_ref) - .await?; + if request.approval_resume.is_some() && request.auth_resume.is_some() { + return Err(AgentLoopHostError::new( + AgentLoopHostErrorKind::InvalidInvocation, + "capability invocation has both approval_resume and auth_resume set; \ + these resume modes are mutually exclusive", + )); + } + let input = match request.approval_resume.as_ref() { + Some(resume) => resume.input.clone(), + None => { + self.input_resolver + .resolve_capability_input(&self.run_context, &request.input_ref) + .await? + } + }; handler .invoke(LocalDevSyntheticCapabilityInvocation { run_context: self.run_context.clone(), diff --git a/crates/ironclaw_reborn_composition/src/runtime/local_dev/tests.rs b/crates/ironclaw_reborn_composition/src/runtime/local_dev/tests.rs index 2ede0d61d1e..62d3fb746f1 100644 --- a/crates/ironclaw_reborn_composition/src/runtime/local_dev/tests.rs +++ b/crates/ironclaw_reborn_composition/src/runtime/local_dev/tests.rs @@ -6,8 +6,11 @@ mod tests { use super::super::*; + use ironclaw_approvals::ApprovalResolver; + use ironclaw_authorization::{CapabilityLeaseStatus, CapabilityLeaseStore}; use ironclaw_host_api::{ - AgentId, EffectKind, MountPermissions, NetworkPolicy, ProjectId, TenantId, ThreadId, + AgentId, CapabilityId, EffectKind, InvocationId, MountPermissions, NetworkPolicy, + ProjectId, TenantId, ThreadId, }; use ironclaw_host_runtime::{ APPLY_PATCH_CAPABILITY_ID, GLOB_CAPABILITY_ID, GREP_CAPABILITY_ID, HTTP_CAPABILITY_ID, @@ -17,19 +20,22 @@ mod tests { WRITE_FILE_CAPABILITY_ID, }; use ironclaw_loop_support::{HostManagedModelMessage, HostSkillContextSource}; + use ironclaw_outbound::CommunicationPreferenceKey; use ironclaw_product_workflow::{ LifecyclePackageKind, LifecyclePackageRef, LifecycleProductAction, LifecycleProductContext, - LifecycleProductFacade, LifecycleProductSurfaceContext, + LifecycleProductFacade, LifecycleProductSurfaceContext, OutboundPreferencesProductFacade, + RebornOutboundDeliveryTargetCapabilities, RebornOutboundDeliveryTargetId, + RebornOutboundDeliveryTargetSummary, RebornServicesError, WebUiAuthenticatedCaller, }; use ironclaw_threads::{ EnsureThreadRequest, InMemorySessionThreadService, MessageKind, ThreadHistoryRequest, ToolResultReferenceEnvelope, ToolResultSafeSummary, }; use ironclaw_turns::{ - AcceptedMessageRef, LoopMessageRef, RunProfileResolutionRequest, RunProfileResolver, - TurnActor, TurnId, TurnRunId, TurnScope, + AcceptedMessageRef, LoopMessageRef, ReplyTargetBindingRef, RunProfileResolutionRequest, + RunProfileResolver, TurnActor, TurnId, TurnRunId, TurnScope, run_profile::{ - CapabilityFailureKind, CapabilityInvocation, CapabilityOutcome, + CapabilityFailureKind, CapabilityInputRef, CapabilityInvocation, CapabilityOutcome, InMemoryLoopHostMilestoneSink, InMemoryRunProfileResolver, ModelProfileId, VisibleCapabilityRequest, }, @@ -39,6 +45,10 @@ mod tests { EXTENSION_ACTIVATE_CAPABILITY_ID, EXTENSION_INSTALL_CAPABILITY_ID, EXTENSION_REMOVE_CAPABILITY_ID, EXTENSION_SEARCH_CAPABILITY_ID, }; + use crate::outbound_preferences::{ + OutboundDeliveryTargetEntry, OutboundDeliveryTargetProvider, + OutboundDeliveryTargetRegistry, RebornOutboundPreferencesFacade, + }; use crate::runtime::local_dev_filesystem_skill_context_source; async fn run_context(label: &str) -> LoopRunContext { @@ -159,6 +169,68 @@ mod tests { provider_tool_call_with_name("builtin_echo", arguments) } + struct StaticOutboundDeliveryTargetProvider { + entry: OutboundDeliveryTargetEntry, + expected_caller: std::sync::Mutex>, + observed_callers: std::sync::Mutex>, + } + + impl StaticOutboundDeliveryTargetProvider { + fn new(entry: OutboundDeliveryTargetEntry) -> Self { + Self { + entry, + expected_caller: std::sync::Mutex::new(None), + observed_callers: std::sync::Mutex::new(Vec::new()), + } + } + + fn expect_caller(&self, caller: WebUiAuthenticatedCaller) { + *self.expected_caller.lock().expect("caller lock") = Some(caller); + } + + fn observed_callers(&self) -> Vec { + self.observed_callers + .lock() + .expect("observed caller lock") + .clone() + } + } + + #[async_trait::async_trait] + impl OutboundDeliveryTargetProvider for StaticOutboundDeliveryTargetProvider { + async fn list_outbound_delivery_targets( + &self, + caller: &WebUiAuthenticatedCaller, + ) -> Result, RebornServicesError> { + self.observed_callers + .lock() + .expect("observed caller lock") + .push(caller.clone()); + if self + .expected_caller + .lock() + .expect("caller lock") + .as_ref() + .is_some_and(|expected| expected != caller) + { + return Ok(Vec::new()); + } + Ok(vec![self.entry.clone()]) + } + } + + fn expected_outbound_delivery_caller( + run_context: &LoopRunContext, + user_id: UserId, + ) -> WebUiAuthenticatedCaller { + WebUiAuthenticatedCaller::new( + run_context.scope.tenant_id.clone(), + user_id, + run_context.scope.agent_id.clone(), + run_context.scope.project_id.clone(), + ) + } + fn skill_md(name: &str, description: &str, prompt: &str) -> String { format!( "---\nname: {name}\ndescription: {description}\nactivation:\n keywords: [\"{name}\"]\n---\n\n{prompt}" @@ -378,6 +450,7 @@ mod tests { Arc::new(UnavailableModelGateway), Arc::new(InMemoryLoopHostMilestoneSink::default()), None, + None, ) .expect("local-dev capability wiring"); @@ -991,6 +1064,10 @@ mod tests { result_writer, milestone_sink: Arc::new(InMemoryLoopHostMilestoneSink::default()), skill_activation_source: Some(Arc::clone(&activation_source)), + outbound_preferences_facade: None, + outbound_delivery_target_set_requires_approval: false, + approval_requests: local_runtime.approval_requests.clone(), + capability_leases: local_runtime.capability_leases.clone(), }; let port = factory .create_capability_port(&run_context) @@ -1104,6 +1181,7 @@ mod tests { Arc::new(UnavailableModelGateway), Arc::new(InMemoryLoopHostMilestoneSink::default()), Some(skill_context.activation_source), + None, ) .expect("capability wiring"); let port = wiring @@ -1120,10 +1198,549 @@ mod tests { surface .descriptors .iter() - .any(|descriptor| descriptor.capability_id.as_str() == SKILL_ACTIVATE_CAPABILITY_ID) + .any(|descriptor| descriptor.capability_id.as_str() == SKILL_ACTIVATE_CAPABILITY_ID) ); } + #[tokio::test] + async fn local_dev_outbound_delivery_capabilities_use_provider_backed_facade() { + let dir = tempfile::tempdir().expect("tempdir"); + let services = crate::build_reborn_services(crate::RebornBuildInput::local_dev( + "local-dev-outbound-delivery-owner", + dir.path().join("local-dev"), + )) + .await + .expect("local-dev services build"); + let runtime = services.host_runtime.clone().expect("host runtime"); + let local_runtime = services + .local_runtime + .as_ref() + .expect("local runtime substrate"); + let slack_target_id = + RebornOutboundDeliveryTargetId::new("slack:test-dm").expect("target id"); + let slack_target_summary = RebornOutboundDeliveryTargetSummary::new( + slack_target_id.clone(), + "slack", + "Slack DM", + Some("Personal Slack direct message".to_string()), + ) + .expect("target summary"); + let slack_target_capabilities = RebornOutboundDeliveryTargetCapabilities { + final_replies: true, + gate_prompts: false, + auth_prompts: false, + }; + let slack_reply_target = + ReplyTargetBindingRef::new("reply:test:slack-dm").expect("reply target"); + let slack_provider = Arc::new(StaticOutboundDeliveryTargetProvider::new( + OutboundDeliveryTargetEntry { + summary: slack_target_summary, + capabilities: slack_target_capabilities, + reply_target_binding_ref: slack_reply_target.clone(), + }, + )); + let slack_provider_delegate: Arc = + slack_provider.clone(); + let target_provider: Arc = + Arc::new(OutboundDeliveryTargetRegistry::new(vec![ + slack_provider_delegate, + ])); + let outbound_preferences_facade: Arc = + Arc::new(RebornOutboundPreferencesFacade::new( + Arc::clone(&local_runtime.outbound_preferences), + target_provider, + )); + let policy = Arc::clone(&local_runtime.capability_policy); + let capability_io = Arc::new(LocalDevCapabilityIo::default()); + let input_resolver: Arc = capability_io.clone(); + let result_writer: Arc = capability_io.clone(); + let fallback_user_id = UserId::new("outbound-delivery-fallback-user").expect("user id"); + let factory = LocalDevLoopCapabilityPortFactory { + runtime, + fallback_user_id: fallback_user_id.clone(), + policy, + workspace_mounts: local_runtime.workspace_mounts.clone(), + memory_mounts: local_runtime.memory_mounts.clone(), + extension_surface_source: LocalDevExtensionSurfaceSource::default(), + input_resolver, + result_writer, + milestone_sink: Arc::new(InMemoryLoopHostMilestoneSink::default()), + skill_activation_source: None, + outbound_preferences_facade: Some(outbound_preferences_facade), + outbound_delivery_target_set_requires_approval: true, + approval_requests: local_runtime.approval_requests.clone(), + capability_leases: local_runtime.capability_leases.clone(), + }; + + let owner_user_id = UserId::new("outbound-delivery-owner").expect("user id"); + let actor_user_id = UserId::new("outbound-delivery-actor").expect("user id"); + let run_context = run_context_with_scope(TurnScope::new_with_owner( + TenantId::new("tenant-outbound-delivery").expect("tenant id"), + Some(AgentId::new("agent-outbound-delivery").expect("agent id")), + Some(ProjectId::new("project-outbound-delivery").expect("project id")), + ThreadId::new("thread-outbound-delivery").expect("thread id"), + Some(owner_user_id.clone()), + )) + .await + .with_actor(TurnActor::new(actor_user_id.clone())); + let expected_provider_caller = + expected_outbound_delivery_caller(&run_context, owner_user_id.clone()); + slack_provider.expect_caller(expected_provider_caller.clone()); + let port = factory + .create_capability_port(&run_context) + .await + .expect("capability port"); + let surface = port + .visible_capabilities(VisibleCapabilityRequest {}) + .await + .expect("visible surface"); + let descriptor_ids = surface + .descriptors + .iter() + .map(|descriptor| descriptor.capability_id.as_str()) + .collect::>(); + assert!(descriptor_ids.contains(&OUTBOUND_DELIVERY_TARGETS_LIST_CAPABILITY_ID)); + assert!(descriptor_ids.contains(&OUTBOUND_DELIVERY_TARGET_SET_CAPABILITY_ID)); + let tool_definition_names = port + .tool_definitions() + .expect("tool definitions") + .into_iter() + .map(|definition| definition.name) + .collect::>(); + assert!(tool_definition_names.contains(&"builtin__outbound_delivery_targets_list".into())); + assert!(tool_definition_names.contains(&"builtin__outbound_delivery_target_set".into())); + + let malformed_list = port + .register_provider_tool_call(provider_tool_call_with_name( + "builtin__outbound_delivery_targets_list", + serde_json::Value::Null, + )) + .await + .expect_err("malformed list input should fail validation"); + assert_eq!( + malformed_list.kind, + AgentLoopHostErrorKind::InvalidInvocation + ); + + let list_candidate = port + .register_provider_tool_call(provider_tool_call_with_name( + "builtin__outbound_delivery_targets_list", + serde_json::json!({ "channel": "slack" }), + )) + .await + .expect("list call stages"); + let list_outcome = port + .invoke_capability(CapabilityInvocation { + surface_version: list_candidate.surface_version, + capability_id: list_candidate.capability_id, + input_ref: list_candidate.input_ref, + approval_resume: None, + auth_resume: None, + }) + .await + .expect("list call invokes"); + let list_result_ref = match list_outcome { + CapabilityOutcome::Completed(message) => message.result_ref, + outcome => panic!("list should complete, got {outcome:?}"), + }; + let list_output = capability_io + .result_output(list_result_ref.as_str()) + .expect("result read succeeds") + .expect("result output exists"); + assert_eq!( + list_output["targets"][0]["target"]["target_id"], + slack_target_id.as_str() + ); + assert_eq!(list_output["targets"][0]["target"]["channel"], "slack"); + assert_eq!( + slack_provider.observed_callers(), + vec![expected_provider_caller.clone()] + ); + + let malformed_set = port + .register_provider_tool_call(provider_tool_call_with_name( + "builtin__outbound_delivery_target_set", + serde_json::json!({ "target_id": "bad\nid" }), + )) + .await + .expect_err("malformed set input should fail validation"); + assert_eq!( + malformed_set.kind, + AgentLoopHostErrorKind::InvalidInvocation + ); + + let owner_preference_key = CommunicationPreferenceKey::personal( + run_context.scope.tenant_id.clone(), + owner_user_id.clone(), + ); + let actor_preference_key = CommunicationPreferenceKey::personal( + run_context.scope.tenant_id.clone(), + actor_user_id.clone(), + ); + let set_candidate = port + .register_provider_tool_call(provider_tool_call_with_name( + "builtin__outbound_delivery_target_set", + serde_json::json!({ "target_id": slack_target_id.as_str() }), + )) + .await + .expect("set call stages"); + let set_surface_version = set_candidate.surface_version.clone(); + let set_capability_id_from_candidate = set_candidate.capability_id.clone(); + let set_input_ref = set_candidate.input_ref.clone(); + let blocked_outcome = port + .invoke_capability(CapabilityInvocation { + surface_version: set_surface_version.clone(), + capability_id: set_capability_id_from_candidate.clone(), + input_ref: set_input_ref.clone(), + approval_resume: None, + auth_resume: None, + }) + .await + .expect("set call reaches approval gate"); + let approval_resume = match blocked_outcome { + CapabilityOutcome::ApprovalRequired { + gate_ref, + approval_resume: Some(resume), + .. + } => { + assert!(gate_ref.as_str().starts_with("gate:approval-")); + resume + } + outcome => panic!("set should require approval, got {outcome:?}"), + }; + assert!( + local_runtime + .outbound_preferences + .load_communication_preference(owner_preference_key.clone()) + .await + .expect("owner preference read before approval") + .is_none() + ); + assert!( + local_runtime + .outbound_preferences + .load_communication_preference(actor_preference_key.clone()) + .await + .expect("actor preference read before approval") + .is_none() + ); + + let set_capability_id = + CapabilityId::new(OUTBOUND_DELIVERY_TARGET_SET_CAPABILITY_ID).expect("capability id"); + let invocation_id = InvocationId::parse(approval_resume.resume_token.as_str()) + .expect("resume token carries invocation id"); + let mut approval_scope = run_context.scope.to_resource_scope(); + approval_scope.user_id = owner_user_id.clone(); + approval_scope.invocation_id = invocation_id; + let approval = local_runtime + .capability_policy + .lease_approval_for( + crate::local_dev_capability_policy::LocalDevApprovalPolicyAction::Dispatch { + capability: &set_capability_id, + }, + &local_runtime.workspace_mounts, + &local_runtime.skill_mounts, + &local_runtime.memory_mounts, + ) + .expect("outbound delivery approval lease terms"); + ApprovalResolver::new( + local_runtime.approval_requests.as_ref(), + local_runtime.capability_leases.as_ref(), + ) + .approve_dispatch( + &approval_scope, + approval_resume.approval_request_id, + approval, + ) + .await + .expect("approval issues dispatch lease"); + + let set_outcome = port + .invoke_capability(CapabilityInvocation { + surface_version: set_surface_version, + capability_id: set_capability_id_from_candidate, + input_ref: CapabilityInputRef::new("input:stale-approval-resume") + .expect("stale input ref"), + approval_resume: Some(approval_resume), + auth_resume: None, + }) + .await + .expect("approved set call invokes"); + let set_result_ref = match set_outcome { + CapabilityOutcome::Completed(message) => message.result_ref, + outcome => panic!("approved set should complete, got {outcome:?}"), + }; + let set_output = capability_io + .result_output(set_result_ref.as_str()) + .expect("set result read succeeds") + .expect("set result output exists"); + assert_eq!( + set_output["final_reply_target"]["target_id"], + slack_target_id.as_str() + ); + let owner_preference = local_runtime + .outbound_preferences + .load_communication_preference(owner_preference_key) + .await + .expect("owner preference read after approval") + .expect("owner preference persisted"); + assert_eq!( + owner_preference + .record + .final_reply_target + .as_ref() + .map(|target| target.as_str()), + Some(slack_reply_target.as_str()) + ); + assert!( + local_runtime + .outbound_preferences + .load_communication_preference(actor_preference_key) + .await + .expect("actor preference read after approval") + .is_none() + ); + let leases = local_runtime + .capability_leases + .leases_for_scope(&approval_scope) + .await; + assert!(leases.iter().any(|lease| { + lease.status == CapabilityLeaseStatus::Consumed + && lease.grant.capability == set_capability_id + })); + let observed_provider_callers = slack_provider.observed_callers(); + assert!( + observed_provider_callers + .iter() + .all(|caller| caller == &expected_provider_caller), + "outbound target provider should be scoped to owner caller: {observed_provider_callers:?}" + ); + assert!( + observed_provider_callers.len() >= 2, + "list and set target resolution should call the outbound target provider" + ); + } + + #[tokio::test] + async fn local_dev_yolo_outbound_delivery_target_set_bypasses_approval_gate() { + let dir = tempfile::tempdir().expect("tempdir"); + let services = crate::build_reborn_services( + crate::RebornBuildInput::local_dev( + "local-yolo-outbound-delivery-owner", + dir.path().join("local-dev"), + ) + .with_runtime_policy(local_dev_minimal_approval_policy()), + ) + .await + .expect("local-dev-yolo services build"); + let local_runtime = services + .local_runtime + .as_ref() + .expect("local runtime substrate"); + let slack_target_id = + RebornOutboundDeliveryTargetId::new("slack:yolo-dm").expect("target id"); + let slack_target_summary = RebornOutboundDeliveryTargetSummary::new( + slack_target_id.clone(), + "slack", + "Slack DM", + Some("Personal Slack direct message".to_string()), + ) + .expect("target summary"); + let slack_reply_target = + ReplyTargetBindingRef::new("reply:test:yolo-slack-dm").expect("reply target"); + let slack_provider = Arc::new(StaticOutboundDeliveryTargetProvider::new( + OutboundDeliveryTargetEntry { + summary: slack_target_summary, + capabilities: RebornOutboundDeliveryTargetCapabilities { + final_replies: true, + gate_prompts: false, + auth_prompts: false, + }, + reply_target_binding_ref: slack_reply_target.clone(), + }, + )); + let slack_provider_delegate: Arc = + slack_provider.clone(); + let target_provider: Arc = + Arc::new(OutboundDeliveryTargetRegistry::new(vec![ + slack_provider_delegate, + ])); + let outbound_preferences_facade: Arc = + Arc::new(RebornOutboundPreferencesFacade::new( + Arc::clone(&local_runtime.outbound_preferences), + target_provider, + )); + let owner_user_id = UserId::new("local-yolo-outbound-owner").expect("user id"); + let actor_user_id = UserId::new("local-yolo-outbound-actor").expect("user id"); + let run_context = run_context_with_scope(TurnScope::new_with_owner( + TenantId::new("tenant-local-yolo-outbound").expect("tenant id"), + Some(AgentId::new("agent-local-yolo-outbound").expect("agent id")), + Some(ProjectId::new("project-local-yolo-outbound").expect("project id")), + ThreadId::new("thread-local-yolo-outbound").expect("thread id"), + Some(owner_user_id.clone()), + )) + .await + .with_actor(TurnActor::new(actor_user_id.clone())); + let expected_provider_caller = + expected_outbound_delivery_caller(&run_context, owner_user_id.clone()); + slack_provider.expect_caller(expected_provider_caller.clone()); + let thread_scope = ThreadScope { + tenant_id: run_context.scope.tenant_id.clone(), + agent_id: run_context.scope.agent_id.clone().expect("agent id"), + project_id: run_context.scope.project_id.clone(), + owner_user_id: Some(owner_user_id.clone()), + mission_id: None, + }; + let wiring = capability_wiring( + &services, + Arc::new(InMemorySessionThreadService::default()), + thread_scope, + UserId::new("local-yolo-outbound-fallback").expect("user id"), + Arc::clone(&local_runtime.capability_policy), + Arc::new(UnavailableModelGateway), + Arc::new(InMemoryLoopHostMilestoneSink::default()), + None, + Some(outbound_preferences_facade), + ) + .expect("capability wiring"); + let port = wiring + .capability_factory + .create_capability_port(&run_context) + .await + .expect("capability port"); + + let owner_preference_key = CommunicationPreferenceKey::personal( + run_context.scope.tenant_id.clone(), + owner_user_id.clone(), + ); + let actor_preference_key = CommunicationPreferenceKey::personal( + run_context.scope.tenant_id.clone(), + actor_user_id.clone(), + ); + let set_candidate = port + .register_provider_tool_call(provider_tool_call_with_name( + "builtin__outbound_delivery_target_set", + serde_json::json!({ "target_id": slack_target_id.as_str() }), + )) + .await + .expect("set call stages"); + let set_outcome = port + .invoke_capability(CapabilityInvocation { + surface_version: set_candidate.surface_version, + capability_id: set_candidate.capability_id, + input_ref: set_candidate.input_ref, + approval_resume: None, + auth_resume: None, + }) + .await + .expect("set call invokes"); + assert!( + matches!(set_outcome, CapabilityOutcome::Completed(_)), + "local-dev-yolo should bypass approval gate, got {set_outcome:?}" + ); + let observed_provider_callers = slack_provider.observed_callers(); + assert!( + !observed_provider_callers.is_empty(), + "set target should resolve through the outbound target provider" + ); + assert!( + observed_provider_callers + .iter() + .all(|caller| caller == &expected_provider_caller), + "outbound target provider should be scoped to owner caller: {observed_provider_callers:?}" + ); + let owner_preference = local_runtime + .outbound_preferences + .load_communication_preference(owner_preference_key) + .await + .expect("owner preference read after direct set") + .expect("owner preference persisted"); + assert_eq!( + owner_preference + .record + .final_reply_target + .as_ref() + .map(|target| target.as_str()), + Some(slack_reply_target.as_str()) + ); + assert!( + local_runtime + .outbound_preferences + .load_communication_preference(actor_preference_key) + .await + .expect("actor preference read after direct set") + .is_none() + ); + } + + #[tokio::test] + async fn local_dev_outbound_delivery_capabilities_hidden_without_provider_facade() { + let dir = tempfile::tempdir().expect("tempdir"); + let services = crate::build_reborn_services(crate::RebornBuildInput::local_dev( + "local-dev-no-outbound-provider-owner", + dir.path().join("local-dev"), + )) + .await + .expect("local-dev services build"); + let runtime = services.host_runtime.clone().expect("host runtime"); + let local_runtime = services + .local_runtime + .as_ref() + .expect("local runtime substrate"); + let policy = Arc::new( + crate::local_dev_capability_policy::local_dev_capability_policy() + .expect("policy parses"), + ); + let capability_io = Arc::new(LocalDevCapabilityIo::default()); + let input_resolver: Arc = capability_io.clone(); + let result_writer: Arc = capability_io; + let factory = LocalDevLoopCapabilityPortFactory { + runtime, + fallback_user_id: UserId::new("outbound-delivery-fallback-user").expect("user id"), + policy, + workspace_mounts: local_runtime.workspace_mounts.clone(), + memory_mounts: local_runtime.memory_mounts.clone(), + extension_surface_source: LocalDevExtensionSurfaceSource::default(), + input_resolver, + result_writer, + milestone_sink: Arc::new(InMemoryLoopHostMilestoneSink::default()), + skill_activation_source: None, + outbound_preferences_facade: None, + outbound_delivery_target_set_requires_approval: false, + approval_requests: local_runtime.approval_requests.clone(), + capability_leases: local_runtime.capability_leases.clone(), + }; + let run_context = run_context("outbound-delivery-hidden") + .await + .with_actor(TurnActor::new( + UserId::new("outbound-delivery-actor").expect("user id"), + )); + let port = factory + .create_capability_port(&run_context) + .await + .expect("capability port"); + let surface = port + .visible_capabilities(VisibleCapabilityRequest {}) + .await + .expect("visible surface"); + let descriptor_ids = surface + .descriptors + .iter() + .map(|descriptor| descriptor.capability_id.as_str()) + .collect::>(); + + assert!(!descriptor_ids.contains(&OUTBOUND_DELIVERY_TARGETS_LIST_CAPABILITY_ID)); + assert!(!descriptor_ids.contains(&OUTBOUND_DELIVERY_TARGET_SET_CAPABILITY_ID)); + let tool_definition_names = port + .tool_definitions() + .expect("tool definitions") + .into_iter() + .map(|definition| definition.name) + .collect::>(); + assert!(!tool_definition_names.contains(&"builtin__outbound_delivery_targets_list".into())); + assert!(!tool_definition_names.contains(&"builtin__outbound_delivery_target_set".into())); + } + #[tokio::test] async fn local_yolo_capability_port_reads_confirmed_host_mount() { let dir = tempfile::tempdir().expect("tempdir"); // safety: test-only setup in #[cfg(test)] module. @@ -1184,6 +1801,10 @@ mod tests { result_writer, milestone_sink: Arc::new(InMemoryLoopHostMilestoneSink::default()), skill_activation_source: None, + outbound_preferences_facade: None, + outbound_delivery_target_set_requires_approval: false, + approval_requests: local_runtime.approval_requests.clone(), + capability_leases: local_runtime.capability_leases.clone(), }; let run_context = run_context("host-mount-read").await; let port = factory @@ -1406,6 +2027,10 @@ mod tests { result_writer, milestone_sink: Arc::new(InMemoryLoopHostMilestoneSink::default()), skill_activation_source: None, + outbound_preferences_facade: None, + outbound_delivery_target_set_requires_approval: false, + approval_requests: local_runtime.approval_requests.clone(), + capability_leases: local_runtime.capability_leases.clone(), }; let run_context = run_context("skill-install-write").await; let port = factory @@ -1498,6 +2123,10 @@ mod tests { result_writer, milestone_sink: Arc::new(InMemoryLoopHostMilestoneSink::default()), skill_activation_source: None, + outbound_preferences_facade: None, + outbound_delivery_target_set_requires_approval: false, + approval_requests: local_runtime.approval_requests.clone(), + capability_leases: local_runtime.capability_leases.clone(), }; let run_context = run_context("no-host-disclosure").await; let port = factory @@ -1658,6 +2287,7 @@ mod tests { Arc::new(UnavailableModelGateway), Arc::new(InMemoryLoopHostMilestoneSink::default()), None, + None, ) .expect("local-dev capability wiring"); assert_github_capabilities_visible(&wiring, &run_context).await; @@ -1693,6 +2323,7 @@ mod tests { Arc::new(UnavailableModelGateway), Arc::new(InMemoryLoopHostMilestoneSink::default()), None, + None, ) .expect("local-dev capability wiring"); let port = wiring @@ -1817,6 +2448,7 @@ mod tests { Arc::new(UnavailableModelGateway), Arc::new(InMemoryLoopHostMilestoneSink::default()), None, + None, ) .expect("local-dev capability wiring"); let port = wiring diff --git a/crates/ironclaw_reborn_composition/src/slack_host_beta.rs b/crates/ironclaw_reborn_composition/src/slack_host_beta.rs index a1636d94116..bb30b358a37 100644 --- a/crates/ironclaw_reborn_composition/src/slack_host_beta.rs +++ b/crates/ironclaw_reborn_composition/src/slack_host_beta.rs @@ -253,7 +253,6 @@ pub struct SlackHostBetaMounts { pub events: PublicRouteMount, pub personal_binding_pairing: SlackPersonalBindingPairingRouteConfig, pub channel_routes: SlackChannelRouteAdminRouteConfig, - /// Internal target-authority handle consumed only by WebUI product-facade composition. pub(crate) outbound_delivery_target_provider: Arc, } @@ -441,11 +440,8 @@ pub fn build_slack_host_beta_mounts( } } - Ok(SlackHostBetaMounts { - events, - personal_binding_pairing: SlackPersonalBindingPairingRouteConfig::new(pairing), - channel_routes, - outbound_delivery_target_provider: Arc::new(SlackHostBetaOutboundTargetProvider::new( + let outbound_delivery_target_provider: Arc = + Arc::new(SlackHostBetaOutboundTargetProvider::new( SlackOutboundTargetProviderConfig { tenant_id: config.tenant_id.clone(), agent_id: config.agent_id.clone(), @@ -465,7 +461,12 @@ pub fn build_slack_host_beta_mounts( }, channel_route_store, Arc::clone(&personal_dm_target_store), - )), + )); + Ok(SlackHostBetaMounts { + events, + personal_binding_pairing: SlackPersonalBindingPairingRouteConfig::new(pairing), + channel_routes, + outbound_delivery_target_provider, }) } diff --git a/docs/reborn/contracts/triggers.md b/docs/reborn/contracts/triggers.md index aee3e3c24a2..8a28d35bd3e 100644 --- a/docs/reborn/contracts/triggers.md +++ b/docs/reborn/contracts/triggers.md @@ -417,7 +417,11 @@ Capability follow-ups before launch: ## 9. Delivery -Trigger delivery is fast-follow. +Trigger delivery target selection is outside trigger identity and goes through +the outbound delivery track. Product-facing outbound preference APIs and +explicit provider-backed target tooling own discovery and approval-gated +selection; durable selection remains product-owned and trigger records still do +not embed delivery targets. - Trigger ingress identity must not include delivery targets. - Trigger record identity must not include delivery targets.