use crate::control::ControlDisposition;
use crate::instance::ConnectorInstanceId;
use monoloop_contracts::{ChannelId, ExternalSessionId, SessionConfig, SessionId, TransactionId};
use secrecy::{ExposeSecret, SecretString};
use std::fmt;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::time::Instant;
use thiserror::Error;
pub trait SessionRoute: Send + Sync {
fn owner(&self) -> &ConnectorInstanceId;
}
pub struct SessionAttachment {
pub owner: ConnectorInstanceId,
pub external_session_id: ExternalSessionId,
pub effective_session_config: SessionConfig,
pub route: Arc<dyn SessionRoute>,
pub create_mode: bool,
pub initial_mcp: Option<McpServerDescriptor>,
}
impl fmt::Debug for SessionAttachment {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("SessionAttachment")
.field("owner", &self.owner)
.field("external_session_id", &"<redacted>")
.field("effective_session_config", &self.effective_session_config)
.field("route_owner", self.route.owner())
.field("create_mode", &self.create_mode)
.field(
"initial_mcp",
&self.initial_mcp.as_ref().map(|_| "<present>"),
)
.finish()
}
}
impl SessionAttachment {
pub fn new(
owner: ConnectorInstanceId,
external_session_id: ExternalSessionId,
effective_session_config: SessionConfig,
route: Arc<dyn SessionRoute>,
) -> Self {
Self {
owner,
external_session_id,
effective_session_config,
route,
create_mode: false,
initial_mcp: None,
}
}
pub fn new_create(
owner: ConnectorInstanceId,
provisional_id: ExternalSessionId,
effective_session_config: SessionConfig,
route: Arc<dyn SessionRoute>,
initial_mcp: Option<McpServerDescriptor>,
) -> Self {
Self {
owner,
external_session_id: provisional_id,
effective_session_config,
route,
create_mode: true,
initial_mcp,
}
}
}
pub struct McpServerDescriptor {
pub server_name: String,
pub protocol_version: String,
capability_url: SecretString,
}
impl McpServerDescriptor {
pub const MAX_NAME_BYTES: usize = 128;
pub const MAX_PROTOCOL_BYTES: usize = 64;
pub const MAX_URL_BYTES: usize = 512;
pub fn try_new(
server_name: impl Into<String>,
protocol_version: impl Into<String>,
capability_url: impl Into<String>,
) -> Result<Self, SessionAttachError> {
let server_name = server_name.into();
let protocol_version = protocol_version.into();
let url = capability_url.into();
if server_name.is_empty()
|| server_name.len() > Self::MAX_NAME_BYTES
|| server_name.chars().any(|c| c.is_control())
{
return Err(SessionAttachError::InvalidMcpDescriptor);
}
if protocol_version.is_empty()
|| protocol_version.len() > Self::MAX_PROTOCOL_BYTES
|| protocol_version.chars().any(|c| c.is_control())
{
return Err(SessionAttachError::InvalidMcpDescriptor);
}
if url.is_empty() || url.len() > Self::MAX_URL_BYTES {
return Err(SessionAttachError::InvalidMcpDescriptor);
}
Ok(Self {
server_name,
protocol_version,
capability_url: SecretString::from(url),
})
}
pub fn expose_capability_url(&self) -> &str {
self.capability_url.expose_secret()
}
}
impl fmt::Debug for McpServerDescriptor {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("McpServerDescriptor")
.field("server_name", &self.server_name)
.field("protocol_version", &self.protocol_version)
.field("capability_url", &"<redacted>")
.finish()
}
}
#[derive(Clone, Debug)]
pub struct SessionAttachRequest {
pub transaction_id: TransactionId,
pub channel_id: ChannelId,
pub requested_session_id: Option<SessionId>,
pub session_config: SessionConfig,
pub initial_mcp: Option<McpServerDescriptor>,
pub deadline: Instant,
}
impl Clone for McpServerDescriptor {
fn clone(&self) -> Self {
Self {
server_name: self.server_name.clone(),
protocol_version: self.protocol_version.clone(),
capability_url: SecretString::from(self.capability_url.expose_secret().to_string()),
}
}
}
pub trait PendingOperationControl: Send + Sync {
fn cancel(&self) -> ControlDisposition;
fn force_terminate(&self) -> ControlDisposition;
}
pub type SessionAttachmentCompletion = Pin<
Box<dyn Future<Output = Result<Arc<SessionAttachment>, SessionAttachError>> + Send + 'static>,
>;
pub struct PendingSessionAttachment {
pub control: Arc<dyn PendingOperationControl>,
pub completion: SessionAttachmentCompletion,
}
pub type SessionConfigurationCompletion =
Pin<Box<dyn Future<Output = Result<(), SessionConfigurationError>> + Send + 'static>>;
pub struct PendingSessionConfiguration {
pub control: Arc<dyn PendingOperationControl>,
pub completion: SessionConfigurationCompletion,
}
pub trait SessionAdapter: Send + Sync {
fn begin_attach(
&self,
request: SessionAttachRequest,
) -> Result<PendingSessionAttachment, SessionAttachError>;
fn begin_refresh_mcp(
&self,
attachment: Arc<SessionAttachment>,
descriptor: Option<McpServerDescriptor>,
) -> Result<PendingSessionConfiguration, SessionConfigurationError>;
}
#[derive(Clone, Debug, Error, PartialEq, Eq)]
pub enum SessionAttachError {
#[error("session attach cancelled")]
Cancelled,
#[error("session attach terminated")]
Terminated,
#[error("session attach deadline exceeded")]
DeadlineExceeded,
#[error("session configuration mismatch")]
ConfigurationMismatch,
#[error("session configuration unsupported")]
UnsupportedConfiguration,
#[error("session operation failed")]
SessionFailed,
#[error("invalid MCP descriptor")]
InvalidMcpDescriptor,
#[error("session attach invariant failed")]
InvariantFailed,
#[error("session capacity exceeded")]
CapacityExceeded,
#[error("session id mismatch")]
SessionIdMismatch,
}
#[derive(Clone, Debug, Error, PartialEq, Eq)]
pub enum SessionConfigurationError {
#[error("session configuration cancelled")]
Cancelled,
#[error("session configuration terminated")]
Terminated,
#[error("session configuration deadline exceeded")]
DeadlineExceeded,
#[error("session attachment owner mismatch")]
OwnerMismatch,
#[error("MCP configuration not supported")]
Unsupported,
#[error("session configuration failed")]
ConfigurationFailed,
#[error("session configuration invariant failed")]
InvariantFailed,
}
pub fn validate_open_attachment_owner(
instance_id: &ConnectorInstanceId,
attachment: Option<&SessionAttachment>,
) -> Result<(), monoloop_contracts::ConnectorError> {
if let Some(att) = attachment {
if &att.owner != instance_id {
return Err(monoloop_contracts::ConnectorError::new(
monoloop_contracts::ConnectorErrorKind::ConfigurationInvalid,
"session attachment owner does not match connector instance",
));
}
if att.route.owner() != instance_id {
return Err(monoloop_contracts::ConnectorError::new(
monoloop_contracts::ConnectorErrorKind::InvariantViolation,
"session route owner does not match attachment owner",
));
}
}
Ok(())
}
pub fn validate_session_id_match(
requested: Option<&SessionId>,
external: &ExternalSessionId,
) -> Result<(), SessionAttachError> {
if let Some(req) = requested {
if req.as_str() != external.as_str() {
return Err(SessionAttachError::SessionIdMismatch);
}
}
Ok(())
}