From 7b89eddf53e2179efbb6e8ead081e97cccc2629d Mon Sep 17 00:00:00 2001 From: Mark Watts Date: Fri, 21 Aug 2026 15:41:42 +0100 Subject: [PATCH 1/3] Add rollup digests to collapse notification bursts per process tree A busy supervised session can fire many individual alerts from one process tree within minutes; NotificationRollup lets the first few notifications per tree root deliver individually within a window, then collapses further ones into one continuously-updated digest, while critical severity always bypasses digesting and never touches the budget. Co-Authored-By: Claude Fable 5 --- Sources/Argus/NotificationManager.swift | 38 +++++ Sources/Argus/NotificationRollup.swift | 103 +++++++++++++ Sources/Argus/ProcessMonitor.swift | 30 +++- .../ArgusTests/NotificationRollupTests.swift | 135 ++++++++++++++++++ 4 files changed, 305 insertions(+), 1 deletion(-) create mode 100644 Sources/Argus/NotificationRollup.swift create mode 100644 Tests/ArgusTests/NotificationRollupTests.swift 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/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)) + } +} From 9fd21428947ebe7cd1b88cefde4ab7f44b27bd3b Mon Sep 17 00:00:00 2001 From: Mark Watts Date: Fri, 21 Aug 2026 15:42:46 +0100 Subject: [PATCH 2/3] Add DockerWatcher: a docker events sensor for container lifecycle activity MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Argus polls the host process table, but containers run inside Docker's Linux VM, so host ps only ever sees the docker CLI and VM helper — container start/exec activity is otherwise invisible. This adds a third, independent sensor (alongside ProcessMonitor and PersistenceWatcher) that subscribes to `docker events --format {{json .}}`, classifying container start (T1610/info) and exec_start (T1609/watch) into synthetic ProcessEvents fed through ProcessMonitor.ingestExternal. Scope is deliberately honest and narrow: only daemon-reported container lifecycle is visible, not processes running inside a container. Inert and silent (beyond one log line) when no docker CLI is found; retries with a fixed 60s backoff if the daemon isn't running yet, so Docker Desktop starting later is picked up without an app restart. Co-Authored-By: Claude Fable 5 --- Sources/Argus/App.swift | 16 ++ Sources/Argus/DockerWatcher.swift | 328 ++++++++++++++++++++++ Tests/ArgusTests/DockerWatcherTests.swift | 255 +++++++++++++++++ 3 files changed, 599 insertions(+) create mode 100644 Sources/Argus/DockerWatcher.swift create mode 100644 Tests/ArgusTests/DockerWatcherTests.swift 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/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() + } +} From 2e4ca87e42f966b275bc1aae3b18d7cd2cf73cde Mon Sep 17 00:00:00 2001 From: Mark Watts Date: Fri, 21 Aug 2026 15:45:41 +0100 Subject: [PATCH 3/3] Document notification rollup digests and Docker container sensor Updated README to cover two merged detection enhancements: notification rollups that collapse busy process trees' sub-critical alerts into digests (keeping critical alerts individual), and a new Docker events subscriber surfacing container lifecycle as synthetic events. Also increased test count from 170 to 210 tests and updated the project layout tree to include the new modules and test suites. Co-Authored-By: Claude Fable 5 --- README.md | 32 ++++++++++++++++++++++++++++---- 1 file changed, 28 insertions(+), 4 deletions(-) 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