Skip to content

Commit e16b37a

Browse files
authored
Merge pull request #22 from NianJiuZst/codex/enforce-session-idle-timeout
fix(daemon): enforce session idle timeout
2 parents d27ef1f + bf6c384 commit e16b37a

6 files changed

Lines changed: 208 additions & 13 deletions

File tree

crates/bsk-cli/src/daemon/mod.rs

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -44,10 +44,16 @@ pub async fn run(
4444
Some(path) => Some(ipc::IpcServer::new(Arc::clone(&state)).bind(path).await?),
4545
None => None,
4646
};
47+
let session_idle_task = start::spawn_session_idle_reaper(Arc::clone(&state));
4748
// Ensure WS/IPC accept loops have polled `shutdown.notified()` before any
4849
// caller can invoke `DaemonHandle::shutdown()` — `Notify::notify_waiters()`
4950
// drops wakeups when nothing is registered yet, which otherwise hangs
5051
// shutdown forever (hit by tests that connect/shutdown immediately).
5152
tokio::task::yield_now().await;
52-
Ok(DaemonHandle::new(state, ws_handle, ipc_handle))
53+
Ok(DaemonHandle::new(
54+
state,
55+
ws_handle,
56+
ipc_handle,
57+
session_idle_task,
58+
))
5359
}

crates/bsk-cli/src/daemon/queue.rs

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -300,7 +300,8 @@ impl ToolQueueRegistry {
300300
queue_state.busy = true;
301301
(entry.sender.clone(), Arc::clone(&entry.state))
302302
};
303-
dispatch_with_sender(
303+
self.sessions.touch(sid);
304+
let outcome = dispatch_with_sender(
304305
sender,
305306
state,
306307
sid.clone(),
@@ -310,7 +311,13 @@ impl ToolQueueRegistry {
310311
false,
311312
inflight,
312313
)
313-
.await
314+
.await;
315+
// A long-running request may outlive the idle threshold. Touching
316+
// after completion prevents the reaper from immediately closing it;
317+
// while it is running, session.stop observes SessionBusy and retries
318+
// on a later sweep rather than interrupting the tool.
319+
self.sessions.touch(sid);
320+
outcome
314321
}
315322

316323
/// Stop accepting new jobs for `sid`, enqueue one final control RPC,

crates/bsk-cli/src/daemon/sessions.rs

Lines changed: 69 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66
use std::collections::HashMap;
77
use std::sync::Arc;
88
use std::sync::Mutex;
9-
use std::time::{SystemTime, UNIX_EPOCH};
9+
use std::time::{Instant, SystemTime, UNIX_EPOCH};
1010

1111
use bsk_protocol::system::{BrowserStatusEntry, SessionStatusEntry};
1212
use bsk_protocol::tools::{
@@ -66,6 +66,9 @@ impl Session {
6666
#[derive(Debug, Default)]
6767
pub struct SessionRegistry {
6868
inner: Mutex<HashMap<SessionId, Session>>,
69+
/// Operational metadata kept outside the public `Session` wire/domain
70+
/// shape so idle enforcement does not break external struct users.
71+
last_activity: Mutex<HashMap<SessionId, Instant>>,
6972
}
7073

7174
impl SessionRegistry {
@@ -100,10 +103,13 @@ impl SessionRegistry {
100103
}
101104

102105
pub fn insert(&self, session: Session) {
103-
self.inner
106+
let session_id = session.id.clone();
107+
let mut sessions = self.inner.lock().expect("session registry poisoned");
108+
sessions.insert(session_id.clone(), session);
109+
self.last_activity
104110
.lock()
105-
.expect("session registry poisoned")
106-
.insert(session.id.clone(), session);
111+
.expect("session activity registry poisoned")
112+
.insert(session_id, Instant::now());
107113
}
108114

109115
/// Reserve a fresh, collision-free [`SessionId`] under the registry
@@ -141,6 +147,10 @@ impl SessionRegistry {
141147
created_at_ms: now_ms_fn(),
142148
},
143149
);
150+
self.last_activity
151+
.lock()
152+
.expect("session activity registry poisoned")
153+
.insert(candidate.clone(), Instant::now());
144154
return Some(candidate);
145155
}
146156
None
@@ -168,13 +178,23 @@ impl SessionRegistry {
168178
.lock()
169179
.expect("session registry poisoned")
170180
.remove(session_id);
181+
self.last_activity
182+
.lock()
183+
.expect("session activity registry poisoned")
184+
.remove(session_id);
171185
}
172186

173187
pub fn remove(&self, id: &SessionId) -> Option<Session> {
174-
self.inner
188+
let removed = self
189+
.inner
175190
.lock()
176191
.expect("session registry poisoned")
177-
.remove(id)
192+
.remove(id);
193+
self.last_activity
194+
.lock()
195+
.expect("session activity registry poisoned")
196+
.remove(id);
197+
removed
178198
}
179199

180200
pub fn get(&self, id: &SessionId) -> Option<Session> {
@@ -185,6 +205,42 @@ impl SessionRegistry {
185205
.cloned()
186206
}
187207

208+
/// Record accepted or completed tool activity for a live session.
209+
pub fn touch(&self, id: &SessionId) -> bool {
210+
if !self
211+
.inner
212+
.lock()
213+
.expect("session registry poisoned")
214+
.contains_key(id)
215+
{
216+
return false;
217+
}
218+
self.last_activity
219+
.lock()
220+
.expect("session activity registry poisoned")
221+
.insert(id.clone(), Instant::now());
222+
true
223+
}
224+
225+
/// Return sessions whose last tool activity is at least `idle_for`
226+
/// old. The caller supplies `now` to keep boundary tests deterministic.
227+
pub fn idle_ids_at(&self, idle_for: Duration, now: Instant) -> Vec<SessionId> {
228+
let sessions = self.inner.lock().expect("session registry poisoned");
229+
let activity = self
230+
.last_activity
231+
.lock()
232+
.expect("session activity registry poisoned");
233+
sessions
234+
.values()
235+
.filter(|session| {
236+
activity
237+
.get(&session.id)
238+
.is_some_and(|last| now.saturating_duration_since(*last) >= idle_for)
239+
})
240+
.map(|session| session.id.clone())
241+
.collect()
242+
}
243+
188244
/// Drop all sessions owned by `browser_id` (e.g. on disconnect).
189245
pub fn purge_browser(&self, browser_id: &BrowserId) -> Vec<Session> {
190246
let mut guard = self.inner.lock().expect("session registry poisoned");
@@ -196,6 +252,13 @@ impl SessionRegistry {
196252
for s in &drained {
197253
guard.remove(&s.id);
198254
}
255+
let mut activity = self
256+
.last_activity
257+
.lock()
258+
.expect("session activity registry poisoned");
259+
for session in &drained {
260+
activity.remove(&session.id);
261+
}
199262
drained
200263
}
201264

crates/bsk-cli/src/daemon/start.rs

Lines changed: 57 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -23,8 +23,11 @@ use tracing::{debug, info, warn};
2323

2424
use crate::cli::daemon::StartArgs;
2525
use crate::daemon::{
26-
browsers::EXTENSION_CONNECT_WAIT, info as daemon_info, ipc, lockfile, paths,
27-
state::DaemonState, ws,
26+
browsers::EXTENSION_CONNECT_WAIT,
27+
info as daemon_info, ipc, lockfile, paths,
28+
sessions::{StopSessionError, forget_session, stop_session},
29+
state::DaemonState,
30+
ws,
2831
};
2932

3033
/// Internal env-var contract: the parent sets this on the spawned child
@@ -196,6 +199,7 @@ pub fn run_foreground(cfg: DaemonConfig) -> Result<()> {
196199
.with_context(|| format!("bind IPC endpoint {}", sock_path.display()))?;
197200

198201
let state = Arc::new(DaemonState::new(cfg.clone()));
202+
let session_idle_task = spawn_session_idle_reaper(Arc::clone(&state));
199203
let ws_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), cfg.ws_port);
200204
let ws_handle = ws::WsServer::new(Arc::clone(&state))
201205
.bind(ws_addr)
@@ -355,6 +359,8 @@ pub fn run_foreground(cfg: DaemonConfig) -> Result<()> {
355359

356360
let _ = ipc_shutdown_tx.send(());
357361
let _ = ipc_task.await;
362+
session_idle_task.abort();
363+
let _ = session_idle_task.await;
358364
ws_handle.shutdown.notify_waiters();
359365
let _ = ws_handle.task.await;
360366

@@ -367,6 +373,55 @@ pub fn run_foreground(cfg: DaemonConfig) -> Result<()> {
367373
Ok(())
368374
}
369375

376+
/// Spawn the cooperative session-idle reaper shared by the production
377+
/// foreground daemon and the test/embed daemon entry point.
378+
pub(crate) fn spawn_session_idle_reaper(state: Arc<DaemonState>) -> tokio::task::JoinHandle<()> {
379+
tokio::spawn(async move {
380+
let session_idle = state.config.session_idle;
381+
let tick = (session_idle / 4)
382+
.max(Duration::from_millis(100))
383+
.min(Duration::from_secs(30));
384+
let mut ticker = tokio::time::interval(tick);
385+
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
386+
// `interval`'s first tick is immediate. Consume it so a zero/very
387+
// short test setting still gets one real inactivity window.
388+
ticker.tick().await;
389+
390+
loop {
391+
ticker.tick().await;
392+
let idle_ids = state.sessions.idle_ids_at(session_idle, Instant::now());
393+
for session_id in idle_ids {
394+
match stop_session(
395+
&state.browsers,
396+
&state.sessions,
397+
&state.tool_queues,
398+
&state.session_interrupts,
399+
&session_id,
400+
Duration::from_secs(10),
401+
)
402+
.await
403+
{
404+
Ok(_) => info!(session = %session_id, "idle session stopped"),
405+
Err(StopSessionError::SessionBusy | StopSessionError::Stopping) => {
406+
debug!(session = %session_id, "idle session still active; retrying later");
407+
}
408+
Err(StopSessionError::NotFound | StopSessionError::BrowserGone) => {
409+
forget_session(
410+
&state.sessions,
411+
&state.tool_queues,
412+
&state.session_interrupts,
413+
&session_id,
414+
);
415+
}
416+
Err(err) => {
417+
warn!(session = %session_id, error = %err, "failed to stop idle session");
418+
}
419+
}
420+
}
421+
}
422+
})
423+
}
424+
370425
fn record_activity(activity: &Arc<Mutex<Instant>>) {
371426
if let Ok(mut a) = activity.lock() {
372427
*a = Instant::now();

crates/bsk-cli/src/daemon/state.rs

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -81,11 +81,22 @@ pub struct DaemonHandle {
8181
state: Arc<DaemonState>,
8282
ws: WsHandle,
8383
ipc: Option<IpcHandle>,
84+
session_idle_task: JoinHandle<()>,
8485
}
8586

8687
impl DaemonHandle {
87-
pub(crate) fn new(state: Arc<DaemonState>, ws: WsHandle, ipc: Option<IpcHandle>) -> Self {
88-
Self { state, ws, ipc }
88+
pub(crate) fn new(
89+
state: Arc<DaemonState>,
90+
ws: WsHandle,
91+
ipc: Option<IpcHandle>,
92+
session_idle_task: JoinHandle<()>,
93+
) -> Self {
94+
Self {
95+
state,
96+
ws,
97+
ipc,
98+
session_idle_task,
99+
}
89100
}
90101

91102
pub fn state(&self) -> Arc<DaemonState> {
@@ -103,6 +114,8 @@ impl DaemonHandle {
103114
/// Stop the WS server (and IPC if running). Returns once both join
104115
/// handles complete.
105116
pub async fn shutdown(self) {
117+
self.session_idle_task.abort();
118+
let _ = await_join(self.session_idle_task).await;
106119
self.ws.shutdown.notify_waiters();
107120
let _ = await_join(self.ws.task).await;
108121
if let Some(ipc) = self.ipc {

crates/bsk-cli/tests/sessions_ipc.rs

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,16 @@ async fn spawn_daemon_with_connect_wait(connect_wait: Duration) -> (daemon::Daem
4646
(handle, sock)
4747
}
4848

49+
async fn spawn_daemon_with_session_idle(session_idle: Duration) -> (daemon::DaemonHandle, PathBuf) {
50+
let port = 0;
51+
52+
let mut config = DaemonConfig::new(port);
53+
config.session_idle = session_idle;
54+
let sock = tempfile_path("bsk-test-ipc");
55+
let handle = daemon::run(config, Some(sock.clone())).await.unwrap();
56+
(handle, sock)
57+
}
58+
4959
async fn connect_ext(
5060
addr: std::net::SocketAddr,
5161
) -> tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>> {
@@ -237,6 +247,47 @@ async fn session_start_stop_round_trip_via_ipc() {
237247
handle.shutdown().await;
238248
}
239249

250+
#[tokio::test]
251+
async fn session_idle_timeout_stops_and_unregisters_session() {
252+
let (handle, _sock) = spawn_daemon_with_session_idle(Duration::from_millis(50)).await;
253+
let mut ws = connect_ext(handle.ws_addr()).await;
254+
let _ = handshake_as_ext(&mut ws).await;
255+
let state = handle.state();
256+
let session_id = bsk::daemon::sessions::SessionId("idle".into());
257+
state.sessions.insert(bsk::daemon::sessions::Session {
258+
id: session_id.clone(),
259+
browser_id: bsk::daemon::browsers::BrowserId(TEST_EXT_ID.into()),
260+
agent_window_id: Some(7),
261+
created_at_ms: 0,
262+
});
263+
state.tool_queues.spawn(session_id);
264+
265+
let request = tokio::time::timeout(Duration::from_secs(1), ws.next())
266+
.await
267+
.expect("idle reaper did not contact extension")
268+
.expect("extension socket closed")
269+
.expect("extension socket failed");
270+
let Message::Text(text) = request else {
271+
panic!("expected text request");
272+
};
273+
let request: RequestFrame = serde_json::from_str(&text).unwrap();
274+
assert_eq!(request.method, Method::ToolSessionStop);
275+
ws.send(Message::Text(
276+
serde_json::to_string(&ResponseFrame {
277+
id: request.id,
278+
body: ResponseBody::Ok(
279+
serde_json::to_value(bsk_protocol::tools::SessionStopResult::default()).unwrap(),
280+
),
281+
})
282+
.unwrap(),
283+
))
284+
.await
285+
.unwrap();
286+
287+
wait_for_no_sessions(&state).await;
288+
handle.shutdown().await;
289+
}
290+
240291
#[tokio::test]
241292
async fn session_start_errors_without_browser() {
242293
let (handle, sock) = spawn_daemon_with_connect_wait(Duration::ZERO).await;

0 commit comments

Comments
 (0)