168 lines
5.9 KiB
Rust
168 lines
5.9 KiB
Rust
//! In-agent socket server (loose-ends v2). Binds the harness-owned
|
|
//! `HIVE_AGENT_SOCKET` and serves the `hive-agent-sock` protocol to the
|
|
//! in-container producers (matrix / bash daemons, forge-notify). Todo ops
|
|
//! hit the harness-local [`Todos`] store; a new-or-changed upsert fires an
|
|
//! in-process [`Notify`] so the serve loop drives a turn — no hive-c0re
|
|
//! round-trip, no broker long-poll, no marker files.
|
|
//!
|
|
//! One request/response line per connection, matching the producers'
|
|
//! existing best-effort JSON-line clients (they just change which socket
|
|
//! they dial, not the payload).
|
|
|
|
use std::path::{Path, PathBuf};
|
|
use std::sync::Arc;
|
|
|
|
use anyhow::{Context, Result};
|
|
use hive_agent_sock::{Request, Response};
|
|
use hive_sh4re::LooseEnd;
|
|
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
|
use tokio::net::{UnixListener, UnixStream};
|
|
use tokio::sync::Notify;
|
|
|
|
use crate::todos::{Todo, Todos};
|
|
|
|
/// Resolve the in-agent socket path from `HIVE_AGENT_SOCKET`. `None` when
|
|
/// unset/empty (e.g. a standalone dev run) — the server then stays off.
|
|
fn socket_path() -> Option<PathBuf> {
|
|
std::env::var_os("HIVE_AGENT_SOCKET")
|
|
.filter(|v| !v.is_empty())
|
|
.map(PathBuf::from)
|
|
}
|
|
|
|
/// Run the in-agent socket server: bind + accept loop, one request/response
|
|
/// line per connection. A no-op (returns `Ok`) when `HIVE_AGENT_SOCKET` is
|
|
/// unset, so a standalone harness without producers just skips it.
|
|
///
|
|
/// # Errors
|
|
///
|
|
/// Returns an error if the socket path is set but can't be bound.
|
|
pub async fn run(store: Arc<Todos>, wake: Arc<Notify>) -> Result<()> {
|
|
let Some(path) = socket_path() else {
|
|
tracing::info!("HIVE_AGENT_SOCKET unset — in-agent todo socket disabled");
|
|
return Ok(());
|
|
};
|
|
let listener = bind(&path)?;
|
|
tracing::info!(socket = %path.display(), "in-agent todo socket listening");
|
|
loop {
|
|
match listener.accept().await {
|
|
Ok((stream, _)) => {
|
|
let store = store.clone();
|
|
let wake = wake.clone();
|
|
tokio::spawn(async move {
|
|
if let Err(e) = handle_conn(stream, &store, &wake).await {
|
|
tracing::warn!(error = ?e, "in-agent todo connection failed");
|
|
}
|
|
});
|
|
}
|
|
Err(e) => tracing::warn!(error = ?e, "in-agent todo accept failed"),
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Bind a `UnixListener` at `path`, creating the parent dir and unlinking a
|
|
/// stale socket left by a prior boot (which would otherwise block `bind`
|
|
/// with `EADDRINUSE`).
|
|
fn bind(path: &Path) -> Result<UnixListener> {
|
|
if let Some(parent) = path.parent() {
|
|
std::fs::create_dir_all(parent)
|
|
.with_context(|| format!("create in-agent socket dir {}", parent.display()))?;
|
|
}
|
|
if path.exists() {
|
|
let _ = std::fs::remove_file(path);
|
|
}
|
|
UnixListener::bind(path).with_context(|| format!("bind in-agent socket {}", path.display()))
|
|
}
|
|
|
|
/// Handle one connection: read a single JSON request line, apply it to the
|
|
/// store, write the JSON response line back.
|
|
async fn handle_conn(stream: UnixStream, store: &Todos, wake: &Notify) -> Result<()> {
|
|
let (read, mut write) = stream.into_split();
|
|
let mut reader = BufReader::new(read);
|
|
let mut line = String::new();
|
|
if reader.read_line(&mut line).await? == 0 {
|
|
return Ok(());
|
|
}
|
|
let resp = match serde_json::from_str::<Request>(line.trim()) {
|
|
Ok(req) => dispatch(req, store, wake),
|
|
Err(e) => Response::Err {
|
|
message: format!("bad request: {e}"),
|
|
},
|
|
};
|
|
let mut out = serde_json::to_string(&resp)?;
|
|
out.push('\n');
|
|
write.write_all(out.as_bytes()).await?;
|
|
write.flush().await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Apply one request to the store, firing `wake` on a new/changed upsert so
|
|
/// the serve loop runs a turn.
|
|
fn dispatch(req: Request, store: &Todos, wake: &Notify) -> Response {
|
|
match req {
|
|
Request::UpsertTodo {
|
|
subsystem,
|
|
key,
|
|
summary,
|
|
source,
|
|
} => match store.upsert(&subsystem, key.as_deref(), &summary, source.as_deref()) {
|
|
Ok((_, changed)) => {
|
|
if changed {
|
|
wake.notify_one();
|
|
}
|
|
Response::Ok
|
|
}
|
|
Err(e) => err(&e),
|
|
},
|
|
Request::ClearTodo {
|
|
subsystem,
|
|
key,
|
|
all,
|
|
} => {
|
|
let result = if all {
|
|
store.clear_subsystem(&subsystem)
|
|
} else {
|
|
store.clear(&subsystem, key.as_deref())
|
|
};
|
|
match result {
|
|
Ok(count) => Response::Acked {
|
|
count: u64::try_from(count).unwrap_or(0),
|
|
},
|
|
Err(e) => err(&e),
|
|
}
|
|
}
|
|
Request::ListTodos { subsystem } => match store.list(subsystem.as_deref()) {
|
|
Ok(todos) => Response::LooseEnds {
|
|
loose_ends: todos.into_iter().map(to_loose_end).collect(),
|
|
},
|
|
Err(e) => err(&e),
|
|
},
|
|
Request::MarkTodoDone { id } => match store.mark_done(id) {
|
|
Ok(count) => Response::Acked {
|
|
count: u64::try_from(count).unwrap_or(0),
|
|
},
|
|
Err(e) => err(&e),
|
|
},
|
|
}
|
|
}
|
|
|
|
/// Format a store error into an `Err` response.
|
|
fn err(e: &anyhow::Error) -> Response {
|
|
Response::Err {
|
|
message: format!("{e:#}"),
|
|
}
|
|
}
|
|
|
|
/// Map a stored [`Todo`] to a [`LooseEnd::Todo`], deriving `age_seconds`
|
|
/// from `updated_at` (saturating so a backwards clock step reads 0).
|
|
fn to_loose_end(t: Todo) -> LooseEnd {
|
|
let now = hive_sh4re::wire_time::now_unix();
|
|
let age = u64::try_from(now.saturating_sub(t.updated_at)).unwrap_or(0);
|
|
LooseEnd::Todo {
|
|
id: t.id,
|
|
subsystem: t.subsystem,
|
|
subsystem_key: t.subsystem_key,
|
|
summary: t.summary,
|
|
source: t.source,
|
|
age_seconds: age,
|
|
}
|
|
}
|