hyperhive/hive-matrix-mcp/src/socket.rs
iris 92b32d06fb fix(clippy): fix all clippy warnings in hive-ag3nt, hive-forge, hive-matrix-mcp, hive-sh4re
Fixes all clippy -D warnings errors in the crates iris owns:

hive-sh4re:
- doc_lazy_continuation: add blank /// separator in priv_proto.rs
- doc_markdown: backtick PRIVATE_NETWORK=0 / PRIVATE_NETWORK=1

hive-matrix-mcp:
- map_unwrap_or: map().unwrap_or_else() -> map_or_else() in paths.rs
- collapsible_if: if-let chains in wake.rs
- doc_markdown: backtick M_UNKNOWN_TOKEN in main.rs
- cast_possible_truncation: usize/u64 -> u32::try_from in handlers.rs
- map_unwrap_or: map_or_else() in handlers.rs
- manual_let_else: match Ok(r) => r, Err => return -> let Ok in handlers.rs
- unused_async: remove async from list_invites; update socket.rs call site

hive-forge:
- doc_markdown: backtick REQUEST_CHANGES / APPROVED / COMMENT in pr_reviews.rs
- unnecessary_wraps: list_reviews_text returns () not Result<()>
- doc_markdown: backtick start_page / last_page in comments.rs
- cast_possible_truncation: PAGE_SIZE u64 -> usize; remove as usize casts

hive-ag3nt:
- collapsible_if: if-let chains in events.rs and mcp.rs
- single_match_else: match -> if let in events.rs and mcp.rs
- items_after_statements: hoist STATUS_MAX_CHARS const in mcp.rs
- map_unwrap_or: map_or_else() in mcp.rs and mcp_loose_ends.rs
- cast_possible_truncation: usize -> u32::try_from in mcp.rs
- doc_markdown: backtick snake_case in mcp.rs, needs_update/deployed_sha
  in web_ui.rs, HISTORY_CAPACITY in web_ui.rs
- identical_match_arms: combine manage_root_agent | query_agent_state
- redundant_closure: |s| s.to_string() -> ToString::to_string in web_ui.rs
- duration_suboptimal_units: from_secs(3600) -> from_hours(1) in turn.rs

Remaining failures in hive-c0re (39), hive-priv (8), hive-bash-mcp (11)
are owned by damocles.
2026-06-05 14:35:12 +02:00

92 lines
3.8 KiB
Rust

//! Unix socket server: the daemon listens here, the stdio MCP bridge
//! `connect()`s on every tool call. One JSON request line in, one
//! JSON response line out. Connections are short-lived (per tool call)
//! so the loop is just accept → dispatch → reply → close.
use std::path::Path;
use std::sync::Arc;
use anyhow::{Context, Result};
use matrix_sdk::Client;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::{UnixListener, UnixStream};
use crate::handlers;
use crate::protocol::{DaemonRequest, DaemonResponse};
/// Start listening on `socket_path` and serve forever. Removes any
/// stale socket file first (daemon restart after a non-clean shutdown
/// would otherwise hit EADDRINUSE).
pub async fn serve(socket_path: &Path, client: Client) -> Result<()> {
let _ = tokio::fs::remove_file(socket_path).await;
if let Some(parent) = socket_path.parent() {
tokio::fs::create_dir_all(parent)
.await
.with_context(|| format!("mkdir {}", parent.display()))?;
}
let listener = UnixListener::bind(socket_path)
.with_context(|| format!("bind unix socket {}", socket_path.display()))?;
tracing::info!(path = %socket_path.display(), "mcp socket listener up");
let client = Arc::new(client);
loop {
let (stream, _) = listener
.accept()
.await
.context("accept connection on mcp socket")?;
let client = client.clone();
tokio::spawn(async move {
if let Err(e) = handle_connection(stream, &client).await {
tracing::warn!(error = %e, "mcp socket connection error");
}
});
}
}
async fn handle_connection(stream: UnixStream, client: &Client) -> Result<()> {
let (reader, mut writer) = stream.into_split();
let mut lines = BufReader::new(reader).lines();
while let Some(line) = lines.next_line().await? {
let response = match serde_json::from_str::<DaemonRequest>(&line) {
Ok(req) => dispatch(req, client).await,
Err(e) => DaemonResponse::error(format!("parse request: {e}")),
};
let mut json = serde_json::to_string(&response)?;
json.push('\n');
writer.write_all(json.as_bytes()).await?;
writer.flush().await?;
}
Ok(())
}
async fn dispatch(req: DaemonRequest, client: &Client) -> DaemonResponse {
match req {
DaemonRequest::Ping => DaemonResponse::ok(&serde_json::json!({"ok": true})),
DaemonRequest::SendMessage { room, body } => {
handlers::send_message(client, &room, &body).await
}
DaemonRequest::SendDm { user_id, body } => handlers::send_dm(client, &user_id, &body).await,
DaemonRequest::SendReaction {
room,
event_id,
key,
} => handlers::send_reaction(client, &room, &event_id, &key).await,
DaemonRequest::SendReply {
room,
event_id,
body,
} => handlers::send_reply(client, &room, &event_id, &body).await,
DaemonRequest::MarkRead { room, event_id } => {
handlers::mark_read(client, &room, &event_id).await
}
DaemonRequest::ListRooms => handlers::list_rooms(client).await,
DaemonRequest::ListInvites => handlers::list_invites(client),
DaemonRequest::JoinRoom { room } => handlers::join_room(client, &room).await,
DaemonRequest::InviteUser { room, user_id } => {
handlers::invite_user(client, &room, &user_id).await
}
DaemonRequest::ListRoomMembers { room } => handlers::list_room_members(client, &room).await,
DaemonRequest::ReadRoom { room, limit } => handlers::read_room(client, &room, limit).await,
DaemonRequest::UnreadCount => handlers::unread_count(client),
DaemonRequest::UnreadSummary => handlers::unread_summary(client).await,
}
}