use std::any::Any;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::completion_feed::CompletionSeq;
use crate::handles::{
CommsDrainHandle, ExternalToolSurfaceHandle, GeneratedAuthLeaseHandle,
GeneratedPeerCommsInstallFactory, InteractionStreamHandle, McpServerLifecycleHandle,
ModelRoutingHandle, PeerCommsHandle, PeerCommsInstallTarget, PeerInteractionHandle,
SessionAdmissionHandle, SessionClaimHandle, SessionContextHandle, TurnStateHandle,
};
use crate::ops_lifecycle::CompletionCursorConsumer;
use crate::ops_lifecycle::OpsLifecycleRegistry;
use crate::tool_scope::GeneratedToolVisibilityOwner;
use crate::types::SessionId;
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct RuntimeEpochId(pub Uuid);
impl RuntimeEpochId {
pub fn new() -> Self {
Self(crate::time_compat::new_uuid_v7())
}
pub fn from_uuid(uuid: Uuid) -> Self {
Self(uuid)
}
}
impl Default for RuntimeEpochId {
fn default() -> Self {
Self::new()
}
}
impl std::fmt::Display for RuntimeEpochId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
pub struct EpochCursorState {
agent_applied_cursor: AtomicU64,
runtime_observed_seq: AtomicU64,
runtime_last_injected_seq: AtomicU64,
}
impl EpochCursorState {
pub fn new() -> Self {
Self {
agent_applied_cursor: AtomicU64::new(0),
runtime_observed_seq: AtomicU64::new(0),
runtime_last_injected_seq: AtomicU64::new(0),
}
}
pub fn from_recovered(
agent_applied_cursor: CompletionSeq,
runtime_observed_seq: CompletionSeq,
runtime_last_injected_seq: CompletionSeq,
) -> Self {
Self {
agent_applied_cursor: AtomicU64::new(agent_applied_cursor),
runtime_observed_seq: AtomicU64::new(runtime_observed_seq),
runtime_last_injected_seq: AtomicU64::new(runtime_last_injected_seq),
}
}
#[doc(hidden)]
pub fn project_authorized_completion_cursor(
&self,
consumer: CompletionCursorConsumer,
cursor: CompletionSeq,
) {
match consumer {
CompletionCursorConsumer::AgentApplied => {
self.agent_applied_cursor.store(cursor, Ordering::Release);
}
CompletionCursorConsumer::RuntimeObserved => {
self.runtime_observed_seq.store(cursor, Ordering::Release);
}
CompletionCursorConsumer::RuntimeInjected => {
self.runtime_last_injected_seq
.store(cursor, Ordering::Release);
}
}
}
}
impl Default for EpochCursorState {
fn default() -> Self {
Self::new()
}
}
impl std::fmt::Debug for EpochCursorState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("EpochCursorState")
.field(
"agent_applied_cursor",
&self.agent_applied_cursor.load(Ordering::Relaxed),
)
.field(
"runtime_observed_seq",
&self.runtime_observed_seq.load(Ordering::Relaxed),
)
.field(
"runtime_last_injected_seq",
&self.runtime_last_injected_seq.load(Ordering::Relaxed),
)
.finish()
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct EpochCursorSnapshot {
pub agent_applied_cursor: CompletionSeq,
pub runtime_observed_seq: CompletionSeq,
pub runtime_last_injected_seq: CompletionSeq,
}
pub struct SessionRuntimeBindings {
session_id: SessionId,
epoch_id: RuntimeEpochId,
ops_lifecycle: Arc<dyn OpsLifecycleRegistry>,
cursor_state: Arc<EpochCursorState>,
tool_visibility_owner: GeneratedToolVisibilityOwner,
turn_state: Arc<dyn TurnStateHandle>,
comms_drain: Arc<dyn CommsDrainHandle>,
external_tool_surface: Arc<dyn ExternalToolSurfaceHandle>,
peer_comms_install: GeneratedPeerCommsInstallFactory,
session_admission: Arc<dyn SessionAdmissionHandle>,
model_routing: Arc<dyn ModelRoutingHandle>,
sticky_model_fallback_commit_coordinator:
Arc<dyn crate::handles::StickyModelFallbackCommitCoordinator>,
auth_lease: GeneratedAuthLeaseHandle,
mcp_server_lifecycle: Arc<dyn McpServerLifecycleHandle>,
peer_interaction: Arc<dyn PeerInteractionHandle>,
session_context: Arc<dyn SessionContextHandle>,
session_claim_handle: Arc<dyn SessionClaimHandle>,
interaction_stream: Arc<dyn InteractionStreamHandle>,
compaction_commit_coordinator: Arc<dyn crate::memory::CompactionCommitCoordinator>,
runtime_authority: Arc<dyn Any + Send + Sync>,
}
impl SessionRuntimeBindings {
#[doc(hidden)]
#[allow(clippy::too_many_arguments)]
pub fn __from_runtime_authority(
session_id: SessionId,
epoch_id: RuntimeEpochId,
ops_lifecycle: Arc<dyn OpsLifecycleRegistry>,
cursor_state: Arc<EpochCursorState>,
tool_visibility_owner: GeneratedToolVisibilityOwner,
turn_state: Arc<dyn TurnStateHandle>,
comms_drain: Arc<dyn CommsDrainHandle>,
external_tool_surface: Arc<dyn ExternalToolSurfaceHandle>,
peer_comms_install: GeneratedPeerCommsInstallFactory,
session_admission: Arc<dyn SessionAdmissionHandle>,
model_routing: Arc<dyn ModelRoutingHandle>,
sticky_model_fallback_commit_coordinator: Arc<
dyn crate::handles::StickyModelFallbackCommitCoordinator,
>,
auth_lease: GeneratedAuthLeaseHandle,
mcp_server_lifecycle: Arc<dyn McpServerLifecycleHandle>,
peer_interaction: Arc<dyn PeerInteractionHandle>,
session_context: Arc<dyn SessionContextHandle>,
session_claim_handle: Arc<dyn SessionClaimHandle>,
interaction_stream: Arc<dyn InteractionStreamHandle>,
compaction_commit_coordinator: Arc<dyn crate::memory::CompactionCommitCoordinator>,
runtime_authority: Arc<dyn Any + Send + Sync>,
) -> Self {
Self {
session_id,
epoch_id,
ops_lifecycle,
cursor_state,
tool_visibility_owner,
turn_state,
comms_drain,
external_tool_surface,
peer_comms_install,
session_admission,
model_routing,
sticky_model_fallback_commit_coordinator,
auth_lease,
mcp_server_lifecycle,
peer_interaction,
session_context,
session_claim_handle,
interaction_stream,
compaction_commit_coordinator,
runtime_authority,
}
}
pub fn session_id(&self) -> &SessionId {
&self.session_id
}
pub fn epoch_id(&self) -> &RuntimeEpochId {
&self.epoch_id
}
pub fn ops_lifecycle(&self) -> &Arc<dyn OpsLifecycleRegistry> {
&self.ops_lifecycle
}
pub fn cursor_state(&self) -> &Arc<EpochCursorState> {
&self.cursor_state
}
pub fn tool_visibility_owner(&self) -> &GeneratedToolVisibilityOwner {
&self.tool_visibility_owner
}
pub fn turn_state(&self) -> &Arc<dyn TurnStateHandle> {
&self.turn_state
}
pub fn comms_drain(&self) -> &Arc<dyn CommsDrainHandle> {
&self.comms_drain
}
pub fn external_tool_surface(&self) -> &Arc<dyn ExternalToolSurfaceHandle> {
&self.external_tool_surface
}
pub fn peer_comms(&self) -> &Arc<dyn PeerCommsHandle> {
self.peer_comms_install.peer_comms_handle()
}
pub fn install_peer_comms_on(
&self,
target: &(dyn PeerCommsInstallTarget + '_),
) -> Result<(), String> {
self.peer_comms_install.install_on_target(target)
}
pub fn session_admission(&self) -> &Arc<dyn SessionAdmissionHandle> {
&self.session_admission
}
pub fn model_routing(&self) -> &Arc<dyn ModelRoutingHandle> {
&self.model_routing
}
pub fn sticky_model_fallback_commit_coordinator(
&self,
) -> &Arc<dyn crate::handles::StickyModelFallbackCommitCoordinator> {
&self.sticky_model_fallback_commit_coordinator
}
pub fn auth_lease(&self) -> &GeneratedAuthLeaseHandle {
&self.auth_lease
}
pub fn mcp_server_lifecycle(&self) -> &Arc<dyn McpServerLifecycleHandle> {
&self.mcp_server_lifecycle
}
pub fn peer_interaction(&self) -> &Arc<dyn PeerInteractionHandle> {
&self.peer_interaction
}
pub fn session_context(&self) -> &Arc<dyn SessionContextHandle> {
&self.session_context
}
pub fn session_claim_handle(&self) -> &Arc<dyn SessionClaimHandle> {
&self.session_claim_handle
}
pub fn interaction_stream(&self) -> &Arc<dyn InteractionStreamHandle> {
&self.interaction_stream
}
pub fn compaction_commit_coordinator(
&self,
) -> &Arc<dyn crate::memory::CompactionCommitCoordinator> {
&self.compaction_commit_coordinator
}
#[doc(hidden)]
pub fn __runtime_authority(&self) -> &(dyn Any + Send + Sync) {
self.runtime_authority.as_ref()
}
}
impl Clone for SessionRuntimeBindings {
fn clone(&self) -> Self {
Self {
session_id: self.session_id.clone(),
epoch_id: self.epoch_id.clone(),
ops_lifecycle: Arc::clone(&self.ops_lifecycle),
cursor_state: Arc::clone(&self.cursor_state),
tool_visibility_owner: self.tool_visibility_owner.clone(),
turn_state: Arc::clone(&self.turn_state),
comms_drain: Arc::clone(&self.comms_drain),
external_tool_surface: Arc::clone(&self.external_tool_surface),
peer_comms_install: self.peer_comms_install.clone(),
session_admission: Arc::clone(&self.session_admission),
model_routing: Arc::clone(&self.model_routing),
sticky_model_fallback_commit_coordinator: Arc::clone(
&self.sticky_model_fallback_commit_coordinator,
),
auth_lease: self.auth_lease.clone(),
mcp_server_lifecycle: Arc::clone(&self.mcp_server_lifecycle),
peer_interaction: Arc::clone(&self.peer_interaction),
session_context: Arc::clone(&self.session_context),
session_claim_handle: Arc::clone(&self.session_claim_handle),
interaction_stream: Arc::clone(&self.interaction_stream),
compaction_commit_coordinator: Arc::clone(&self.compaction_commit_coordinator),
runtime_authority: Arc::clone(&self.runtime_authority),
}
}
}
impl std::fmt::Debug for SessionRuntimeBindings {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SessionRuntimeBindings")
.field("session_id", &self.session_id)
.field("epoch_id", &self.epoch_id)
.field("ops_lifecycle", &"<dyn OpsLifecycleRegistry>")
.field("cursor_state", &self.cursor_state)
.field("tool_visibility_owner", &"<dyn ToolVisibilityOwner>")
.field("turn_state", &"<dyn TurnStateHandle>")
.field("comms_drain", &"<dyn CommsDrainHandle>")
.field("external_tool_surface", &"<dyn ExternalToolSurfaceHandle>")
.field("peer_comms", &"<dyn PeerCommsHandle>")
.field("session_admission", &"<dyn SessionAdmissionHandle>")
.field("model_routing", &"<dyn ModelRoutingHandle>")
.field(
"sticky_model_fallback_commit_coordinator",
&"<dyn StickyModelFallbackCommitCoordinator>",
)
.field("auth_lease", &"<dyn AuthLeaseHandle>")
.field("mcp_server_lifecycle", &"<dyn McpServerLifecycleHandle>")
.field("peer_interaction", &"<dyn PeerInteractionHandle>")
.field("session_context", &"<dyn SessionContextHandle>")
.field("session_claim_handle", &"<dyn SessionClaimHandle>")
.field("interaction_stream", &"<dyn InteractionStreamHandle>")
.field(
"compaction_commit_coordinator",
&"<dyn CompactionCommitCoordinator>",
)
.finish()
}
}
#[allow(clippy::large_enum_variant)]
pub enum RuntimeBuildMode {
StandaloneEphemeral,
SessionOwned(SessionRuntimeBindings),
}
impl Clone for RuntimeBuildMode {
fn clone(&self) -> Self {
match self {
Self::StandaloneEphemeral => Self::StandaloneEphemeral,
Self::SessionOwned(b) => Self::SessionOwned(b.clone()),
}
}
}
impl std::fmt::Debug for RuntimeBuildMode {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::StandaloneEphemeral => write!(f, "StandaloneEphemeral"),
Self::SessionOwned(b) => f
.debug_tuple("SessionOwned")
.field(&b.session_id)
.field(&b.epoch_id)
.finish(),
}
}
}