use monoloop_contracts::{
Bytes, CanonicalMessage, DialectDescriptor, EncodedExchange, EncodingError,
ExchangeInputPolicy, ExtensionKey, InitialEncodeRequest, OutboundDialectEncoder,
ToolContinuationEncodeRequest, VersionedExtension,
};
use serde_json::{json, Map, Value};
use std::collections::BTreeMap;
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub enum AcpPromptWireShape {
#[default]
ParamsObject,
MethodEnvelope,
PlainText,
}
#[derive(Clone, Debug)]
pub struct AcpPromptEncoder {
pub shape: AcpPromptWireShape,
pub dialect: DialectDescriptor,
pub max_encoded_bytes: usize,
}
impl Default for AcpPromptEncoder {
fn default() -> Self {
Self {
shape: AcpPromptWireShape::ParamsObject,
dialect: DialectDescriptor::acp_json_rpc("1"),
max_encoded_bytes: 1024 * 1024,
}
}
}
impl AcpPromptEncoder {
pub fn grok() -> Self {
Self {
shape: AcpPromptWireShape::ParamsObject,
dialect: DialectDescriptor::acp_json_rpc("1"),
max_encoded_bytes: 1024 * 1024,
}
}
pub fn cursor() -> Self {
Self {
shape: AcpPromptWireShape::PlainText,
dialect: DialectDescriptor::cursor_acp("1"),
max_encoded_bytes: 1024 * 1024,
}
}
pub fn codex() -> Self {
Self {
shape: AcpPromptWireShape::PlainText,
dialect: DialectDescriptor::codex_acp("1"),
max_encoded_bytes: 1024 * 1024,
}
}
pub fn agy() -> Self {
Self {
shape: AcpPromptWireShape::PlainText,
dialect: DialectDescriptor::agy_acp("1"),
max_encoded_bytes: 1024 * 1024,
}
}
fn collect_user_text(messages: &[CanonicalMessage]) -> Result<String, EncodingError> {
let mut text = String::new();
for msg in messages {
match msg {
CanonicalMessage::User { content, .. }
| CanonicalMessage::System { content, .. } => {
for part in content {
if !text.is_empty() {
text.push('\n');
}
text.push_str(part.text());
}
}
CanonicalMessage::Assistant { content, .. } => {
for part in content {
if !text.is_empty() {
text.push('\n');
}
text.push_str(part.text());
}
}
CanonicalMessage::Tool { content, .. } => {
for part in content {
if !text.is_empty() {
text.push('\n');
}
text.push_str(part.text());
}
}
}
}
if text.trim().is_empty() {
return Err(EncodingError::UnrepresentableInput);
}
Ok(text)
}
}
impl OutboundDialectEncoder for AcpPromptEncoder {
fn encode_initial(
&self,
request: InitialEncodeRequest<'_>,
) -> Result<EncodedExchange, EncodingError> {
let text = Self::collect_user_text(request.input.messages())?;
if !request.tools.is_empty() {
return Err(EncodingError::Unsupported(
"non-empty tools require MCP gateway for external-agent ACP profiles",
));
}
let meta = encode_acp_extension_meta(&request.config.extensions)?;
let bytes = match self.shape {
AcpPromptWireShape::PlainText => {
if meta.is_some() {
return Err(EncodingError::Unsupported(
"plain-text ACP shape cannot encode extensions",
));
}
Bytes::from(text.into_bytes())
}
AcpPromptWireShape::ParamsObject => {
let mut params = json!({
"prompt": [{ "type": "text", "text": text }]
});
if let Some(m) = meta {
params
.as_object_mut()
.ok_or(EncodingError::UnrepresentableInput)?
.insert("_meta".into(), m);
}
Bytes::from(
serde_json::to_vec(¶ms).map_err(|_| EncodingError::UnrepresentableInput)?,
)
}
AcpPromptWireShape::MethodEnvelope => {
let mut params = json!({
"prompt": [{ "type": "text", "text": text }]
});
if let Some(m) = meta {
params
.as_object_mut()
.ok_or(EncodingError::UnrepresentableInput)?
.insert("_meta".into(), m);
}
let v = json!({
"method": "session/prompt",
"params": params
});
Bytes::from(
serde_json::to_vec(&v).map_err(|_| EncodingError::UnrepresentableInput)?,
)
}
};
if bytes.len() > self.max_encoded_bytes {
return Err(EncodingError::LimitExceeded);
}
Ok(EncodedExchange {
bytes,
required_input_dialect: self.dialect.clone(),
input_policy: ExchangeInputPolicy::SendAndFinish,
})
}
fn encode_tool_continuation(
&self,
_request: ToolContinuationEncodeRequest<'_>,
) -> Result<EncodedExchange, EncodingError> {
Err(EncodingError::Unsupported(
"ACP external agents do not encode model tool continuations",
))
}
}
fn encode_acp_extension_meta(
extensions: &BTreeMap<ExtensionKey, VersionedExtension>,
) -> Result<Option<Value>, EncodingError> {
if extensions.is_empty() {
return Ok(None);
}
let mut meta = Map::new();
for (key, ext) in extensions {
let field = key
.as_str()
.strip_prefix("acp.meta.")
.or_else(|| key.as_str().strip_prefix("x.ai."))
.ok_or(EncodingError::Unsupported("non-acp extension"))?;
if meta.contains_key(field) {
return Err(EncodingError::Unsupported("duplicate extension meta key"));
}
meta.insert(field.to_string(), ext.value.clone());
}
Ok(Some(Value::Object(meta)))
}
#[derive(Clone, Debug)]
pub struct HeadlessPromptEncoder {
pub dialect: DialectDescriptor,
pub max_encoded_bytes: usize,
}
impl HeadlessPromptEncoder {
pub fn zai() -> Self {
Self {
dialect: DialectDescriptor::zai_cli("1"),
max_encoded_bytes: 1024 * 1024,
}
}
pub fn claude() -> Self {
Self {
dialect: DialectDescriptor::claude_code("1"),
max_encoded_bytes: 1024 * 1024,
}
}
}
impl OutboundDialectEncoder for HeadlessPromptEncoder {
fn encode_initial(
&self,
request: InitialEncodeRequest<'_>,
) -> Result<EncodedExchange, EncodingError> {
if !request.tools.is_empty() {
return Err(EncodingError::Unsupported(
"headless CLI profiles reject Monoloop-linked tools (MCP None)",
));
}
if !request.config.extensions.is_empty() {
return Err(EncodingError::Unsupported(
"headless CLI encoder cannot encode extensions",
));
}
let text = AcpPromptEncoder::collect_user_text(request.input.messages())?;
let bytes = Bytes::from(text.into_bytes());
if bytes.len() > self.max_encoded_bytes {
return Err(EncodingError::LimitExceeded);
}
Ok(EncodedExchange {
bytes,
required_input_dialect: self.dialect.clone(),
input_policy: ExchangeInputPolicy::SendAndFinish,
})
}
fn encode_tool_continuation(
&self,
_request: ToolContinuationEncodeRequest<'_>,
) -> Result<EncodedExchange, EncodingError> {
Err(EncodingError::Unsupported(
"headless CLI has no tool continuation encoding",
))
}
}
#[cfg(test)]
mod tests {
use super::*;
use monoloop_contracts::{user_text_input, EffectiveConfig, ExchangeId, TransactionId};
fn bare_cfg() -> EffectiveConfig {
EffectiveConfig {
model: None,
temperature: None,
reasoning_effort: None,
max_output_tokens: None,
stop: vec![],
response_format: None,
continuation_policy: Default::default(),
deadline: None,
extensions: Default::default(),
session: Default::default(),
}
}
#[test]
fn grok_params_object_has_no_prompt_on_argv_shape() {
let enc = AcpPromptEncoder::grok();
let input = user_text_input("hello agent").unwrap();
let tid = TransactionId::generate();
let eid = ExchangeId::generate();
let encoded = enc
.encode_initial(InitialEncodeRequest {
transaction_id: &tid,
exchange_id: &eid,
input: &input,
config: &bare_cfg(),
tools: &[],
})
.unwrap();
let v: serde_json::Value = serde_json::from_slice(&encoded.bytes).unwrap();
assert!(v.get("prompt").is_some());
assert!(v.get("method").is_none());
}
#[test]
fn plain_text_cursor_shape() {
let enc = AcpPromptEncoder::cursor();
let input = user_text_input("hi").unwrap();
let tid = TransactionId::generate();
let eid = ExchangeId::generate();
let encoded = enc
.encode_initial(InitialEncodeRequest {
transaction_id: &tid,
exchange_id: &eid,
input: &input,
config: &bare_cfg(),
tools: &[],
})
.unwrap();
assert_eq!(&encoded.bytes[..], b"hi");
}
}