(
&mut self,
idle_timeout_secs: u64,
absolute_timeout_secs: u64,
)
| 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); |
no test coverage detected