MCPcopy Create free account
hub / github.com/NodeDB-Lab/nodedb / run_inner

Method run_inner

nodedb/src/control/server/session.rs:272–359  ·  view source on GitHub ↗
(
        &mut self,
        idle_timeout_secs: u64,
        absolute_timeout_secs: u64,
    )

Source from the content-addressed store, hash-verified

270 }
271
272 async fn run_inner(
273 &mut self,
274 idle_timeout_secs: u64,
275 absolute_timeout_secs: u64,
276 ) -> crate::Result<()> {
277 loop {
278 // Hard-revoke check: bus consumer sent kill signal.
279 if self.is_killed() {
280 let msg = r#"{"status":"error","sqlstate":"57P01","error":"session revoked by administrator"}"#;
281 let resp_len = (msg.len() as u32).to_be_bytes();
282 let _ = self.stream.write_all(&resp_len).await;
283 let _ = self.stream.write_all(msg.as_bytes()).await;
284 return Ok(());
285 }
286
287 // Enforce absolute session lifetime (SQLSTATE 57P01 "admin shutdown").
288 if absolute_timeout_secs > 0
289 && self.connected_at.elapsed().as_secs() >= absolute_timeout_secs
290 {
291 debug!(
292 "session absolute timeout ({}s), closing connection",
293 absolute_timeout_secs
294 );
295 let msg = r#"{"status":"error","sqlstate":"57P01","error":"session timeout: absolute lifetime exceeded"}"#;
296 let resp_len = (msg.len() as u32).to_be_bytes();
297 let _ = self.stream.write_all(&resp_len).await;
298 let _ = self.stream.write_all(msg.as_bytes()).await;
299 return Ok(());
300 }
301
302 // Read length prefix with idle timeout.
303 let mut len_buf = [0u8; 4];
304 let read_result: std::io::Result<usize> = if idle_timeout_secs > 0 {
305 match tokio::time::timeout(
306 Duration::from_secs(idle_timeout_secs),
307 self.stream.read_exact(&mut len_buf),
308 )
309 .await
310 {
311 Ok(result) => result,
312 Err(_) => {
313 debug!("session idle timeout ({}s)", idle_timeout_secs);
314 return Ok(());
315 }
316 }
317 } else {
318 self.stream.read_exact(&mut len_buf).await
319 };
320 match read_result {
321 Ok(_) => {}
322 Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => {
323 debug!("client disconnected");
324 return Ok(());
325 }
326 Err(e) => return Err(e.into()),
327 }
328
329 let payload_len = u32::from_be_bytes(len_buf);

Callers 1

runMethod · 0.80

Calls 8

is_killedMethod · 0.80
elapsedMethod · 0.80
kindMethod · 0.80
handle_frameMethod · 0.80
lenMethod · 0.45
as_bytesMethod · 0.45
read_exactMethod · 0.45
next_request_idMethod · 0.45

Tested by

no test coverage detected