Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 12 additions & 16 deletions Sources/Cloud/VMClient+ResourceStats.swift
Original file line number Diff line number Diff line change
Expand Up @@ -3,22 +3,18 @@ import Foundation
extension VMClient {
/// All callers share revisioned resource state, including CLI and sidebar reads.
func stats(id: String) async throws -> VMStats {
try Task.checkCancellation()
let task = await resourceStats.read(machineID: id) {
try await self.fetchStats(id: id)
}
// A disappearing consumer must not cancel another panel's shared read.
let stats = try await task.value
try Task.checkCancellation()
return stats
}

private func fetchStats(id: String) async throws -> VMStats {
try await withOperation(.stats, foreground: false) {
let encodedID = try pathSegment(id, fieldName: "vm id")
let (data, http) = try await request("GET", path: "/api/vm/\(encodedID)/stats", timeoutSeconds: 30)
try ensureOK(http, data: data)
return VMStats(json: try decodeJSONObject(data))
let read = await resourceStats.beginRead(machineID: id)
do {
let stats = try await withOperation(.stats, foreground: false) {
let encodedID = try pathSegment(id, fieldName: "vm id")
let (data, http) = try await request("GET", path: "/api/vm/\(encodedID)/stats", timeoutSeconds: 30)
try ensureOK(http, data: data)
return VMStats(json: try decodeJSONObject(data))
}
return await resourceStats.finishRead(read, stats: stats)
} catch {
_ = await resourceStats.finishRead(read, stats: nil)
throw error
}
}

Expand Down
1 change: 0 additions & 1 deletion Sources/Cloud/VMResourceStatsStore+Entry.swift
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,5 @@ extension VMResourceStatsStore {
var readSequence: UInt64 = 0
var acceptedSequence: UInt64 = 0
var stats: VMStats?
var readTask: Task<VMStats, Error>?
}
}
27 changes: 0 additions & 27 deletions Sources/Cloud/VMResourceStatsStore.swift
Original file line number Diff line number Diff line change
Expand Up @@ -19,31 +19,6 @@ final class VMResourceStatsStore {

func stats(for machineID: String) -> VMStats? { entries[machineID]?.stats }

/// Concurrent consumers share only an active request, never a cached reading.
/// Resize, reset, and removal discard the entry's task before another read joins.
func read(
machineID: String,
fetch: @escaping @Sendable () async throws -> VMStats
) -> Task<VMStats, Error> {
if let task = entries[machineID]?.readTask { return task }
let request = beginRead(machineID: machineID)
let task = Task { @MainActor in
defer {
if entries[machineID]?.revision == request.revision {
entries[machineID]?.readTask = nil
}
}
do {
return finishRead(request, stats: try await fetch())
} catch {
finishRead(request, stats: nil)
throw error
}
}
entries[machineID]?.readTask = task
return task
}

func beginRead(machineID: String) -> Request {
var entry = entry(for: machineID)
entry.readSequence &+= 1
Expand Down Expand Up @@ -76,7 +51,6 @@ final class VMResourceStatsStore {
func beginResize(machineID: String) -> Request {
var entry = entry(for: machineID)
entry.revision = UUID()
entry.readTask = nil
entry.resizing = true
entry.stats = .unavailable(at: now())
entries[machineID] = entry
Expand All @@ -88,7 +62,6 @@ final class VMResourceStatsStore {
guard var entry = entries[request.machineID], entry.revision == request.revision else { return }
// Also fence reads started while the resize was in progress.
entry.revision = UUID()
entry.readTask = nil
entry.resizing = false
entry.stats = stats ?? .unavailable(at: now())
entries[request.machineID] = entry
Expand Down
98 changes: 0 additions & 98 deletions cmuxTests/VMResourceStatsStoreTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -222,102 +222,4 @@ struct VMResourceStatsStoreTests {
}
#expect(store.stats(for: "new") != nil)
}

@Test func concurrentConsumersShareOneFetchButLaterReadsAreFresh() async throws {
let store = VMResourceStatsStore(now: { self.time })
let reading = stats(memory: 8192, disk: 32768)
let first = store.read(machineID: "vm") { reading }
let second = store.read(machineID: "vm") {
Issue.record("A concurrent consumer started a duplicate fetch")
return reading
}
#expect(try await first.value == reading)
#expect(try await second.value == reading)
let fresh = stats(memory: 16384, disk: 65536)
let next = store.read(machineID: "vm") { fresh }
#expect(try await next.value == fresh)
}

@Test func differentMachinesFetchIndependently() async throws {
let store = VMResourceStatsStore(now: { self.time })
let firstReading = stats(memory: 8192, disk: 32768)
let secondReading = stats(memory: 16384, disk: 65536)
let first = store.read(machineID: "a") { firstReading }
let second = store.read(machineID: "b") { secondReading }
#expect(try await first.value == firstReading)
#expect(try await second.value == secondReading)
}

@Test func failedSharedFetchDoesNotPoisonTheNextRead() async throws {
let store = VMResourceStatsStore(now: { self.time })
let reading = stats(memory: 8192, disk: 32768)
store.finishRead(store.beginRead(machineID: "vm"), stats: reading)
let first = store.read(machineID: "vm") { throw URLError(.badServerResponse) }
let second = store.read(machineID: "vm") {
Issue.record("A concurrent consumer started a duplicate failed fetch")
return reading
}
for task in [first, second] {
do {
_ = try await task.value
Issue.record("The fetch should fail for each consumer")
} catch let error as URLError {
#expect(error.code == .badServerResponse)
}
}
#expect(store.stats(for: "vm")?.memoryTotalMb == reading.memoryTotalMb)
#expect(store.stats(for: "vm")?.cpuPercent == nil)
let retry = store.read(machineID: "vm") { reading }
#expect(try await retry.value == reading)
}

@Test(arguments: ["reset", "removal", "resize"])
func invalidationSeparatesRequestsAndOldCompletionCannotClearNewFetch(_ invalidation: String) async throws {
let store = VMResourceStatsStore(now: { self.time })
let oldGate = AsyncStream<VMStats>.makeStream()
let newGate = AsyncStream<VMStats>.makeStream()
let oldReading = stats(memory: 8192, disk: 32768)
let newReading = stats(memory: 16384, disk: 65536)
let old = store.read(machineID: "vm") {
var iterator = oldGate.stream.makeAsyncIterator()
return await iterator.next()!
}
switch invalidation {
case "reset": store.reset()
case "removal": store.retain(machineIDs: [], token: store.beginRetention())
default: store.finishResize(store.beginResize(machineID: "vm"), stats: newReading)
}
let current = store.read(machineID: "vm") {
var iterator = newGate.stream.makeAsyncIterator()
return await iterator.next()!
}
oldGate.continuation.yield(oldReading)
_ = try await old.value
#expect(store.stats(for: "vm") != oldReading)
let joined = store.read(machineID: "vm") {
Issue.record("An obsolete fetch cleared the current shared request")
return oldReading
}
newGate.continuation.yield(newReading)
#expect(try await current.value == newReading)
#expect(try await joined.value == newReading)
}

@Test func readsDuringResizeAreNotReusedAfterResizeFinishes() async throws {
let store = VMResourceStatsStore(now: { self.time })
let oldGate = AsyncStream<VMStats>.makeStream()
let before = stats(memory: 8192, disk: 32768)
let after = stats(memory: 16384, disk: 65536)
let mutation = store.beginResize(machineID: "vm")
let during = store.read(machineID: "vm") {
var iterator = oldGate.stream.makeAsyncIterator()
return await iterator.next()!
}
store.finishResize(mutation, stats: after)
let fresh = store.read(machineID: "vm") { after }
#expect(try await fresh.value == after)
oldGate.continuation.yield(before)
#expect(try await during.value == after)
#expect(store.stats(for: "vm") == after)
}
}
Loading