diff --git a/.yoi/tickets/00001KY1ZKCPC/item.md b/.yoi/tickets/00001KY1ZKCPC/item.md index 5f92f1cb..82c9096e 100644 --- a/.yoi/tickets/00001KY1ZKCPC/item.md +++ b/.yoi/tickets/00001KY1ZKCPC/item.md @@ -2,7 +2,7 @@ title: 'Backend runtime経由の操作をprotocol transportへ統一しTUI同等にする' state: 'planning' created_at: '2026-07-21T09:20:07Z' -updated_at: '2026-07-22T02:11:40Z' +updated_at: '2026-07-22T02:54:28Z' assignee: null readiness: 'draft' --- diff --git a/.yoi/tickets/00001KY1ZKCPC/thread.md b/.yoi/tickets/00001KY1ZKCPC/thread.md index 8193fa08..5b134e97 100644 --- a/.yoi/tickets/00001KY1ZKCPC/thread.md +++ b/.yoi/tickets/00001KY1ZKCPC/thread.md @@ -198,3 +198,40 @@ Verification: - `git diff --check` --- + + + +## Implementation report + +Correction after review: + +The previous fix synthesized `Event::UserMessage` in `Runtime::send_protocol_method(...)`. That was the wrong layer. `UserMessage` is part of the normal Worker protocol/log output path and should not be constructed by the runtime transport layer. + +Root cause refined: +- The embedded Worker backend bridge only subscribed to `WorkerHandle::event_tx`. +- Local socket clients also subscribe to `SegmentLogSink` live entries; `LogEntry::UserInput` is mapped to `Event::UserMessage` there. +- Runtime observation/protocol streaming missed that `SegmentLogSink` live-entry path, so user input echo events were absent when using runtime protocol transport. + +Fix: +- Removed the runtime-side synthetic protocol-method echo from `Runtime::send_protocol_method(...)`. +- Updated `WorkerRuntimeExecutionBackend::connect_handle(...)` to bridge both: + - direct Worker protocol events from `WorkerHandle::subscribe()` + - live session log entries from `handle.sink.subscribe_with_snapshot()` +- Added runtime bridge mapping matching local IPC behavior: + - `LogEntry::SegmentStart` -> `Event::SegmentRotated` + - `LogEntry::UserInput` -> `Event::UserMessage` + - `LogEntry::SystemItem` -> `Event::SystemItem` + - `LogEntry::Invoke` -> `Event::InvokeStart` +- Removed the fake protocol-WS echo test that depended on synthetic runtime behavior. +- Added a unit test that locks the normal bridge mapping: live `LogEntry::UserInput` maps to `Event::UserMessage`. + +Verification: +- `nix develop -c cargo fmt -- --check` +- `nix develop -c cargo check -p worker-runtime -p yoi-workspace-server -p client -p tui` +- `nix develop -c cargo test -p worker-runtime --features ws-server --lib runtime_bridge_maps_live_user_input_log_entry_to_user_message` +- `nix develop -c cargo test -p worker-runtime --features ws-server --lib protocol_ws` +- `nix develop -c cargo test -p yoi-workspace-server protocol_ws --lib` +- `nix develop -c cargo test -p client backend_runtime --lib` +- `git diff --check` + +--- diff --git a/crates/worker-runtime/src/http_server.rs b/crates/worker-runtime/src/http_server.rs index 6bbb72d3..252d94a0 100644 --- a/crates/worker-runtime/src/http_server.rs +++ b/crates/worker-runtime/src/http_server.rs @@ -1339,34 +1339,6 @@ mod ws_tests { )); } - #[tokio::test] - async fn protocol_ws_echoes_accepted_run_as_user_message() { - let (_runtime, _worker_ref, url) = spawn_runtime_server().await; - let (mut stream, _) = connect_async(&url).await.unwrap(); - let _ = next_frame(&mut stream).await; - - stream - .send(Message::Text( - serde_json::to_string(&protocol::Method::Run { - input: vec![protocol::Segment::text("hello from protocol")], - }) - .unwrap() - .into(), - )) - .await - .unwrap(); - - match next_frame(&mut stream).await { - protocol::Event::UserMessage { segments } => { - assert_eq!( - protocol::Segment::flatten_to_text(&segments), - "hello from protocol" - ); - } - event => panic!("expected user message echo, got {event:?}"), - } - } - #[tokio::test] async fn protocol_ws_reports_malformed_cursor_and_method_frame() { let (_runtime, _worker_ref, url) = spawn_runtime_server().await; diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index cf813f2d..be362126 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -576,8 +576,6 @@ impl Runtime { return Ok(vec![Event::Completions { kind, entries }]); } - let observation_payload = protocol_method_observation_event(&method); - let (backend, handle) = { let mut state = self.lock()?; state.ensure_running()?; @@ -624,9 +622,6 @@ impl Runtime { } self.record_execution_result(worker_ref, dispatch_result)?; - if let Some(payload) = observation_payload { - self.record_protocol_method_observation(worker_ref, payload)?; - } Ok(Vec::new()) } @@ -963,7 +958,12 @@ impl Runtime { worker_ref: &WorkerRef, input: WorkerInput, ) -> Result<(), RuntimeError> { - self.record_protocol_method_observation(worker_ref, input_protocol_event(&input)) + let mut state = self.lock()?; + state.ensure_worker_ref(worker_ref)?; + let event = + state.push_worker_observation_event(worker_ref.clone(), input_protocol_event(&input)); + state.persist_worker_observation_event(&event)?; + Ok(()) } #[cfg(not(feature = "ws-server"))] @@ -975,28 +975,6 @@ impl Runtime { Ok(()) } - #[cfg(feature = "ws-server")] - fn record_protocol_method_observation( - &self, - worker_ref: &WorkerRef, - payload: Event, - ) -> Result<(), RuntimeError> { - let mut state = self.lock()?; - state.ensure_worker_ref(worker_ref)?; - let event = state.push_worker_observation_event(worker_ref.clone(), payload); - state.persist_worker_observation_event(&event)?; - Ok(()) - } - - #[cfg(not(feature = "ws-server"))] - fn record_protocol_method_observation( - &self, - _worker_ref: &WorkerRef, - _payload: Event, - ) -> Result<(), RuntimeError> { - Ok(()) - } - fn transition_worker( &self, worker_ref: &WorkerRef, @@ -1858,48 +1836,6 @@ fn validate_worker_input(input: &WorkerInput) -> Result<(), RuntimeError> { Ok(()) } -#[cfg(feature = "ws-server")] -fn protocol_method_observation_event(method: &Method) -> Option { - match method { - Method::Run { input } => Some(Event::UserMessage { - segments: input.clone(), - }), - Method::Notify { message, .. } => Some(Event::SystemItem { - item: serde_json::json!({ - "kind": "embedded_worker_system_input", - "content": message, - }), - }), - Method::RegisterPeer { name } => Some(Event::SystemItem { - item: serde_json::json!({ - "kind": "embedded_worker_command_input", - "command": "register_peer", - "content": name, - }), - }), - Method::Compact => Some(Event::SystemItem { - item: serde_json::json!({ - "kind": "embedded_worker_command_input", - "command": "compact", - "content": "", - }), - }), - Method::ListRewindTargets => Some(Event::SystemItem { - item: serde_json::json!({ - "kind": "embedded_worker_command_input", - "command": "list_rewind_targets", - "content": "", - }), - }), - _ => None, - } -} - -#[cfg(not(feature = "ws-server"))] -fn protocol_method_observation_event(_method: &Method) -> Option { - None -} - #[cfg(feature = "ws-server")] fn input_protocol_event(input: &WorkerInput) -> protocol::Event { match input.kind { @@ -1941,7 +1877,6 @@ mod tests { WorkerExecutionRestoreRequest, WorkerExecutionRunState, }; use std::collections::BTreeMap; - use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; fn task_request(_objective: &str) -> CreateWorkerRequest { diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index 35f12934..89fa049d 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -30,9 +30,9 @@ use crate::working_directory::{ }; use async_trait::async_trait; use manifest::paths; -use protocol::{Method, Segment, WorkerStatus}; +use protocol::{Event, Method, Segment, WorkerStatus}; use session_store::FsStore; -use session_store::{CombinedStore, FsWorkerStore}; +use session_store::{CombinedStore, FsWorkerStore, LogEntry}; use tokio::runtime::Runtime; #[cfg(feature = "ws-server")] use tokio::sync::broadcast; @@ -667,19 +667,35 @@ where #[cfg(feature = "ws-server")] { let mut events = handle.subscribe(); + let (_entries, mut entry_events) = handle.sink.subscribe_with_snapshot(); let bridge_handle = handle.clone(); let bridge_busy = busy.clone(); if let Err(message) = self.spawn_on_adapter_runtime(async move { loop { - match events.recv().await { - Ok(event) => { - let _ = bridge_context.publish_protocol_event(event); - if bridge_handle.shared_state.get_status() == WorkerStatus::Idle { - bridge_busy.store(false, Ordering::SeqCst); + tokio::select! { + event = events.recv() => { + match event { + Ok(event) => { + let _ = bridge_context.publish_protocol_event(event); + if bridge_handle.shared_state.get_status() == WorkerStatus::Idle { + bridge_busy.store(false, Ordering::SeqCst); + } + } + Err(broadcast::error::RecvError::Lagged(_)) => continue, + Err(broadcast::error::RecvError::Closed) => break, + } + } + entry = entry_events.recv() => { + match entry { + Ok(entry) => { + if let Some(event) = live_log_entry_event(entry) { + let _ = bridge_context.publish_protocol_event(event); + } + } + Err(broadcast::error::RecvError::Lagged(_)) => continue, + Err(broadcast::error::RecvError::Closed) => break, } } - Err(broadcast::error::RecvError::Lagged(_)) => continue, - Err(broadcast::error::RecvError::Closed) => break, } } }) { @@ -722,6 +738,22 @@ impl Drop for WorkerRuntimeExecutionBackend { } } +fn live_log_entry_event(entry: LogEntry) -> Option { + match entry { + LogEntry::SegmentStart { .. } => { + let value = serde_json::to_value(&entry).expect("LogEntry is Serialize"); + Some(Event::SegmentRotated { entry: value }) + } + LogEntry::UserInput { segments, .. } => Some(Event::UserMessage { segments }), + LogEntry::SystemItem { item, .. } => { + let value = serde_json::to_value(&item).expect("SystemItem is Serialize"); + Some(Event::SystemItem { item: value }) + } + LogEntry::Invoke { trigger, .. } => Some(Event::InvokeStart { kind: trigger }), + _ => None, + } +} + fn method_starts_turn(method: &Method) -> bool { matches!( method, @@ -1208,6 +1240,21 @@ mod tests { use llm_engine::llm_client::{ClientError, LlmClient, Request}; use manifest::{Scope, WorkerManifest}; + #[test] + fn runtime_bridge_maps_live_user_input_log_entry_to_user_message() { + let segments = vec![Segment::text("hello through normal bridge")]; + let event = live_log_entry_event(LogEntry::UserInput { + ts: session_store::segment_log::now_millis(), + segments: segments.clone(), + }) + .expect("UserInput must be live-relevant"); + + match event { + Event::UserMessage { segments: echoed } => assert_eq!(echoed, segments), + other => panic!("expected UserMessage, got {other:?}"), + } + } + #[derive(Clone)] struct MockClient { responses: Arc>>,