diff --git a/.opengrep/agentx-ifstack-rules.yaml b/.opengrep/agentx-ifstack-rules.yaml index 3de9baf..c4788c7 100644 --- a/.opengrep/agentx-ifstack-rules.yaml +++ b/.opengrep/agentx-ifstack-rules.yaml @@ -57,16 +57,23 @@ rules: # A relationship built inside the equality branch is never a rejection, # even when a terminating statement follows it. - patterns: - - pattern: | - StackRelationship { higher: $HIGHER, lower: $LOWER } + # Rust accepts the fields in either order, so the rule matches both. + - pattern-either: + - pattern: | + StackRelationship { higher: $HIGHER, lower: $LOWER } + - pattern: | + StackRelationship { lower: $LOWER, higher: $HIGHER } - pattern-either: - pattern-inside: | if $HIGHER == $LOWER { ... } - pattern-inside: | if $LOWER == $HIGHER { ... } - patterns: - - pattern: | - StackRelationship { higher: $HIGHER, lower: $LOWER } + - pattern-either: + - pattern: | + StackRelationship { higher: $HIGHER, lower: $LOWER } + - pattern: | + StackRelationship { lower: $LOWER, higher: $HIGHER } # The guard may sit in an enclosing block, so the trailing ellipsis spans # nested blocks. Anchor the terminator as the equality branch's last # statement. The leading ellipsis then covers any direct observations but diff --git a/.opengrep/tests/agentx-stack-relationship-without-self-guard.rs b/.opengrep/tests/agentx-stack-relationship-without-self-guard.rs index addbb42..e564939 100644 --- a/.opengrep/tests/agentx-stack-relationship-without-self-guard.rs +++ b/.opengrep/tests/agentx-stack-relationship-without-self-guard.rs @@ -378,3 +378,58 @@ fn guard_misplaced_inside_the_filter( } Ok(()) } + +// Rust accepts the fields in either order, so the rule must see both. A reversed +// construction without a guard is the same defect as the forward one. +fn unguarded_with_reversed_field_order( + peer: u32, + link: &Link, + relationships: &mut BTreeSet, +) { + // ruleid: agentx-stack-relationship-without-self-guard + relationships.insert(StackRelationship { + lower: peer, + higher: link.index, + }); +} + +fn unguarded_with_reversed_shorthand_fields( + lower: u32, + higher: u32, + relationships: &mut BTreeSet, +) { + // ruleid: agentx-stack-relationship-without-self-guard + relationships.insert(StackRelationship { lower, higher }); +} + +fn reversed_field_order_inside_the_equality_branch( + peer: u32, + link: &Link, + relationships: &mut BTreeSet, +) -> Result<()> { + if link.index == peer { + // ruleid: agentx-stack-relationship-without-self-guard + relationships.insert(StackRelationship { + lower: peer, + higher: link.index, + }); + return Err(invalid("self")); + } + Ok(()) +} + +fn guarded_with_reversed_field_order( + peer: u32, + link: &Link, + relationships: &mut BTreeSet, +) -> Result<()> { + if link.index == peer { + return Err(invalid("self")); + } + // ok: agentx-stack-relationship-without-self-guard + relationships.insert(StackRelationship { + lower: peer, + higher: link.index, + }); + Ok(()) +} diff --git a/README.md b/README.md index 841aabf..573774e 100644 --- a/README.md +++ b/README.md @@ -48,7 +48,12 @@ The packaged unit runs as root by default because `/var/agentx` is normally root-owned with mode 0700. It has no capabilities and writes no files. Its systemd sandbox permits Unix and netlink sockets in the host network namespace. The unit orders itself after `snmpd.service` without pulling that service in. -A missing master causes connection retries, not startup failure. +A missing master causes connection retries. The unit is `Type=notify` and reports +`READY=1` only once it registers the table, so `systemctl start` reports success +when the AgentX registration completes, and fails at `TimeoutStartSec` when no +master answers the retries. Registration does not mean the table serves rows yet. +systemd then restarts the unit, and the start limit stops it after three attempts +rather than retrying for ever. To run the subagent without root, create a system group and user: diff --git a/packaging/agentx-ifstack.service b/packaging/agentx-ifstack.service index 69bcec4..a27c2e4 100644 --- a/packaging/agentx-ifstack.service +++ b/packaging/agentx-ifstack.service @@ -1,11 +1,18 @@ [Unit] Description=AgentX ifStackTable subagent After=snmpd.service -StartLimitIntervalSec=60s +# Wide enough to hold the three start attempts it counts. A window shorter +# than an attempt resets between them, and the limit never stops the retries. +StartLimitIntervalSec=10min StartLimitBurst=3 [Service] -Type=simple +# The subagent reports READY=1 once it registers the table, so a start that +# never registers fails instead of reporting a healthy service that serves +# nothing. The timeout allows for the reconnect backoff. +Type=notify +NotifyAccess=main +TimeoutStartSec=60s User=root ExecStart=/usr/bin/agentx-ifstack Restart=on-failure diff --git a/packaging/test_policy.py b/packaging/test_policy.py index 2327b44..90be618 100644 --- a/packaging/test_policy.py +++ b/packaging/test_policy.py @@ -224,6 +224,17 @@ def documented_shell_commands(path): yield from shell_commands(span) +def seconds(span): + """Read a systemd time span. Bare numbers are seconds.""" + units = {"us": 0.000001, "ms": 0.001, "s": 1, "min": 60, "h": 3600, "d": 86400} + total = 0.0 + for value, unit in re.findall(r"(\d+(?:\.\d+)?)\s*([a-z]*)", span): + if not value: + continue + total += float(value) * units[unit or "s"] + return total + + def runs_real_namespace_suite(step): """Return whether a workflow step runs the suite with a compatible iproute2.""" return list(shell_commands(str(step.get("run", "")))) == REAL_NAMESPACE_COMMANDS @@ -578,9 +589,45 @@ def test_permanent_errors_do_not_restart_and_crashes_are_limited(self): self.assertEqual(unit["Service"]["Restart"], "on-failure") self.assertEqual(unit["Service"]["RestartPreventExitStatus"], "1") self.assertEqual(unit["Service"]["RestartSec"], "5s") - self.assertEqual(unit["Unit"]["StartLimitIntervalSec"], "60s") + self.assertEqual(unit["Unit"]["StartLimitIntervalSec"], "10min") self.assertEqual(unit["Unit"].getint("StartLimitBurst"), 3) + def test_a_start_that_never_registers_stops_retrying(self): + """Attempts spaced wider than the limit window reset it, so it never bounds them.""" + unit = configparser.ConfigParser(interpolation=None) + unit.read(ROOT / "packaging/agentx-ifstack.service") + attempt = seconds(unit["Service"]["TimeoutStartSec"]) + seconds( + unit["Service"]["RestartSec"] + ) + attempts = unit["Unit"].getint("StartLimitBurst") + self.assertGreaterEqual( + seconds(unit["Unit"]["StartLimitIntervalSec"]), + attempts * attempt, + "the start limit window must cover every attempt it counts", + ) + + def test_the_unit_reports_active_only_once_the_table_is_registered(self): + """A subagent that never registers serves nothing, and must not look healthy.""" + unit = configparser.ConfigParser(interpolation=None) + unit.read(ROOT / "packaging/agentx-ifstack.service") + self.assertEqual(unit["Service"]["Type"], "notify") + self.assertEqual(unit["Service"]["NotifyAccess"], "main") + self.assertEqual(unit["Service"]["TimeoutStartSec"], "60s") + + def test_readiness_documentation_describes_table_registration(self): + """READY=1 follows registration, which precedes every served read.""" + session = (ROOT / "src/session.rs").read_text() + registered = session.index("registered();") + self.assertLess( + session.index("acknowledge(&mut stream, ®ister.header)?"), registered + ) + self.assertLess(registered, session.index("loop {")) + self.assertNotIn("inventory", session[:registered]) + prose = " ".join((ROOT / "README.md").read_text().split()) + self.assertIn("`READY=1` only once it registers the table", prose) + self.assertNotIn("reports success when the subagent serves rows", prose) + self.assertIn("Registration does not mean the table serves rows", prose) + def test_third_party_actions_are_pinned_to_full_commit_shas(self): """A movable tag lets a compromised action change what CI and releases run.""" unpinned = [] @@ -1774,6 +1821,20 @@ def test_shell_commands_keep_quoted_operators_as_arguments(self): self.assertEqual(shell_options(commands[0][3:], {"--json", "--jq"}), (["1"], [("--json", "title,body"), ("--jq", ".[] | .title")])) + def test_readiness_guard_rejects_row_availability_wording(self): + path = ROOT / "README.md" + text = path.read_text() + start = text.index("A missing master causes connection retries.") + end = text.index("To run the subagent without root") + path.write_text(text[:start] + ( + "A missing master causes connection retries. The unit is `Type=notify` " + "and reports\n`READY=1` only once it registers the table, so " + "`systemctl start` reports success\nwhen the subagent serves rows, and " + "fails at `TimeoutStartSec` when no master\nanswers the retries.\n\n" + ) + text[end:]) + with self.assertRaises(AssertionError): + self.policy.test_readiness_documentation_describes_table_registration() + if __name__ == "__main__": unittest.main() diff --git a/src/main.rs b/src/main.rs index 37cfe53..6ea57cc 100644 --- a/src/main.rs +++ b/src/main.rs @@ -3,6 +3,7 @@ mod link; mod mib; mod monitor; mod netlink; +mod notify; mod session; use std::path::Path; @@ -54,7 +55,7 @@ fn main() -> ExitCode { let mut backoff = Duration::from_secs(1); loop { let started = Instant::now(); - match session::run(&config, &tables) { + match session::run(&config, &tables, ¬ify::ready) { Ok(()) => log::warn!("AgentX master closed the session"), Err(error) => log::warn!("AgentX session ended: {error}"), } diff --git a/src/notify.rs b/src/notify.rs new file mode 100644 index 0000000..7bf4136 --- /dev/null +++ b/src/notify.rs @@ -0,0 +1,75 @@ +use std::ffi::OsStr; +use std::io::{Error, ErrorKind, Result}; +use std::os::linux::net::SocketAddrExt; +use std::os::unix::ffi::OsStrExt; +use std::os::unix::net::{SocketAddr, UnixDatagram}; +use std::path::Path; + +const SOCKET: &str = "NOTIFY_SOCKET"; +const READY: &[u8] = b"READY=1"; + +/// Tell the service manager the subagent serves the table, so a Type=notify +/// unit reaches "active" on registration rather than on process start. An +/// absent socket means the process runs outside systemd, which is not an error. +pub fn ready() { + let Some(socket) = std::env::var_os(SOCKET) else { + return; + }; + // A lost notification costs the unit its start timeout. It must not cost + // the subagent a session it has already registered. + if let Err(error) = send(&socket, READY) { + log::warn!("Cannot report readiness to the service manager: {error}"); + } +} + +fn send(socket: &OsStr, message: &[u8]) -> Result<()> { + let datagram = UnixDatagram::unbound()?; + datagram.send_to_addr(message, &address(socket)?)?; + Ok(()) +} + +// systemd passes an abstract socket with a leading @, which is not a filesystem +// path: reading it as one addresses a socket that does not exist. +fn address(socket: &OsStr) -> Result { + match socket.as_bytes() { + [b'@', name @ ..] if !name.is_empty() => SocketAddr::from_abstract_name(name), + [b'/', ..] => SocketAddr::from_pathname(Path::new(socket)), + _ => Err(Error::new( + ErrorKind::InvalidInput, + format!("{SOCKET} is neither an absolute path nor an abstract name"), + )), + } +} + +#[cfg(test)] +mod tests { + use std::ffi::OsString; + + use super::*; + + #[test] + fn an_absolute_path_addresses_a_socket_file() { + let address = address(OsStr::new("/run/systemd/notify")).expect("path address"); + + assert_eq!( + address.as_pathname(), + Some(Path::new("/run/systemd/notify")) + ); + } + + #[test] + fn a_leading_at_sign_addresses_the_abstract_namespace() { + let address = address(OsStr::new("@systemd/notify")).expect("abstract address"); + + assert_eq!(address.as_abstract_name(), Some(&b"systemd/notify"[..])); + } + + #[test] + fn an_unusable_socket_value_is_an_error() { + for value in ["", "@", "run/systemd/notify"] { + let error = address(&OsString::from(value)).expect_err(value); + + assert_eq!(error.kind(), ErrorKind::InvalidInput); + } + } +} diff --git a/src/session.rs b/src/session.rs index 8443693..56822b7 100644 --- a/src/session.rs +++ b/src/session.rs @@ -16,7 +16,7 @@ const NOT_WRITABLE: u16 = 17; const COMMIT_FAILED: u16 = 14; const UNDO_FAILED: u16 = 15; -pub fn run(config: &Config, tables: &TableReader) -> Result<()> { +pub fn run(config: &Config, tables: &TableReader, registered: &dyn Fn()) -> Result<()> { log::info!("Connecting to AgentX master at {}", config.socket.display()); let mut stream = UnixStream::connect(&config.socket)?; stream.set_read_timeout(Some(IO_TIMEOUT))?; @@ -40,6 +40,7 @@ pub fn run(config: &Config, tables: &TableReader) -> Result<()> { "AgentX session {} registered ifStackTable", opened.session_id ); + registered(); loop { let (header, bytes) = receive(&mut stream, IO_TIMEOUT)?; diff --git a/tests/real_namespace.rs b/tests/real_namespace.rs index c8a0244..86759b5 100644 --- a/tests/real_namespace.rs +++ b/tests/real_namespace.rs @@ -14,11 +14,13 @@ use agentx::pdu::{self, Header, ResError, Response, Type}; use serde::Deserialize; use tempfile::TempDir; -mod support; +#[path = "support/master.rs"] +mod master; +#[path = "support/requests.rs"] +mod requests; -use support::agentx::{ - AgentxMaster, NETWORK_ORDER, STATUS, exchange, header, oid, range, read_frame, -}; +use master::{AgentxMaster, NETWORK_ORDER, read_frame}; +use requests::{STATUS, exchange, header, oid, range}; const TEST_UID: u32 = 65_534; const TEST_GID: u32 = 65_534; diff --git a/tests/support/agentx.rs b/tests/support/master.rs similarity index 74% rename from tests/support/agentx.rs rename to tests/support/master.rs index 3000288..0eae244 100644 --- a/tests/support/agentx.rs +++ b/tests/support/master.rs @@ -5,11 +5,10 @@ use std::path::Path; use std::str::FromStr; use std::time::{Duration, Instant}; -use agentx::encodings::{ID, SearchRange, SearchRangeList}; +use agentx::encodings::ID; use agentx::pdu::{self, Header, Response, Type}; pub const NETWORK_ORDER: u8 = 1 << pdu::NETWORK_BYTE_ORDER; -pub const STATUS: &str = "1.3.6.1.2.1.31.1.2.1.3"; pub struct AgentxMaster { listener: UnixListener, @@ -91,32 +90,3 @@ pub fn read_frame(stream: &mut UnixStream) -> Vec { .expect("read AgentX payload"); bytes } - -pub fn header(ty: Type, flags: u8, session_id: u32, packet_id: u32) -> Header { - let mut header = Header::new(ty); - header.flags = flags; - header.session_id = session_id; - header.transaction_id = 0x12345678; - header.packet_id = packet_id; - header -} - -pub fn exchange(stream: &mut UnixStream, bytes: &[u8]) -> Response { - let request = Header::from_bytes(bytes).expect("decode request header"); - stream.write_all(bytes).expect("send AgentX request"); - let response = Response::from_bytes(&read_frame(stream)).expect("decode AgentX response"); - assert_eq!(response.header.ty, Type::Response); - assert_eq!(response.header.flags, request.flags & NETWORK_ORDER); - assert_eq!(response.header.session_id, request.session_id); - assert_eq!(response.header.transaction_id, request.transaction_id); - assert_eq!(response.header.packet_id, request.packet_id); - response -} - -pub fn oid(suffix: &str) -> ID { - ID::from_str(&format!("{STATUS}{suffix}")).expect("valid test OID") -} - -pub fn range(name: ID) -> SearchRangeList { - SearchRangeList(vec![SearchRange::new(name, ID::default())]) -} diff --git a/tests/support/mod.rs b/tests/support/mod.rs deleted file mode 100644 index 38456a4..0000000 --- a/tests/support/mod.rs +++ /dev/null @@ -1 +0,0 @@ -pub mod agentx; diff --git a/tests/support/requests.rs b/tests/support/requests.rs new file mode 100644 index 0000000..91a05fb --- /dev/null +++ b/tests/support/requests.rs @@ -0,0 +1,39 @@ +use std::io::Write; +use std::os::unix::net::UnixStream; +use std::str::FromStr; + +use agentx::encodings::{ID, SearchRange, SearchRangeList}; +use agentx::pdu::{Header, Response, Type}; + +use crate::master::{NETWORK_ORDER, read_frame}; + +pub const STATUS: &str = "1.3.6.1.2.1.31.1.2.1.3"; + +pub fn header(ty: Type, flags: u8, session_id: u32, packet_id: u32) -> Header { + let mut header = Header::new(ty); + header.flags = flags; + header.session_id = session_id; + header.transaction_id = 0x12345678; + header.packet_id = packet_id; + header +} + +pub fn exchange(stream: &mut UnixStream, bytes: &[u8]) -> Response { + let request = Header::from_bytes(bytes).expect("decode request header"); + stream.write_all(bytes).expect("send AgentX request"); + let response = Response::from_bytes(&read_frame(stream)).expect("decode AgentX response"); + assert_eq!(response.header.ty, Type::Response); + assert_eq!(response.header.flags, request.flags & NETWORK_ORDER); + assert_eq!(response.header.session_id, request.session_id); + assert_eq!(response.header.transaction_id, request.transaction_id); + assert_eq!(response.header.packet_id, request.packet_id); + response +} + +pub fn oid(suffix: &str) -> ID { + ID::from_str(&format!("{STATUS}{suffix}")).expect("valid test OID") +} + +pub fn range(name: ID) -> SearchRangeList { + SearchRangeList(vec![SearchRange::new(name, ID::default())]) +} diff --git a/tests/systemd_readiness.rs b/tests/systemd_readiness.rs new file mode 100644 index 0000000..9c14c8b --- /dev/null +++ b/tests/systemd_readiness.rs @@ -0,0 +1,120 @@ +use std::io::ErrorKind; +use std::os::linux::net::SocketAddrExt; +use std::os::unix::net::{SocketAddr, UnixDatagram}; +use std::path::Path; +use std::process::{Child, Command, Stdio}; +use std::time::Duration; + +use tempfile::TempDir; + +#[path = "support/master.rs"] +mod master; + +use master::{AgentxMaster, NETWORK_ORDER}; + +// systemd holds a Type=notify unit in "activating" until READY=1 arrives, so +// this datagram is what makes "started" mean "registered" for the deployment. +const READY: &[u8] = b"READY=1"; + +struct Subagent { + child: Child, +} + +impl Drop for Subagent { + fn drop(&mut self) { + let _ = self.child.kill(); + let _ = self.child.wait(); + } +} + +fn spawn(directory: &Path, socket: &Path, notify: &str) -> Subagent { + let config = directory.join("config.toml"); + std::fs::write(&config, format!("socket = {socket:?}\n")).expect("write test configuration"); + let child = Command::new(env!("CARGO_BIN_EXE_agentx-ifstack")) + .arg("--config") + .arg(config) + .env("NOTIFY_SOCKET", notify) + .env_remove("RUST_LOG") + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .expect("start agentx-ifstack test process"); + Subagent { child } +} + +fn receiver(address: &SocketAddr) -> UnixDatagram { + let datagram = UnixDatagram::bind_addr(address).expect("bind notify socket"); + datagram + .set_read_timeout(Some(Duration::from_secs(10))) + .expect("bound the readiness wait"); + datagram +} + +fn read(datagram: &UnixDatagram) -> Vec { + let mut message = vec![0; 64]; + let read = datagram.recv(&mut message).expect("receive readiness"); + message.truncate(read); + message +} + +#[test] +fn the_subagent_reports_readiness_once_it_registers_the_table() { + let directory = TempDir::new().expect("create test directory"); + let socket = directory.path().join("master"); + let notify = directory.path().join("notify"); + let datagram = receiver(&SocketAddr::from_pathname(¬ify).expect("notify socket path")); + let master = AgentxMaster::bind(&socket); + let _subagent = spawn( + directory.path(), + &socket, + notify.to_str().expect("notify path is UTF-8"), + ); + + let _session = master.connect(NETWORK_ORDER, 100, 127); + + assert_eq!(read(&datagram), READY); +} + +// systemd passes an abstract socket as a leading @, which is not a filesystem +// path. Reading it as one silently loses every notification. +#[test] +fn an_abstract_notify_socket_receives_the_same_message() { + let directory = TempDir::new().expect("create test directory"); + let socket = directory.path().join("master"); + let name = format!("agentx-ifstack-test-{}", std::process::id()); + let datagram = + receiver(&SocketAddr::from_abstract_name(name.as_bytes()).expect("abstract notify socket")); + let master = AgentxMaster::bind(&socket); + let _subagent = spawn(directory.path(), &socket, &format!("@{name}")); + + let _session = master.connect(NETWORK_ORDER, 100, 127); + + assert_eq!(read(&datagram), READY); +} + +// A subagent that is up but has never registered serves nothing. Reporting +// readiness on start would hide exactly the failure the unit must catch. +#[test] +fn a_subagent_without_a_master_reports_nothing() { + let directory = TempDir::new().expect("create test directory"); + let socket = directory.path().join("master"); + let notify = directory.path().join("notify"); + let datagram = receiver(&SocketAddr::from_pathname(¬ify).expect("notify socket path")); + datagram + .set_read_timeout(Some(Duration::from_secs(3))) + .expect("bound the negative wait"); + let _subagent = spawn( + directory.path(), + &socket, + notify.to_str().expect("notify path is UTF-8"), + ); + + let error = datagram + .recv(&mut [0; 64]) + .expect_err("readiness without a registered session"); + + assert!( + matches!(error.kind(), ErrorKind::WouldBlock | ErrorKind::TimedOut), + "unexpected error: {error}" + ); +}