From c9c46fa419839b1fcb887445ef8395afc6fed854 Mon Sep 17 00:00:00 2001 From: jsflax Date: Sat, 19 Sep 2026 08:23:52 -0400 Subject: [PATCH] Bound advice output and capture stalled native test diagnostics --- .github/workflows/release.yml | 16 +- CHANGELOG.md | 2 + Sources/EngramKit/MemoryTools+Core.swift | 57 +++- Sources/EngramKit/MemoryTools+Service.swift | 45 ++- Sources/EngramKit/MemoryTools.swift | 6 + .../MemoryServiceContractTests.swift | 79 ++++- scripts/run_native_tests.py | 291 ++++++++++++++++++ scripts/test_codex_native_tests.py | 202 ++++++++++++ scripts/test_codex_release_workflow.py | 13 + 9 files changed, 687 insertions(+), 24 deletions(-) create mode 100644 scripts/run_native_tests.py create mode 100644 scripts/test_codex_native_tests.py diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 81022da..7c4beac 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -128,8 +128,13 @@ jobs: # tension recall assertions from the revival. ~345 correctness tests # remain as the release bar. - name: Run tests - timeout-minutes: 30 + id: native_tests + # Keep the full 30-minute test budget, plus bounded diagnostics/cleanup. + timeout-minutes: 32 run: > + python3 -B scripts/run_native_tests.py + --diagnostics-dir build/native-test-diagnostics + --timeout-seconds 1800 --silence-seconds 300 -- swift test --force-resolved-versions --skip-build --filter "EngramTests|EngramMemoryCoreTests|EngramRealityKitTests|PositionVersionTests" --skip "PerfTests" --skip "keyBERTKeywordExtraction" @@ -139,6 +144,15 @@ jobs: --skip "recall_statementBudget" --skip "clusters_statementBudget" + - name: Upload native test diagnostics + if: ${{ (failure() || cancelled()) && steps.native_tests.outcome != 'skipped' }} + uses: actions/upload-artifact@v4 + with: + name: native-test-diagnostics-${{ github.run_id }}-${{ github.run_attempt }} + path: build/native-test-diagnostics/ + if-no-files-found: ignore + retention-days: 7 + # Statement-budget regressions read the process-global SQL statement # counter; suites running in parallel pollute every measurement window # (sustained, not bursty — min-of-N can't rescue it). They run alone in diff --git a/CHANGELOG.md b/CHANGELOG.md index 8901944..0401687 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,8 @@ All notable changes to Engram are documented in this file. failure is reported instead of continuing after a failed transaction start. - Include the actual advice query as JSON-quoted provenance, preventing a multiline query from introducing extra result headings. +- Keep advice within its overall character budget, including query provenance, + headings, and truncation markers, and report only memory IDs retained in it. - Persist project changes made through the typed Swift memory service, including project-only updates and updates that also edit a topic. - Preserve plugin-owned Codex configuration during native app/CLI installation; diff --git a/Sources/EngramKit/MemoryTools+Core.swift b/Sources/EngramKit/MemoryTools+Core.swift index 8b04f64..af97b21 100644 --- a/Sources/EngramKit/MemoryTools+Core.swift +++ b/Sources/EngramKit/MemoryTools+Core.swift @@ -388,9 +388,27 @@ extension MemoryTools { // MARK: - recall + /// Capture provenance while rendering, never by parsing recalled content. + nonisolated static func appendRecallRowMarker(_ id: UUID, to output: inout String, + rows: inout [RecallRowBoundary]) { + output += "[id:\(id.uuidString)] " + rows.append(RecallRowBoundary(id: id, markerEnd: output.count)) + } + func handleRecall(_ args: [String: Value]?) async throws -> CallTool.Result { lastRecallHits = [] lastRecallMode = .vector + lastRecallRows = [] + var recalledHits: [RecallHit] = [] + var recalledMode: RecallMode = .vector + var renderedRows: [RecallRowBoundary] = [] + // Keep capture local across embedding awaits, then publish it with + // this invocation's result, including empty/error returns. + defer { + lastRecallHits = recalledHits + lastRecallMode = recalledMode + lastRecallRows = renderedRows + } let a = try args.decode(RecallArgs.self) guard !a.query.isEmpty else { throw MCPError.invalidParams("'query' is required") @@ -556,7 +574,7 @@ extension MemoryTools { // access-stat bump — the bump's writes invalidate the row caches, // so a post-bump read would re-issue one SELECT per field. // Traversal hits append below. - lastRecallHits = filtered.compactMap { hit in + recalledHits = filtered.compactMap { hit in guard hit.object.globalId != nil else { return nil } return RecallHit(memory: record(from: hit.object), distance: hit.distance, @@ -583,7 +601,7 @@ extension MemoryTools { return isHubResident(gid) } - let lines = filtered.compactMap { match -> String? in + let lines = filtered.compactMap { match -> (id: UUID, body: String)? in let m = match.object guard let mGid = m.globalId else { return nil } let dist = String(format: "%.3f", match.distance) @@ -609,15 +627,20 @@ extension MemoryTools { // escape-hardened indentation fence. let body = (isForeign && fenceForeignContent) ? Self.fencedForeignContent(m.content) : m.content - return "[id:\(mGid.uuidString)] [\(m.project)/\(m.topic)]\(badge)\(via) (distance: \(dist)\(impInfo)\(expires)\(created)) \(body)" + return (mGid, "[\(m.project)/\(m.topic)]\(badge)\(via) (distance: \(dist)\(impInfo)\(expires)\(created)) \(body)") } - var output = lines.joined(separator: "\n\n") + var output = "" // Knowledge gap detection — signal when recall results are weak let avgDistance = filtered.map(\.distance).reduce(0, +) / Double(max(filtered.count, 1)) if avgDistance > 1.05 { // v2: relevant query→memory hits measure 0.86–1.05; beyond = weak - output = "⚠️ Weak recall (avg distance: \(String(format: "%.3f", avgDistance)), count: \(filtered.count)). Results may not be closely related to the query.\n\n" + output + output = "⚠️ Weak recall (avg distance: \(String(format: "%.3f", avgDistance)), count: \(filtered.count)). Results may not be closely related to the query.\n\n" + } + for (index, line) in lines.enumerated() { + if index > 0 { output += "\n\n" } + Self.appendRecallRowMarker(line.id, to: &output, rows: &renderedRows) + output += line.body } log("[recall] Output formatted, \(lines.count) lines") @@ -679,7 +702,7 @@ extension MemoryTools { for mem in connected { let m = mem.memory guard let memGlobalId = m.globalId else { continue } - lastRecallHits.append(RecallHit(memory: record(from: m), + recalledHits.append(RecallHit(memory: record(from: m), distance: 0, depth: mem.depth, isForeign: isForeignAuthored(m))) @@ -699,17 +722,19 @@ extension MemoryTools { ? " [by:\(GroupDirectory.badgeName(for: m.authorUserId))]" : "" let connVia = viaMarker(for: memGlobalId) // Small memories shown in full; large ones get a compact preview + output += "\n\n" + Self.appendRecallRowMarker(memGlobalId, to: &output, rows: &renderedRows) if isForeignAuthored(m) && fenceForeignContent { - output += "\n\n[id:\(memGlobalId.uuidString)] [\(m.project)/\(m.topic)]\(connBadge)\(connVia)\(expires)\(edgeInfo) \(Self.fencedForeignContent(m.content))" + output += "[\(m.project)/\(m.topic)]\(connBadge)\(connVia)\(expires)\(edgeInfo) \(Self.fencedForeignContent(m.content))" } else if m.content.count <= 500 { - output += "\n\n[id:\(memGlobalId.uuidString)] [\(m.project)/\(m.topic)]\(connBadge)\(connVia)\(expires)\(edgeInfo) \(m.content)" + output += "[\(m.project)/\(m.topic)]\(connBadge)\(connVia)\(expires)\(edgeInfo) \(m.content)" } else { let firstLine = m.content.split(separator: "\n", maxSplits: 1).first.map(String.init) ?? m.content let preview = String(firstLine.prefix(120)) let charCount = m.content.count let sectionCount = m.content.components(separatedBy: "\n").filter { $0.hasPrefix("## ") || $0.hasPrefix("### ") }.count let sizeInfo = sectionCount > 0 ? "\(sectionCount) sections, \(charCount) chars" : "\(charCount) chars" - output += "\n\n[id:\(memGlobalId.uuidString)] [\(m.project)/\(m.topic)]\(connBadge)\(connVia) (\(sizeInfo)\(expires))\(edgeInfo) \(preview)\(charCount > 120 ? "..." : "")" + output += "[\(m.project)/\(m.topic)]\(connBadge)\(connVia) (\(sizeInfo)\(expires))\(edgeInfo) \(preview)\(charCount > 120 ? "..." : "")" } } } @@ -740,7 +765,7 @@ extension MemoryTools { sessionLog("[recall] DONE, returning \(output.count) chars") return CallTool.Result(content: [.text(output)], isError: false) } else { - lastRecallMode = .fullText + recalledMode = .fullText log("[recall] No embedding available, falling back to FTS5") // Degraded mode: FTS5 full-text search (no embedding model loaded) let contentWords = Self.extractContentWords(from: query) @@ -763,7 +788,7 @@ extension MemoryTools { } let ftsResults = results.matching(ftsQuery, on: \.content, limit: limit) - var lines: [String] = [] + var output = "" for match in ftsResults { let m = match.object m.materialize() // hydrated by the FTS query — format for free @@ -773,7 +798,7 @@ extension MemoryTools { let expires = m.expiresAt == .distantFuture ? "" : ", expires: \(Self.dateFormatter.string(from: m.expiresAt))" let created = hasTemporalFilter ? ", created: \(Self.dateFormatter.string(from: m.createdAt))" : "" guard let mGid = m.globalId else { continue } - lastRecallHits.append(RecallHit(memory: record(from: m), + recalledHits.append(RecallHit(memory: record(from: m), distance: 0, depth: 0, isForeign: isForeign)) @@ -782,12 +807,14 @@ extension MemoryTools { let via = viaMarker(for: mGid) let body = (isForeign && fenceForeignContent) ? Self.fencedForeignContent(m.content) : m.content - lines.append("[id:\(mGid.uuidString)] [\(m.project)/\(m.topic)]\(badge)\(via)\(ftsInfo)\(expires)\(created) \(body)") + if !output.isEmpty { output += "\n\n" } + Self.appendRecallRowMarker(mGid, to: &output, rows: &renderedRows) + output += "[\(m.project)/\(m.topic)]\(badge)\(via)\(ftsInfo)\(expires)\(created) \(body)" } - if lines.isEmpty { + if output.isEmpty { return CallTool.Result(content: [.text("No memories found.")], isError: false) } - return CallTool.Result(content: [.text(lines.joined(separator: "\n\n"))], isError: false) + return CallTool.Result(content: [.text(output)], isError: false) } } diff --git a/Sources/EngramKit/MemoryTools+Service.swift b/Sources/EngramKit/MemoryTools+Service.swift index d08a1d1..4ed8941 100644 --- a/Sources/EngramKit/MemoryTools+Service.swift +++ b/Sources/EngramKit/MemoryTools+Service.swift @@ -25,6 +25,12 @@ extension MemoryTools: MemoryService { // MARK: Core reads public func recall(_ request: RecallRequest) async throws -> RecallResult { + let captured = try await recallWithRows(request) + return captured.result + } + + private func recallWithRows(_ request: RecallRequest) async throws + -> (result: RecallResult, rows: [RecallRowBoundary]) { var args: [String: Value] = [ "query": .string(request.query), "depth": .int(request.depth), @@ -32,9 +38,11 @@ extension MemoryTools: MemoryService { ] if let project = request.project { args["project"] = .string(project) } let result = try await handleRecall(args) - return RecallResult(hits: lastRecallHits, - mode: lastRecallMode, - renderedText: Self.text(from: result)) + // Snapshot both captures together before returning across another + // await; advice must never consult a later call's actor state. + return (RecallResult(hits: lastRecallHits, + mode: lastRecallMode, + renderedText: Self.text(from: result)), lastRecallRows) } public func advise(_ request: AdviseRequest) async throws -> AdviseResult { @@ -44,19 +52,42 @@ extension MemoryTools: MemoryService { // analytics loop is the tuner (plan §advise). let words = MemoryTools.extractContentWords(from: request.prompt) let query = words.isEmpty ? request.prompt : words.joined(separator: " ") - let recallResult = try await recall(RecallRequest( + let captured = try await recallWithRows(RecallRequest( query: query, project: request.project, depth: 1, limit: 5)) + return Self.boundedAdvice(captured.result, rows: captured.rows, + query: query, budget: request.budget) + } + + /// Keep the exact query provenance and already-fenced recall prefix inside + /// the overall block budget. Partial or omitted row IDs are not feedback. + nonisolated static func boundedAdvice(_ recallResult: RecallResult, + rows: [RecallRowBoundary], + query: String, budget: Int) -> AdviseResult { guard !recallResult.hits.isEmpty, recallResult.renderedText != "No memories found." else { return AdviseResult(block: nil, memoryIds: [], mode: recallResult.mode) } + let budget = max(0, budget) + let overhead = AdviseAssembly.memorySection(renderedRecall: "", query: query).count + guard budget > overhead else { + return AdviseResult(block: nil, memoryIds: [], mode: recallResult.mode) + } + let available = budget - overhead var rendered = recallResult.renderedText - if rendered.count > request.budget { - rendered = String(rendered.prefix(request.budget)) + "\n… (truncated)" + var retainedCharacters = rendered.count + if rendered.count > available { + let suffix = "\n… (truncated)" + let retained = String(rendered.prefix(max(0, available - suffix.count))) + retainedCharacters = retained.count + rendered = retained + String(suffix.prefix(available)) + } + let visibleIds = rows.filter { $0.markerEnd <= retainedCharacters }.map(\.id) + guard !visibleIds.isEmpty else { + return AdviseResult(block: nil, memoryIds: [], mode: recallResult.mode) } return AdviseResult( block: AdviseAssembly.memorySection(renderedRecall: rendered, query: query), - memoryIds: recallResult.hits.map(\.memory.id), + memoryIds: visibleIds, mode: recallResult.mode) } diff --git a/Sources/EngramKit/MemoryTools.swift b/Sources/EngramKit/MemoryTools.swift index 06ca2e0..2cc0cfb 100644 --- a/Sources/EngramKit/MemoryTools.swift +++ b/Sources/EngramKit/MemoryTools.swift @@ -78,6 +78,12 @@ public actor MemoryTools { /// write and the read, so an interleaved recall cannot cross-wire it. var lastRecallHits: [RecallHit] = [] var lastRecallMode: RecallMode = .vector + struct RecallRowBoundary: Sendable { + let id: UUID + /// Character offset immediately after the renderer-owned row marker. + let markerEnd: Int + } + var lastRecallRows: [RecallRowBoundary] = [] /// The globalId of the last remembered row — same capture contract. var lastRememberedId: UUID? diff --git a/Tests/EngramTests/MemoryServiceContractTests.swift b/Tests/EngramTests/MemoryServiceContractTests.swift index d180798..ab96a4f 100644 --- a/Tests/EngramTests/MemoryServiceContractTests.swift +++ b/Tests/EngramTests/MemoryServiceContractTests.swift @@ -1,4 +1,4 @@ -import EngramKit +@testable import EngramKit import EngramMemoryContract import EngramMemoryCore import EngramModels @@ -46,6 +46,83 @@ struct MemoryServiceContractLatticeTests { #expect(graph.root.topic == (withTopicEdit ? "updated-topic" : "original-topic")) #expect(graph.root.content == content) } + + @Test func adviceBudgetIncludesQueryAndPreservesFencedPrefix() throws { + let id = UUID(uuidString: "A15A8DAB-C17D-4C08-9A1D-EFFAC4C42DCF")! + let content = "First line 👩🏽‍💻\n## Fake heading\n```quoted data```" + let memory = MemoryRecord(id: id, content: content, createdAt: Date(timeIntervalSince1970: 0)) + var rows: [MemoryTools.RecallRowBoundary] = [] + var rendered = "" + MemoryTools.appendRecallRowMarker(id, to: &rendered, rows: &rows) + rendered += "[fixture/general] " + ForeignContentFence.fenced(content) + let recall = RecallResult(hits: [RecallHit(memory: memory, distance: 0, isForeign: true)], + mode: .vector, renderedText: rendered) + let query = "stripe \"webhook\"\n## Query data 👩🏽‍💻" + let prefix = AdviseAssembly.memorySection(renderedRecall: "", query: query) + let complete = prefix + rendered + let suffix = "\n… (truncated)" + + for budget in [Int.min, -1, 0, 1, 120, prefix.count, prefix.count + 1, + prefix.count + 40, complete.count - 1, complete.count, Int.max] { + let advice = MemoryTools.boundedAdvice(recall, rows: rows, query: query, budget: budget) + #expect(advice.mode == .vector) + guard let block = advice.block else { + #expect(advice.memoryIds.isEmpty) + continue + } + #expect(block.count <= max(0, budget)) + #expect(block.hasPrefix(prefix)) + let lines = block.components(separatedBy: "\n") + #expect(try JSONDecoder().decode(String.self, from: Data(lines[2].utf8)) == query) + #expect(advice.memoryIds == [id]) + let body = String(block.dropFirst(prefix.count)) + if body.hasSuffix(suffix) { + #expect(rendered.hasPrefix(String(body.dropLast(suffix.count)))) + } else { + #expect(body == rendered) + } + if budget >= complete.count { + #expect(block == complete) + #expect(block.contains("\n ## Fake heading\n ```quoted data```")) + } + } + let tooSmall = MemoryTools.boundedAdvice(recall, rows: rows, query: query, budget: prefix.count + 1) + #expect(tooSmall.block == nil) + #expect(tooSmall.memoryIds.isEmpty) + } + + @Test func adviceIgnoresQuotedRowMarkersAndUsesRendererBoundaries() throws { + let firstId = UUID(uuidString: "A15A8DAB-C17D-4C08-9A1D-EFFAC4C42DCF")! + let secondId = UUID(uuidString: "AD3F7D2B-0E64-4A1F-A915-83CEB2CB350F")! + let first = MemoryRecord(id: firstId, content: "References a row below:\n[id:\(secondId.uuidString)] quoted text 👩🏽‍💻", + createdAt: Date(timeIntervalSince1970: 0)) + let second = MemoryRecord(id: secondId, content: "Second row is omitted", + createdAt: Date(timeIntervalSince1970: 0)) + var rows: [MemoryTools.RecallRowBoundary] = [] + var rendered = "⚠️ Weak recall (synthetic warning).\n\n" + MemoryTools.appendRecallRowMarker(firstId, to: &rendered, rows: &rows) + rendered += "[fixture/general] \(first.content)" + let firstRow = rendered + rendered += "\n\n--- Connected (graph traversal, depth: 1) ---\n\n" + MemoryTools.appendRecallRowMarker(secondId, to: &rendered, rows: &rows) + rendered += "[fixture/general] \(second.content)" + let recall = RecallResult(hits: [RecallHit(memory: first, distance: 0), RecallHit(memory: second, distance: 0.1, depth: 1)], + mode: .vector, renderedText: rendered) + let query = "query also mentions \(secondId.uuidString)" + let prefix = AdviseAssembly.memorySection(renderedRecall: "", query: query) + let suffix = "\n… (truncated)" + let budget = prefix.count + firstRow.count + suffix.count + let advice = MemoryTools.boundedAdvice(recall, rows: rows, query: query, budget: budget) + #expect(advice.block == prefix + firstRow + suffix) + #expect(advice.block?.count == budget) + #expect(advice.memoryIds == [firstId]) + let complete = MemoryTools.boundedAdvice(recall, rows: rows, query: query, budget: Int.max) + #expect(complete.memoryIds == [firstId, secondId]) + let cutMarker = MemoryTools.boundedAdvice(recall, rows: rows, query: query, + budget: prefix.count + rows[0].markerEnd - 1 + suffix.count) + #expect(cutMarker.block == nil) + #expect(cutMarker.memoryIds.isEmpty) + } } /// Builds `MemoryTools` over throwaway sqlite files. Peers share the diff --git a/scripts/run_native_tests.py b/scripts/run_native_tests.py new file mode 100644 index 0000000..7972d7b --- /dev/null +++ b/scripts/run_native_tests.py @@ -0,0 +1,291 @@ +#!/usr/bin/env python3 +"""Run the unchanged native test command with bounded stall diagnostics.""" + +import argparse +import json +import os +from pathlib import Path +import signal +import subprocess +import sys +import threading +import time + + +def processes(): + # comm excludes arguments and environments (which can contain credentials). + output = subprocess.check_output( + ["/bin/ps", "-axo", "pid=,ppid=,pgid=,lstart=,stat=,comm="], + text=True, timeout=2) + result = {} + for line in output.splitlines(): + fields = line.split(None, 9) + if len(fields) == 10: + pid, parent, group = map(int, fields[:3]) + result[pid] = {"pid": pid, "parent": parent, "group": group, + "started": " ".join(fields[3:8]), "state": fields[8], + "command": fields[9]} + elif line.strip(): + raise RuntimeError("Unparseable process inventory row") + if not result: + raise RuntimeError("Empty process inventory") + return result + + +def identity(process): + # Executing another binary preserves process identity. + return process["pid"], process["started"] + + +class Runner: + def __init__(self, command, directory, *, timeout=1800, silence=300, + log_limit=64 * 1024**2, sample_tool="/usr/bin/sample", + sample_timeout=10, sample_seconds=5, cleanup_grace=5, + console=None): + self.command = command + self.directory = Path(directory) + self.timeout, self.silence = timeout, silence + self.log_limit = log_limit + self.sample_tool = sample_tool + self.sample_timeout, self.sample_seconds = sample_timeout, sample_seconds + self.cleanup_grace = cleanup_grace + self.console = sys.stdout.buffer if console is None else console + self.known = {} + self.errors, self.diagnostics, self.signals = [], [], [] + self.output_bytes = self.log_bytes = 0 + self.reader_error = None + self.group_closed = False + + def refresh(self): + current = {pid: item for pid, item in processes().items() if "Z" not in item["state"]} + # Remember descendants even after reparenting, but never reuse a stale + # PID after its start identity changes. + owned = {pid: item for pid, item in current.items() + if pid in self.known and identity(item) == identity(self.known[pid])} + # Popen created this private group. Stop following its number forever + # once the original group disappears, so a later group cannot be adopted. + root = current.get(self.child.pid) + if root and self.child.pid in self.known and identity(root) != identity(self.known[self.child.pid]): + self.group_closed = True + group = [item for item in current.values() if item["group"] == self.child.pid] + if not self.group_closed: + if group: + owned.update((item["pid"], item) for item in group) + elif self.child.poll() is not None: + self.group_closed = True + while True: + children = {pid: item for pid, item in current.items() + if item["parent"] in owned and pid not in owned} + if not children: + break + owned.update(children) + self.known.update(owned) + return owned + + def drain(self): + try: + with (self.directory / "output.log").open("xb") as output: + while True: + chunk = os.read(self.child.stdout.fileno(), 65536) + if not chunk: + break + self.last_output = time.monotonic() + self.output_bytes += len(chunk) + kept = chunk[:max(0, self.log_limit - self.log_bytes)] + output.write(kept) + output.flush() + self.log_bytes += len(kept) + if self.console is not None: + try: + self.console.write(chunk) + self.console.flush() + except (BrokenPipeError, OSError): + # A closed console must not stop draining the child. + self.console = None + except Exception as error: + self.reader_error = repr(error) + + def capture(self, reason, deadline): + record = {"reason": reason, "elapsed_seconds": time.monotonic() - self.started, + "samples": [], "errors": []} + self.diagnostics.append(record) + prefix = self.directory / (str(len(self.diagnostics)) + "-" + reason) + try: + owned = self.refresh() + (prefix.with_suffix(".json")).write_text(json.dumps( + {"processes": list(owned.values())[:256]}, indent=2) + "\n") + # Prefer native test executables over SwiftPM and helper processes. + candidates = sorted(owned.values(), key=lambda item: ( + not any(name in item["command"].lower() + for name in (".xctest", "xctest", "engrampackagetests", "swiftpm-testing")), + item["pid"] == self.child.pid, item["pid"]))[:3] + for item in candidates: + remaining = deadline - time.monotonic() + if remaining <= 0 or self.signals: + break + current = processes().get(item["pid"]) + if current is None or identity(current) != identity(item): + continue + path = Path(str(prefix) + "-" + str(item["pid"]) + ".sample.txt") + sample = {"process": item, "path": path.name} + record["samples"].append(sample) + try: + result = subprocess.run( + [self.sample_tool, str(item["pid"]), str(self.sample_seconds), + "10", "-file", str(path)], + stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + timeout=min(self.sample_timeout, max(0.01, remaining))) + sample["exit_code"] = result.returncode + except (OSError, subprocess.TimeoutExpired) as error: + sample["error"] = repr(error) + finally: + if path.exists(): + # Keep each stack artifact bounded as well as output.log. + with path.open("r+b") as output: + if output.seek(0, os.SEEK_END) > 4 * 1024**2: + output.truncate(4 * 1024**2) + sample["truncated"] = True + except Exception as error: + record["errors"].append(repr(error)) + + def cleanup(self, deadline): + receipt = {"signals": [], "remaining": [], "errors": []} + for number in (signal.SIGTERM, signal.SIGKILL): + if time.monotonic() >= deadline: + break + try: + owned = self.refresh() + # Children first, then their launcher. Revalidate every PID + # against a fresh snapshot before each signaling phase. + for pid, item in sorted(owned.items(), key=lambda pair: pair[0] == self.child.pid): + try: + os.kill(pid, number) + receipt["signals"].append({"pid": pid, "started": item["started"], + "signal": signal.Signals(number).name}) + except ProcessLookupError: + pass + until = min(deadline, time.monotonic() + self.cleanup_grace) + while time.monotonic() < until: + self.child.poll() # Reap our direct child before inspecting ps. + if not self.refresh(): + break + time.sleep(0.05) + except Exception as error: + receipt["errors"].append(repr(error)) + # An unreadable inventory removes authority to target unknown + # descendants, but Popen still owns its unreaped direct child. + if self.child.poll() is None: + try: + if number == signal.SIGTERM: + self.child.terminate() + else: + self.child.kill() + receipt["signals"].append({"pid": self.child.pid, + "signal": signal.Signals(number).name, + "direct_child_fallback": True}) + self.child.wait(timeout=max(0.01, min( + self.cleanup_grace, deadline - time.monotonic()))) + except subprocess.TimeoutExpired: + pass + except Exception as fallback_error: + receipt["errors"].append(repr(fallback_error)) + try: + self.child.wait(timeout=max(0.01, min(1, deadline - time.monotonic()))) + except subprocess.TimeoutExpired: + receipt["errors"].append("Direct child did not exit") + try: + receipt["remaining"] = list(self.refresh().values()) + except Exception as error: + receipt["errors"].append(repr(error)) + return receipt + + def run(self): + self.directory.mkdir(parents=True, exist_ok=False) + self.started = self.last_output = time.monotonic() + runtime_deadline = self.started + self.timeout + # Independent of the workflow's 32-minute timeout; leave room for upload. + final_deadline = runtime_deadline + 90 + self.child = subprocess.Popen(self.command, stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, stderr=subprocess.STDOUT, + start_new_session=True) + reader = threading.Thread(target=self.drain, daemon=True) + reader.start() + previous = {} + for number in (signal.SIGINT, signal.SIGTERM): + previous[number] = signal.signal(number, lambda number, frame: self.signals.append(number)) + reason, code, child_code = "exit", None, None + silence_captured = False + next_snapshot = 0 + try: + while True: + now = time.monotonic() + if now >= next_snapshot: + self.refresh() + next_snapshot = now + 1 + child_code = self.child.poll() + if child_code is not None and child_code != 0: + code = child_code if child_code >= 0 else 128 - child_code + break + if self.signals: + reason, code = "signal", 128 + self.signals[0] + break + if self.reader_error: + raise RuntimeError("Output reader failed: " + self.reader_error) + now = time.monotonic() + if now >= runtime_deadline: + reason, code = "timeout", 124 + break + if child_code == 0: + code = 0 + break + if not silence_captured and now - self.last_output >= self.silence: + silence_captured = True + self.capture("silence", min(now + 35, runtime_deadline)) + time.sleep(0.25) + if code and not self.signals: + self.capture(reason, min(time.monotonic() + 35, final_deadline - 15)) + except Exception as error: + reason, code = "supervisor_error", 125 + self.errors.append(repr(error)) + finally: + cleanup = self.cleanup(min(time.monotonic() + 15, final_deadline)) + reader.join(timeout=2) + reader_complete = not reader.is_alive() + self.child.stdout.close() + for number, handler in previous.items(): + signal.signal(number, handler) + if not code and self.signals: + reason, code = "signal", 128 + self.signals[0] + if not code and (cleanup["remaining"] or cleanup["errors"] or + not reader_complete or self.reader_error): + code = 125 + receipt = {"command": self.command, "runtime_budget_seconds": self.timeout, + "silence_threshold_seconds": self.silence, "reason": reason, + "child_exit_code": self.child.returncode, "exit_code": code, + "elapsed_seconds": time.monotonic() - self.started, + "output_bytes": self.output_bytes, "saved_log_bytes": self.log_bytes, + "log_truncated": self.output_bytes > self.log_bytes, + "reader_complete": reader_complete, "reader_error": self.reader_error, + "diagnostics": self.diagnostics, "cleanup": cleanup, + "signals_received": self.signals, "errors": self.errors} + (self.directory / "receipt.json").write_text(json.dumps(receipt, indent=2) + "\n") + return code + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--diagnostics-dir", type=Path, required=True) + parser.add_argument("--timeout-seconds", type=float, default=1800) + parser.add_argument("--silence-seconds", type=float, default=300) + parser.add_argument("command", nargs=argparse.REMAINDER) + args = parser.parse_args() + command = args.command[1:] if args.command[:1] == ["--"] else args.command + if not command or args.timeout_seconds <= 0 or args.silence_seconds <= 0: + parser.error("A command and positive runtime/silence budgets are required") + return Runner(command, args.diagnostics_dir, timeout=args.timeout_seconds, + silence=args.silence_seconds).run() + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/scripts/test_codex_native_tests.py b/scripts/test_codex_native_tests.py new file mode 100644 index 0000000..ca63205 --- /dev/null +++ b/scripts/test_codex_native_tests.py @@ -0,0 +1,202 @@ +"""Exercise native-test supervision using disposable Python child processes.""" + +import importlib.util +import io +import json +import os +from pathlib import Path +import signal +import subprocess +import sys +import tempfile +import time +import unittest +from unittest.mock import Mock, patch + + +SCRIPT = Path(__file__).with_name("run_native_tests.py") +SPEC = importlib.util.spec_from_file_location("native_test_runner", SCRIPT) +HARNESS = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(HARNESS) + + +class NativeTestRunnerTests(unittest.TestCase): + def setUp(self): + self.temporary = tempfile.TemporaryDirectory(prefix="native-test-supervisor-") + self.addCleanup(self.temporary.cleanup) + self.root = Path(self.temporary.name) + self.sample = self.root / "sample" + self.sample.write_text("#!" + sys.executable + "\n" + + "import sys\nfrom pathlib import Path\n" + "Path(sys.argv[-1]).write_text('fixture stack for ' + sys.argv[1])\n") + self.sample.chmod(0o700) + + def runner(self, program, **overrides): + options = dict(timeout=5, silence=60, sample_tool=str(self.sample), + sample_timeout=0.3, sample_seconds=0.01, + cleanup_grace=0.2, console=io.BytesIO()) + options.update(overrides) + return HARNESS.Runner([sys.executable, "-c", program], + self.root / "diagnostics", **options) + + def receipt(self): + return json.loads((self.root / "diagnostics/receipt.json").read_text()) + + def test_cli_preserves_exact_arguments_output_and_failure_exit(self): + command = [sys.executable, "-c", + "import json,sys; print(json.dumps(sys.argv[1:])); sys.exit(7)", + "argument with spaces", "--filter", "A|B", "--skip", "C"] + result = subprocess.run( + [sys.executable, str(SCRIPT), "--diagnostics-dir", str(self.root / "diagnostics"), + "--timeout-seconds", "5", "--", *command], + capture_output=True, timeout=10) + self.assertEqual(result.returncode, 7, result.stderr) + self.assertEqual(json.loads(result.stdout), command[3:]) + receipt = self.receipt() + self.assertEqual(receipt["command"], command) + self.assertEqual(receipt["child_exit_code"], 7) + self.assertEqual(receipt["cleanup"]["remaining"], []) + + def test_capped_log_keeps_draining_and_streaming_all_output(self): + runner = self.runner("import sys; sys.stdout.write('x' * (2 * 1024**2) + 'FINISHED')", + log_limit=127) + self.assertEqual(runner.run(), 0) + receipt = self.receipt() + self.assertEqual(receipt["saved_log_bytes"], 127) + self.assertEqual((self.root / "diagnostics/output.log").stat().st_size, 127) + self.assertEqual(receipt["output_bytes"], 2 * 1024**2 + 8) + self.assertTrue(receipt["log_truncated"]) + self.assertTrue(runner.console.getvalue().endswith(b"FINISHED")) + self.assertTrue(receipt["reader_complete"]) + + def test_silence_collects_diagnostics_but_eventual_success_still_passes(self): + runner = self.runner("import time; time.sleep(0.7); print('completed')", silence=0.1) + self.assertEqual(runner.run(), 0) + receipt = self.receipt() + self.assertEqual(receipt["reason"], "exit") + self.assertEqual([item["reason"] for item in receipt["diagnostics"]], ["silence"]) + self.assertTrue(list((self.root / "diagnostics").glob("*.sample.txt"))) + self.assertEqual(receipt["cleanup"]["signals"], []) + + def test_timeout_cleans_detached_descendant_that_ignores_term(self): + child_program = "import signal,time; signal.signal(signal.SIGTERM,signal.SIG_IGN); time.sleep(30)" + program = ("import subprocess,sys,time; " + f"child=subprocess.Popen([sys.executable,'-c',{child_program!r}],start_new_session=True); " + "print(child.pid,flush=True); time.sleep(30)") + runner = self.runner(program, timeout=1.2) + self.assertEqual(runner.run(), 124) + receipt = self.receipt() + child_pid = int(runner.console.getvalue().strip()) + self.assertIn({"pid": child_pid, "signal": "SIGKILL"}, + [{"pid": entry["pid"], "signal": entry["signal"]} + for entry in receipt["cleanup"]["signals"]]) + self.assertEqual(receipt["cleanup"]["remaining"], []) + current = HARNESS.processes().get(child_pid) + self.assertTrue(current is None or "Z" in current["state"], current) + self.assertTrue(receipt["reader_complete"]) + + def test_unresponsive_sampler_is_bounded_and_does_not_change_timeout_exit(self): + self.sample.write_text("#!" + sys.executable + "\nimport time\ntime.sleep(30)\n") + runner = self.runner("import time; time.sleep(30)", timeout=0.3, sample_timeout=0.15) + start = time.monotonic() + self.assertEqual(runner.run(), 124) + self.assertLess(time.monotonic() - start, 3) + self.assertIn("TimeoutExpired", self.receipt()["diagnostics"][0]["samples"][0]["error"]) + self.assertEqual(self.receipt()["cleanup"]["remaining"], []) + + def test_reused_pid_is_never_signaled(self): + runner = self.runner("pass") + runner.child = Mock(pid=100) + runner.child.poll.return_value = 0 + old = {"pid": 200, "parent": 1, "group": 200, "started": "old", "state": "S", "command": "old"} + reused = dict(old, started="new", command="unrelated") + runner.known = {200: old} + with patch.object(HARNESS, "processes", return_value={200: reused}), \ + patch.object(HARNESS.os, "kill") as kill: + receipt = runner.cleanup(time.monotonic() + 2) + kill.assert_not_called() + self.assertEqual(receipt["remaining"], []) + + def test_observed_descendant_remains_owned_after_reparenting_and_exec(self): + runner = self.runner("pass") + runner.child = Mock(pid=100) + runner.child.poll.return_value = 0 + old = {"pid": 200, "parent": 100, "group": 200, "started": "same", "state": "S", "command": "launcher"} + current = dict(old, parent=1, command="swiftpm-testing") + runner.known = {200: old} + with patch.object(HARNESS, "processes", return_value={200: current}): + self.assertEqual(runner.refresh(), {200: current}) + + def test_unreadable_inventory_cannot_report_success(self): + runner = self.runner("import time; time.sleep(30)") + with patch.object(HARNESS, "processes", side_effect=OSError("fixture ps unavailable")): + self.assertEqual(runner.run(), 125) + receipt = self.receipt() + self.assertEqual(receipt["reason"], "supervisor_error") + self.assertTrue(receipt["cleanup"]["errors"]) + self.assertIsNotNone(runner.child.returncode) + with self.assertRaises(ProcessLookupError): + os.kill(runner.child.pid, 0) + self.assertTrue(any(item.get("direct_child_fallback") + for item in receipt["cleanup"]["signals"])) + + def test_empty_or_malformed_inventory_is_an_error(self): + for output in ("", "unexpected process format\n"): + with self.subTest(output=output), \ + patch.object(HARNESS.subprocess, "check_output", return_value=output), \ + self.assertRaises(RuntimeError): + HARNESS.processes() + + def test_cancellation_propagates_and_cleans_owned_child(self): + output = self.root / "diagnostics" + wrapper = subprocess.Popen( + [sys.executable, str(SCRIPT), "--diagnostics-dir", str(output), + "--timeout-seconds", "5", "--", sys.executable, "-c", + "import time; print('ready',flush=True); time.sleep(30)"], + stdout=subprocess.PIPE, stderr=subprocess.PIPE) + try: + self.assertEqual(wrapper.stdout.readline().strip(), b"ready") + wrapper.send_signal(signal.SIGTERM) + _, stderr = wrapper.communicate(timeout=10) + self.assertEqual(wrapper.returncode, 128 + signal.SIGTERM, stderr) + receipt = self.receipt() + self.assertEqual(receipt["reason"], "signal") + self.assertEqual(receipt["cleanup"]["remaining"], []) + finally: + if wrapper.poll() is None: + wrapper.kill() + wrapper.wait(timeout=2) + + def test_signal_recorded_during_success_cleanup_cannot_report_success(self): + runner = self.runner("pass") + cleanup = runner.cleanup + + def interrupted_cleanup(deadline): + runner.signals.append(signal.SIGTERM) + return cleanup(deadline) + + with patch.object(runner, "cleanup", side_effect=interrupted_cleanup): + self.assertEqual(runner.run(), 128 + signal.SIGTERM) + self.assertEqual(self.receipt()["reason"], "signal") + self.assertEqual(self.receipt()["child_exit_code"], 0) + + def test_zero_exit_first_observed_after_runtime_budget_cannot_pass(self): + runner = self.runner("pass", timeout=0.1) + refresh = runner.refresh + first = True + + def delayed_snapshot(): + nonlocal first + if first: + first = False + time.sleep(0.25) + return refresh() + + with patch.object(runner, "refresh", side_effect=delayed_snapshot): + self.assertEqual(runner.run(), 124) + self.assertEqual(self.receipt()["reason"], "timeout") + self.assertEqual(self.receipt()["child_exit_code"], 0) + + +if __name__ == "__main__": + unittest.main() diff --git a/scripts/test_codex_release_workflow.py b/scripts/test_codex_release_workflow.py index 8b1e88f..350b960 100644 --- a/scripts/test_codex_release_workflow.py +++ b/scripts/test_codex_release_workflow.py @@ -3,6 +3,7 @@ import os from pathlib import Path import re +import shutil import subprocess import sys import tempfile @@ -51,6 +52,9 @@ def test_xcode_and_package_sdk_mirrors_match_owned_fork(self): def test_release_and_linux_commands_enforce_locks_and_xcode_version(self): with tempfile.TemporaryDirectory(prefix='engram-workflow-') as temporary: root = Path(temporary) + (root / 'scripts').mkdir() + shutil.copyfile(ROOT / 'scripts/run_native_tests.py', + root / 'scripts/run_native_tests.py') bin_dir = root / 'bin' bin_dir.mkdir() log = root / 'calls.jsonl' @@ -95,6 +99,15 @@ def test_release_and_linux_commands_enforce_locks_and_xcode_version(self): self.assertEqual(len(tests), 3) self.assertTrue(any('recall_statementBudget' in call for call in tests)) self.assertTrue(any('clusters_statementBudget' in call for call in tests)) + self.assertEqual(tests[0], [ + 'swift', 'test', '--force-resolved-versions', '--skip-build', '--filter', + 'EngramTests|EngramMemoryCoreTests|EngramRealityKitTests|PositionVersionTests', + '--skip', 'PerfTests', '--skip', 'keyBERTKeywordExtraction', + '--skip', 'recall_semanticRelevanceOrdering', + '--skip', 'recall_connectedMemory_showsEdgeRelation', + '--skip', 'recall_graphTraversal_relatesToDoesNotLeakViaUnrelatedStructuralEdge', + '--skip', 'recall_statementBudget', '--skip', 'clusters_statementBudget', + ]) def test_native_mcp_gate_includes_persistence_and_existing_cases(self): body = run_block('release.yml', 'Verify real MCP lifecycle and database lock regressions')