historyを返すプロトコル
This commit is contained in:
@@ -231,6 +231,10 @@ impl PodController {
|
||||
message: "Pod is not running".into(),
|
||||
});
|
||||
}
|
||||
|
||||
// GetHistory is handled at the socket layer (direct response).
|
||||
// If it somehow reaches the controller, ignore it.
|
||||
Method::GetHistory => {}
|
||||
}
|
||||
}
|
||||
});
|
||||
@@ -287,6 +291,9 @@ where
|
||||
message: "Pod is already executing a turn".into(),
|
||||
});
|
||||
}
|
||||
Some(Method::GetHistory) => {
|
||||
// Handled at socket layer; ignore here.
|
||||
}
|
||||
None => {
|
||||
let _ = cancel_tx.try_send(());
|
||||
shared_state.set_status(PodStatus::Idle);
|
||||
|
||||
@@ -49,6 +49,10 @@ impl PodSharedState {
|
||||
self.status.read().map(|s| *s).unwrap_or(PodStatus::Idle)
|
||||
}
|
||||
|
||||
pub fn history(&self) -> Vec<Item> {
|
||||
self.history.read().map(|h| h.clone()).unwrap_or_default()
|
||||
}
|
||||
|
||||
pub fn update_history(&self, items: Vec<Item>) {
|
||||
if let Ok(mut h) = self.history.write() {
|
||||
*h = items;
|
||||
|
||||
@@ -66,31 +66,44 @@ async fn handle_connection(stream: tokio::net::UnixStream, handle: PodHandle) {
|
||||
let mut writer = JsonLineWriter::new(writer);
|
||||
let mut rx = handle.subscribe();
|
||||
|
||||
// Event writer: broadcast events → socket
|
||||
let write_task = tokio::spawn(async move {
|
||||
while let Ok(event) = rx.recv().await {
|
||||
if writer.write(&event).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Method reader: socket → controller
|
||||
loop {
|
||||
match reader.next::<Method>().await {
|
||||
Ok(Some(method)) => {
|
||||
let _ = handle.send(method).await;
|
||||
tokio::select! {
|
||||
// Broadcast events → this client
|
||||
event = rx.recv() => {
|
||||
match event {
|
||||
Ok(event) => {
|
||||
if writer.write(&event).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
Ok(None) => break,
|
||||
Err(e) => {
|
||||
let _ = handle.send_event(Event::Error {
|
||||
code: protocol::ErrorCode::Internal,
|
||||
message: format!("invalid method: {e}"),
|
||||
});
|
||||
// Client methods → handle or forward to controller
|
||||
method = reader.next::<Method>() => {
|
||||
match method {
|
||||
Ok(Some(Method::GetHistory)) => {
|
||||
let items = handle.shared_state.history();
|
||||
let values = items
|
||||
.iter()
|
||||
.map(|item| serde_json::to_value(item).expect("Item is Serialize"))
|
||||
.collect();
|
||||
if writer.write(&Event::History { items: values }).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Ok(Some(method)) => {
|
||||
let _ = handle.send(method).await;
|
||||
}
|
||||
Ok(None) => break,
|
||||
Err(e) => {
|
||||
let _ = handle.send_event(Event::Error {
|
||||
code: protocol::ErrorCode::Internal,
|
||||
message: format!("invalid method: {e}"),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Client disconnected — stop the write task
|
||||
write_task.abort();
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ pub enum Method {
|
||||
Run { input: String },
|
||||
Resume,
|
||||
Cancel,
|
||||
GetHistory,
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -63,6 +64,9 @@ pub enum Event {
|
||||
code: ErrorCode,
|
||||
message: String,
|
||||
},
|
||||
History {
|
||||
items: Vec<serde_json::Value>,
|
||||
},
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -138,6 +142,25 @@ mod tests {
|
||||
assert_eq!(parsed["data"]["result"], "limit_reached");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn method_get_history() {
|
||||
let json = r#"{"method":"get_history"}"#;
|
||||
let method: Method = serde_json::from_str(json).unwrap();
|
||||
assert!(matches!(method, Method::GetHistory));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn event_history_format() {
|
||||
let event = Event::History {
|
||||
items: vec![serde_json::json!({"type": "message", "role": "user"})],
|
||||
};
|
||||
let json = serde_json::to_string(&event).unwrap();
|
||||
let parsed: serde_json::Value = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(parsed["event"], "history");
|
||||
assert!(parsed["data"]["items"].is_array());
|
||||
assert_eq!(parsed["data"]["items"][0]["role"], "user");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn event_error_format() {
|
||||
let event = Event::Error {
|
||||
|
||||
@@ -135,6 +135,7 @@ impl App {
|
||||
self.push_status(format!("[run end] {result:?}"));
|
||||
}
|
||||
Event::ToolCallArgsDelta { .. } => {}
|
||||
Event::History { .. } => {}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user