use crate::canonical::{InterpretationId, ToolActionId, UnitId};
use crate::id::{ConnectionId, ExternalSessionId, MonoloopRunId};
use serde::{Deserialize, Serialize};
#[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct LoopId(String);
impl LoopId {
pub fn new(value: impl Into<String>) -> Self {
Self(value.into())
}
pub fn generate() -> Self {
Self(uuid::Uuid::new_v4().to_string())
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct ToolExecutionId(String);
impl ToolExecutionId {
pub fn new(value: impl Into<String>) -> Self {
Self(value.into())
}
pub fn generate() -> Self {
Self(uuid::Uuid::new_v4().to_string())
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Clone, Debug)]
pub struct LoopLimits {
pub max_tool_actions: usize,
pub max_concurrent_executions: usize,
pub max_queued_ready: usize,
pub max_output_queue: usize,
pub max_dedup_entries: usize,
}
impl Default for LoopLimits {
fn default() -> Self {
Self {
max_tool_actions: 1024,
max_concurrent_executions: 32,
max_queued_ready: 256,
max_output_queue: 4096,
max_dedup_entries: 4096,
}
}
}
#[derive(Clone, Debug)]
pub struct LoopScope {
pub monoloop_run_id: MonoloopRunId,
pub loop_id: LoopId,
pub accepted_interpretation_ids: Vec<InterpretationId>,
pub accepted_connection_ids: Vec<ConnectionId>,
pub accepted_external_session_ids: Vec<ExternalSessionId>,
pub accept_all_in_run: bool,
}
impl Default for LoopScope {
fn default() -> Self {
Self {
monoloop_run_id: MonoloopRunId::generate(),
loop_id: LoopId::generate(),
accepted_interpretation_ids: Vec::new(),
accepted_connection_ids: Vec::new(),
accepted_external_session_ids: Vec::new(),
accept_all_in_run: true,
}
}
}
impl LoopScope {
pub fn single(
run_id: MonoloopRunId,
loop_id: LoopId,
interpretation_id: InterpretationId,
connection_id: ConnectionId,
external_session_id: Option<ExternalSessionId>,
) -> Self {
let mut sessions = Vec::new();
if let Some(s) = external_session_id {
sessions.push(s);
}
Self {
monoloop_run_id: run_id,
loop_id,
accepted_interpretation_ids: vec![interpretation_id],
accepted_connection_ids: vec![connection_id],
accepted_external_session_ids: sessions,
accept_all_in_run: false,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum ToolUnavailableReason {
NoRegisteredTool,
NotFound,
Denied,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum OutboundToolOutcome {
Success,
ToolUnavailable,
DispatchRejected,
ExecutionFailed,
Cancelled,
ExecutionLost,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct OutboundToolResult {
pub outbound_result_id: String,
pub monoloop_run_id: MonoloopRunId,
pub loop_id: LoopId,
pub source_interpretation_id: InterpretationId,
pub source_connection_id: ConnectionId,
pub external_session_id: Option<ExternalSessionId>,
pub tool_action_id: ToolActionId,
pub request_generation: u64,
pub tool_execution_id: Option<ToolExecutionId>,
pub outcome: OutboundToolOutcome,
pub payload: String,
pub source_unit_id: UnitId,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum LoopOutputEvent {
ToolDispatchRequested {
tool_action_id: ToolActionId,
request_generation: u64,
},
ToolUnavailable {
tool_action_id: ToolActionId,
reason: ToolUnavailableReason,
},
OutboundToolResult(OutboundToolResult),
Diagnostic {
message: String,
},
LoopEnded(LoopEnd),
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct LoopEnd {
pub monoloop_run_id: MonoloopRunId,
pub loop_id: LoopId,
pub kind: LoopEndKind,
pub delivery_events_received: u64,
pub duplicate_events: u64,
pub tools_unavailable: u64,
pub outbound_results_emitted: u64,
pub safe_diagnostics: Vec<String>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum LoopEndKind {
Drained,
Cancelled,
SubscriptionLost,
OutputFailed,
InvariantFailed,
ConfigurationFailed,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum LoopErrorKind {
EventOutOfScope,
DeliverySequenceGap,
UnitIdentityConflict,
ToolRequestIncomplete,
LimitExceeded,
Cancelled,
InvariantViolation,
ConfigurationInvalid,
OutputBackpressure,
}
#[derive(Clone, Debug, thiserror::Error, PartialEq, Eq)]
#[error("{kind:?}: {message}")]
pub struct LoopError {
pub kind: LoopErrorKind,
pub message: String,
}
impl LoopError {
pub fn new(kind: LoopErrorKind, message: impl Into<String>) -> Self {
Self {
kind,
message: message.into(),
}
}
pub fn cancelled() -> Self {
Self::new(LoopErrorKind::Cancelled, "loop cancelled")
}
pub fn gap() -> Self {
Self::new(
LoopErrorKind::DeliverySequenceGap,
"subscription gap detected",
)
}
pub fn limit(message: impl Into<String>) -> Self {
Self::new(LoopErrorKind::LimitExceeded, message)
}
}