Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 58 additions & 9 deletions crates/buzz-acp/src/acp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -646,10 +646,13 @@ impl AcpClient {
/// - `Some(SystemPromptTransport::ClaudeMeta(text))` — `_meta.systemPrompt`
/// as `{"append": text}`, keeping claude-agent-acp's native preset intact.
///
/// `session_title` rides in `_meta.sessionTitle` when `Some`; `_meta` is
/// omitted entirely otherwise, since adapters may distinguish an absent
/// member from a null one. When both `ClaudeMeta` and `session_title` are
/// present the two `_meta` members are merged into a single object.
/// `session_title` rides in `_meta.sessionTitle` when `Some`.
/// `session_key` rides in `_meta.sessionKey` when `Some` — OpenClaw's ACP
/// mapper uses that as the Gateway store key (else the process `--session`
/// default). `_meta` is omitted entirely when neither is set, since adapters
/// may distinguish an absent member from a null one. When `ClaudeMeta`,
/// `session_title`, and/or `session_key` are present the `_meta` members
/// are merged into a single object.
///
/// Callers use [`extract_model_config_options`] and [`extract_model_state`]
/// to pull model info from the raw result.
Expand All @@ -659,6 +662,7 @@ impl AcpClient {
mcp_servers: Vec<McpServer>,
system_prompt: Option<SystemPromptTransport<'_>>,
session_title: Option<&str>,
session_key: Option<&str>,
) -> Result<SessionNewResponse, AcpError> {
let mut params = serde_json::json!({
"cwd": cwd,
Expand All @@ -669,7 +673,7 @@ impl AcpClient {
params["systemPrompt"] = serde_json::Value::String(sp.to_owned());
}
Some(SystemPromptTransport::ClaudeMeta(sp)) => {
// Merge into _meta so sessionTitle (set below) is not clobbered.
// Merge into _meta so sessionTitle / sessionKey (set below) are not clobbered.
params["_meta"]["systemPrompt"] = serde_json::json!({ "append": sp });
}
None => {}
Expand All @@ -678,6 +682,9 @@ impl AcpClient {
// Merge — _meta may already carry systemPrompt from ClaudeMeta above.
params["_meta"]["sessionTitle"] = serde_json::Value::String(title.to_owned());
}
if let Some(key) = session_key {
params["_meta"]["sessionKey"] = serde_json::Value::String(key.to_owned());
}
let result = self.send_request("session/new", params).await?;
let session_id = result["sessionId"]
.as_str()
Expand All @@ -700,9 +707,10 @@ impl AcpClient {
mcp_servers: Vec<McpServer>,
system_prompt: Option<SystemPromptTransport<'_>>,
session_title: Option<&str>,
session_key: Option<&str>,
) -> Result<String, AcpError> {
Ok(self
.session_new_full(cwd, mcp_servers, system_prompt, session_title)
.session_new_full(cwd, mcp_servers, system_prompt, session_title, session_key)
.await?
.session_id)
}
Expand Down Expand Up @@ -3492,6 +3500,7 @@ mod tests {
vec![],
Some(SystemPromptTransport::Field("Custom system prompt")),
None,
None,
)
.await
.expect("session_new_full should succeed");
Expand Down Expand Up @@ -3577,7 +3586,7 @@ mod tests {
.expect("initialize should succeed");

let resp = client
.session_new_full("/tmp", vec![], None, None)
.session_new_full("/tmp", vec![], None, None, None)
.await
.expect("session_new_full should succeed");

Expand Down Expand Up @@ -3605,7 +3614,7 @@ mod tests {
.expect("initialize should succeed");

let resp = client
.session_new_full("/tmp", vec![], None, Some("Fizz · #buzz-dev"))
.session_new_full("/tmp", vec![], None, Some("Fizz · #buzz-dev"), None)
.await
.expect("session_new_full should succeed");

Expand Down Expand Up @@ -3633,7 +3642,7 @@ mod tests {
.expect("initialize should succeed");

let resp = client
.session_new_full("/tmp", vec![], None, None)
.session_new_full("/tmp", vec![], None, None, None)
.await
.expect("session_new_full should succeed");

Expand Down Expand Up @@ -3669,6 +3678,7 @@ mod tests {
vec![],
Some(SystemPromptTransport::ClaudeMeta("Be concise")),
None,
None,
)
.await
.expect("session_new_full should succeed");
Expand Down Expand Up @@ -3708,6 +3718,7 @@ mod tests {
vec![],
Some(SystemPromptTransport::ClaudeMeta("Be concise")),
Some("Fizz · #buzz-dev"),
None,
)
.await
.expect("session_new_full should succeed");
Expand All @@ -3725,6 +3736,44 @@ mod tests {
);
}

#[tokio::test]
async fn session_new_full_sends_session_key_and_title_in_meta() {
let script = r#"
read -t 2 _init
echo '{"jsonrpc":"2.0","id":0,"result":{"protocolVersion":1,"agentCapabilities":{}}}'
read -t 2 REQ
echo '{"jsonrpc":"2.0","id":1,"result":{"sessionId":"ses_key","_receivedRequest":'"$REQ"'}}'
sleep 1
"#;
let mut client = spawn_script(script).await;
client
.initialize()
.await
.expect("initialize should succeed");

let resp = client
.session_new_full(
"/tmp",
vec![],
None,
Some("Captain · #general"),
Some("agent:captain:buzz:channel:11111111-1111-1111-1111-111111111111"),
)
.await
.expect("session_new_full should succeed");

let received = &resp.raw["_receivedRequest"];
assert_eq!(
received["params"]["_meta"]["sessionTitle"].as_str(),
Some("Captain · #general"),
);
assert_eq!(
received["params"]["_meta"]["sessionKey"].as_str(),
Some("agent:captain:buzz:channel:11111111-1111-1111-1111-111111111111"),
"OpenClaw Gateway key must ride in _meta.sessionKey"
);
}

// ── Goose-native steer scaffold (PR follow-up to #1160) ──────────────

/// Helper: spawn an inert `cat` subprocess so we have a real AcpClient
Expand Down
95 changes: 55 additions & 40 deletions crates/buzz-acp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ mod engram_fetch;
mod filter;
mod last_mile;
mod observer;
mod openclaw_session;
mod pool;
mod pool_lifecycle;
mod queue;
Expand Down Expand Up @@ -2238,6 +2239,7 @@ async fn tokio_main() -> Result<()> {
memory_enabled: config.memory_enabled,
harness_name: crate::config::normalize_agent_command_identity(&config.agent_command),
relay_url: config.relay_url.clone(),
openclaw_agent_id: crate::openclaw_session::parse_openclaw_agent_id(&config.agent_args),
});

if !config.memory_enabled {
Expand Down Expand Up @@ -4047,15 +4049,20 @@ fn handle_prompt_result(
// The task may have invalidated this session before returning. Never
// resurrect delivery state for a dead session; its replacement must
// receive fresh standing context and history.
if let Some(live_session_id) = result.agent.state.sessions.get(channel_id).cloned() {
let conversation = result
.batch
.as_ref()
.map(crate::openclaw_session::ConversationKey::from_batch)
.unwrap_or_else(|| crate::openclaw_session::ConversationKey::channel(*channel_id));
if let Some(live_session_id) = result.agent.state.sessions.get(&conversation).cloned() {
let event_ids = successful_steer_deliveries
.into_iter()
.filter(|delivery| delivery.session_id == live_session_id)
.map(|delivery| delivery.event_id);
result
.agent
.state
.mark_channel_delivery_success(*channel_id, false, event_ids);
.mark_channel_delivery_success(conversation, false, event_ids);
}
}

Expand Down Expand Up @@ -5088,7 +5095,9 @@ async fn run_models(args: ModelsArgs) -> Result<()> {
// so shutdown() runs on all paths (success, error, timeout).
let protocol_result = tokio::time::timeout(MODELS_TIMEOUT, async {
let init = client.initialize().await?;
let session = client.session_new_full(&cwd, vec![], None, None).await?;
let session = client
.session_new_full(&cwd, vec![], None, None, None)
.await?;
Ok::<_, acp::AcpError>((init, session))
})
.await;
Expand Down Expand Up @@ -7524,14 +7533,14 @@ mod error_outcome_emission_tests {
let channel_id = Uuid::new_v4();
let steer_event_id = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
let mut agent = dummy_agent(0).await;
agent
.state
.sessions
.insert(channel_id, "live-session".into());
agent
.state
.deliveries
.insert(channel_id, Default::default());
agent.state.sessions.insert(
crate::openclaw_session::ConversationKey::channel(channel_id),
"live-session".into(),
);
agent.state.deliveries.insert(
crate::openclaw_session::ConversationKey::channel(channel_id),
Default::default(),
);

let mut pool = AgentPool::from_slots(vec![None]);
let task_id = pool.join_set.spawn(async {}).id();
Expand Down Expand Up @@ -7587,7 +7596,8 @@ mod error_outcome_emission_tests {
);

let returned = pool.agents_mut()[0].as_ref().expect("returned agent");
assert!(returned.state.deliveries[&channel_id]
assert!(returned.state.deliveries
[&crate::openclaw_session::ConversationKey::channel(channel_id)]
.delivered_event_ids
.contains(steer_event_id));
}
Expand All @@ -7596,14 +7606,14 @@ mod error_outcome_emission_tests {
async fn in_flight_stale_native_steer_ack_cannot_update_replacement_session() {
let channel_id = Uuid::new_v4();
let mut agent = dummy_agent(0).await;
agent
.state
.sessions
.insert(channel_id, "replacement-session".into());
agent
.state
.deliveries
.insert(channel_id, Default::default());
agent.state.sessions.insert(
crate::openclaw_session::ConversationKey::channel(channel_id),
"replacement-session".into(),
);
agent.state.deliveries.insert(
crate::openclaw_session::ConversationKey::channel(channel_id),
Default::default(),
);

let mut pool = AgentPool::from_slots(vec![None]);
let task_id = pool.join_set.spawn(async {}).id();
Expand Down Expand Up @@ -7659,7 +7669,8 @@ mod error_outcome_emission_tests {
);

let returned = pool.agents_mut()[0].as_ref().expect("returned agent");
assert!(returned.state.deliveries[&channel_id]
assert!(returned.state.deliveries
[&crate::openclaw_session::ConversationKey::channel(channel_id)]
.delivered_event_ids
.is_empty());
}
Expand All @@ -7669,14 +7680,14 @@ mod error_outcome_emission_tests {
let channel_id = Uuid::new_v4();
let steer_event_id = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb";
let mut agent = dummy_agent(0).await;
agent
.state
.sessions
.insert(channel_id, "live-session".into());
agent
.state
.deliveries
.insert(channel_id, Default::default());
agent.state.sessions.insert(
crate::openclaw_session::ConversationKey::channel(channel_id),
"live-session".into(),
);
agent.state.deliveries.insert(
crate::openclaw_session::ConversationKey::channel(channel_id),
Default::default(),
);
let mut pool = AgentPool::from_slots(vec![Some(agent)]);

assert!(pool.record_successful_steer(
Expand All @@ -7685,7 +7696,8 @@ mod error_outcome_emission_tests {
"live-session".into(),
));
let returned = pool.agents_mut()[0].as_ref().expect("idle returned agent");
assert!(returned.state.deliveries[&channel_id]
assert!(returned.state.deliveries
[&crate::openclaw_session::ConversationKey::channel(channel_id)]
.delivered_event_ids
.contains(steer_event_id));
}
Expand All @@ -7694,14 +7706,14 @@ mod error_outcome_emission_tests {
async fn late_native_steer_ack_cannot_update_replacement_session() {
let channel_id = Uuid::new_v4();
let mut agent = dummy_agent(0).await;
agent
.state
.sessions
.insert(channel_id, "replacement-session".into());
agent
.state
.deliveries
.insert(channel_id, Default::default());
agent.state.sessions.insert(
crate::openclaw_session::ConversationKey::channel(channel_id),
"replacement-session".into(),
);
agent.state.deliveries.insert(
crate::openclaw_session::ConversationKey::channel(channel_id),
Default::default(),
);
let mut pool = AgentPool::from_slots(vec![Some(agent)]);

assert!(!pool.record_successful_steer(
Expand All @@ -7710,7 +7722,8 @@ mod error_outcome_emission_tests {
"old-session".into(),
));
let returned = pool.agents_mut()[0].as_ref().expect("replacement agent");
assert!(returned.state.deliveries[&channel_id]
assert!(returned.state.deliveries
[&crate::openclaw_session::ConversationKey::channel(channel_id)]
.delivered_event_ids
.is_empty());
}
Expand Down Expand Up @@ -7773,7 +7786,9 @@ mod error_outcome_emission_tests {
);

let returned = pool.agents_mut()[0].as_ref().expect("returned agent");
assert!(!returned.state.deliveries.contains_key(&channel_id));
assert!(!returned.state.deliveries.contains_key(
&crate::openclaw_session::ConversationKey::channel(channel_id)
));
}

/// Drive one error outcome through `handle_prompt_result` and return how
Expand Down
Loading