From 5c82112d95af25e839448f1f05771f1cafad910b Mon Sep 17 00:00:00 2001 From: Bryan Font Date: Mon, 13 Jul 2026 15:06:25 -0400 Subject: [PATCH] Stream Codex lineage families --- .../Generated/CodexParserHash.generated.swift | 2 +- .../Providers/Codex/CodexLineageEngine.swift | 168 +++++++++++++ .../Codex/CodexLineageTwoPassDiscovery.swift | 227 ++++++++++++++++++ .../Vendored/CostUsage/CostUsageScanner.swift | 67 +++++- .../CodexLineageTwoPassDiscoveryTests.swift | 152 ++++++++++++ 5 files changed, 602 insertions(+), 14 deletions(-) create mode 100644 Sources/CodexBarCore/Providers/Codex/CodexLineageTwoPassDiscovery.swift create mode 100644 Tests/CodexBarTests/CodexLineageTwoPassDiscoveryTests.swift diff --git a/Sources/CodexBarCore/Generated/CodexParserHash.generated.swift b/Sources/CodexBarCore/Generated/CodexParserHash.generated.swift index 90a5d3ff38..c0e0dd2e96 100644 --- a/Sources/CodexBarCore/Generated/CodexParserHash.generated.swift +++ b/Sources/CodexBarCore/Generated/CodexParserHash.generated.swift @@ -1,5 +1,5 @@ // Generated by Scripts/regenerate-codex-parser-hash.sh. Do not edit by hand. enum CodexParserHash { - static let value = "4c426503aec9cb49" + static let value = "84c52421265e9fc8" } diff --git a/Sources/CodexBarCore/Providers/Codex/CodexLineageEngine.swift b/Sources/CodexBarCore/Providers/Codex/CodexLineageEngine.swift index 8e9e3c921c..22bd558301 100644 --- a/Sources/CodexBarCore/Providers/Codex/CodexLineageEngine.swift +++ b/Sources/CodexBarCore/Providers/Codex/CodexLineageEngine.swift @@ -24,6 +24,14 @@ enum CodexLineageEngine { let observationCount: Int } + struct PreparedDescriptorFamily: Equatable, Sendable { + let stableID: String + let inputFingerprint: Fingerprint + let descriptors: [CodexLineageTwoPassDiscovery.Descriptor] + let unresolvedParents: Set + let observationCount: Int + } + struct FamilyResult: Equatable, Sendable { let stableID: String let inputFingerprint: Fingerprint @@ -117,6 +125,134 @@ enum CodexLineageEngine { return families.sorted { $0.stableID < $1.stableID } } + static func prepareDescriptorFamilies( + descriptors: [CodexLineageTwoPassDiscovery.Descriptor], + unresolvedParents: Set = [], + checkCancellation: CostUsageScanner.CancellationCheck? = nil) throws -> [PreparedDescriptorFamily] + { + var graph = DisjointSet() + for descriptor in descriptors { + try checkCancellation?() + let owner = Self.scoped(descriptor.ownerID, scopeID: descriptor.scopeID) + graph.insert(owner) + if let metadata = Self.nonEmpty(descriptor.metadataSessionID) { + graph.union(owner, Self.scoped(metadata, scopeID: descriptor.scopeID)) + } + if let parent = Self.nonEmpty(descriptor.parentSessionID) { + graph.union(owner, Self.scoped(parent, scopeID: descriptor.scopeID)) + } + } + var descriptorsByRoot: [String: [CodexLineageTwoPassDiscovery.Descriptor]] = [:] + for descriptor in descriptors { + try checkCancellation?() + let root = graph.find(Self.scoped(descriptor.ownerID, scopeID: descriptor.scopeID)) + descriptorsByRoot[root, default: []].append(descriptor) + } + var families: [PreparedDescriptorFamily] = [] + for var familyDescriptors in descriptorsByRoot.values { + try checkCancellation?() + familyDescriptors.sort { Self.descriptorKey($0) < Self.descriptorKey($1) } + let identities = Set(familyDescriptors.flatMap { descriptor in + [descriptor.ownerID, descriptor.metadataSessionID, descriptor.parentSessionID].compactMap { value in + Self.nonEmpty(value).map { Self.scoped($0, scopeID: descriptor.scopeID) } + } + }) + let familyUnresolved = unresolvedParents.filter { + identities.contains(Self.scoped($0.sessionID, scopeID: $0.scopeID)) + } + let fingerprint = Self.descriptorFamilyFingerprint( + descriptors: familyDescriptors, + unresolvedParents: familyUnresolved) + families.append(.init( + stableID: identities.min() ?? "", + inputFingerprint: fingerprint, + descriptors: familyDescriptors, + unresolvedParents: familyUnresolved, + observationCount: familyDescriptors.reduce(0) { $0 + $1.observationCount })) + } + return families.sorted { $0.stableID < $1.stableID } + } + + static func reconcileStreaming( + families: [PreparedDescriptorFamily], + previousCache: Cache? = nil, + localTimeZone: TimeZone, + checkCancellation: CostUsageScanner.CancellationCheck? = nil, + loadDocument: ((CodexLineageTwoPassDiscovery.Descriptor) throws -> CodexLineageLedger.Document)? = nil) throws + -> Result + { + let reusable = previousCache?.algorithmVersion == Self.algorithmVersion + ? previousCache?.familiesByInputFingerprint ?? [:] + : [:] + var results: [FamilyResult] = [] + var recomputed = 0 + var reused = 0 + var peakLoadedObservations = 0 + var peakAccepted = 0 + for family in families.sorted(by: { $0.stableID < $1.stableID }) { + try checkCancellation?() + let cacheFingerprint = Self.cacheFingerprint( + input: family.inputFingerprint, + localTimeZone: localTimeZone) + if let cached = reusable[cacheFingerprint], cached.stableID == family.stableID { + results.append(cached) + reused += 1 + peakAccepted = max(peakAccepted, cached.report.acceptedObservationCount) + continue + } + var documents: [CodexLineageLedger.Document] = [] + documents.reserveCapacity(family.descriptors.count) + var loadedObservations = 0 + for descriptor in family.descriptors { + try checkCancellation?() + let document: CodexLineageLedger.Document = if let loadDocument { + try loadDocument(descriptor) + } else { + try CodexLineageTwoPassDiscovery.loadDocument( + descriptor, + checkCancellation: checkCancellation) + } + loadedObservations += document.observations.count + documents.append(document) + } + peakLoadedObservations = max(peakLoadedObservations, loadedObservations) + let conservative = try CodexLineageLedger.reconcileConservatively( + documents: documents, + unresolvedParents: family.unresolvedParents, + localTimeZone: localTimeZone, + checkCancellation: checkCancellation) + let quality = conservative.families.first?.quality ?? .primary + let result = FamilyResult( + stableID: family.stableID, + inputFingerprint: family.inputFingerprint, + familyFingerprint: Self.familyFingerprint(input: cacheFingerprint, quality: quality), + quality: quality, + report: conservative.primary) + results.append(result) + recomputed += 1 + peakAccepted = max(peakAccepted, result.report.acceptedObservationCount) + } + results.sort { $0.stableID < $1.stableID } + let report = try Self.compose(results.map(\.report), checkCancellation: checkCancellation) + let candidate = Cache( + algorithmVersion: Self.algorithmVersion, + familiesByInputFingerprint: Dictionary(uniqueKeysWithValues: results.map { + (Self.cacheFingerprint(input: $0.inputFingerprint, localTimeZone: localTimeZone), $0) + })) + try checkCancellation?() + return Result( + report: report, + families: results, + candidateCache: candidate, + diagnostics: .init( + familyCount: families.count, + recomputedFamilyCount: recomputed, + reusedFamilyCount: reused, + observationCount: families.reduce(0) { $0 + $1.observationCount }, + peakFamilyObservationCount: peakLoadedObservations, + peakAcceptedFingerprintCount: peakAccepted)) + } + static func reconcile( families: [PreparedFamily], previousCache: Cache? = nil, @@ -300,6 +436,38 @@ enum CodexLineageEngine { return digest.finalize() } + private static func descriptorFamilyFingerprint( + descriptors: [CodexLineageTwoPassDiscovery.Descriptor], + unresolvedParents: Set) -> Fingerprint + { + var digest = DigestBuilder() + digest.append(contentsOf: ["codex-lineage-descriptor-family", String(Self.algorithmVersion)]) + for descriptor in descriptors { + digest.append(contentsOf: [ + "descriptor", descriptor.scopeID, Self.canonical(descriptor.ownerID), + descriptor.metadataSessionID.map(Self.canonical) ?? "", + descriptor.parentSessionID.map(Self.canonical) ?? "", + String(descriptor.incompleteObservationCount), String(descriptor.observationCount), + String(descriptor.signature.size), String(descriptor.signature.modifiedMilliseconds), + descriptor.signature.contentSHA256, + ]) + } + for parent in unresolvedParents.sorted(by: { ($0.scopeID, $0.sessionID) < ($1.scopeID, $1.sessionID) }) { + digest.append(contentsOf: ["unresolved", parent.scopeID, Self.canonical(parent.sessionID)]) + } + return digest.finalize() + } + + private static func descriptorKey(_ descriptor: CodexLineageTwoPassDiscovery.Descriptor) -> String { + [ + descriptor.scopeID, self.canonical(descriptor.ownerID), + descriptor.metadataSessionID.map(self.canonical) ?? "", + descriptor.parentSessionID.map(self.canonical) ?? "", + descriptor.signature.contentSHA256, + descriptor.fileURL.standardizedFileURL.path, + ].joined(separator: "\u{0}") + } + private static func familyFingerprint( input: Fingerprint, quality: CodexLineageLedger.FamilyQuality) -> Fingerprint diff --git a/Sources/CodexBarCore/Providers/Codex/CodexLineageTwoPassDiscovery.swift b/Sources/CodexBarCore/Providers/Codex/CodexLineageTwoPassDiscovery.swift new file mode 100644 index 0000000000..16edd16396 --- /dev/null +++ b/Sources/CodexBarCore/Providers/Codex/CodexLineageTwoPassDiscovery.swift @@ -0,0 +1,227 @@ +#if canImport(CryptoKit) +import CryptoKit +#else +import Crypto +#endif +import Foundation + +/// Pass-one lineage index. It retains file identity and quality evidence, never token observations. +enum CodexLineageTwoPassDiscovery { + enum DiscoveryError: Error, Equatable { + case fileChangedDuringScan + } + + struct FileSignature: Equatable, Hashable, Sendable { + let size: Int64 + let modifiedMilliseconds: Int64 + let contentSHA256: String + } + + struct Descriptor: Equatable, Sendable { + let fileURL: URL + let ownerID: String + let metadataSessionID: String? + let parentSessionID: String? + let scopeID: String + let incompleteObservationCount: Int + let observationCount: Int + let signature: FileSignature + } + + struct Report: Equatable, Sendable { + let descriptors: [Descriptor] + let referencedParentDocumentCount: Int + let unresolvedParents: Set + let peakRetainedObservationCount: Int + } + + static func discover( + includedFiles: [URL], + roots: [URL], + checkCancellation: CostUsageScanner.CancellationCheck? = nil) throws -> Report + { + var locator = ParentLocator(roots: roots, checkCancellation: checkCancellation) + var descriptors: [Descriptor] = [] + var known: Set = [] + var pending: [ScopedIdentity] = [] + var unresolved: Set = [] + var seenPaths: Set = [] + var referencedParents = 0 + + func remember(_ descriptor: Descriptor) { + descriptors.append(descriptor) + known.insert(.init(scopeID: descriptor.scopeID, sessionID: Self.canonical(descriptor.ownerID))) + if let metadata = Self.nonEmpty(descriptor.metadataSessionID) { + known.insert(.init(scopeID: descriptor.scopeID, sessionID: Self.canonical(metadata))) + } + if let parent = Self.nonEmpty(descriptor.parentSessionID) { + pending.append(.init(scopeID: descriptor.scopeID, sessionID: Self.canonical(parent))) + } + } + + for fileURL in includedFiles.sorted(by: { $0.path < $1.path }) { + try checkCancellation?() + guard seenPaths.insert(fileURL.standardizedFileURL.path).inserted else { continue } + try remember(Self.describe(fileURL: fileURL, checkCancellation: checkCancellation)) + } + + var nextParent = 0 + while nextParent < pending.count { + try checkCancellation?() + let identity = pending[nextParent] + nextParent += 1 + let unresolvedIdentity = CodexLineageLedger.ParentIdentity( + scopeID: identity.scopeID, + sessionID: identity.sessionID) + guard !known.contains(identity), !unresolved.contains(unresolvedIdentity) else { continue } + guard let matches = try locator.fileURLs(for: identity) else { + unresolved.insert(unresolvedIdentity) + continue + } + var found = false + for fileURL in matches { + try checkCancellation?() + guard seenPaths.insert(fileURL.standardizedFileURL.path).inserted else { continue } + let descriptor = try Self.describe(fileURL: fileURL, checkCancellation: checkCancellation) + let owner = Self.canonical(descriptor.ownerID) + let metadata = descriptor.metadataSessionID.map(Self.canonical) + guard owner == identity.sessionID || metadata == identity.sessionID else { continue } + remember(descriptor) + referencedParents += 1 + found = true + } + if !found { + unresolved.insert(unresolvedIdentity) + } + } + + return Report( + descriptors: descriptors.sorted { $0.fileURL.path < $1.fileURL.path }, + referencedParentDocumentCount: referencedParents, + unresolvedParents: unresolved, + peakRetainedObservationCount: 0) + } + + static func describe( + fileURL: URL, + checkCancellation: CostUsageScanner.CancellationCheck? = nil) throws -> Descriptor + { + // Validate both sides of the parse. A post-parse signature alone can bless a + // mixed read when a rollout is replaced while the parser has the file open. + let initialSignature = try Self.signature(fileURL: fileURL, checkCancellation: checkCancellation) + let summary = try CostUsageScanner.parseCodexLineageDocumentSummary( + fileURL: fileURL, + checkCancellation: checkCancellation) + let finalSignature = try Self.signature(fileURL: fileURL, checkCancellation: checkCancellation) + guard initialSignature == finalSignature else { throw DiscoveryError.fileChangedDuringScan } + return Descriptor( + fileURL: fileURL, + ownerID: summary.ownerID, + metadataSessionID: summary.metadataSessionID, + parentSessionID: summary.parentSessionID, + scopeID: summary.scopeID, + incompleteObservationCount: summary.incompleteObservationCount, + observationCount: summary.observationCount, + signature: finalSignature) + } + + static func loadDocument( + _ descriptor: Descriptor, + checkCancellation: CostUsageScanner.CancellationCheck? = nil) throws -> CodexLineageLedger.Document + { + guard try self.signature(fileURL: descriptor.fileURL, checkCancellation: checkCancellation) == descriptor + .signature + else { throw DiscoveryError.fileChangedDuringScan } + let document = try CostUsageScanner.parseCodexLineageDocument( + fileURL: descriptor.fileURL, + checkCancellation: checkCancellation) + guard try Self.signature(fileURL: descriptor.fileURL, checkCancellation: checkCancellation) == descriptor + .signature + else { throw DiscoveryError.fileChangedDuringScan } + return document + } + + private static func signature( + fileURL: URL, + checkCancellation: CostUsageScanner.CancellationCheck?) throws -> FileSignature + { + let values = try fileURL.resourceValues(forKeys: [.fileSizeKey, .contentModificationDateKey]) + let handle = try FileHandle(forReadingFrom: fileURL) + defer { try? handle.close() } + var hasher = SHA256() + while true { + try checkCancellation?() + let data = try handle.read(upToCount: 256 * 1024) ?? Data() + guard !data.isEmpty else { break } + hasher.update(data: data) + } + let digest = hasher.finalize().map { String(format: "%02x", $0) }.joined() + return FileSignature( + size: Int64(values.fileSize ?? 0), + modifiedMilliseconds: Int64((values.contentModificationDate?.timeIntervalSince1970 ?? 0) * 1000), + contentSHA256: digest) + } + + private struct ScopedIdentity: Equatable, Hashable { + let scopeID: String + let sessionID: String + } + + private struct ParentLocator { + let roots: [URL] + let checkCancellation: CostUsageScanner.CancellationCheck? + var indexed = false + var filesByIdentity: [ScopedIdentity: Set] = [:] + + mutating func fileURLs(for identity: ScopedIdentity) throws -> [URL]? { + if !self.indexed { + try self.index() + } + guard let matches = self.filesByIdentity[identity], !matches.isEmpty else { return nil } + let owners = Set(matches.compactMap(CostUsageScanner.codexRolloutOwnerID(fileURL:))) + guard matches.count == 1 || owners.count == 1 else { return nil } + return matches.sorted { $0.path < $1.path } + } + + private mutating func index() throws { + self.indexed = true + for root in self.roots { + try self.checkCancellation?() + guard let enumerator = FileManager.default.enumerator( + at: root, + includingPropertiesForKeys: [.isRegularFileKey], + options: [.skipsHiddenFiles, .skipsPackageDescendants]) + else { continue } + while let fileURL = enumerator.nextObject() as? URL { + try self.checkCancellation?() + guard fileURL.pathExtension.lowercased() == "jsonl" else { continue } + let scopeID = CostUsageScanner.codexLineageScopeID(fileURL: fileURL) + if let owner = CostUsageScanner.codexRolloutOwnerID(fileURL: fileURL) { + self.filesByIdentity[ + .init(scopeID: scopeID, sessionID: CodexLineageTwoPassDiscovery.canonical(owner)), + default: [], + ].insert(fileURL) + } + if let metadata = try CostUsageScanner.parseCodexSessionIdentifier( + fileURL: fileURL, + checkCancellation: self.checkCancellation) + { + self.filesByIdentity[ + .init(scopeID: scopeID, sessionID: CodexLineageTwoPassDiscovery.canonical(metadata)), + default: [], + ].insert(fileURL) + } + } + } + } + } + + private static func canonical(_ value: String) -> String { + UUID(uuidString: value)?.uuidString.lowercased() ?? value + } + + private static func nonEmpty(_ value: String?) -> String? { + guard let value, !value.isEmpty else { return nil } + return value + } +} diff --git a/Sources/CodexBarCore/Vendored/CostUsage/CostUsageScanner.swift b/Sources/CodexBarCore/Vendored/CostUsage/CostUsageScanner.swift index 8514a9bcb2..88f724d30e 100644 --- a/Sources/CodexBarCore/Vendored/CostUsage/CostUsageScanner.swift +++ b/Sources/CodexBarCore/Vendored/CostUsage/CostUsageScanner.swift @@ -104,6 +104,16 @@ enum CostUsageScanner { let snapshots: [CodexTimestampedTotals] let observations: [CodexLineageLedger.Observation] let incompleteObservationCount: Int + let observationCount: Int + } + + struct CodexLineageDocumentSummary: Equatable, Sendable { + let ownerID: String + let metadataSessionID: String? + let parentSessionID: String? + let scopeID: String + let incompleteObservationCount: Int + let observationCount: Int } enum CodexForkBaseline { @@ -1689,6 +1699,8 @@ enum CostUsageScanner { // swiftlint:disable:next cyclomatic_complexity function_body_length private static func parseCodexTokenSnapshots( fileURL: URL, + retainEvidence: Bool = true, + suppressScanErrors: Bool = true, checkCancellation: CancellationCheck? = nil) throws -> CodexParsedTokenEvidence { var sessionId: String? @@ -1698,6 +1710,7 @@ enum CostUsageScanner { var snapshots: [CodexTimestampedTotals] = [] var observations: [CodexLineageLedger.Observation] = [] var incompleteObservationCount = 0 + var observationCount = 0 var warnedAboutUnparsedTimestamp = false func parsedSnapshotDate(timestamp: String) -> Date? { @@ -1721,19 +1734,24 @@ enum CostUsageScanner { if last == nil || total == nil { incompleteObservationCount += 1 } - let counted = accumulator.apply(last: last, total: total) - snapshots.append(CodexTimestampedTotals( - timestamp: timestamp, - date: parsedSnapshotDate(timestamp: timestamp), - totals: counted)) - if let last, let total { - observations.append(CodexLineageLedger.Observation( + if retainEvidence { + let counted = accumulator.apply(last: last, total: total) + snapshots.append(CodexTimestampedTotals( timestamp: timestamp, - model: Self.codexModelEvidence(model) - ?? Self.codexModelEvidence(currentModel) - ?? CostUsagePricing.codexUnattributedModel, - last: Self.lineageTotals(last), - total: Self.lineageTotals(total))) + date: parsedSnapshotDate(timestamp: timestamp), + totals: counted)) + } + if let last, let total { + observationCount += 1 + if retainEvidence { + observations.append(CodexLineageLedger.Observation( + timestamp: timestamp, + model: Self.codexModelEvidence(model) + ?? Self.codexModelEvidence(currentModel) + ?? CostUsagePricing.codexUnattributedModel, + last: Self.lineageTotals(last), + total: Self.lineageTotals(total))) + } } } @@ -1845,6 +1863,9 @@ enum CostUsageScanner { } catch is CancellationError { throw CancellationError() } catch { + if !suppressScanErrors { + throw error + } self.log.warning( "Codex cost usage failed while scanning parent token snapshots", metadata: ["path": fileURL.path, "error": error.localizedDescription]) @@ -1855,7 +1876,8 @@ enum CostUsageScanner { forkedFromId: forkedFromId, snapshots: snapshots, observations: observations, - incompleteObservationCount: incompleteObservationCount) + incompleteObservationCount: incompleteObservationCount, + observationCount: observationCount) } static func parseCodexLineageDocument( @@ -1864,6 +1886,7 @@ enum CostUsageScanner { { let parsed = try Self.parseCodexTokenSnapshots( fileURL: fileURL, + suppressScanErrors: false, checkCancellation: checkCancellation) return CodexLineageLedger.Document( ownerID: Self.codexRolloutOwnerID(fileURL: fileURL) ?? parsed.sessionId ?? fileURL.standardizedFileURL.path, @@ -1874,6 +1897,24 @@ enum CostUsageScanner { incompleteObservationCount: parsed.incompleteObservationCount) } + static func parseCodexLineageDocumentSummary( + fileURL: URL, + checkCancellation: CancellationCheck? = nil) throws -> CodexLineageDocumentSummary + { + let parsed = try Self.parseCodexTokenSnapshots( + fileURL: fileURL, + retainEvidence: false, + suppressScanErrors: false, + checkCancellation: checkCancellation) + return CodexLineageDocumentSummary( + ownerID: Self.codexRolloutOwnerID(fileURL: fileURL) ?? parsed.sessionId ?? fileURL.standardizedFileURL.path, + metadataSessionID: parsed.sessionId, + parentSessionID: parsed.forkedFromId, + scopeID: Self.codexLineageScopeID(fileURL: fileURL), + incompleteObservationCount: parsed.incompleteObservationCount, + observationCount: parsed.observationCount) + } + static func codexLineageScopeID(fileURL: URL) -> String { let standardized = fileURL.standardizedFileURL let components = standardized.pathComponents diff --git a/Tests/CodexBarTests/CodexLineageTwoPassDiscoveryTests.swift b/Tests/CodexBarTests/CodexLineageTwoPassDiscoveryTests.swift new file mode 100644 index 0000000000..11a750c2bf --- /dev/null +++ b/Tests/CodexBarTests/CodexLineageTwoPassDiscoveryTests.swift @@ -0,0 +1,152 @@ +import Foundation +import Testing +@testable import CodexBarCore + +struct CodexLineageTwoPassDiscoveryTests { + @Test + func `pass one retains descriptors without observations and follows exceptional archived parents`() throws { + let environment = try CostUsageTestEnvironment() + defer { environment.cleanup() } + let parentID = "11111111-1111-4111-8111-111111111111" + let childID = "22222222-2222-4222-8222-222222222222" + _ = try Self.writeRollout( + root: environment.codexArchivedSessionsRoot, + ownerID: parentID, + observations: 3) + let child = try Self.writeRollout( + root: environment.codexSessionsRoot.appendingPathComponent("2026/07/09"), + ownerID: childID, + parentID: parentID, + observations: 2) + + let report = try CodexLineageTwoPassDiscovery.discover( + includedFiles: [child], + roots: [environment.codexSessionsRoot, environment.codexArchivedSessionsRoot]) + + #expect(report.descriptors.count == 2) + #expect(report.descriptors.reduce(0) { $0 + $1.observationCount } == 5) + #expect(report.referencedParentDocumentCount == 1) + #expect(report.unresolvedParents.isEmpty) + #expect(report.peakRetainedObservationCount == 0) + } + + @Test + func `streaming reconciliation loads one family at a time and warm reuse loads none`() throws { + let environment = try CostUsageTestEnvironment() + defer { environment.cleanup() } + let first = try Self.writeRollout(root: environment.codexSessionsRoot, ownerID: Self.uuid(1), observations: 50) + let second = try Self.writeRollout(root: environment.codexSessionsRoot, ownerID: Self.uuid(2), observations: 20) + let discovery = try CodexLineageTwoPassDiscovery.discover(includedFiles: [first, second], roots: []) + let families = try CodexLineageEngine.prepareDescriptorFamilies(descriptors: discovery.descriptors) + var coldLoads = 0 + let cold = try CodexLineageEngine.reconcileStreaming( + families: families, + localTimeZone: .gmt, + loadDocument: { descriptor in + coldLoads += 1 + return try CodexLineageTwoPassDiscovery.loadDocument(descriptor) + }) + var warmLoads = 0 + let warm = try CodexLineageEngine.reconcileStreaming( + families: families, + previousCache: cold.candidateCache, + localTimeZone: .gmt, + loadDocument: { descriptor in + warmLoads += 1 + return try CodexLineageTwoPassDiscovery.loadDocument(descriptor) + }) + + #expect(coldLoads == 2) + #expect(cold.diagnostics.observationCount == 70) + #expect(cold.diagnostics.peakFamilyObservationCount == 50) + #expect(warmLoads == 0) + #expect(warm.diagnostics.reusedFamilyCount == 2) + #expect(warm.report == cold.report) + } + + @Test + func `descriptor and family fingerprints are deterministic under file permutation`() throws { + let environment = try CostUsageTestEnvironment() + defer { environment.cleanup() } + let first = try Self.writeRollout(root: environment.codexSessionsRoot, ownerID: Self.uuid(3), observations: 2) + let second = try Self.writeRollout(root: environment.codexSessionsRoot, ownerID: Self.uuid(4), observations: 2) + let forward = try CodexLineageTwoPassDiscovery.discover(includedFiles: [first, second], roots: []) + let reverse = try CodexLineageTwoPassDiscovery.discover(includedFiles: [second, first], roots: []) + let forwardFamilies = try CodexLineageEngine.prepareDescriptorFamilies(descriptors: forward.descriptors) + let reverseFamilies = try CodexLineageEngine.prepareDescriptorFamilies(descriptors: reverse.descriptors) + + #expect(forwardFamilies.map(\.inputFingerprint) == reverseFamilies.map(\.inputFingerprint)) + #expect(forwardFamilies.map(\.stableID) == reverseFamilies.map(\.stableID)) + } + + @Test + func `cancellation before streaming completion produces no candidate`() throws { + let environment = try CostUsageTestEnvironment() + defer { environment.cleanup() } + let file = try Self.writeRollout(root: environment.codexSessionsRoot, ownerID: Self.uuid(5), observations: 100) + let discovery = try CodexLineageTwoPassDiscovery.discover(includedFiles: [file], roots: []) + let families = try CodexLineageEngine.prepareDescriptorFamilies(descriptors: discovery.descriptors) + var checks = 0 + #expect(throws: CancellationError.self) { + _ = try CodexLineageEngine.reconcileStreaming( + families: families, + localTimeZone: .gmt, + checkCancellation: { + checks += 1 + if checks == 4 { + throw CancellationError() + } + }) + } + } + + @Test + func `pass two rejects a rollout changed after discovery`() throws { + let environment = try CostUsageTestEnvironment() + defer { environment.cleanup() } + let file = try Self.writeRollout( + root: environment.codexSessionsRoot, + ownerID: Self.uuid(6), + observations: 2) + let discovery = try CodexLineageTwoPassDiscovery.discover(includedFiles: [file], roots: []) + let descriptor = try #require(discovery.descriptors.first) + try "\n".append(to: file) + + #expect(throws: CodexLineageTwoPassDiscovery.DiscoveryError.fileChangedDuringScan) { + _ = try CodexLineageTwoPassDiscovery.loadDocument(descriptor) + } + } + + private static func writeRollout( + root: URL, + ownerID: String, + parentID: String? = nil, + observations: Int) throws -> URL + { + try FileManager.default.createDirectory(at: root, withIntermediateDirectories: true) + let file = root.appendingPathComponent("rollout-2026-07-09T00-00-00-\(ownerID).jsonl") + var lines = [#"{"type":"session_meta","payload":{"id":"\#(ownerID)""# + + (parentID.map { #", "forked_from_id":"\#($0)""# } ?? "") + "}}"] + lines += (0.. String { + String(format: "00000000-0000-4000-8000-%012d", value) + } +} + +extension String { + fileprivate func append(to fileURL: URL) throws { + let handle = try FileHandle(forWritingTo: fileURL) + defer { try? handle.close() } + try handle.seekToEnd() + try handle.write(contentsOf: Data(self.utf8)) + } +}