Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 11 additions & 4 deletions .opengrep/agentx-ifstack-rules.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
55 changes: 55 additions & 0 deletions .opengrep/tests/agentx-stack-relationship-without-self-guard.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<StackRelationship>,
) {
// 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<StackRelationship>,
) {
// 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<StackRelationship>,
) -> 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<StackRelationship>,
) -> 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(())
}
7 changes: 6 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
11 changes: 9 additions & 2 deletions packaging/agentx-ifstack.service
Original file line number Diff line number Diff line change
@@ -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
Expand Down
63 changes: 62 additions & 1 deletion packaging/test_policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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, &register.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 = []
Expand Down Expand Up @@ -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()
3 changes: 2 additions & 1 deletion src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ mod link;
mod mib;
mod monitor;
mod netlink;
mod notify;
mod session;

use std::path::Path;
Expand Down Expand Up @@ -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, &notify::ready) {
Ok(()) => log::warn!("AgentX master closed the session"),
Err(error) => log::warn!("AgentX session ended: {error}"),
}
Expand Down
75 changes: 75 additions & 0 deletions src/notify.rs
Original file line number Diff line number Diff line change
@@ -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<SocketAddr> {
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);
}
}
}
3 changes: 2 additions & 1 deletion src/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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))?;
Expand All @@ -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)?;
Expand Down
10 changes: 6 additions & 4 deletions tests/real_namespace.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
32 changes: 1 addition & 31 deletions tests/support/agentx.rs → tests/support/master.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -91,32 +90,3 @@ pub fn read_frame(stream: &mut UnixStream) -> Vec<u8> {
.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())])
}
1 change: 0 additions & 1 deletion tests/support/mod.rs

This file was deleted.

Loading
Loading