Merge nucleic/sleek-thistle-egret-fyej into dev
This commit is contained in:
@@ -5,3 +5,8 @@ edition = "2021"
|
||||
description = "nash-observe — event model, taps, batching, and transport for nash"
|
||||
license = "MIT"
|
||||
publish = false
|
||||
|
||||
[dependencies]
|
||||
brush-core = { path = "../../third_party/brush/brush-core" }
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
|
||||
+496
-4
@@ -1,5 +1,497 @@
|
||||
//! nash-observe — event model, redirection/pipe taps, batching, spool, and
|
||||
//! transport for nash (see docs/NASH.md §4–§6).
|
||||
//! nash-observe — the observation layer for nash (docs/NASH.md §4–§6).
|
||||
//!
|
||||
//! M0 placeholder: this crate gains the `Gate` implementation (allow-all),
|
||||
//! the shell-batch event schema, and the unix-socket HTTP transport in M1.
|
||||
//! Implements the allow-all [`brush_core::gate::Gate`], the shell-batch event
|
||||
//! model, near-live batching (200 events / 500 ms / exit), and the transport
|
||||
//! chain: HTTP over a unix socket (`NUCLEIC_SHELL_SOCKET`), then HTTP over TCP
|
||||
//! (`NUCLEIC_SHELL_HOOK_URL`), then a JSONL disk spool (`NUCLEIC_SHELL_SPOOL`).
|
||||
//!
|
||||
//! Failure doctrine (docs/NASH.md §5.5): nothing in this crate may affect a
|
||||
//! command's outcome. Hook bodies are panic-caught, the event channel drops
|
||||
//! (with a `dropped` counter) rather than blocks, and transport is
|
||||
//! fire-and-forget with short timeouts and spool fallback.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::io::{Read, Write};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
||||
use std::sync::mpsc;
|
||||
use std::sync::{Mutex, OnceLock};
|
||||
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
||||
|
||||
use serde::Serialize;
|
||||
|
||||
const FLUSH_MAX_EVENTS: usize = 200;
|
||||
const FLUSH_MAX_AGE: Duration = Duration::from_millis(500);
|
||||
const CHANNEL_BOUND: usize = 4096;
|
||||
const IO_TIMEOUT: Duration = Duration::from_millis(1500);
|
||||
const EXIT_FLUSH_TIMEOUT: Duration = Duration::from_millis(2000);
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Event model (docs/NASH.md §6.1)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
#[derive(Serialize)]
|
||||
#[serde(tag = "kind", rename_all = "camelCase")]
|
||||
enum Event {
|
||||
#[serde(rename = "exec", rename_all = "camelCase")]
|
||||
Exec {
|
||||
seq: u64,
|
||||
ts: u64,
|
||||
argv: Vec<String>,
|
||||
cwd: String,
|
||||
exit_code: i32,
|
||||
duration_ms: u64,
|
||||
},
|
||||
#[serde(rename = "cd", rename_all = "camelCase")]
|
||||
Cd { seq: u64, ts: u64, from: String, to: String },
|
||||
#[serde(rename = "export", rename_all = "camelCase")]
|
||||
Export {
|
||||
seq: u64,
|
||||
ts: u64,
|
||||
op: String,
|
||||
names: Vec<String>,
|
||||
},
|
||||
#[serde(rename = "fallback", rename_all = "camelCase")]
|
||||
Fallback {
|
||||
seq: u64,
|
||||
ts: u64,
|
||||
reason: String,
|
||||
input: String,
|
||||
},
|
||||
#[serde(rename = "dropped", rename_all = "camelCase")]
|
||||
Dropped { seq: u64, ts: u64, count: u64 },
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
struct Batch<'a> {
|
||||
#[serde(rename = "type")]
|
||||
batch_type: &'static str,
|
||||
source: &'static str,
|
||||
session_id: Option<&'a str>,
|
||||
shell_id: &'a str,
|
||||
parent_shell: Option<&'a str>,
|
||||
events: &'a [Event],
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Config
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
struct Config {
|
||||
socket: Option<PathBuf>,
|
||||
hook_url: Option<HookUrl>,
|
||||
token: Option<String>,
|
||||
session_id: Option<String>,
|
||||
spool: Option<PathBuf>,
|
||||
}
|
||||
|
||||
struct HookUrl {
|
||||
host: String,
|
||||
port: u16,
|
||||
path: String,
|
||||
}
|
||||
|
||||
fn parse_hook_url(url: &str) -> Option<HookUrl> {
|
||||
let rest = url.strip_prefix("http://")?;
|
||||
let (authority, path) = match rest.find('/') {
|
||||
Some(idx) => (&rest[..idx], rest[idx..].to_string()),
|
||||
None => (rest, "/shell-event".to_string()),
|
||||
};
|
||||
let (host, port) = match authority.rsplit_once(':') {
|
||||
Some((h, p)) => (h.to_string(), p.parse().ok()?),
|
||||
None => (authority.to_string(), 80),
|
||||
};
|
||||
Some(HookUrl { host, port, path })
|
||||
}
|
||||
|
||||
impl Config {
|
||||
fn from_env() -> Option<Self> {
|
||||
if std::env::var("NUCLEIC_SHELL_CAPTURE").as_deref() == Ok("off") {
|
||||
return None;
|
||||
}
|
||||
let socket = std::env::var_os("NUCLEIC_SHELL_SOCKET").map(PathBuf::from);
|
||||
let hook_url = std::env::var("NUCLEIC_SHELL_HOOK_URL")
|
||||
.ok()
|
||||
.and_then(|u| parse_hook_url(&u));
|
||||
let spool = std::env::var_os("NUCLEIC_SHELL_SPOOL").map(PathBuf::from);
|
||||
if socket.is_none() && hook_url.is_none() && spool.is_none() {
|
||||
return None;
|
||||
}
|
||||
Some(Self {
|
||||
socket,
|
||||
hook_url,
|
||||
token: std::env::var("NUCLEIC_HOOK_TOKEN").ok(),
|
||||
session_id: std::env::var("NUCLEIC_SESSION_ID").ok(),
|
||||
spool,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Observer state
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
enum Msg {
|
||||
Event(Event),
|
||||
FlushSync(mpsc::SyncSender<()>),
|
||||
}
|
||||
|
||||
struct PendingExec {
|
||||
argv: Vec<String>,
|
||||
cwd: PathBuf,
|
||||
started: Instant,
|
||||
}
|
||||
|
||||
struct Observer {
|
||||
tx: mpsc::SyncSender<Msg>,
|
||||
next_token: AtomicU64,
|
||||
next_seq: AtomicU64,
|
||||
dropped: AtomicU64,
|
||||
pending: Mutex<HashMap<u64, PendingExec>>,
|
||||
}
|
||||
|
||||
static OBSERVER: OnceLock<Observer> = OnceLock::new();
|
||||
|
||||
impl Observer {
|
||||
fn seq(&self) -> u64 {
|
||||
self.next_seq.fetch_add(1, Ordering::Relaxed)
|
||||
}
|
||||
|
||||
fn send(&self, event: Event) {
|
||||
match self.tx.try_send(Msg::Event(event)) {
|
||||
Ok(()) => {
|
||||
// Piggyback a dropped-counter event whenever capacity returns.
|
||||
let dropped = self.dropped.swap(0, Ordering::Relaxed);
|
||||
if dropped > 0 {
|
||||
let _ = self.tx.try_send(Msg::Event(Event::Dropped {
|
||||
seq: self.seq(),
|
||||
ts: now_millis(),
|
||||
count: dropped,
|
||||
}));
|
||||
}
|
||||
}
|
||||
Err(_) => {
|
||||
self.dropped.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn now_millis() -> u64 {
|
||||
SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.map(|d| u64::try_from(d.as_millis()).unwrap_or(0))
|
||||
.unwrap_or(0)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Gate implementation (allow-all; docs/NASH.md §4.2)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
struct RecordingGate;
|
||||
|
||||
impl brush_core::gate::Gate for RecordingGate {
|
||||
fn on_exec(&self, ev: brush_core::gate::ExecStart) -> brush_core::gate::Verdict {
|
||||
let token = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
|
||||
let obs = OBSERVER.get()?;
|
||||
let token = obs.next_token.fetch_add(1, Ordering::Relaxed);
|
||||
if let Ok(mut pending) = obs.pending.lock() {
|
||||
pending.insert(
|
||||
token,
|
||||
PendingExec {
|
||||
argv: ev.argv,
|
||||
cwd: ev.cwd,
|
||||
started: Instant::now(),
|
||||
},
|
||||
);
|
||||
}
|
||||
Some(token)
|
||||
}))
|
||||
.ok()
|
||||
.flatten()
|
||||
.unwrap_or(0);
|
||||
brush_core::gate::Verdict::Allow(token)
|
||||
}
|
||||
|
||||
fn on_exit(&self, token: u64, ev: brush_core::gate::ExecEnd) {
|
||||
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
|
||||
let Some(obs) = OBSERVER.get() else { return };
|
||||
let Some(pending) = obs.pending.lock().ok().and_then(|mut p| p.remove(&token))
|
||||
else {
|
||||
return;
|
||||
};
|
||||
let ts = now_millis();
|
||||
let cwd_str = pending.cwd.to_string_lossy().into_owned();
|
||||
|
||||
// Derived, names-only state-change events (docs/NASH.md §5.3–5.4).
|
||||
match pending.argv.first().map(String::as_str) {
|
||||
Some("cd") if ev.exit_code == 0 => {
|
||||
let to = pending
|
||||
.argv
|
||||
.get(1)
|
||||
.map_or_else(|| "~".to_string(), Clone::clone);
|
||||
obs.send(Event::Cd {
|
||||
seq: obs.seq(),
|
||||
ts,
|
||||
from: cwd_str.clone(),
|
||||
to,
|
||||
});
|
||||
}
|
||||
Some(op @ ("export" | "declare" | "typeset" | "unset" | "readonly" | "local"))
|
||||
if ev.exit_code == 0 =>
|
||||
{
|
||||
let names: Vec<String> = pending
|
||||
.argv
|
||||
.iter()
|
||||
.skip(1)
|
||||
.filter(|a| !a.starts_with('-'))
|
||||
.map(|a| a.split('=').next().unwrap_or(a).to_string())
|
||||
.collect();
|
||||
if !names.is_empty() {
|
||||
obs.send(Event::Export {
|
||||
seq: obs.seq(),
|
||||
ts,
|
||||
op: op.to_string(),
|
||||
names,
|
||||
});
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
|
||||
obs.send(Event::Exec {
|
||||
seq: obs.seq(),
|
||||
ts,
|
||||
argv: pending.argv,
|
||||
cwd: cwd_str,
|
||||
exit_code: ev.exit_code,
|
||||
duration_ms: u64::try_from(pending.started.elapsed().as_millis()).unwrap_or(0),
|
||||
});
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Transport
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Set once the host rejects our token; we stop posting and spool instead.
|
||||
static UNAUTHORIZED: AtomicBool = AtomicBool::new(false);
|
||||
|
||||
fn http_request(path: &str, token: Option<&str>, body: &[u8]) -> Vec<u8> {
|
||||
let mut req = format!(
|
||||
"POST {path} HTTP/1.1\r\nHost: nucleic\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n",
|
||||
body.len()
|
||||
);
|
||||
if let Some(token) = token {
|
||||
req.push_str(&format!("Authorization: Bearer {token}\r\n"));
|
||||
}
|
||||
req.push_str("\r\n");
|
||||
let mut bytes = req.into_bytes();
|
||||
bytes.extend_from_slice(body);
|
||||
bytes
|
||||
}
|
||||
|
||||
fn read_status<R: Read>(mut stream: R) -> Option<u16> {
|
||||
let mut buf = [0u8; 64];
|
||||
let n = stream.read(&mut buf).ok()?;
|
||||
let text = String::from_utf8_lossy(&buf[..n]);
|
||||
let status = text.split_whitespace().nth(1)?;
|
||||
status.parse().ok()
|
||||
}
|
||||
|
||||
fn check_status(status: Option<u16>) -> Result<(), ()> {
|
||||
match status {
|
||||
Some(code) if (200..300).contains(&code) => Ok(()),
|
||||
Some(401) => {
|
||||
UNAUTHORIZED.store(true, Ordering::Relaxed);
|
||||
Err(())
|
||||
}
|
||||
// Treat unreadable/other responses as delivered-at-best-effort: the
|
||||
// host 202s known tokens, so anything else is not retryable.
|
||||
Some(_) | None => Err(()),
|
||||
}
|
||||
}
|
||||
|
||||
fn post_unix(socket: &Path, token: Option<&str>, body: &[u8]) -> Result<(), ()> {
|
||||
let stream = std::os::unix::net::UnixStream::connect(socket).map_err(|_| ())?;
|
||||
let _ = stream.set_write_timeout(Some(IO_TIMEOUT));
|
||||
let _ = stream.set_read_timeout(Some(IO_TIMEOUT));
|
||||
let mut stream = stream;
|
||||
stream
|
||||
.write_all(&http_request("/shell-event", token, body))
|
||||
.map_err(|_| ())?;
|
||||
check_status(read_status(&mut stream))
|
||||
}
|
||||
|
||||
fn post_tcp(url: &HookUrl, token: Option<&str>, body: &[u8]) -> Result<(), ()> {
|
||||
let addr = (url.host.as_str(), url.port);
|
||||
let stream = std::net::TcpStream::connect_timeout(
|
||||
&std::net::ToSocketAddrs::to_socket_addrs(&addr)
|
||||
.ok()
|
||||
.and_then(|mut a| a.next())
|
||||
.ok_or(())?,
|
||||
IO_TIMEOUT,
|
||||
)
|
||||
.map_err(|_| ())?;
|
||||
let _ = stream.set_write_timeout(Some(IO_TIMEOUT));
|
||||
let _ = stream.set_read_timeout(Some(IO_TIMEOUT));
|
||||
let mut stream = stream;
|
||||
stream
|
||||
.write_all(&http_request(&url.path, token, body))
|
||||
.map_err(|_| ())?;
|
||||
check_status(read_status(&mut stream))
|
||||
}
|
||||
|
||||
fn spool(dir: &Path, shell_id: &str, body: &[u8]) {
|
||||
let _ = std::fs::create_dir_all(dir);
|
||||
let path = dir.join(format!("{shell_id}.jsonl"));
|
||||
if let Ok(mut f) = std::fs::OpenOptions::new().create(true).append(true).open(path) {
|
||||
let _ = f.write_all(body);
|
||||
let _ = f.write_all(b"\n");
|
||||
}
|
||||
}
|
||||
|
||||
fn deliver(cfg: &Config, shell_id: &str, body: &[u8]) {
|
||||
let token = cfg.token.as_deref();
|
||||
if !UNAUTHORIZED.load(Ordering::Relaxed) {
|
||||
if let Some(socket) = &cfg.socket {
|
||||
if post_unix(socket, token, body).is_ok() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
if let Some(url) = &cfg.hook_url {
|
||||
if post_tcp(url, token, body).is_ok() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Some(dir) = &cfg.spool {
|
||||
spool(dir, shell_id, body);
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Flusher thread
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn flusher(cfg: Config, shell_id: String, parent_shell: Option<String>, rx: mpsc::Receiver<Msg>) {
|
||||
let mut buf: Vec<Event> = Vec::new();
|
||||
let mut first_at: Option<Instant> = None;
|
||||
|
||||
let flush = |buf: &mut Vec<Event>| {
|
||||
if buf.is_empty() {
|
||||
return;
|
||||
}
|
||||
let batch = Batch {
|
||||
batch_type: "shell-batch",
|
||||
source: "nash",
|
||||
session_id: cfg.session_id.as_deref(),
|
||||
shell_id: &shell_id,
|
||||
parent_shell: parent_shell.as_deref(),
|
||||
events: buf,
|
||||
};
|
||||
if let Ok(body) = serde_json::to_vec(&batch) {
|
||||
deliver(&cfg, &shell_id, &body);
|
||||
}
|
||||
buf.clear();
|
||||
};
|
||||
|
||||
loop {
|
||||
let timeout = match first_at {
|
||||
Some(t) => FLUSH_MAX_AGE.saturating_sub(t.elapsed()),
|
||||
None => Duration::from_secs(3600),
|
||||
};
|
||||
match rx.recv_timeout(timeout) {
|
||||
Ok(Msg::Event(ev)) => {
|
||||
if buf.is_empty() {
|
||||
first_at = Some(Instant::now());
|
||||
}
|
||||
buf.push(ev);
|
||||
if buf.len() >= FLUSH_MAX_EVENTS {
|
||||
flush(&mut buf);
|
||||
first_at = None;
|
||||
}
|
||||
}
|
||||
Ok(Msg::FlushSync(ack)) => {
|
||||
flush(&mut buf);
|
||||
first_at = None;
|
||||
let _ = ack.send(());
|
||||
}
|
||||
Err(mpsc::RecvTimeoutError::Timeout) => {
|
||||
flush(&mut buf);
|
||||
first_at = None;
|
||||
}
|
||||
Err(mpsc::RecvTimeoutError::Disconnected) => {
|
||||
flush(&mut buf);
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Public API
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Reads `NUCLEIC_SHELL_*` config from the environment; when observation is
|
||||
/// configured, installs the recording gate into brush-core and starts the
|
||||
/// flusher thread. Returns whether observation is active.
|
||||
///
|
||||
/// Also threads shell lineage: reads `NUCLEIC_SHELL_PARENT` as this shell's
|
||||
/// parent and re-exports it as this shell's own id for child shells.
|
||||
pub fn install_from_env() -> bool {
|
||||
let Some(cfg) = Config::from_env() else {
|
||||
return false;
|
||||
};
|
||||
|
||||
let shell_id = format!("{:x}-{:x}", std::process::id(), now_millis());
|
||||
let parent_shell = std::env::var("NUCLEIC_SHELL_PARENT").ok();
|
||||
// SAFETY: called from `main` before any other threads exist.
|
||||
unsafe { std::env::set_var("NUCLEIC_SHELL_PARENT", &shell_id) };
|
||||
|
||||
let (tx, rx) = mpsc::sync_channel(CHANNEL_BOUND);
|
||||
let observer = Observer {
|
||||
tx,
|
||||
next_token: AtomicU64::new(1),
|
||||
next_seq: AtomicU64::new(1),
|
||||
dropped: AtomicU64::new(0),
|
||||
pending: Mutex::new(HashMap::new()),
|
||||
};
|
||||
if OBSERVER.set(observer).is_err() {
|
||||
return false;
|
||||
}
|
||||
|
||||
std::thread::Builder::new()
|
||||
.name("nash-observe-flush".into())
|
||||
.spawn(move || flusher(cfg, shell_id, parent_shell, rx))
|
||||
.ok();
|
||||
|
||||
brush_core::gate::set_gate(Box::new(RecordingGate));
|
||||
true
|
||||
}
|
||||
|
||||
/// Flushes any buffered events, waiting up to a short deadline. Intended for
|
||||
/// the process-exit hook; a no-op when observation is not installed.
|
||||
pub fn flush_sync() {
|
||||
let Some(obs) = OBSERVER.get() else { return };
|
||||
let (ack_tx, ack_rx) = mpsc::sync_channel(1);
|
||||
if obs.tx.try_send(Msg::FlushSync(ack_tx)).is_ok() {
|
||||
let _ = ack_rx.recv_timeout(EXIT_FLUSH_TIMEOUT);
|
||||
}
|
||||
}
|
||||
|
||||
/// Records a compat-fallback event (docs/NASH.md §4.1): nash is about to
|
||||
/// re-exec the real bash because it could not handle `input`.
|
||||
pub fn report_fallback(reason: &str, input: &str) {
|
||||
let Some(obs) = OBSERVER.get() else { return };
|
||||
let truncated: String = input.chars().take(4096).collect();
|
||||
obs.send(Event::Fallback {
|
||||
seq: obs.seq(),
|
||||
ts: now_millis(),
|
||||
reason: reason.to_string(),
|
||||
input: truncated,
|
||||
});
|
||||
flush_sync();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user