use crate::paths::AgentsHome;
use crate::protocol::{write_response, ErrorCode, Request, Response};
use futures_util::{SinkExt, StreamExt};
use serde_json::{json, Value};
use std::io::{Read, Seek, SeekFrom};
use std::os::unix::fs::MetadataExt;
use std::path::{Path, PathBuf};
use std::time::Duration;
use tokio::net::UnixStream;
use tokio_tungstenite::tungstenite::Message;
const UPGRADE_TIMEOUT: Duration = Duration::from_secs(10);
const POLL: Duration = Duration::from_millis(500);
const KEEPALIVE: Duration = Duration::from_secs(30);
pub async fn handle_logs(home: &AgentsHome, req: &Request, mut stream: UnixStream) {
let name = req
.params
.get("name")
.and_then(|v| v.as_str())
.unwrap_or("");
let log_path = registry_log_path(&home.registry_json(), name);
let log_path = match log_path {
Some(p) => PathBuf::from(p),
None => {
let resp = Response::err(
req.id,
ErrorCode::AgentNotFound,
format!("agent `{name}` has no tee'd log path to follow"),
);
let _ = write_response(&mut stream, &resp).await;
return;
}
};
let ack = Response::ok(req.id, json!({"streaming": true}));
if write_response(&mut stream, &ack).await.is_err() {
return;
}
let ws = match tokio::time::timeout(UPGRADE_TIMEOUT, tokio_tungstenite::accept_async(stream))
.await
{
Ok(Ok(ws)) => ws,
Ok(Err(_)) | Err(_) => return,
};
stream_follow(ws, &log_path).await;
}
fn registry_log_path(registry_path: &Path, name: &str) -> Option<String> {
let bytes = std::fs::read(registry_path).ok()?;
let raw: Value = serde_json::from_str(&String::from_utf8_lossy(&bytes)).ok()?;
let rows = raw
.get("agents")
.or_else(|| raw.get("entries"))?
.as_array()?;
rows.iter()
.find(|e| e.get("name").and_then(Value::as_str) == Some(name))
.and_then(|e| e.get("log_path").and_then(Value::as_str))
.map(str::to_string)
}
fn stat(path: &Path) -> Option<(u64, u64)> {
std::fs::metadata(path).ok().map(|m| (m.len(), m.ino()))
}
fn end_frame(msg: String) -> Message {
Message::Text(json!({"t": "end", "msg": msg}).to_string().into())
}
async fn stream_follow(ws: tokio_tungstenite::WebSocketStream<UnixStream>, path: &Path) {
let (mut sink, mut source) = ws.split();
let mut file = match std::fs::File::open(path) {
Ok(f) => f,
Err(_) => {
let _ = sink
.send(end_frame(format!(
"log file disappeared: {}\n",
path.display()
)))
.await;
return;
}
};
let (mut offset, inode) = match stat(path) {
Some(s) => s,
None => {
let _ = sink
.send(end_frame(format!(
"log file disappeared: {}\n",
path.display()
)))
.await;
return;
}
};
let mut last_size = offset;
let mut poll = tokio::time::interval(POLL);
poll.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut keepalive = tokio::time::interval(KEEPALIVE);
keepalive.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
keepalive.tick().await;
loop {
tokio::select! {
_ = keepalive.tick() => {
if sink.send(Message::Ping(Vec::new().into())).await.is_err() {
break; }
}
_ = poll.tick() => {
let (size, ino) = match stat(path) {
Some(s) => s,
None => {
let _ = sink.send(end_frame(
format!("log file disappeared: {}\n", path.display()))).await;
break;
}
};
if ino != inode {
let _ = sink.send(end_frame(
format!("log file rotated (inode changed): {}\n", path.display()))).await;
break;
}
if size < offset {
let _ = sink.send(end_frame(
format!("log file truncated: {}\n", path.display()))).await;
break;
}
if size < last_size {
let _ = sink.send(end_frame(
format!("log file truncated (size shrank across poll): {}\n", path.display()))).await;
break;
}
last_size = size;
if size > offset {
match read_from(&mut file, offset) {
Ok(bytes) if !bytes.is_empty() => {
offset += bytes.len() as u64;
if sink.send(Message::Binary(bytes.into())).await.is_err() {
break; }
}
Ok(_) => {}
Err(_) => {
}
}
}
}
msg = source.next() => {
match msg {
Some(Ok(Message::Text(_))) => break,
Some(Ok(Message::Close(_))) | None => break,
Some(Err(_)) => break,
_ => {}
}
}
}
}
let _ = sink.close().await;
}
fn read_from(f: &mut std::fs::File, offset: u64) -> std::io::Result<Vec<u8>> {
f.seek(SeekFrom::Start(offset))?;
let mut buf = Vec::new();
f.read_to_end(&mut buf)?;
Ok(buf)
}