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
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
39 changes: 39 additions & 0 deletions app/Sources/MenuBarCore/PollingCoordinator.swift
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ public actor PollingCoordinator {
/// Attempt time, distinct from success time: a persistently failing quota endpoint
/// must still respect the slower cadence.
private var lastQuotaAttempt: Date?
private var lastUsageAttempt: Date?

public init(client: ProxyClient, endpoint: ProxyEndpoint) {
self.client = client
Expand Down Expand Up @@ -174,6 +175,9 @@ public actor PollingCoordinator {
}

if popoverOpen {
let previousOpenCodeGoActivity = Set(snapshot.activity?.activities.compactMap { activity in
activity.provider?.lowercased() == "opencode-go" ? activity.id : nil
} ?? [])
await refreshActivity(cycle: cycle)
guard isCurrent(cycle) else {
refreshInFlight = false
Expand All @@ -184,6 +188,12 @@ public actor PollingCoordinator {
// tick would turn a rarely changing endpoint into a two-second poller.
if includeHeavy { await refreshOnOpen(cycle: cycle) }

let currentOpenCodeGoActivity = Set(snapshot.activity?.activities.compactMap { activity in
activity.provider?.lowercased() == "opencode-go" ? activity.id : nil
} ?? [])
let completedOpenCodeGoRequest = snapshot.activityLoaded
&& !previousOpenCodeGoActivity.isSubset(of: currentOpenCodeGoActivity)

// Rate-limit on ATTEMPT, not success, so a persistently failing endpoint is
// not retried on every two-second activity cycle.
let quotaDue = forceQuotaRefresh || (lastQuotaAttempt.map {
Expand All @@ -193,6 +203,12 @@ public actor PollingCoordinator {
lastQuotaAttempt = Date()
await refreshQuotas(cycle: cycle, forceRefresh: forceQuotaRefresh)
}
let usageDue = includeHeavy || forceQuotaRefresh || completedOpenCodeGoRequest || (lastUsageAttempt.map {
Date().timeIntervalSince($0) >= Self.heavyInterval
} ?? true)
if usageDue, isCurrent(cycle) {
await refreshOpenCodeGoUsage(cycle: cycle)
}
}

if cycle == generation { publish() }
Expand Down Expand Up @@ -292,6 +308,29 @@ public actor PollingCoordinator {
}
}

/// Usage is a completed-request ledger. Fetch on open, on a finished Go
/// request, and at the heavy cadence; never poll it every two seconds.
private func refreshOpenCodeGoUsage(cycle: Int) async {
guard isCurrent(cycle) else { return }
guard snapshot.providersLoaded else { return }
guard snapshot.providers.contains(where: { $0.name.lowercased() == "opencode-go" && $0.isEnabled }) else {
snapshot.openCodeGoModelUsage = []
snapshot.openCodeGoUsageLoaded = false
snapshot.openCodeGoUsageUpdatedAt = nil
return
}
lastUsageAttempt = Date()
guard let usage = try? await client.usageLast30Days(), isCurrent(cycle) else { return }
snapshot.openCodeGoModelUsage = usage.models.filter {
$0.provider.lowercased() == "opencode-go"
&& !$0.model.isEmpty
&& $0.requests >= 0
&& $0.totalTokens >= 0
}
snapshot.openCodeGoUsageLoaded = true
snapshot.openCodeGoUsageUpdatedAt = Date()
}

/// Readiness is an orthogonal observation. Once authenticated startup health has
/// succeeded, a public-probe transport or contract failure means "unavailable"—it
/// must not be allowed to rewrite the live process as stopped or degraded.
Expand Down
43 changes: 38 additions & 5 deletions app/Sources/MenuBarCore/ProxyClient.swift
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,28 @@ public actor ProxyClient {
try await authenticatedGet("api/codex-routing")
}

public func codexCatalogStatus() async throws -> CodexCatalogStatus {
try await authenticatedGet("api/codex-catalog/status")
}

/// Explicit panel refresh uses the same non-disruptive sync as `ccx sync`.
/// Catalog gathering can exceed ordinary management mutation timeouts.
public func syncCodexCatalog() async throws -> CodexCatalogSyncOutcome {
let data = try await authenticatedSend(
method: "POST",
path: "api/sync",
body: nil as EmptyBody?,
timeout: 90
)
let response: CodexCatalogSyncResponse
do {
response = try JSONDecoder().decode(CodexCatalogSyncResponse.self, from: data)
} catch {
throw ProxyError.decoding
}
return response.outcome
}

/// Public post-startup readiness. This request intentionally carries no management
/// credential, and accepts the endpoint's contractually meaningful 503 response for
/// `pending` and `failed` observations.
Expand Down Expand Up @@ -356,6 +378,13 @@ public actor ProxyClient {
)
}

public func usageLast30Days() async throws -> ProviderUsageEnvelope {
try await authenticatedGet(
"api/usage",
query: [URLQueryItem(name: "range", value: "30d")]
)
}

public func restart() async throws -> RestartAccepted {
let data = try await authenticatedSend(
method: "POST",
Expand Down Expand Up @@ -494,7 +523,8 @@ public actor ProxyClient {
method: String,
path: String,
query: [URLQueryItem] = [],
body: Body?
body: Body?,
timeout: TimeInterval? = nil
) async throws -> Data {
// First attempt: fresh descriptors + identity check immediately before the
// credential-bearing request.
Expand All @@ -505,7 +535,8 @@ public actor ProxyClient {
method: method,
path: path,
query: query,
body: body
body: body,
timeout: timeout
)
} catch ProxyError.unauthorized {
// Exactly one retry. Rediscovery may pick up a rotated token or replacement
Expand All @@ -516,7 +547,8 @@ public actor ProxyClient {
method: method,
path: path,
query: query,
body: body
body: body,
timeout: timeout
)
}
}
Expand All @@ -526,7 +558,8 @@ public actor ProxyClient {
method: String,
path: String,
query: [URLQueryItem],
body: Body?
body: Body?,
timeout: TimeInterval?
) async throws -> Data {
guard installation.credential?.isEmpty == false else {
throw ProxyError.authenticationUnavailable
Expand Down Expand Up @@ -558,7 +591,7 @@ public actor ProxyClient {
query: query,
body: body,
credential: credential,
timeout: method == "GET" ? 4 : 8
timeout: timeout ?? (method == "GET" ? 4 : 8)
)
}

Expand Down
107 changes: 104 additions & 3 deletions app/Sources/MenuBarCore/ProxyModels.swift
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,89 @@ public struct CodexRouteStatus: Decodable, Equatable, Sendable {
}
}

/// Read-only activation observation from `GET /api/codex-catalog/status`.
/// A restart prompt requires both a published catalog and confirmed stale workers;
/// an unknown or pending status must not tell the user a restart will fix it.
public enum CatalogReloadStatus: Equatable, Sendable {
case unknown
case current
case restartRequired(staleWorkerCount: Int)
}

/// Closed result contract from authenticated `POST /api/sync`.
public enum CodexCatalogSyncStatus: String, Decodable, Equatable, Sendable {
case applied
case skipped
case refused
}

public enum CodexCatalogSyncSkipReason: String, Decodable, Equatable, Sendable {
case desiredDisabled = "desired_disabled"
case externalProvider = "external_provider"
}

public struct CodexCatalogSyncResponse: Decodable, Equatable, Sendable {
public let status: CodexCatalogSyncStatus
public let ok: Bool
public let skippedReason: CodexCatalogSyncSkipReason?
public let warning: String?

public var outcome: CodexCatalogSyncOutcome {
guard ok else { return .failed }
switch status {
case .applied:
return .applied(warning: warning)
case .skipped:
guard let skippedReason else { return .failed }
return .skipped(skippedReason)
case .refused:
return .failed
}
}
}

public enum CodexCatalogSyncOutcome: Equatable, Sendable {
case applied(warning: String?)
case skipped(CodexCatalogSyncSkipReason)
case failed
}

public struct CodexCatalogStatus: Decodable, Sendable {
public struct Activation: Decodable, Sendable {
public struct State: Decodable, Sendable {
public let status: String
}

public struct Workers: Decodable, Sendable {
public let status: String
public let staleCount: Int
}

public let schemaVersion: Int
public let catalog: State
public let routing: State
public let workers: Workers
}

public let activation: Activation

public var reloadStatus: CatalogReloadStatus {
guard activation.schemaVersion == 1,
activation.catalog.status == "current",
activation.routing.status == "current"
else { return .unknown }

switch activation.workers.status {
case "reload_required" where activation.workers.staleCount > 0:
return .restartRequired(staleWorkerCount: activation.workers.staleCount)
case "current", "not_running":
return .current
default:
return .unknown
}
}
}

/// `GET /api/startup-health`
public struct StartupHealth: Decodable, Equatable, Sendable {
public let status: String
Expand Down Expand Up @@ -265,13 +348,13 @@ public struct QuotaWindow: Decodable, Equatable, Sendable {
public let resetAt: Double?
}

/// A published provider cap paired with observations from this local CodexCommander usage
/// log. This is reference data, not a provider-reported balance or remaining percent.
/// Observations from this local CodexCommander usage log. Older proxy versions may
/// include a published cap, but it is not a provider-reported balance or current limit.
public struct QuotaReferenceWindow: Decodable, Equatable, Sendable {
public let id: String
public let label: String
public let windowSeconds: Double
public let publishedLimitUsd: Double
public let publishedLimitUsd: Double?
public let observedSpendUsd: Double?
public let observedTokens: Int64
public let observedRequests: Int
Expand Down Expand Up @@ -381,6 +464,24 @@ public struct ProviderQuotaEnvelope: Decodable, Equatable, Sendable {
public let availability: [ProviderQuotaAvailability]
}

/// Completed local requests in the management usage log. These are observations,
/// not an OpenCode Go account balance or a live token counter.
public struct ProviderModelUsage: Decodable, Equatable, Sendable {
public let provider: String
public let model: String
public let requests: Int
public let measuredRequests: Int
public let totalTokens: Int64
public let inputTokens: Int64
public let outputTokens: Int64
public let estimatedCostUsd: Double?
}

public struct ProviderUsageEnvelope: Decodable, Equatable, Sendable {
public let generatedAt: Double
public let models: [ProviderModelUsage]
}

public enum AgentActivityRole: String, Decodable, Equatable, Sendable {
case primary
case subagent
Expand Down
10 changes: 10 additions & 0 deletions app/Sources/MenuBarCore/ProxySnapshot.swift
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,10 @@ public struct ProxySnapshot: Equatable, Sendable {
public var quotaAvailability: [ProviderQuotaAvailability]
public var activity: AgentActivitySnapshot?
public var providers: [ProviderSummary]
/// Completed OpenCode Go requests over the last 30 days, from /api/usage.
public var openCodeGoModelUsage: [ProviderModelUsage]
public var openCodeGoUsageLoaded: Bool
public var openCodeGoUsageUpdatedAt: Date?
public var lastUpdated: Date?
public var consecutiveFailures: Int
/// Remembered from the last successful health read, so a stopped proxy can still
Expand All @@ -149,6 +153,9 @@ public struct ProxySnapshot: Equatable, Sendable {
quotaAvailability: [ProviderQuotaAvailability] = [],
activity: AgentActivitySnapshot? = nil,
providers: [ProviderSummary] = [],
openCodeGoModelUsage: [ProviderModelUsage] = [],
openCodeGoUsageLoaded: Bool = false,
openCodeGoUsageUpdatedAt: Date? = nil,
lastUpdated: Date? = nil,
consecutiveFailures: Int = 0,
lastKnownStartCommand: String? = nil,
Expand All @@ -166,6 +173,9 @@ public struct ProxySnapshot: Equatable, Sendable {
self.quotaAvailability = quotaAvailability
self.activity = activity
self.providers = providers
self.openCodeGoModelUsage = openCodeGoModelUsage
self.openCodeGoUsageLoaded = openCodeGoUsageLoaded
self.openCodeGoUsageUpdatedAt = openCodeGoUsageUpdatedAt
self.lastUpdated = lastUpdated
self.consecutiveFailures = consecutiveFailures
self.lastKnownStartCommand = lastKnownStartCommand
Expand Down
Loading
Loading