Files
mcp-memory/server/src/handlers/logs.rs
T

183 lines
6.2 KiB
Rust

use crate::router::McpTool;
use crate::state::MemoryState;
use crate::tools::{ProcessLogAction, ProcessLogsTool};
use async_trait::async_trait;
use serde_json::Value;
use std::fs::File;
use std::io::{Read, Seek, SeekFrom};
use std::sync::Arc;
pub struct ProcessLogsHandler;
#[async_trait]
impl McpTool for ProcessLogsHandler {
fn name(&self) -> &'static str {
"process_logs"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<ProcessLogsTool>(
"process_logs",
"Monitor, tail, and manage process logs: watch a log file, tail recent output, or clear log files.",
)
}
async fn execute(&self, args: Value, _state: Arc<MemoryState>) -> crate::error::Result<String> {
let tool_args: ProcessLogsTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let safe_path = crate::handlers::utils::validate_safe_path(&tool_args.file_path)?;
match tool_args.action {
ProcessLogAction::Watch => {
if !safe_path.exists() {
return Err(crate::error::AppError::Internal(format!(
"File does not exist: {}",
tool_args.file_path
)));
}
Ok(format!("Started watching logs for {}", tool_args.file_path))
}
ProcessLogAction::Get => {
let max_lines = tool_args.max_lines.unwrap_or(100);
let result =
tokio::task::spawn_blocking(move || -> crate::error::Result<String> {
let mut file = File::open(&safe_path).map_err(|e| {
crate::error::AppError::Internal(format!("Failed to open file: {}", e))
})?;
let len = file.metadata().map_err(|e| e.to_string())?.len();
let read_size = std::cmp::min(16 * 1024, len);
file.seek(SeekFrom::End(-(read_size as i64)))
.map_err(|e| e.to_string())?;
let mut vec_buf = Vec::new();
file.read_to_end(&mut vec_buf).map_err(|e| e.to_string())?;
let buffer = String::from_utf8_lossy(&vec_buf).to_string();
let lines: Vec<&str> = buffer.lines().collect();
let recent_lines = if lines.len() > max_lines {
lines[lines.len() - max_lines..].join("\n")
} else {
buffer
};
Ok(recent_lines)
})
.await
.map_err(|e| {
crate::error::AppError::Internal(format!("Task panic: {}", e))
})??;
Ok(result)
}
ProcessLogAction::Clear => {
if safe_path.exists() {
std::fs::write(&safe_path, "").map_err(|e| {
crate::error::AppError::Internal(format!("Failed to clear log file: {}", e))
})?;
}
Ok(format!("Cleared logs in {}", tool_args.file_path))
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use std::sync::Arc;
use tempfile::tempdir;
#[tokio::test]
async fn test_process_logs_watch() {
let dir = tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = ProcessLogsHandler;
let log_file = dir.path().join("test.log");
std::fs::write(&log_file, "line1\nline2").unwrap();
let args = json!({
"action": "watch",
"file_path": log_file.to_str().unwrap()
});
let result = handler
.execute(args, state)
.await
.map_err(|e| format!("Failed to watch logs: {}", e))
.unwrap();
assert!(result.contains("Started watching logs"));
}
#[tokio::test]
async fn test_process_logs_get() {
let dir = tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = ProcessLogsHandler;
let log_file = dir.path().join("test_recent.log");
std::fs::write(&log_file, "line1\nline2\nline3").unwrap();
let args = json!({
"action": "get",
"file_path": log_file.to_str().unwrap()
});
let result = handler
.execute(args, state)
.await
.map_err(|e| format!("Failed to get recent logs: {}", e))
.unwrap();
assert!(result.contains("line1"));
assert!(result.contains("line3"));
}
#[tokio::test]
async fn test_process_logs_clear() {
let dir = tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = ProcessLogsHandler;
let log_file = dir.path().join("test_clear.log");
std::fs::write(&log_file, "line1\nline2\nline3").unwrap();
let args = json!({
"action": "clear",
"file_path": log_file.to_str().unwrap()
});
let result = handler.execute(args, state.clone()).await.unwrap();
assert!(result.contains("Cleared logs"));
let content = std::fs::read_to_string(&log_file).unwrap();
assert_eq!(content, "");
}
#[tokio::test]
async fn test_process_logs_get_with_large_file() {
let dir = tempfile::tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = ProcessLogsHandler;
let log_file = dir.path().join("large_test.log");
let mut buffer = String::new();
for _ in 0..1000 {
buffer.push_str("line\n");
}
std::fs::write(&log_file, buffer).unwrap();
let args = serde_json::json!({
"action": "get",
"file_path": log_file.to_str().unwrap()
});
let result = handler
.execute(args, state)
.await
.map_err(|e| format!("Failed to get recent logs: {}", e))
.unwrap();
assert!(result.contains("line"));
}
}