プロトコル経由のshutdow経路
This commit is contained in:
@@ -43,7 +43,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
let pod = pod::Pod::from_manifest_toml(&toml, store).await?;
|
||||
|
||||
let runtime_tmp = tempfile::tempdir()?;
|
||||
let handle = PodController::spawn(pod, runtime_tmp.path()).await?;
|
||||
let (handle, _shutdown_rx) = PodController::spawn(pod, runtime_tmp.path()).await?;
|
||||
|
||||
// Check initial status via shared state
|
||||
println!("[shared_state] {}", handle.shared_state.status_json());
|
||||
|
||||
@@ -4,7 +4,7 @@ use std::sync::Arc;
|
||||
use llm_worker::WorkerError;
|
||||
use llm_worker::llm_client::client::LlmClient;
|
||||
use session_store::Store;
|
||||
use tokio::sync::{broadcast, mpsc};
|
||||
use tokio::sync::{broadcast, mpsc, oneshot};
|
||||
|
||||
use crate::notifier::Notifier;
|
||||
use crate::pod::{Pod, PodError, PodRunResult};
|
||||
@@ -50,17 +50,20 @@ impl PodHandle {
|
||||
// PodController — actor that owns a Pod
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
pub type ShutdownReceiver = oneshot::Receiver<()>;
|
||||
|
||||
pub struct PodController;
|
||||
|
||||
impl PodController {
|
||||
pub async fn spawn<C, St>(
|
||||
mut pod: Pod<C, St>,
|
||||
runtime_base: &Path,
|
||||
) -> Result<PodHandle, std::io::Error>
|
||||
) -> Result<(PodHandle, ShutdownReceiver), std::io::Error>
|
||||
where
|
||||
C: LlmClient + 'static,
|
||||
St: Store + 'static,
|
||||
{
|
||||
let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
|
||||
let (method_tx, mut method_rx) = mpsc::channel::<Method>(32);
|
||||
let (event_tx, _) = broadcast::channel::<Event>(256);
|
||||
let notifier = Notifier::new(event_tx.clone());
|
||||
@@ -225,7 +228,7 @@ impl PodController {
|
||||
shared_state.set_status(PodStatus::Running);
|
||||
let _ = runtime_dir.write_status(&shared_state).await;
|
||||
|
||||
let new_status = run_with_cancel_support(
|
||||
let (new_status, shutdown) = run_with_cancel_support(
|
||||
pod.run(&input),
|
||||
&mut method_rx,
|
||||
&event_tx,
|
||||
@@ -234,7 +237,6 @@ impl PodController {
|
||||
)
|
||||
.await;
|
||||
|
||||
// Proactive post-run compaction (best-effort).
|
||||
if new_status == PodStatus::Idle {
|
||||
if let Err(e) = pod.try_post_run_compact().await {
|
||||
tracing::warn!(error = %e, "Post-run compaction error");
|
||||
@@ -251,6 +253,11 @@ impl PodController {
|
||||
shared_state.set_status(new_status);
|
||||
let _ = runtime_dir.write_status(&shared_state).await;
|
||||
let _ = runtime_dir.write_history(&shared_state).await;
|
||||
|
||||
if shutdown {
|
||||
let _ = event_tx.send(Event::Shutdown);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
Method::Resume => {
|
||||
@@ -264,7 +271,7 @@ impl PodController {
|
||||
shared_state.set_status(PodStatus::Running);
|
||||
let _ = runtime_dir.write_status(&shared_state).await;
|
||||
|
||||
let new_status = run_with_cancel_support(
|
||||
let (new_status, shutdown) = run_with_cancel_support(
|
||||
pod.resume(),
|
||||
&mut method_rx,
|
||||
&event_tx,
|
||||
@@ -273,7 +280,6 @@ impl PodController {
|
||||
)
|
||||
.await;
|
||||
|
||||
// Proactive post-run compaction (best-effort).
|
||||
if new_status == PodStatus::Idle {
|
||||
if let Err(e) = pod.try_post_run_compact().await {
|
||||
tracing::warn!(error = %e, "Post-run compaction error");
|
||||
@@ -290,6 +296,11 @@ impl PodController {
|
||||
shared_state.set_status(new_status);
|
||||
let _ = runtime_dir.write_status(&shared_state).await;
|
||||
let _ = runtime_dir.write_history(&shared_state).await;
|
||||
|
||||
if shutdown {
|
||||
let _ = event_tx.send(Event::Shutdown);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
Method::Cancel => {
|
||||
@@ -299,30 +310,39 @@ impl PodController {
|
||||
});
|
||||
}
|
||||
|
||||
Method::Shutdown => {
|
||||
let _ = event_tx.send(Event::Shutdown);
|
||||
break;
|
||||
}
|
||||
|
||||
// GetHistory is handled at the socket layer (direct response).
|
||||
// If it somehow reaches the controller, ignore it.
|
||||
Method::GetHistory => {}
|
||||
}
|
||||
}
|
||||
|
||||
let _ = shutdown_tx.send(());
|
||||
});
|
||||
|
||||
Ok(handle)
|
||||
Ok((handle, shutdown_rx))
|
||||
}
|
||||
}
|
||||
|
||||
/// Runs a Pod future while concurrently processing incoming methods.
|
||||
/// Only `Cancel` is handled during execution; `Run` and `Resume` get errors.
|
||||
///
|
||||
/// Returns `(final_status, shutdown_requested)`.
|
||||
async fn run_with_cancel_support<F>(
|
||||
pod_future: F,
|
||||
method_rx: &mut mpsc::Receiver<Method>,
|
||||
event_tx: &broadcast::Sender<Event>,
|
||||
cancel_tx: &mpsc::Sender<()>,
|
||||
shared_state: &Arc<PodSharedState>,
|
||||
) -> PodStatus
|
||||
) -> (PodStatus, bool)
|
||||
where
|
||||
F: std::future::Future<Output = Result<PodRunResult, PodError>>,
|
||||
{
|
||||
tokio::pin!(pod_future);
|
||||
let mut shutdown_requested = false;
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
@@ -335,7 +355,7 @@ where
|
||||
PodRunResult::LimitReached => (PodStatus::Idle, RunResult::LimitReached),
|
||||
};
|
||||
let _ = event_tx.send(Event::RunEnd { result: run_result });
|
||||
status
|
||||
(status, shutdown_requested)
|
||||
}
|
||||
Err(e) => {
|
||||
let code = worker_error_code(&e);
|
||||
@@ -343,7 +363,7 @@ where
|
||||
code,
|
||||
message: e.to_string(),
|
||||
});
|
||||
PodStatus::Idle
|
||||
(PodStatus::Idle, shutdown_requested)
|
||||
}
|
||||
};
|
||||
}
|
||||
@@ -352,19 +372,21 @@ where
|
||||
Some(Method::Cancel) => {
|
||||
let _ = cancel_tx.try_send(());
|
||||
}
|
||||
Some(Method::Shutdown) => {
|
||||
shutdown_requested = true;
|
||||
let _ = cancel_tx.try_send(());
|
||||
}
|
||||
Some(Method::Run { .. } | Method::Resume) => {
|
||||
let _ = event_tx.send(Event::Error {
|
||||
code: ErrorCode::AlreadyRunning,
|
||||
message: "Pod is already executing a turn".into(),
|
||||
});
|
||||
}
|
||||
Some(Method::GetHistory) => {
|
||||
// Handled at socket layer; ignore here.
|
||||
}
|
||||
Some(Method::GetHistory) => {}
|
||||
None => {
|
||||
let _ = cancel_tx.try_send(());
|
||||
shared_state.set_status(PodStatus::Idle);
|
||||
return PodStatus::Idle;
|
||||
return (PodStatus::Idle, false);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -19,7 +19,7 @@ mod usage_tracker;
|
||||
|
||||
pub use token_counter::{EstimateSource, SplitPoint, TokenEstimate};
|
||||
|
||||
pub use controller::{PodController, PodHandle};
|
||||
pub use controller::{PodController, PodHandle, ShutdownReceiver};
|
||||
pub use factory::{FactoryError, PodFactory};
|
||||
pub use notifier::Notifier;
|
||||
pub use hook::{Hook, HookEventKind, HookRegistryBuilder};
|
||||
|
||||
@@ -157,8 +157,8 @@ async fn main() -> ExitCode {
|
||||
return ExitCode::FAILURE;
|
||||
}
|
||||
};
|
||||
let handle = match PodController::spawn(pod, &runtime_base).await {
|
||||
Ok(h) => h,
|
||||
let (handle, shutdown_rx) = match PodController::spawn(pod, &runtime_base).await {
|
||||
Ok(pair) => pair,
|
||||
Err(e) => {
|
||||
eprintln!("error: failed to start pod controller: {e}");
|
||||
return ExitCode::FAILURE;
|
||||
@@ -170,13 +170,12 @@ async fn main() -> ExitCode {
|
||||
handle.runtime_dir.socket_path()
|
||||
);
|
||||
|
||||
// Wait for shutdown signal
|
||||
match tokio::signal::ctrl_c().await {
|
||||
Ok(()) => {
|
||||
eprintln!("pod: {pod_name} shutting down");
|
||||
tokio::select! {
|
||||
_ = tokio::signal::ctrl_c() => {
|
||||
eprintln!("pod: {pod_name} shutting down (signal)");
|
||||
}
|
||||
Err(e) => {
|
||||
eprintln!("error: failed to listen for signal: {e}");
|
||||
_ = shutdown_rx => {
|
||||
eprintln!("pod: {pod_name} shutting down (client request)");
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -111,9 +111,9 @@ use pod::PodHandle;
|
||||
async fn spawn_controller(pod: Pod<MockClient, FsStore>) -> PodHandle {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let runtime_base = tmp.path().to_owned();
|
||||
// Leak tempdir so it survives the test
|
||||
std::mem::forget(tmp);
|
||||
PodController::spawn(pod, &runtime_base).await.unwrap()
|
||||
let (handle, _shutdown_rx) = PodController::spawn(pod, &runtime_base).await.unwrap();
|
||||
handle
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
@@ -12,6 +12,7 @@ pub enum Method {
|
||||
Run { input: String },
|
||||
Resume,
|
||||
Cancel,
|
||||
Shutdown,
|
||||
GetHistory,
|
||||
}
|
||||
|
||||
@@ -69,6 +70,7 @@ pub enum Event {
|
||||
greeting: Greeting,
|
||||
},
|
||||
Notification(Notification),
|
||||
Shutdown,
|
||||
}
|
||||
|
||||
/// User-facing notification emitted from the Pod layer.
|
||||
|
||||
@@ -12,6 +12,7 @@ pub struct App {
|
||||
pub input: String,
|
||||
pub cursor: usize,
|
||||
pub quit: bool,
|
||||
pub shutdown_confirm: Option<std::time::Instant>,
|
||||
/// Lines waiting to be flushed to terminal via insert_before.
|
||||
pub output_queue: Vec<OutputItem>,
|
||||
/// Partial streaming text not yet terminated by newline.
|
||||
@@ -55,6 +56,7 @@ impl App {
|
||||
input: String::new(),
|
||||
cursor: 0,
|
||||
quit: false,
|
||||
shutdown_confirm: None,
|
||||
output_queue: Vec::new(),
|
||||
pending_text: String::new(),
|
||||
}
|
||||
@@ -193,6 +195,9 @@ impl App {
|
||||
self.output_queue.insert(1, OutputItem::Blank);
|
||||
}
|
||||
}
|
||||
Event::Shutdown => {
|
||||
self.quit = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -168,6 +168,9 @@ fn handle_key(app: &mut App, key: KeyEvent) -> Option<Method> {
|
||||
}
|
||||
KeyCode::Char('r') if key.modifiers.contains(KeyModifiers::CONTROL) => Some(Method::Resume),
|
||||
KeyCode::Char('x') if key.modifiers.contains(KeyModifiers::CONTROL) => Some(Method::Cancel),
|
||||
KeyCode::Char('d') if key.modifiers.contains(KeyModifiers::CONTROL) => {
|
||||
return handle_shutdown(app);
|
||||
}
|
||||
KeyCode::Enter => app.submit_input(),
|
||||
KeyCode::Backspace => {
|
||||
app.delete_char_before();
|
||||
@@ -200,3 +203,23 @@ fn handle_key(app: &mut App, key: KeyEvent) -> Option<Method> {
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
const SHUTDOWN_CONFIRM_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3);
|
||||
|
||||
fn handle_shutdown(app: &mut App) -> Option<Method> {
|
||||
if !app.running {
|
||||
return Some(Method::Shutdown);
|
||||
}
|
||||
if let Some(t) = app.shutdown_confirm {
|
||||
if t.elapsed() < SHUTDOWN_CONFIRM_TIMEOUT {
|
||||
app.shutdown_confirm = None;
|
||||
return Some(Method::Shutdown);
|
||||
}
|
||||
}
|
||||
app.shutdown_confirm = Some(std::time::Instant::now());
|
||||
app.output_queue.push(app::OutputItem::Padded(
|
||||
app::MessageKind::Error,
|
||||
"Turn is running. Press Ctrl-D again to cancel and shut down.".into(),
|
||||
));
|
||||
None
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user