Skip to main content

monoloop_connector_codex/
channel_binding.rs

1//! WP-11: Codex ACP ChannelBinding + ConnectorFactory.
2
3use crate::CodexConnector;
4use monoloop_connector::{
5    Connector, ConnectorBuildError, ConnectorFactory, ConnectorInstance, ConnectorInstanceId,
6    ControlDisposition, McpServerDescriptor, PendingOperationControl, PendingSessionAttachment,
7    PendingSessionConfiguration, SessionAdapter, SessionAttachError, SessionAttachRequest,
8    SessionAttachment, SessionAttachmentCompletion, SessionConfigurationError, SessionRoute,
9};
10use monoloop_contracts::{
11    ChannelCapabilities, ChannelDefaults, ChannelId, ChannelKind, ChannelLimits,
12    ContinuationPolicy, DialectDescriptor, ExchangeMode, ExternalSessionId,
13    McpConfigurationCapability, McpReachability, SessionId, SessionMode, ToolExecutionMode,
14};
15use std::collections::{BTreeSet, HashMap};
16use std::sync::{Arc, Mutex};
17
18#[derive(Default)]
19/// ConnectorFactory for this profile.
20pub struct CodexConnectorFactory;
21impl CodexConnectorFactory {
22    /// Create a default factory.
23    /// Build a ChannelBinding for StartedRuntime / ChannelRegistry composition.
24    pub fn new() -> Self {
25        Self
26    }
27}
28impl ConnectorFactory for CodexConnectorFactory {
29    fn create(&self) -> Result<ConnectorInstance, ConnectorBuildError> {
30        let instance_id = ConnectorInstanceId::generate();
31        let connector = Arc::new(CodexConnector::new());
32        let sessions = Arc::new(ProfileSessionAdapter::new(instance_id.clone()));
33        Ok(ConnectorInstance::new(
34            instance_id,
35            connector as Arc<dyn Connector>,
36            Some(sessions as Arc<dyn SessionAdapter>),
37        ))
38    }
39}
40
41struct ProfileRoute {
42    owner: ConnectorInstanceId,
43}
44impl SessionRoute for ProfileRoute {
45    fn owner(&self) -> &ConnectorInstanceId {
46        &self.owner
47    }
48}
49struct PendingCtrl {
50    cancel: std::sync::atomic::AtomicBool,
51    terminate: std::sync::atomic::AtomicBool,
52}
53impl PendingOperationControl for PendingCtrl {
54    fn cancel(&self) -> ControlDisposition {
55        if self.cancel.swap(true, std::sync::atomic::Ordering::SeqCst) {
56            ControlDisposition::AlreadyRequested
57        } else {
58            ControlDisposition::Accepted
59        }
60    }
61    fn force_terminate(&self) -> ControlDisposition {
62        if self
63            .terminate
64            .swap(true, std::sync::atomic::Ordering::SeqCst)
65        {
66            ControlDisposition::AlreadyRequested
67        } else {
68            ControlDisposition::Accepted
69        }
70    }
71}
72struct ProfileSessionAdapter {
73    owner: ConnectorInstanceId,
74    known: Arc<Mutex<HashMap<String, monoloop_contracts::SessionConfig>>>,
75}
76impl ProfileSessionAdapter {
77    fn new(owner: ConnectorInstanceId) -> Self {
78        Self {
79            owner,
80            known: Arc::new(Mutex::new(HashMap::new())),
81        }
82    }
83}
84impl SessionAdapter for ProfileSessionAdapter {
85    fn begin_attach(
86        &self,
87        request: SessionAttachRequest,
88    ) -> Result<PendingSessionAttachment, SessionAttachError> {
89        let control = Arc::new(PendingCtrl {
90            cancel: std::sync::atomic::AtomicBool::new(false),
91            terminate: std::sync::atomic::AtomicBool::new(false),
92        });
93        let control_api: Arc<dyn PendingOperationControl> = control.clone();
94        let owner = self.owner.clone();
95        let known = Arc::clone(&self.known);
96
97        let completion: SessionAttachmentCompletion = Box::pin(async move {
98            if control.terminate.load(std::sync::atomic::Ordering::SeqCst) {
99                return Err(SessionAttachError::Terminated);
100            }
101            if control.cancel.load(std::sync::atomic::Ordering::SeqCst) {
102                return Err(SessionAttachError::Cancelled);
103            }
104            let route = Arc::new(ProfileRoute {
105                owner: owner.clone(),
106            });
107            if let Some(ref sid) = request.requested_session_id {
108                let map = known
109                    .lock()
110                    .map_err(|_| SessionAttachError::SessionFailed)?;
111                let cfg = map
112                    .get(sid.as_str())
113                    .cloned()
114                    .unwrap_or_else(|| request.session_config.clone());
115                let ext = ExternalSessionId::try_new(sid.as_str())
116                    .map_err(|_| SessionAttachError::SessionFailed)?;
117                monoloop_connector::validate_session_id_match(Some(sid), &ext)
118                    .map_err(|_| SessionAttachError::SessionIdMismatch)?;
119                Ok(Arc::new(SessionAttachment::new(owner, ext, cfg, route)))
120            } else {
121                let provisional = SessionId::generate();
122                let ext = ExternalSessionId::try_new(provisional.as_str())
123                    .map_err(|_| SessionAttachError::SessionFailed)?;
124                known
125                    .lock()
126                    .map_err(|_| SessionAttachError::SessionFailed)?
127                    .insert(
128                        provisional.as_str().to_string(),
129                        request.session_config.clone(),
130                    );
131                Ok(Arc::new(SessionAttachment::new_create(
132                    owner,
133                    ext,
134                    request.session_config,
135                    route,
136                    request.initial_mcp,
137                )))
138            }
139        });
140        Ok(PendingSessionAttachment {
141            control: control_api,
142            completion,
143        })
144    }
145
146    fn begin_refresh_mcp(
147        &self,
148        attachment: Arc<SessionAttachment>,
149        _: Option<McpServerDescriptor>,
150    ) -> Result<PendingSessionConfiguration, SessionConfigurationError> {
151        if attachment.owner != self.owner {
152            return Err(SessionConfigurationError::OwnerMismatch);
153        }
154        Err(SessionConfigurationError::Unsupported)
155    }
156}
157
158/// Build a ChannelBinding for StartedRuntime / ChannelRegistry composition.
159pub fn codex_channel_binding(
160    id: impl AsRef<str>,
161    endpoint_ref: impl Into<String>,
162    encoder: Arc<dyn monoloop_contracts::OutboundDialectEncoder>,
163    interpreter: Arc<dyn monoloop_interpreter::InterpreterFactory>,
164) -> monoloop_loop::ChannelBinding {
165    let d = DialectDescriptor::codex_acp("1");
166    monoloop_loop::ChannelBinding {
167        id: ChannelId::try_new(id.as_ref()).expect("channel id"),
168        kind: ChannelKind::ExternalAgent,
169        tool_mode: ToolExecutionMode::McpGateway,
170        connector_factory: Arc::new(CodexConnectorFactory::new()),
171        encoder,
172        interpreter,
173        endpoint_ref: endpoint_ref.into(),
174        credential_ref: None,
175        defaults: ChannelDefaults::default(),
176        capabilities: ChannelCapabilities {
177            session_mode: SessionMode::External,
178            mcp_configuration: McpConfigurationCapability::CreationOnly,
179            mcp_reachability: McpReachability::SameLoopbackNamespace,
180            exchange_mode: ExchangeMode::Bidirectional,
181            continuation_policies: BTreeSet::from([ContinuationPolicy::CallerControlled]),
182            supports_distinct_session_concurrency: true,
183            input_dialect: d.clone(),
184            output_dialect: d,
185            option_policy: monoloop_contracts::OptionPolicy::external_agent(),
186        },
187        limits: ChannelLimits::default(),
188    }
189}
190
191#[cfg(test)]
192mod tests {
193    use super::*;
194    #[test]
195    fn codex_binding_validates() {
196        use monoloop_interpreter::DefaultInterpreterFactory;
197        use monoloop_loop::AcpPromptEncoder;
198        let b = codex_channel_binding(
199            "codex-1",
200            "codex:stdio",
201            Arc::new(AcpPromptEncoder::codex()),
202            Arc::new(DefaultInterpreterFactory::new()),
203        );
204        assert!(b.descriptor().validate().is_ok());
205    }
206}