From e48c3f6840d66266cff7c5d835adf268c0971331 Mon Sep 17 00:00:00 2001 From: Gronod Date: Fri, 11 Sep 2026 11:36:50 +0100 Subject: [PATCH] test(process): complete ProcessManager edge contracts (#84) Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../ICCeryCore/Process/ProcessManager.swift | 20 +- .../ICCeryCoreTests/ProcessManagerTests.swift | 192 ++++++++++++++++++ 2 files changed, 206 insertions(+), 6 deletions(-) diff --git a/Packages/ICCeryCore/Sources/ICCeryCore/Process/ProcessManager.swift b/Packages/ICCeryCore/Sources/ICCeryCore/Process/ProcessManager.swift index a084908..f4ba148 100644 --- a/Packages/ICCeryCore/Sources/ICCeryCore/Process/ProcessManager.swift +++ b/Packages/ICCeryCore/Sources/ICCeryCore/Process/ProcessManager.swift @@ -145,9 +145,7 @@ public actor ProcessManager { ) let process = prepared.process - AppLogger(category: "process").debug( - "spawn \(id): \(binary.path) \(LogSanitizer.sanitizeArgs(arguments))" - ) + logSpawn(id: id, binary: binary, arguments: arguments, captured: false) children[id] = RunningChild( process: process, @@ -214,9 +212,7 @@ public actor ProcessManager { let stdoutPipe = prepared.stdoutPipe let stderrPipe = prepared.stderrPipe - AppLogger(category: "process").debug( - "spawn(captured) \(id): \(binary.path) \(LogSanitizer.sanitizeArgs(arguments))" - ) + logSpawn(id: id, binary: binary, arguments: arguments, captured: true) // Register and set up the termination hand-off before run() so // a very fast exit is never missed (#50, #52). @@ -466,6 +462,18 @@ public actor ProcessManager { ) } + private nonisolated func logSpawn( + id: String, + binary: URL, + arguments: [String], + captured: Bool + ) { + let prefix = captured ? "spawn(captured)" : "spawn" + AppLogger(category: "process").debug( + "\(prefix) \(id): \(binary.path) \(LogSanitizer.sanitizeArgs(arguments))" + ) + } + /// `terminationHandler` can lose a fast-exit race on a loaded host; /// `waitUntilExit` on a detached thread is the fallback (#50, #52). /// The handler is attached before `run()`; the wait thread starts diff --git a/Tests/ICCeryCoreTests/ProcessManagerTests.swift b/Tests/ICCeryCoreTests/ProcessManagerTests.swift index 083d5f7..716bf92 100644 --- a/Tests/ICCeryCoreTests/ProcessManagerTests.swift +++ b/Tests/ICCeryCoreTests/ProcessManagerTests.swift @@ -58,6 +58,59 @@ struct ProcessManagerTests { func finish() -> Bool { lock.lock(); defer { lock.unlock() }; if finished { return false }; finished = true; return true } } + /// Subscribes synchronously (registration happens inside `events()`) + /// then records every event for `id` until the task is cancelled. + /// Unlike `collect`, observation continues past `.exit` so tests can + /// prove exactly-once exit emission. + private func observe( + _ manager: ProcessManager, + id: String, + into box: Box + ) -> Task { + let stream = manager.events() + return Task { + for await event in stream { + guard event.id == id else { continue } + box.append(event) + } + } + } + + private func exitCount(in box: Box) -> Int { + box.events.filter { if case .exit = $0 { return true }; return false }.count + } + + private func waitForExit(in box: Box, timeout: TimeInterval = 10) async -> Bool { + let deadline = Date().addingTimeInterval(timeout) + while Date() < deadline { + if exitCount(in: box) > 0 { return true } + try? await Task.sleep(for: .milliseconds(10)) + } + return false + } + + private func waitForFile(_ url: URL, timeout: TimeInterval = 5) async -> Bool { + let deadline = Date().addingTimeInterval(timeout) + while Date() < deadline { + if FileManager.default.fileExists(atPath: url.path) { return true } + try? await Task.sleep(for: .milliseconds(10)) + } + return false + } + + private func waitForRunning( + _ manager: ProcessManager, + id: String, + timeout: TimeInterval = 5 + ) async -> Bool { + let deadline = Date().addingTimeInterval(timeout) + while Date() < deadline { + if await manager.isRunning(id) { return true } + try? await Task.sleep(for: .milliseconds(10)) + } + return false + } + // MARK: - Tests @Test func streamsStdoutAndEmitsExit() async throws { @@ -214,6 +267,145 @@ struct ProcessManagerTests { try await pm.sendStdin(id: "nope", text: "d\n") } } + + @Test func explicitPartialFlushEmitsRowColorsJSON() async throws { + let pm = ProcessManager() + let marker = Self.fixtureDir + .appendingPathComponent("partial-row-ready-\(UUID().uuidString)") + let bin = try script( + "partial-row.sh", + "#!/bin/sh\nprintf 'ROW_COLORS_JSON: {\"row\":9}'\ntouch \"$1\"\nsleep 30\n" + ) + let box = Box() + let observer = observe(pm, id: "t11", into: box) + try await pm.runStreaming(id: "t11", binary: bin, arguments: [marker.path]) + #expect(await waitForFile(marker)) + // Retry the flush so the pipe-ingest task can win the actor race + // on a loaded host; the first successful flush emits the row. + var flushed = false + for _ in 0..<50 { + await pm.flushPartialLine(id: "t11") + if box.events.contains(where: { if case .jsonRow = $0 { return true }; return false }) { + flushed = true + break + } + try await Task.sleep(for: .milliseconds(20)) + } + #expect(flushed) + await pm.kill(id: "t11") + #expect(await waitForExit(in: box)) + observer.cancel() + let events = box.events + let rows = events.compactMap { e -> String? in + if case .jsonRow(_, let d) = e { return String(decoding: d, as: UTF8.self) } + return nil + } + #expect(rows == ["{\"row\":9}"]) + // Prefixed tails must not leak into stdout, even via finalize. + #expect(!events.contains(.stdout(id: "t11", line: "ROW_COLORS_JSON: {\"row\":9}"))) + #expect(exitCount(in: box) == 1) + } + + @Test func unterminatedRowTailFinalizesAsJSONRow() async throws { + let pm = ProcessManager() + let bin = try script( + "row-tail.sh", + "#!/bin/sh\nprintf 'ROW_COLORS_JSON: {\"row\":42}'\n" + ) + let box = Box() + let observer = observe(pm, id: "t12", into: box) + try await pm.runStreaming(id: "t12", binary: bin, arguments: []) + #expect(await waitForExit(in: box)) + observer.cancel() + let events = box.events + let rows = events.compactMap { e -> String? in + if case .jsonRow(_, let d) = e { return String(decoding: d, as: UTF8.self) } + return nil + } + #expect(rows == ["{\"row\":42}"]) + #expect(!events.contains(.stdout(id: "t12", line: "ROW_COLORS_JSON: {\"row\":42}"))) + let rowIndex = events.firstIndex { + if case .jsonRow = $0 { return true }; return false + } + let exitIndexes = events.indices.filter { + if case .exit = events[$0] { return true }; return false + } + #expect(exitIndexes.count == 1) + if let rowIndex, let exitIndex = exitIndexes.first { + #expect(rowIndex < exitIndex) + } else { + Issue.record("expected a jsonRow before the exit event") + } + } + + @Test func fastStreamingExitEmitsExactlyOneExit() async throws { + let pm = ProcessManager() + let bin = try script("fast-stream.sh", "#!/bin/sh\nexit 0\n") + let box = Box() + let observer = observe(pm, id: "t13", into: box) + try await pm.runStreaming(id: "t13", binary: bin, arguments: []) + #expect(await waitForExit(in: box)) + // The grace window must outlast the 2 s finalize watchdog so a + // duplicate emission from it would be observed. + try await Task.sleep(for: .milliseconds(2500)) + observer.cancel() + #expect(box.events == [.exit(id: "t13", code: 0)]) + } + + @Test func fastCapturedExitEmitsExactlyOneExit() async throws { + let pm = ProcessManager() + let bin = try script("fast-cap.sh", "#!/bin/sh\nexit 7\n") + let box = Box() + let observer = observe(pm, id: "t14", into: box) + let result = try await pm.runCaptured(id: "t14", binary: bin, arguments: []) + #expect(result.exitCode == 7) + // Both the termination handler and the waitUntilExit watchdog + // resume the same box; give the slower path time to fire. + try await Task.sleep(for: .milliseconds(500)) + observer.cancel() + #expect(box.events == [.exit(id: "t14", code: 7)]) + } + + @Test func capturedRunSetsArgyllNotInteractive() async throws { + let pm = ProcessManager() + let bin = try script( + "cap-env.sh", + "#!/bin/sh\necho \"ANI=$ARGYLL_NOT_INTERACTIVE\"\n" + ) + let result = try await pm.runCaptured(id: "t15", binary: bin, arguments: []) + #expect(result.stdout == "ANI=1\n") + } + + @Test func killAllTerminatesStreamingAndCapturedChildren() async throws { + let pm = ProcessManager() + let marker = Self.fixtureDir + .appendingPathComponent("mixed-cap-ready-\(UUID().uuidString)") + let slowBin = try script("mixed-slow.sh", "#!/bin/sh\nsleep 30\n") + let capBin = try script("mixed-cap.sh", "#!/bin/sh\ntouch \"$1\"\nsleep 30\n") + let streamBox = Box() + let capBox = Box() + let streamObserver = observe(pm, id: "t16", into: streamBox) + let capObserver = observe(pm, id: "t17", into: capBox) + try await pm.runStreaming(id: "t16", binary: slowBin, arguments: []) + let capTask = Task { + try await pm.runCaptured(id: "t17", binary: capBin, arguments: [marker.path]) + } + #expect(await waitForFile(marker)) + #expect(await waitForRunning(pm, id: "t16")) + #expect(await waitForRunning(pm, id: "t17")) + #expect(await pm.killAll() == 2) + _ = try await capTask.value + #expect(await waitForExit(in: streamBox)) + #expect(await waitForExit(in: capBox)) + // Grace window outlasts the streaming finalize watchdog. + try await Task.sleep(for: .milliseconds(2500)) + streamObserver.cancel() + capObserver.cancel() + #expect(!(await pm.isRunning("t16"))) + #expect(!(await pm.isRunning("t17"))) + #expect(exitCount(in: streamBox) == 1) + #expect(exitCount(in: capBox) == 1) + } } @Suite("ProcessLineDecoder") -- 2.39.5