diff --git a/README.md b/README.md index 93abd26..aa79dfa 100644 --- a/README.md +++ b/README.md @@ -55,7 +55,13 @@ without requiring enterprise EDR tooling or kernel entitlements. or above a threshold you choose (off by default beyond critical-only). Notifications carry "Show in Argus" and "Allowlist…" action buttons — the allowlist action goes through the exact same Touch ID gate as the in-app path. - Authorization is requested once on first launch. + Authorization is requested once on first launch. Busy process trees (e.g. + an AI-agent session doing extensive automation) no longer produce unbounded + notification spam: per process-tree root, the first 3 sub-critical + notifications in a rolling 5-minute window deliver individually; after that + they collapse into one continuously-updated digest notification ("15 alerts, + techniques: privilege escalation, persistence… since 14:02"). Critical + alerts always deliver individually and never count against the budget. - **Event export** — right-click an event to "Copy as JSON"; from the History panel, export the full event history as JSON or CSV (RFC 4180-quoted so command lines with embedded commas and quotes don't misalign spreadsheets). @@ -120,6 +126,19 @@ sample it. Allowlist filtering deliberately doesn't apply to these events — a persistence artifact change is a different thing than an allowlisted process. +### Docker container activity + +Containers run inside Docker's Linux VM, invisible to the host's `ps`. A +`docker events` subscriber surfaces container lifecycle activity as synthetic +process events tagged with provenance "docker": container starts (info level, +T1610 Deploy Container) and execs into running containers (watch level, T1609 +Container Administration Command — a common post-compromise lateral step). +The watcher is resilient to the Docker daemon restarting (60-second retry +loop) and inert if Docker CLI is not installed. Note the scope limitation: +only container operations the daemon reports; processes *inside* containers +remain invisible to host-level detection — in-container detection is a +different tool's job. + ### Tamper evidence The app records HMAC-SHA256 MACs of `rules-state.json` and `allowlist.json` @@ -281,7 +300,7 @@ sidecar records HMAC-SHA256 MACs of the security-relevant JSON files swift test ``` -170 tests. The Sigma engine is validated two ways: `SigmaEngineTests` checks +210 tests. The Sigma engine is validated two ways: `SigmaEngineTests` checks the YAML parser, condition evaluator, and matcher against real rule text fetched from SigmaHQ (including the `N of selection_*` quantifier, `base64`/ `base64offset`/`cased` modifiers, and keyword matching), and `BundledRulesTests` @@ -294,8 +313,9 @@ covered: `RuleStoreTests` (enable/disable persistence, user-rule loading), `ChainCorrelatorTests` (sequence detection), `PersistenceWatcherTests` (artifact diffing and event generation), `IntegrityGuardTests` (MAC recording and verification), `ProvenanceClassifierTests` (ancestry classification), -and `AgentActivityPolicyTests` (agent-attributed event escalation and notification -quieting). +`AgentActivityPolicyTests` (agent-attributed event escalation and notification +quieting), `NotificationRollupTests` (digest throttling and event batching), +and `DockerWatcherTests` (container lifecycle event synthesis). ## Project layout @@ -317,6 +337,7 @@ Sources/Argus/ AgentActivityPolicy.swift escalation and notification quieting for agent-attributed events ChainCorrelator.swift correlates techniques in same process tree (10-min window) PersistenceWatcher.swift event-driven monitoring of LaunchAgents/Daemons/periodic + DockerWatcher.swift docker events subscriber for container lifecycle events IntegrityGuard.swift HMAC verification of rules-state.json and allowlist.json AllowlistStore.swift persisted (rule, executable, provenance scope) suppression EventStore.swift persisted event history (events.jsonl) @@ -325,6 +346,7 @@ Sources/Argus/ HistoryStats.swift day-bucketing + technique-frequency aggregation AppSettings.swift tunable poll interval, decay, notification threshold NotificationManager.swift UNUserNotificationCenter wrapper + NotificationRollup.swift per-tree-root notification throttling and digesting NotificationResponder.swift notification actions (Show in Argus, Allowlist) DiagnosticsLog.swift on-disk activity log Theme.swift color/type tokens @@ -348,6 +370,8 @@ Tests/ArgusTests/ AgentActivityPolicyTests.swift escalation rules and notification quieting logic ChainCorrelatorTests.swift sequence detection in process trees PersistenceWatcherTests.swift artifact diffing and event generation + DockerWatcherTests.swift container lifecycle event synthesis and resilience + NotificationRollupTests.swift per-tree throttling, digest batching, budget enforcement IntegrityGuardTests.swift MAC recording and verification Resources/ Info.plist app bundle metadata diff --git a/Sources/Argus/App.swift b/Sources/Argus/App.swift index 4ad259a..3fda0df 100644 --- a/Sources/Argus/App.swift +++ b/Sources/Argus/App.swift @@ -38,6 +38,10 @@ struct ArgusApp: App { /// alive — `PersistenceWatcher` isn't observed by any view, so nothing /// else in the view hierarchy retains it. private let persistenceWatcher: PersistenceWatcher + /// Held for the app's lifetime purely to keep its child `docker events` + /// process and retry timer alive — `DockerWatcher` isn't observed by any + /// view either, so nothing else in the view hierarchy retains it. + private let dockerWatcher: DockerWatcher /// Held for the app's lifetime purely to keep it alive — see its own /// doc comment. `UNUserNotificationCenter.delegate` is a weak reference, /// so nothing else retains this object. @@ -89,6 +93,17 @@ struct ArgusApp: App { } watcher.start() + // Third, independent sensor: containers run inside Docker's Linux VM, + // so ProcessMonitor's host `ps` poll never sees anything running + // inside one. This subscribes to the Docker daemon's own event + // stream to surface container start/exec activity instead — inert + // and silent (beyond one log line) on machines without Docker. + let dockerEventWatcher = DockerWatcher() + dockerEventWatcher.onEvent = { [weak m] event in + m?.ingestExternal(event) + } + dockerEventWatcher.start() + // Verify rules-state.json/allowlist.json against their last // authenticated-write MAC — asynchronously, off this init: the // Keychain key fetch can present a consent prompt (seen in practice @@ -119,6 +134,7 @@ struct ArgusApp: App { _ruleStore = StateObject(wrappedValue: rules) eventStore = events persistenceWatcher = watcher + dockerWatcher = dockerEventWatcher notificationResponder = responder } diff --git a/Sources/Argus/DockerWatcher.swift b/Sources/Argus/DockerWatcher.swift new file mode 100644 index 0000000..cabb3a6 --- /dev/null +++ b/Sources/Argus/DockerWatcher.swift @@ -0,0 +1,328 @@ +import Foundation + +/// One `docker events --format '{{json .}}'` line, reduced to the fields the +/// classifier needs. Field names match the Docker CLI's JSON event schema +/// exactly (capitalized `Type`/`Action`/`Actor`, nested `Attributes`) so +/// decoding requires no remapping. +struct DockerEventLine: Decodable { + struct Actor: Decodable { + struct Attributes: Decodable { + let image: String? + let name: String? + } + let Attributes: Attributes? + } + + let `Type`: String + let Action: String + let Actor: Actor? + let time: Int? +} + +/// The container-lifecycle actions Argus reports on. Deliberately narrow: +/// `die`/`stop`/`destroy` and friends are lifecycle noise with no detection +/// value here — a container exiting isn't itself suspicious, and reporting +/// every stop/destroy would just be churn in the feed. `start` marks a new +/// container coming up; `exec` marks something being run inside one after +/// the fact, which is the more interesting signal (see `DockerActionKind`'s +/// technique mapping below). +enum DockerActionKind: Equatable { + case start + case execCreate + case execStart +} + +/// Pure classification of one already-decoded `DockerEventLine` into an +/// action Argus cares about, or nil for everything else (non-container +/// types, and container actions outside the lifecycle list above). No +/// Process/filesystem access — this is exercised directly in tests against +/// literal `DockerEventLine` values. +enum DockerEventClassifier { + /// `exec_create`/`exec_start` actions arrive from `docker events` as + /// `"exec_create: "` / `"exec_start: "` — the command + /// is appended after a colon-space rather than living in its own field. + /// Splitting on the *first* ": " keeps the rest of a command that itself + /// contains ": " intact. + static func splitAction(_ action: String) -> (verb: String, detail: String?) { + guard let range = action.range(of: ": ") else { return (action, nil) } + let verb = String(action[action.startIndex.. (kind: DockerActionKind, execCommand: String?)? { + guard line.Type == "container" else { return nil } + let (verb, detail) = splitAction(line.Action) + switch verb { + case "start": + return (.start, nil) + case "exec_create": + return (.execCreate, detail) + case "exec_start": + return (.execStart, detail) + default: + return nil + } + } +} + +/// Builds the synthetic `ProcessEvent` for one classified Docker action. +/// Pure (no Process/filesystem access), mirroring `PersistenceEventBuilder`, +/// so the severity/technique/explanation mapping is directly testable. +/// +/// These events have no real pid — they represent something the Docker +/// daemon reported, not a process Argus itself observed — so `pid`/`ppid` +/// are 0 and `executable`/`command` describe the container instead. +enum DockerEventBuilder { + /// `exec_create` and `exec_start` are two halves of the same exec — the + /// daemon fires both for one `docker exec` invocation. Emitting on both + /// would double-count every exec, so only `exec_start` (the point the + /// command actually ran) produces an event; `exec_create` is classified + /// but intentionally dropped here. + static func makeEvent(kind: DockerActionKind, execCommand: String?, image: String?, name: String?, timestamp: Date) -> ProcessEvent? { + let containerLabel = name ?? image ?? "unknown" + let imageLabel = image ?? "unknown" + + switch kind { + case .start: + let command = "docker start — image=\(imageLabel) name=\(containerLabel)" + let rule = MatchedRule( + name: "Docker container started", + severity: .info, + technique: "T1610", + explanation: "A new container (\(containerLabel), image \(imageLabel)) came up. Deploying a container is routine, but " + + "it's also a documented way to stand up throwaway compute an attacker controls — worth a glance at what image it's running." + ) + return ProcessEvent(pid: 0, ppid: 0, executable: containerLabel, command: command, rules: [rule], timestamp: timestamp, provenance: ["docker"]) + + case .execCreate: + return nil + + case .execStart: + let commandDetail = execCommand?.trimmingCharacters(in: .whitespaces) + let commandSummary = (commandDetail?.isEmpty == false) ? commandDetail! : "unknown command" + let command = "docker exec — \(commandSummary) in \(containerLabel)" + let rule = MatchedRule( + name: "Docker exec into running container", + severity: .watch, + technique: "T1609", + explanation: "Something was executed inside the running container \(containerLabel) (\(commandSummary)). " + + "Exec-ing into a live container is a common post-compromise/lateral-movement step — usually benign (debugging), " + + "but worth a glance." + ) + return ProcessEvent(pid: 0, ppid: 0, executable: containerLabel, command: command, rules: [rule], timestamp: timestamp, provenance: ["docker"]) + } + } +} + +/// Splits a stream of arbitrary byte chunks (as delivered by +/// `FileHandle.readabilityHandler`, which has no notion of line boundaries) +/// into complete lines, holding back any trailing partial line until more +/// data arrives. Pure and state-carrying via an inout buffer rather than a +/// class, so a chunk boundary landing mid-line is directly testable without +/// standing up a real pipe. +enum LineBuffer { + /// Appends `chunk` to `buffer`, returns every complete (newline-terminated) + /// line found, and leaves any trailing partial line in `buffer` for the + /// next call. + static func consume(_ chunk: String, buffer: inout String) -> [String] { + buffer += chunk + var lines: [String] = [] + while let newlineRange = buffer.range(of: "\n") { + let line = String(buffer[buffer.startIndex.. ProcessEvent? { + let trimmed = raw.trimmingCharacters(in: .whitespacesAndNewlines) + guard !trimmed.isEmpty, let data = trimmed.data(using: .utf8) else { return nil } + guard let line = try? decoder.decode(DockerEventLine.self, from: data) else { return nil } + guard let (kind, execCommand) = DockerEventClassifier.classify(line) else { return nil } + + let timestamp = line.time.map { Date(timeIntervalSince1970: TimeInterval($0)) } ?? Date() + return DockerEventBuilder.makeEvent( + kind: kind, + execCommand: execCommand, + image: line.Actor?.Attributes?.image, + name: line.Actor?.Attributes?.name, + timestamp: timestamp + ) + } +} + +/// Third, independent sensor alongside `ProcessMonitor` and +/// `PersistenceWatcher`: containers run inside Docker's Linux VM, so the +/// host `ps` table `ProcessMonitor` polls only ever sees the `docker` CLI +/// and the VM's own helper process — every process running *inside* a +/// container is invisible to it. This watcher subscribes to the Docker +/// daemon's own event stream instead, which the daemon reports regardless +/// of host-side process visibility. +/// +/// Honest scope: this surfaces only container lifecycle (start) and exec +/// (exec_start) events exactly as the Docker daemon reports them — it has no +/// visibility into what runs *inside* a container once it's up (a shell +/// spawned by that execced process, a payload downloaded and run within the +/// container's own PID namespace, etc.). In-container process monitoring is +/// a different tool's job (e.g. something running an agent inside the +/// container, or an EDR with container-aware sensors); this is the +/// container-boundary signal only. +final class DockerWatcher { + /// Invoked on the main actor for every detected Docker event. Set + /// before calling `start()`. + var onEvent: (@MainActor (ProcessEvent) -> Void)? + + private let candidatePaths = ["/usr/local/bin/docker", "/opt/homebrew/bin/docker", "/usr/bin/docker"] + private let retryInterval: TimeInterval + private let queue = DispatchQueue(label: "com.argus.dockerwatcher") + + private var process: Process? + private var lineBuffer = "" + private var retryWorkItem: DispatchWorkItem? + private var stopped = false + private var hasLoggedMissingCLI = false + private var hasLoggedExit = false + + init(retryInterval: TimeInterval = 60) { + self.retryInterval = retryInterval + } + + /// Resolves the docker CLI once at call time (Docker Desktop's install + /// location doesn't move while the app is running) and, if found, spawns + /// `docker events`. If no CLI is found anywhere in `candidatePaths`, logs + /// one line and stays permanently inert — there is nothing to retry when + /// the binary itself isn't installed. + func start() { + guard !stopped else { return } + guard let dockerPath = resolveDockerPath() else { + if !hasLoggedMissingCLI { + DiagnosticsLog.write("docker watcher — no docker CLI found, staying inert") + hasLoggedMissingCLI = true + } + return + } + launch(dockerPath: dockerPath) + } + + /// Terminates the child process and cancels any pending retry. Safe to + /// call multiple times. + func stop() { + stopped = true + retryWorkItem?.cancel() + retryWorkItem = nil + process?.terminationHandler = nil + if process?.isRunning == true { + process?.terminate() + } + process = nil + } + + deinit { + process?.terminationHandler = nil + if process?.isRunning == true { + process?.terminate() + } + } + + private func resolveDockerPath() -> String? { + candidatePaths.first { FileManager.default.isExecutableFile(atPath: $0) } + } + + /// Spawns `docker events --format '{{json .}}'` as a long-lived child and + /// wires its stdout through a `readabilityHandler` on `queue` so lines + /// are consumed incrementally rather than buffered to EOF (which would + /// never come for a subscription that runs forever). If the daemon isn't + /// running (or Docker Desktop later quits), the process exits quickly; + /// that's treated as a transition to retry from, not a fatal error — + /// Docker Desktop starting later must be picked up without restarting + /// Argus. + private func launch(dockerPath: String) { + let proc = Process() + proc.executableURL = URL(fileURLWithPath: dockerPath) + proc.arguments = ["events", "--format", "{{json .}}"] + let outPipe = Pipe() + proc.standardOutput = outPipe + proc.standardError = Pipe() + + lineBuffer = "" + + outPipe.fileHandleForReading.readabilityHandler = { [weak self] handle in + let data = handle.availableData + // Empty data means EOF. The pipe's fd stays "readable" at EOF, so + // without clearing the handler here it would keep firing in a + // tight spin until `handleExit` (driven separately by + // `terminationHandler`) gets around to clearing it. + guard !data.isEmpty else { + handle.readabilityHandler = nil + return + } + guard let chunk = String(data: data, encoding: .utf8) else { return } + self?.handle(chunk: chunk) + } + + proc.terminationHandler = { [weak self] _ in + self?.queue.async { + self?.handleExit(outPipe: outPipe) + } + } + + do { + try proc.run() + process = proc + hasLoggedExit = false + } catch { + DiagnosticsLog.write("docker watcher — failed to launch docker events: \(error)") + scheduleRetry() + } + } + + private func handle(chunk: String) { + queue.async { [weak self] in + guard let self else { return } + let lines = LineBuffer.consume(chunk, buffer: &self.lineBuffer) + for line in lines { + guard let event = DockerLineProcessor.processLine(line) else { continue } + Task { @MainActor in + self.onEvent?(event) + } + } + } + } + + /// Runs on `queue`. A missing daemon and a Desktop quit look identical + /// from here (the child just exits) — both are logged once and retried + /// after a fixed backoff rather than distinguished, since the recovery + /// action (retry later) is the same either way. + private func handleExit(outPipe: Pipe) { + outPipe.fileHandleForReading.readabilityHandler = nil + process = nil + guard !stopped else { return } + if !hasLoggedExit { + DiagnosticsLog.write("docker watcher — docker events exited (daemon not running?), retrying in \(Int(retryInterval))s") + hasLoggedExit = true + } + scheduleRetry() + } + + private func scheduleRetry() { + guard !stopped else { return } + let work = DispatchWorkItem { [weak self] in + self?.start() + } + retryWorkItem = work + queue.asyncAfter(deadline: .now() + retryInterval, execute: work) + } +} diff --git a/Sources/Argus/NotificationManager.swift b/Sources/Argus/NotificationManager.swift index fb7b9b2..585f2f5 100644 --- a/Sources/Argus/NotificationManager.swift +++ b/Sources/Argus/NotificationManager.swift @@ -54,4 +54,42 @@ enum NotificationManager { let request = UNNotificationRequest(identifier: event.id.uuidString, content: content, trigger: nil) UNUserNotificationCenter.current().add(request) } + + /// Stable per-root identifier for a rollup digest notification. Reusing + /// the same identifier on every call for a given root is what makes + /// `notifyDigest` replace the previous digest in place rather than + /// stacking a new banner per folded-in event — see that method's doc + /// comment. + static func rollupIdentifier(rootPID: Int32) -> String { "argus.rollup.\(rootPID)" } + + private static let digestTimeFormatter: DateFormatter = { + let f = DateFormatter() + f.dateFormat = "HH:mm" + return f + }() + + /// Posts, or on a later call for the same `rootPID`, replaces, the + /// digest notification summarizing a busy process tree's rolled-up + /// alerts. `UNUserNotificationCenter.add` replaces any pending/delivered + /// request that shares an identifier, so reusing + /// `rollupIdentifier(rootPID:)` across calls updates one notification in + /// place instead of stacking a new banner per event `NotificationRollup` + /// folds into the digest — the whole point of rolling up. + /// + /// Carries no `ruleNameKey`/`executableKey` userInfo — unlike a single + /// `notify(event:)` call, a digest spans many rules and processes, not + /// one, so there's no single rule its "Allowlist…" action could act on. + /// It still sets `categoryIdentifier`, so tapping the notification body + /// opens Argus via the same default-action handling `NotificationResponder` + /// already gives every other event notification. + static func notifyDigest(rootPID: Int32, count: Int, techniques: [String], since: Date) { + let content = UNMutableNotificationContent() + content.title = "Busy process tree (pid \(rootPID))" + content.body = "\(count) alerts, techniques: \(techniques.joined(separator: ", ")) since \(digestTimeFormatter.string(from: since))" + content.sound = .default + content.categoryIdentifier = categoryIdentifier + + let request = UNNotificationRequest(identifier: rollupIdentifier(rootPID: rootPID), content: content, trigger: nil) + UNUserNotificationCenter.current().add(request) + } } diff --git a/Sources/Argus/NotificationRollup.swift b/Sources/Argus/NotificationRollup.swift new file mode 100644 index 0000000..ff645c7 --- /dev/null +++ b/Sources/Argus/NotificationRollup.swift @@ -0,0 +1,103 @@ +import Foundation + +/// Collapses a burst of notifications from one process tree into a single, +/// continuously-updated digest, so a busy supervised session (e.g. an agent +/// tearing through many LOLBin-adjacent commands in minutes) can't fire more +/// than a handful of individual banners before the rest of that burst +/// folds into one digest — addressing notification fatigue without ever +/// hiding a genuinely critical alert. +/// +/// Pure and free of any `UserNotifications` import so the windowing/budget +/// logic is fully testable without a running app bundle: `ProcessMonitor` is +/// the only caller of `record`, and `NotificationManager.notifyDigest` is the +/// only place that turns a `.digest` decision into an actual OS notification. +/// +/// Keys strictly on the process-tree **root pid**, never on provenance label +/// (e.g. "claude"). Two independent agent sessions both attributed to the +/// same supervisor label are still two unrelated process trees — sharing one +/// digest between them would mix one session's alert count and techniques +/// into another's, which is exactly the kind of cross-session confusion this +/// package exists to avoid, not cause. Root pid is the one identifier that's +/// guaranteed distinct per tree for as long as that tree is active. +final class NotificationRollup { + /// The outcome of folding one matched event into this root's rollup state. + enum Decision: Equatable { + /// Deliver an individual notification as today. + case deliver + /// Fold into the root's running digest instead of delivering + /// individually. + /// + /// - `count`: total events this root has produced within the current + /// window, including the ones already delivered individually — + /// not just the events that themselves became digests. This is the + /// number a digest notification's body should show. + /// - `techniques`: the union of technique IDs seen this window, + /// deduplicated, in first-seen order. + /// - `since`: the window's start timestamp (the first event's + /// timestamp that opened this window), for the digest body's + /// "since HH:mm" text. + /// - `isFirstDigest`: true only for the single event that flips this + /// root from individual delivery into digesting. Callers should + /// write a diagnostics line exactly when this is true — every + /// later digest update for the same window is a silent replace. + case digest(count: Int, techniques: [String], since: Date, isFirstDigest: Bool) + } + + private struct RootState { + let windowStart: Date + var deliveredCount: Int = 0 + var totalCount: Int = 0 + var techniques: [String] = [] + var techniquesSeen: Set = [] + var hasDigested: Bool = false + } + + private let window: TimeInterval + private let budget: Int + private var state: [Int32: RootState] = [:] + + init(window: TimeInterval = 300, budget: Int = 3) { + self.window = window + self.budget = budget + } + + /// Records one matched event for `rootPID` at `timestamp` and returns + /// whether it should be delivered individually or folded into the + /// root's running digest. + /// + /// Critical-severity events always return `.deliver` and do not consume + /// the per-window budget. A rollup exists to cut UI noise, not to make a + /// genuinely critical alert reachable only by a user opening a digest + /// notification they might otherwise dismiss unread — so `.critical` + /// bypasses digesting entirely, even mid-window, even after the budget + /// for this root is already exhausted. + func record(rootPID: Int32, techniques: [String], severity: Severity, timestamp: Date) -> Decision { + // Drop any root's state once its window has fully elapsed relative + // to this event's timestamp — including `rootPID`'s own state, which + // is what makes window expiry the mechanism that resets a busy tree + // back to individual delivery: once the window is gone, the next + // event below starts a brand new window with a fresh budget. + state = state.filter { timestamp.timeIntervalSince($0.value.windowStart) < window } + + var root = state[rootPID] ?? RootState(windowStart: timestamp) + root.totalCount += 1 + for id in techniques where root.techniquesSeen.insert(id).inserted { + root.techniques.append(id) + } + + let decision: Decision + if severity == .critical { + decision = .deliver + } else if root.deliveredCount < budget { + root.deliveredCount += 1 + decision = .deliver + } else { + let isFirstDigest = !root.hasDigested + root.hasDigested = true + decision = .digest(count: root.totalCount, techniques: root.techniques, since: root.windowStart, isFirstDigest: isFirstDigest) + } + + state[rootPID] = root + return decision + } +} diff --git a/Sources/Argus/ProcessMonitor.swift b/Sources/Argus/ProcessMonitor.swift index d8cc182..14810da 100644 --- a/Sources/Argus/ProcessMonitor.swift +++ b/Sources/Argus/ProcessMonitor.swift @@ -181,6 +181,12 @@ final class ProcessMonitor: ObservableObject { private var tickIndex = 0 private var samplingHealth = SamplingHealthTracker() private let chainCorrelator = ChainCorrelator() + /// Collapses a burst of individually-noteworthy notifications from one + /// busy process tree into a single, continuously-updated digest — see + /// `NotificationRollup`'s doc comment. Consulted only for events that + /// already passed the threshold/agent-quieting checks below; a quieted + /// event never reaches it. + private let notificationRollup = NotificationRollup() func configure(allowlist: AllowlistStore) { self.allowlist = allowlist @@ -366,7 +372,29 @@ final class ProcessMonitor: ObservableObject { if settings.quietAgentNotifications, assessment == .routine { agentQuietedNotificationCount += 1 } else { - NotificationManager.notify(event: event) + // The rollup key is the furthest-known ancestor pid + // — the process-tree root — not `proc.ppid` alone + // and not provenance label. See `NotificationRollup`'s + // doc comment for why two distinct agent sessions + // must never share one digest just because they + // carry the same supervisor label. + let rootPID = ancestors.last?.pid ?? proc.ppid + switch notificationRollup.record(rootPID: rootPID, techniques: eventRules.map(\.technique), severity: event.topSeverity, timestamp: event.timestamp) { + case .deliver: + NotificationManager.notify(event: event) + case .digest(let count, let techniques, let since, let isFirstDigest): + NotificationManager.notifyDigest(rootPID: rootPID, count: count, techniques: techniques, since: since) + // Only the event that flips this root from + // individual delivery into digesting gets a + // diagnostics line — every later update to the + // same window's digest replaces the notification + // in place silently, or the log would be as + // noisy as the notifications this feature exists + // to quiet. + if isFirstDigest { + DiagnosticsLog.write("notification rollup — pid \(rootPID) tree exceeded its per-window budget, digesting further alerts") + } + } } } riskScore = min(100, riskScore + event.topSeverity.weight) diff --git a/Tests/ArgusTests/DockerWatcherTests.swift b/Tests/ArgusTests/DockerWatcherTests.swift new file mode 100644 index 0000000..79555fe --- /dev/null +++ b/Tests/ArgusTests/DockerWatcherTests.swift @@ -0,0 +1,255 @@ +import XCTest +@testable import Argus + +/// Real-shaped fixtures for `docker events --format '{{json .}}'` lines. +/// Field names/casing/nesting match Docker's actual JSON event schema. +private enum Fixture { + static let containerStart = """ + {"status":"start","id":"abc123","from":"nginx:latest","Type":"container","Action":"start","Actor":{"ID":"abc123","Attributes":{"image":"nginx:latest","name":"web1"}},"scope":"local","time":1700000000,"timeNano":1700000000000000000} + """ + + static let execCreate = """ + {"status":"exec_create: /bin/sh -c ls","id":"abc123","from":"nginx:latest","Type":"container","Action":"exec_create: /bin/sh -c ls","Actor":{"ID":"abc123","Attributes":{"image":"nginx:latest","name":"web1"}},"scope":"local","time":1700000001,"timeNano":1700000001000000000} + """ + + static let execStart = """ + {"status":"exec_start: /bin/sh -c ls","id":"abc123","from":"nginx:latest","Type":"container","Action":"exec_start: /bin/sh -c ls","Actor":{"ID":"abc123","Attributes":{"image":"nginx:latest","name":"web1"}},"scope":"local","time":1700000002,"timeNano":1700000002000000000} + """ + + static let networkConnect = """ + {"status":"connect","id":"net1","Type":"network","Action":"connect","Actor":{"ID":"net1","Attributes":{"container":"abc123","name":"bridge"}},"scope":"local","time":1700000003} + """ + + static let imagePull = """ + {"status":"pull","id":"nginx:latest","Type":"image","Action":"pull","Actor":{"ID":"nginx:latest","Attributes":{"name":"nginx:latest"}},"scope":"local","time":1700000004} + """ + + static let volumeCreate = """ + {"status":"create","id":"vol1","Type":"volume","Action":"create","Actor":{"ID":"vol1","Attributes":{}},"scope":"local","time":1700000005} + """ + + static let containerDie = """ + {"status":"die","id":"abc123","from":"nginx:latest","Type":"container","Action":"die","Actor":{"ID":"abc123","Attributes":{"image":"nginx:latest","name":"web1"}},"scope":"local","time":1700000006} + """ + + static let missingAttributes = """ + {"status":"start","id":"abc123","Type":"container","Action":"start","Actor":{"ID":"abc123"},"scope":"local","time":1700000007} + """ + + static let missingActor = """ + {"status":"start","id":"abc123","Type":"container","Action":"start","scope":"local","time":1700000008} + """ +} + +final class DockerEventClassifierTests: XCTestCase { + func testSplitActionOnSimpleVerb() { + let (verb, detail) = DockerEventClassifier.splitAction("start") + XCTAssertEqual(verb, "start") + XCTAssertNil(detail) + } + + func testSplitActionSplitsOnFirstColonSpace() { + let (verb, detail) = DockerEventClassifier.splitAction("exec_start: /bin/sh -c echo hi: there") + XCTAssertEqual(verb, "exec_start") + XCTAssertEqual(detail, "/bin/sh -c echo hi: there") + } + + func testClassifyIgnoresNonContainerType() { + let line = try! JSONDecoder().decode(DockerEventLine.self, from: Data(Fixture.networkConnect.utf8)) + XCTAssertNil(DockerEventClassifier.classify(line)) + } + + func testClassifyIgnoresDieAction() { + let line = try! JSONDecoder().decode(DockerEventLine.self, from: Data(Fixture.containerDie.utf8)) + XCTAssertNil(DockerEventClassifier.classify(line)) + } + + func testClassifyStart() { + let line = try! JSONDecoder().decode(DockerEventLine.self, from: Data(Fixture.containerStart.utf8)) + let result = DockerEventClassifier.classify(line) + XCTAssertEqual(result?.kind, .start) + XCTAssertNil(result?.execCommand) + } + + func testClassifyExecStartExtractsCommand() { + let line = try! JSONDecoder().decode(DockerEventLine.self, from: Data(Fixture.execStart.utf8)) + let result = DockerEventClassifier.classify(line) + XCTAssertEqual(result?.kind, .execStart) + XCTAssertEqual(result?.execCommand, "/bin/sh -c ls") + } + + func testClassifyExecCreateExtractsCommand() { + let line = try! JSONDecoder().decode(DockerEventLine.self, from: Data(Fixture.execCreate.utf8)) + let result = DockerEventClassifier.classify(line) + XCTAssertEqual(result?.kind, .execCreate) + XCTAssertEqual(result?.execCommand, "/bin/sh -c ls") + } +} + +final class DockerEventBuilderTests: XCTestCase { + func testStartProducesInfoT1610WithNameAndImage() { + let event = DockerEventBuilder.makeEvent(kind: .start, execCommand: nil, image: "nginx:latest", name: "web1", timestamp: Date()) + XCTAssertNotNil(event) + XCTAssertEqual(event?.pid, 0) + XCTAssertEqual(event?.ppid, 0) + XCTAssertEqual(event?.executable, "web1") + XCTAssertTrue(event!.command.contains("nginx:latest")) + XCTAssertTrue(event!.command.contains("web1")) + XCTAssertEqual(event?.rules.count, 1) + XCTAssertEqual(event?.rules.first?.severity, .info) + XCTAssertEqual(event?.rules.first?.technique, "T1610") + XCTAssertEqual(event?.provenance, ["docker"]) + } + + func testExecCreateIsIgnored() { + let event = DockerEventBuilder.makeEvent(kind: .execCreate, execCommand: "/bin/sh -c ls", image: "nginx:latest", name: "web1", timestamp: Date()) + XCTAssertNil(event) + } + + func testExecStartProducesWatchT1609WithCommand() { + let event = DockerEventBuilder.makeEvent(kind: .execStart, execCommand: "/bin/sh -c ls", image: "nginx:latest", name: "web1", timestamp: Date()) + XCTAssertNotNil(event) + XCTAssertEqual(event?.rules.first?.severity, .watch) + XCTAssertEqual(event?.rules.first?.technique, "T1609") + XCTAssertTrue(event!.command.contains("/bin/sh -c ls")) + XCTAssertTrue(event!.command.contains("web1")) + XCTAssertEqual(event?.provenance, ["docker"]) + } + + func testFallsBackToImageWhenNameMissing() { + let event = DockerEventBuilder.makeEvent(kind: .start, execCommand: nil, image: "nginx:latest", name: nil, timestamp: Date()) + XCTAssertEqual(event?.executable, "nginx:latest") + } + + func testUnknownWhenBothNameAndImageMissing() { + let event = DockerEventBuilder.makeEvent(kind: .start, execCommand: nil, image: nil, name: nil, timestamp: Date()) + XCTAssertEqual(event?.executable, "unknown") + } +} + +final class DockerLineProcessorTests: XCTestCase { + func testContainerStartLineProducesEvent() { + let event = DockerLineProcessor.processLine(Fixture.containerStart) + XCTAssertNotNil(event) + XCTAssertEqual(event?.executable, "web1") + XCTAssertEqual(event?.rules.first?.technique, "T1610") + XCTAssertEqual(event?.provenance, ["docker"]) + } + + func testExecStartLineProducesEvent() { + let event = DockerLineProcessor.processLine(Fixture.execStart) + XCTAssertNotNil(event) + XCTAssertEqual(event?.rules.first?.technique, "T1609") + XCTAssertTrue(event!.command.contains("/bin/sh -c ls")) + } + + func testExecCreateLineIsIgnored() { + XCTAssertNil(DockerLineProcessor.processLine(Fixture.execCreate)) + } + + func testNetworkTypeLineIsIgnored() { + XCTAssertNil(DockerLineProcessor.processLine(Fixture.networkConnect)) + } + + func testImageTypeLineIsIgnored() { + XCTAssertNil(DockerLineProcessor.processLine(Fixture.imagePull)) + } + + func testVolumeTypeLineIsIgnored() { + XCTAssertNil(DockerLineProcessor.processLine(Fixture.volumeCreate)) + } + + func testContainerDieLineIsIgnored() { + XCTAssertNil(DockerLineProcessor.processLine(Fixture.containerDie)) + } + + func testNonJSONGarbageLineDoesNotCrash() { + XCTAssertNil(DockerLineProcessor.processLine("this is not json at all {{{")) + } + + func testEmptyLineIsIgnored() { + XCTAssertNil(DockerLineProcessor.processLine("")) + XCTAssertNil(DockerLineProcessor.processLine(" \n")) + } + + func testMissingAttributesIsTolerated() { + let event = DockerLineProcessor.processLine(Fixture.missingAttributes) + XCTAssertNotNil(event) + XCTAssertEqual(event?.executable, "unknown") + } + + func testMissingActorIsTolerated() { + let event = DockerLineProcessor.processLine(Fixture.missingActor) + XCTAssertNotNil(event) + XCTAssertEqual(event?.executable, "unknown") + } + + func testMalformedPartialJSONLineIsIgnored() { + let partial = String(Fixture.containerStart.prefix(40)) + XCTAssertNil(DockerLineProcessor.processLine(partial)) + } +} + +final class LineBufferTests: XCTestCase { + func testSingleCompleteLine() { + var buffer = "" + let lines = LineBuffer.consume("hello\n", buffer: &buffer) + XCTAssertEqual(lines, ["hello"]) + XCTAssertEqual(buffer, "") + } + + func testPartialLineIsHeldBack() { + var buffer = "" + let lines = LineBuffer.consume("hel", buffer: &buffer) + XCTAssertEqual(lines, []) + XCTAssertEqual(buffer, "hel") + } + + func testChunkBoundaryMidLineIsReassembled() { + var buffer = "" + // Simulates a readabilityHandler delivering an NDJSON line split + // across two separate chunks, as a real pipe read can do. + let firstLines = LineBuffer.consume("{\"Type\":\"cont", buffer: &buffer) + XCTAssertEqual(firstLines, []) + XCTAssertEqual(buffer, "{\"Type\":\"cont") + + let secondLines = LineBuffer.consume("ainer\"}\n", buffer: &buffer) + XCTAssertEqual(secondLines, ["{\"Type\":\"container\"}"]) + XCTAssertEqual(buffer, "") + } + + func testMultipleLinesInOneChunk() { + var buffer = "" + let lines = LineBuffer.consume("line1\nline2\nline3\n", buffer: &buffer) + XCTAssertEqual(lines, ["line1", "line2", "line3"]) + XCTAssertEqual(buffer, "") + } + + func testTrailingPartialAfterMultipleCompleteLines() { + var buffer = "" + let lines = LineBuffer.consume("line1\nline2\npartial", buffer: &buffer) + XCTAssertEqual(lines, ["line1", "line2"]) + XCTAssertEqual(buffer, "partial") + } +} + +/// Verifies the watcher itself never touches a real `docker` binary and +/// stays silent (beyond the diagnostics line) when none is found — the CI +/// environment and most dev machines have no Docker installed at all. +final class DockerWatcherTests: XCTestCase { + func testStartIsInertWithoutRealDockerCLI() { + // DockerWatcher only ever probes the three fixed candidate paths; + // on a machine (like CI) without any of them present this must not + // throw, hang, or spawn anything — just log and return. + let watcher = DockerWatcher() + watcher.start() + watcher.stop() + // No assertion beyond "did not crash/hang" — there is no real docker + // CLI in this environment to assert against, by design. + } + + func testStopBeforeStartIsSafe() { + let watcher = DockerWatcher() + watcher.stop() + } +} diff --git a/Tests/ArgusTests/NotificationRollupTests.swift b/Tests/ArgusTests/NotificationRollupTests.swift new file mode 100644 index 0000000..b0ff854 --- /dev/null +++ b/Tests/ArgusTests/NotificationRollupTests.swift @@ -0,0 +1,135 @@ +import XCTest +@testable import Argus + +final class NotificationRollupTests: XCTestCase { + private let base = Date(timeIntervalSince1970: 1_700_000_000) + + private func rollup(window: TimeInterval = 300, budget: Int = 3) -> NotificationRollup { + NotificationRollup(window: window, budget: budget) + } + + // MARK: - budget + + func testFirstBudgetEventsDeliverIndividually() { + let r = rollup(budget: 3) + for i in 0..<3 { + let decision = r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base.addingTimeInterval(Double(i))) + XCTAssertEqual(decision, .deliver, "event \(i) should still be within budget") + } + } + + func testEventAfterBudgetBecomesDigestWithCorrectCountAndTechniques() { + let r = rollup(budget: 3) + for i in 0..<3 { + _ = r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base.addingTimeInterval(Double(i))) + } + let decision = r.record(rootPID: 100, techniques: ["T1553"], severity: .watch, timestamp: base.addingTimeInterval(3)) + XCTAssertEqual(decision, .digest(count: 4, techniques: ["T1059", "T1553"], since: base, isFirstDigest: true)) + } + + func testSubsequentDigestEventsAreNotFirstDigestAndAccumulateCount() { + let r = rollup(budget: 3) + for i in 0..<3 { + _ = r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base.addingTimeInterval(Double(i))) + } + let first = r.record(rootPID: 100, techniques: [], severity: .watch, timestamp: base.addingTimeInterval(3)) + XCTAssertEqual(first, .digest(count: 4, techniques: ["T1059"], since: base, isFirstDigest: true)) + + let second = r.record(rootPID: 100, techniques: [], severity: .watch, timestamp: base.addingTimeInterval(4)) + XCTAssertEqual(second, .digest(count: 5, techniques: ["T1059"], since: base, isFirstDigest: false)) + } + + // MARK: - critical always delivers + + func testCriticalAlwaysDeliversAndDoesNotConsumeBudget() { + let r = rollup(budget: 1) + // Exhaust the budget with a non-critical event. + XCTAssertEqual(r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base), .deliver) + // Now digesting. + if case .digest = r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base.addingTimeInterval(1)) { + // expected + } else { + XCTFail("expected digest once budget exhausted") + } + // A critical event still delivers individually even though the + // budget is exhausted and this root is already digesting. + let criticalDecision = r.record(rootPID: 100, techniques: ["T1562"], severity: .critical, timestamp: base.addingTimeInterval(2)) + XCTAssertEqual(criticalDecision, .deliver) + } + + func testCriticalDeliveryDoesNotConsumeBudgetForLaterEvents() { + let r = rollup(budget: 1) + // Critical events, however many, never consume the single-slot budget. + for i in 0..<5 { + let decision = r.record(rootPID: 100, techniques: ["T1562"], severity: .critical, timestamp: base.addingTimeInterval(Double(i))) + XCTAssertEqual(decision, .deliver, "critical event \(i) must always deliver") + } + // The budget is still untouched, so the next non-critical event + // still delivers individually. + let watchDecision = r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base.addingTimeInterval(5)) + XCTAssertEqual(watchDecision, .deliver) + } + + // MARK: - window expiry + + func testWindowExpiryResetsToIndividualDelivery() { + let r = rollup(window: 300, budget: 1) + XCTAssertEqual(r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base), .deliver) + if case .digest = r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base.addingTimeInterval(1)) { + // expected: digesting within the same window + } else { + XCTFail("expected digest before window expiry") + } + + // Past the window boundary — a brand new window opens with a fresh budget. + let afterExpiry = r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base.addingTimeInterval(301)) + XCTAssertEqual(afterExpiry, .deliver) + } + + // MARK: - independent roots + + func testTwoDistinctRootsTrackIndependently() { + let r = rollup(budget: 1) + XCTAssertEqual(r.record(rootPID: 100, techniques: ["T1059"], severity: .watch, timestamp: base), .deliver) + // A different root's first event still delivers, unaffected by root + // 100 already having used its budget. + XCTAssertEqual(r.record(rootPID: 200, techniques: ["T1059"], severity: .watch, timestamp: base), .deliver) + + // Root 100's second event digests... + if case .digest = r.record(rootPID: 100, techniques: [], severity: .watch, timestamp: base.addingTimeInterval(1)) { + // expected + } else { + XCTFail("expected root 100 to be digesting") + } + // ...while root 200's second event still digests independently, + // with its own isFirstDigest and its own count, not root 100's. + let root200Second = r.record(rootPID: 200, techniques: ["T1553"], severity: .watch, timestamp: base.addingTimeInterval(1)) + XCTAssertEqual(root200Second, .digest(count: 2, techniques: ["T1059", "T1553"], since: base, isFirstDigest: true)) + } + + // MARK: - isFirstDigest exactly once per window per root + + func testIsFirstDigestTrueExactlyOncePerWindow() { + let r = rollup(budget: 2) + _ = r.record(rootPID: 100, techniques: [], severity: .watch, timestamp: base) + _ = r.record(rootPID: 100, techniques: [], severity: .watch, timestamp: base.addingTimeInterval(1)) + + var firstDigestCount = 0 + for i in 2..<10 { + let decision = r.record(rootPID: 100, techniques: [], severity: .watch, timestamp: base.addingTimeInterval(Double(i))) + if case .digest(_, _, _, let isFirst) = decision, isFirst { + firstDigestCount += 1 + } + } + XCTAssertEqual(firstDigestCount, 1) + } + + // MARK: - technique dedup order + + func testTechniquesDedupPreservesFirstSeenOrder() { + let r = rollup(budget: 1) + _ = r.record(rootPID: 100, techniques: ["T1059", "T1553"], severity: .watch, timestamp: base) + let decision = r.record(rootPID: 100, techniques: ["T1553", "T1547", "T1059"], severity: .watch, timestamp: base.addingTimeInterval(1)) + XCTAssertEqual(decision, .digest(count: 2, techniques: ["T1059", "T1553", "T1547"], since: base, isFirstDigest: true)) + } +}