315 lines
16 KiB
Swift
315 lines
16 KiB
Swift
import Foundation
|
|
import Testing
|
|
import NucleicProtocol
|
|
|
|
@testable import NucleicCore
|
|
|
|
/// Validates the real authority wiring: a live `AppStore` standing in as `SyncHostBridge`,
|
|
/// projected to a `SyncClient` over an in-memory pipe. Complements `SyncHostTests` (which
|
|
/// uses a fake bridge) by exercising the actual session→wire mapping and snapshot path.
|
|
@MainActor
|
|
@Suite("AppStore sync bridge", .serialized, .isolatedContainerSettings)
|
|
struct AppStoreSyncBridgeTests {
|
|
private func makeStore(repo: GitTestRepo) -> AppStore {
|
|
let database = try! GRDBMetadataStore(path: nil)
|
|
let worktrees = GitWorktreeManager(now: { Date(timeIntervalSince1970: 1_700_000_000) })
|
|
let transcriptsDir = URL(fileURLWithPath: (repo.container as NSString)
|
|
.appendingPathComponent("transcripts"))
|
|
return AppStore(
|
|
database: database, worktrees: worktrees, transcriptsDir: transcriptsDir,
|
|
now: { Date(timeIntervalSince1970: 1_700_000_000) }
|
|
) { session in
|
|
let wtPath = session.worktreePath ?? ""
|
|
return ScriptedBackend { e, _ in
|
|
e.emit(.sessionStarted(SessionStarted(
|
|
backendSessionID: "be", model: "m", cwd: wtPath, toolNames: ["Write"])))
|
|
e.emit(.assistantText(TextChunk(messageID: "m1", text: "on it", isPartial: false)))
|
|
e.emit(.turnCompleted(TurnCompleted(stopReason: "end_turn", usage: nil)))
|
|
e.emit(.runFinished(RunFinished(outcome: .completed)))
|
|
}
|
|
}
|
|
}
|
|
|
|
@Test func phoneSeesRealSessionsAndSnapshot() async throws {
|
|
let repo = try await GitTestRepo()
|
|
defer { repo.cleanup() }
|
|
let store = makeStore(repo: repo)
|
|
store.syncHostName = "Andrew's Mac"
|
|
|
|
let project = await store.addProject(name: "demo", rootPath: repo.root, defaultBranch: "main")
|
|
let projectID = try #require(project?.id)
|
|
let sessionID = try await store.createSession(
|
|
in: store.project(projectID)!, title: "make a file", prompt: "go")
|
|
store.openSessionID = sessionID
|
|
await store.awaitOpenSessionSettled()
|
|
|
|
// Stand the sync server up with the live store as its authority.
|
|
let hostIdentity = DeviceIdentity()
|
|
let host = SyncHost(identity: hostIdentity, bridge: store, store: InMemoryPairedDeviceStore())
|
|
let listener = ManualListener()
|
|
try await host.start(listener: listener)
|
|
let (clientCh, serverCh) = MemoryChannel.pair()
|
|
listener.push(serverCh)
|
|
let payload = await host.beginPairing()
|
|
|
|
let client = SyncClient(
|
|
channel: clientCh, identity: DeviceIdentity(), hostStaticKey: payload.hostStaticKey,
|
|
mode: .pair(secret: payload.pairingSecret), deviceID: "iphone-1", deviceLabel: "iPhone")
|
|
let recorder = EventRecorder()
|
|
await recorder.consume(client.start())
|
|
|
|
let ready = await recorder.waitFor { if case .ready = $0 { return true } else { return false } }
|
|
guard case .ready(let welcome) = ready else { Issue.record("no welcome"); return }
|
|
#expect(welcome.host.hostName == "Andrew's Mac")
|
|
|
|
await client.send(.listSessions)
|
|
let list = await recorder.waitFor { if case .sessionList = $0 { return true } else { return false } }
|
|
guard case .sessionList(let summaries) = list else { Issue.record("no list"); return }
|
|
#expect(summaries.count == 1)
|
|
let summary = try #require(summaries.first)
|
|
#expect(summary.sessionID == sessionID)
|
|
#expect(summary.projectName == "demo")
|
|
#expect(summary.title == "make a file")
|
|
#expect(summary.status == .awaitingInput)
|
|
|
|
await client.send(.subscribe(Subscribe(sessionID: sessionID, sinceSeq: nil, verbosity: .full)))
|
|
let snap = await recorder.waitFor { if case .snapshot = $0 { return true } else { return false } }
|
|
guard case .snapshot(let snapshot) = snap else { Issue.record("no snapshot"); return }
|
|
#expect(snapshot.summary.sessionID == sessionID)
|
|
// The transcript tail came across: the agent's assistant text is in there.
|
|
#expect(snapshot.recentEvents.contains {
|
|
if case .assistantText(let c) = $0.kind { return c.text == "on it" } else { return false }
|
|
})
|
|
|
|
await host.stop()
|
|
}
|
|
|
|
// MARK: Full-transcript mesh sync — serving side
|
|
|
|
@Test func bridgeAdvertisesTranscriptSyncCapability() async throws {
|
|
let repo = try await GitTestRepo()
|
|
defer { repo.cleanup() }
|
|
let store = makeStore(repo: repo)
|
|
#expect(store.capabilities.canSyncTranscripts == true)
|
|
}
|
|
|
|
@Test func transcriptEventsServesFullHistoryAndHeader() async throws {
|
|
let repo = try await GitTestRepo()
|
|
defer { repo.cleanup() }
|
|
let store = makeStore(repo: repo)
|
|
let project = try #require(await store.addProject(name: "demo", rootPath: repo.root, defaultBranch: "main"))
|
|
let sessionID = try await store.createSession(in: project, title: "make a file", prompt: "go")
|
|
store.openSessionID = sessionID
|
|
await store.awaitOpenSessionSettled()
|
|
|
|
// afterSeq: 0 ⇒ the WHOLE transcript (not a 200 tail) + the header, so a requester can
|
|
// build a faithful mirror.
|
|
let full = try #require(await store.transcriptEvents(sessionID, afterSeq: 0))
|
|
#expect(full.header != nil)
|
|
#expect(full.header?.sessionID == sessionID)
|
|
#expect(full.events.count >= 2)
|
|
#expect(full.events.contains {
|
|
if case .assistantText(let c) = $0.kind { return c.text == "on it" } else { return false }
|
|
})
|
|
// Events are canonically seq-ordered and contiguous from 1.
|
|
#expect(full.events.map(\.seq) == Array(1...UInt64(full.events.count)))
|
|
|
|
// A warm cursor returns only the tail after it, and no header.
|
|
let tipSeq = try #require(full.events.last?.seq)
|
|
let afterFirst = try #require(await store.transcriptEvents(sessionID, afterSeq: 1))
|
|
#expect(afterFirst.header == nil)
|
|
#expect(afterFirst.events.allSatisfy { $0.seq > 1 })
|
|
#expect(afterFirst.events.last?.seq == tipSeq)
|
|
|
|
// Unknown session ⇒ nil ⇒ the handler replies `notHeld`.
|
|
#expect(await store.transcriptEvents(SessionID(rawValue: "ghost"), afterSeq: 0) == nil)
|
|
}
|
|
|
|
// MARK: Session transfer — destination side (mesh P5)
|
|
|
|
@Test func bridgeAdvertisesTransferCapability() async throws {
|
|
let repo = try await GitTestRepo()
|
|
defer { repo.cleanup() }
|
|
let store = makeStore(repo: repo)
|
|
#expect(store.capabilities.canReceiveSessionTransfer == true)
|
|
}
|
|
|
|
@Test func offerForUnknownProjectIsRejected() async throws {
|
|
let repo = try await GitTestRepo()
|
|
defer { repo.cleanup() }
|
|
let store = makeStore(repo: repo) // no projects added → nothing resolves
|
|
|
|
let offer = TransferOffer(
|
|
transferID: "x",
|
|
record: SessionTransferRecord(
|
|
sessionID: SessionID(rawValue: "s1"), backend: .claudeCode, backendSessionID: nil,
|
|
title: "t", branch: "nucleic/feature", baseSHA: "abc",
|
|
model: nil, effort: nil, lastSeq: 0, createdAt: Date(timeIntervalSince1970: 1)),
|
|
project: ProjectDescriptor(
|
|
projectID: ProjectID(rawValue: "unknown"), name: "ghost",
|
|
rootCommitSHA: "nope", normalizedRemote: nil, defaultBranch: "main"),
|
|
transcriptHeaderVersion: SessionHeader.currentVersion,
|
|
availableBaseSHAs: ["abc"], knownItems: [])
|
|
let replies = await store.receiveTransferOffer(offer, from: "host-A")
|
|
guard case .transferReject(let reject) = replies.first else {
|
|
Issue.record("expected a reject, got \(replies)"); return
|
|
}
|
|
#expect(reject.reason == .projectNotFound)
|
|
}
|
|
|
|
@Test func offerForAKnownProjectIsAcceptedWithHaveSHAs() async throws {
|
|
let repo = try await GitTestRepo()
|
|
defer { repo.cleanup() }
|
|
let store = makeStore(repo: repo)
|
|
let project = try #require(await store.addProject(name: "demo", rootPath: repo.root, defaultBranch: "main"))
|
|
let baseSHA = try await repo.revParse("HEAD")
|
|
|
|
let offer = TransferOffer(
|
|
transferID: "x",
|
|
record: SessionTransferRecord(
|
|
sessionID: SessionID(rawValue: "s1"), backend: .claudeCode, backendSessionID: nil,
|
|
title: "t", branch: "nucleic/never-seen", baseSHA: baseSHA,
|
|
model: nil, effort: nil, lastSeq: 0, createdAt: Date(timeIntervalSince1970: 1)),
|
|
project: ProjectDescriptor(
|
|
projectID: project.id, name: "demo", rootCommitSHA: nil,
|
|
normalizedRemote: nil, defaultBranch: "main"),
|
|
transcriptHeaderVersion: SessionHeader.currentVersion,
|
|
availableBaseSHAs: [baseSHA], knownItems: [
|
|
TransferItemDescriptor(kind: .transcript, byteCount: 10, sha256: "z", chunkCount: 1)])
|
|
let replies = await store.receiveTransferOffer(offer, from: "host-A")
|
|
guard case .transferAccept(let accept) = replies.first else {
|
|
Issue.record("expected accept, got \(replies)"); return
|
|
}
|
|
#expect(accept.resolvedProjectID == project.id)
|
|
// The destination shares the base commit (its own HEAD) → it's offered as a have.
|
|
#expect(accept.haveSHAs.contains(baseSHA))
|
|
// Clean up the staging the offer opened.
|
|
await store.receiveTransferCancel("x")
|
|
}
|
|
|
|
// MARK: Relaunch recovery (mesh P5)
|
|
|
|
/// `recoverInterruptedTransfers` reconciles locks a crash left behind: pre-tombstone outbound
|
|
/// locks, orphaned inbound staging, and an unrecoverable inbound `ready` (no manifest on disk)
|
|
/// are all cleared, while a tombstoned source with no reachable peer is preserved for a later
|
|
/// attempt.
|
|
@Test func recoverInterruptedTransfersCleansUpAbandonedLocks() async throws {
|
|
let repo = try await GitTestRepo()
|
|
defer { repo.cleanup() }
|
|
let db = try GRDBMetadataStore(path: nil)
|
|
let worktrees = GitWorktreeManager(now: { Date(timeIntervalSince1970: 1_700_000_000) })
|
|
let transcriptsDir = URL(fileURLWithPath: (repo.container as NSString)
|
|
.appendingPathComponent("transcripts"))
|
|
let store = AppStore(
|
|
database: db, worktrees: worktrees, transcriptsDir: transcriptsDir,
|
|
now: { Date(timeIntervalSince1970: 1_700_000_000) }
|
|
) { _ in ScriptedBackend { e, _ in e.emit(.runFinished(RunFinished(outcome: .completed))) } }
|
|
|
|
let t = Date(timeIntervalSince1970: 1_700_000_000)
|
|
func seed(_ id: String, _ session: String, _ dir: TransferDirection, _ state: SessionTransferState) async throws {
|
|
try await db.beginTransfer(SessionTransferLock(
|
|
transferID: id, sessionID: SessionID(rawValue: session), direction: dir,
|
|
peerDeviceID: "peer-B", state: state, createdAt: t, updatedAt: t))
|
|
}
|
|
try await seed("o-offer", "s1", .outbound, .offering)
|
|
try await seed("o-stream", "s2", .outbound, .streaming)
|
|
try await seed("o-await", "s3", .outbound, .awaitingReady)
|
|
try await seed("o-tomb", "s4", .outbound, .tombstoned)
|
|
try await seed("i-stage", "s5", .inbound, .staging)
|
|
try await seed("i-ready", "s6", .inbound, .ready)
|
|
|
|
await store.recoverInterruptedTransfers()
|
|
|
|
// Abandoned pre-tombstone outbound locks: cleared (the sessions live here again).
|
|
#expect(try await db.transfer(transferID: "o-offer") == nil)
|
|
#expect(try await db.transfer(transferID: "o-stream") == nil)
|
|
#expect(try await db.transfer(transferID: "o-await") == nil)
|
|
// Orphaned inbound partial receive: cleared.
|
|
#expect(try await db.transfer(transferID: "i-stage") == nil)
|
|
// Tombstoned with no reachable peer (no PeerClient here): preserved for a later attempt.
|
|
#expect(try await db.transfer(transferID: "o-tomb")?.state == .tombstoned)
|
|
// Inbound `ready` with no manifest on disk is unrecoverable → its orphan lock is cleared.
|
|
#expect(try await db.transfer(transferID: "i-ready") == nil)
|
|
#expect(store.pendingArrivedTransfers.isEmpty)
|
|
}
|
|
|
|
/// A stranded inbound `ready` transfer *with* its on-disk manifest is surfaced as a pending
|
|
/// arrival the user can activate or discard (mesh P5 relaunch recovery).
|
|
@Test func recoverInterruptedTransfersSurfacesStrandedArrival() async throws {
|
|
let repo = try await GitTestRepo()
|
|
defer { repo.cleanup() }
|
|
let db = try GRDBMetadataStore(path: nil)
|
|
let worktrees = GitWorktreeManager(now: { Date(timeIntervalSince1970: 1_700_000_000) })
|
|
let transcriptsDir = URL(fileURLWithPath: (repo.container as NSString)
|
|
.appendingPathComponent("transcripts"))
|
|
let store = AppStore(
|
|
database: db, worktrees: worktrees, transcriptsDir: transcriptsDir,
|
|
now: { Date(timeIntervalSince1970: 1_700_000_000) }
|
|
) { _ in ScriptedBackend { e, _ in e.emit(.runFinished(RunFinished(outcome: .completed))) } }
|
|
|
|
let t = Date(timeIntervalSince1970: 1_700_000_000)
|
|
let sessionID = SessionID(rawValue: "arr-session")
|
|
try await db.beginTransfer(SessionTransferLock(
|
|
transferID: "arr", sessionID: sessionID, direction: .inbound,
|
|
peerDeviceID: "laptop-A", state: .ready, createdAt: t, updatedAt: t))
|
|
// Write the staged-session manifest where the importer looks (…/transfers/<id>/).
|
|
let stagingRoot = transcriptsDir.deletingLastPathComponent()
|
|
.appendingPathComponent("transfers", isDirectory: true)
|
|
let dir = stagingRoot.appendingPathComponent("arr", isDirectory: true)
|
|
try FileManager.default.createDirectory(at: dir, withIntermediateDirectories: true)
|
|
let staged = Session(
|
|
id: sessionID, projectID: ProjectID(rawValue: "p1"), backend: .claudeCode,
|
|
title: "ported chat", status: .awaitingInput, worktreePath: "/tmp/wt",
|
|
branch: "nucleic/ported", baseSHA: "abc", lastSeq: 1,
|
|
transcriptPath: dir.appendingPathComponent("t.jsonl").path,
|
|
arrivedFromDeviceID: "laptop-A", createdAt: t, updatedAt: t)
|
|
try JSONEncoder().encode(staged).write(
|
|
to: dir.appendingPathComponent("staged-session.json"), options: .atomic)
|
|
|
|
await store.recoverInterruptedTransfers()
|
|
|
|
#expect(store.pendingArrivedTransfers.count == 1)
|
|
let prompt = try #require(store.pendingArrivedTransfers.first)
|
|
#expect(prompt.transferID == "arr")
|
|
#expect(prompt.title == "ported chat")
|
|
#expect(prompt.sourceDeviceID == "laptop-A")
|
|
// Still recoverable (the lock stays until the user resolves it).
|
|
#expect(try await db.transfer(transferID: "arr")?.state == .ready)
|
|
}
|
|
|
|
// MARK: Bulk hand-off (mesh P5)
|
|
|
|
/// `transferableSessions` lists the eligible chats (excluding archived ones), and a bulk
|
|
/// `moveSessionsToPeer` with no reachable peer fails every move, sets a rollup `lastError`, and
|
|
/// restores the sessions (nothing is lost when a hand-off can't reach the destination).
|
|
@Test func bulkHandoffListsEligibleAndRollsUpFailures() async throws {
|
|
let repo = try await GitTestRepo()
|
|
defer { repo.cleanup() }
|
|
let store = makeStore(repo: repo)
|
|
let project = try #require(await store.addProject(name: "demo", rootPath: repo.root, defaultBranch: "main"))
|
|
let s1 = try await store.createSession(in: store.project(project.id)!, title: "one", prompt: "go")
|
|
let s2 = try await store.createSession(in: store.project(project.id)!, title: "two", prompt: "go")
|
|
|
|
// Wait for both turns to quiesce (a mid-turn session isn't transferable).
|
|
var eligible = await store.transferableSessions()
|
|
var tries = 0
|
|
while eligible.count < 2, tries < 200 {
|
|
try await Task.sleep(for: .milliseconds(20))
|
|
eligible = await store.transferableSessions()
|
|
tries += 1
|
|
}
|
|
#expect(Set(eligible.map(\.id)) == [s1, s2])
|
|
|
|
// Archiving one drops it from the eligible list.
|
|
await store.setSessionArchived(s1, true)
|
|
#expect(await store.transferableSessions().map(\.id) == [s2])
|
|
|
|
// Hand-off to a peer that isn't connected: the move fails, but the session is restored.
|
|
let result = await store.moveSessionsToPeer([s2], to: "ghost-device", label: "Studio")
|
|
#expect(result.moved == 0)
|
|
#expect(result.failed == 1)
|
|
#expect(store.lastError != nil)
|
|
#expect(store.summaries(for: project.id).map(\.id).contains(s2))
|
|
}
|
|
}
|