Skip to content
Draft
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
89 changes: 51 additions & 38 deletions Sources/ContextCore/LocalObservabilityClient.swift
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ public struct LocalObservabilityClient: Sendable {
agents: AgentRuntime.allCases.map { runtime in
switch runtime {
case .codex:
AgentTelemetry(runtime: .codex, connected: true, source: "Local OTEL · Prometheus · Tempo", signals: signalResult.signals)
AgentTelemetry(runtime: .codex, connected: signalResult.signals.values.contains { $0.value != nil }, source: "Local OTEL · Prometheus · Tempo", signals: signalResult.signals)
case .claude:
claudeResult.agent
default:
Expand All @@ -58,42 +58,44 @@ public struct LocalObservabilityClient: Sendable {
sourceSchema: "codex-native-otel+claude-code-metrics/prometheus-tempo-v1"
)
} catch {
return .unavailable(reason: "Local telemetry error: \(error.localizedDescription)")
return .unavailable(reason: "Local telemetry error: \(Self.safeError(error))")
}
}

private func loadAggregateSignalsOutcome() async -> SignalsOutcome {
do { return SignalsOutcome(signals: try await loadAggregateSignals(), error: nil) }
catch {
let reason = error.localizedDescription
let reason = Self.safeError(error)
return SignalsOutcome(signals: unavailableSignals(reason: reason), error: reason)
}
}

private func loadAggregateSignals() async throws -> [TelemetrySignal: TelemetryValue] {
async let contextTokens = queryScalar("sum(increase(codex_turn_token_usage_sum{token_type=\"total\"}[24h])) or vector(0)")
async let nativeToolCalls = queryScalar("sum(increase(codex_tool_call_total[24h])) or vector(0)")
async let apiRequests = queryScalar("sum(increase(codex_api_request_total[24h])) or vector(0)")
async let websocketRequests = queryScalar("sum(increase(codex_websocket_request_total[24h])) or vector(0)")
async let proxyRequests = queryScalar("sum(increase(codex_turn_network_proxy_total[24h])) or vector(0)")
let measuredNetworkCalls = try await apiRequests + websocketRequests + proxyRequests
async let tokens = sampleOutcome(expression: "sum(increase(codex_turn_token_usage_sum{token_type=\"total\"}[24h]))")
async let tools = sampleOutcome(expression: "sum(increase(codex_tool_call_total[24h]))")
async let api = sampleOutcome(expression: "sum(increase(codex_api_request_total[24h]))")
async let websocket = sampleOutcome(expression: "sum(increase(codex_websocket_request_total[24h]))")
async let proxy = sampleOutcome(expression: "sum(increase(codex_turn_network_proxy_total[24h]))")
let (tokenResult, toolResult, apiResult, websocketResult, proxyResult) = await (tokens, tools, api, websocket, proxy)
func metric(_ result: SamplesOutcome, unit: String, note: String) -> TelemetryValue {
let value = result.samples.first?.value
return TelemetryValue(value: value?.rounded(), unit: unit, quality: value == nil ? .unavailable : .derived,
note: result.error ?? (value == nil ? "No range samples returned; missing or one-sample series cannot establish zero activity." : note))
}
let network = [apiResult, websocketResult, proxyResult]
let networkValues = network.compactMap { $0.samples.first?.value }
let completeNetwork = networkValues.count == network.count
let networkTotal = networkValues.reduce(0, +)
return [
.contextTokens: TelemetryValue(
value: try await contextTokens.rounded(), unit: "estimated tokens / 24h", quality: .derived,
note: "PromQL increase of the total-token counter. Short-lived one-sample series can be missed; components are shown separately."
),
.toolCalls: TelemetryValue(
value: try await nativeToolCalls.rounded(), unit: "estimated invocations / 24h", quality: .derived,
note: "PromQL increase of Codex tool calls. Short-lived series can be missed; MCP calls can overlap."
),
.networkCalls: TelemetryValue(
value: measuredNetworkCalls.rounded(), unit: "estimated logical events / 24h", quality: .derived,
note: "PromQL increase across API, websocket-request, and proxy events. Categories can overlap and short-lived series can be missed."
),
.internetUsage: TelemetryValue(
value: nil, unit: "bytes", quality: .unavailable,
note: "The current OTLP stream does not measure bytes transferred."
),
.contextTokens: metric(tokenResult, unit: "estimated tokens / 24h",
note: "PromQL increase of the total-token counter. Short-lived one-sample series can be missed; components are shown separately."),
.toolCalls: metric(toolResult, unit: "estimated invocations / 24h",
note: "PromQL increase of Codex tool calls. Short-lived series can be missed; MCP calls can overlap."),
.networkCalls: TelemetryValue(value: completeNetwork && networkTotal.isFinite ? networkTotal.rounded() : nil,
unit: "estimated logical events / 24h", quality: completeNetwork && networkTotal.isFinite ? .derived : .unavailable,
note: completeNetwork ? "PromQL increase across API, websocket-request, and proxy events. Categories can overlap." : "One or more logical-event series are unavailable; a complete total cannot be established."),
.internetUsage: TelemetryValue(value: nil, unit: "bytes", quality: .unavailable,
note: "The current OTLP stream does not measure bytes transferred."),
]
}

Expand All @@ -117,7 +119,7 @@ public struct LocalObservabilityClient: Sendable {
sessions: try await sessions, costs: try await costs,
presence: try await presence)
} catch {
let reason = error.localizedDescription
let reason = Self.safeError(error)
return ClaudeLoad(agent: adapter(for: .claude), sections: [], error: reason)
}
}
Expand Down Expand Up @@ -239,7 +241,7 @@ public struct LocalObservabilityClient: Sendable {

private func sampleOutcome(expression: String) async -> SamplesOutcome {
do { return SamplesOutcome(samples: try await querySamples(expression), error: nil) }
catch { return SamplesOutcome(samples: [], error: error.localizedDescription) }
catch { return SamplesOutcome(samples: [], error: Self.safeError(error)) }
}

private func items(_ outcome: SamplesOutcome?, label: String, unit: String,
Expand Down Expand Up @@ -283,7 +285,7 @@ public struct LocalObservabilityClient: Sendable {
URLQueryItem(name: "end", value: String(Int(now.timeIntervalSince1970))),
]
let (data, response) = try await session.data(from: components.url!)
try requireSuccess(response: response, data: data, source: "Tempo session search")
try requireSuccess(response: response, source: "Tempo session search")
let payload = try JSONDecoder().decode(TempoSearchResponse.self, from: data)
let runs = payload.traces.compactMap { trace -> OTelRun? in
guard let nanos = Double(trace.startTimeUnixNano) else { return nil }
Expand All @@ -298,7 +300,7 @@ public struct LocalObservabilityClient: Sendable {
}.sorted { $0.startedAt > $1.startedAt }
return RunsOutcome(runs: runs, error: nil)
} catch {
return RunsOutcome(runs: [], error: "Tempo recent sessions unavailable: \(error.localizedDescription)")
return RunsOutcome(runs: [], error: "Tempo recent sessions unavailable: \(Self.safeError(error))")
}
}

Expand All @@ -314,22 +316,32 @@ public struct LocalObservabilityClient: Sendable {
)!
components.queryItems = [URLQueryItem(name: "query", value: expression)]
let (data, response) = try await session.data(from: components.url!)
try requireSuccess(response: response, data: data, source: "Prometheus query")
try requireSuccess(response: response, source: "Prometheus query")
let payload = try JSONDecoder().decode(PrometheusResponse.self, from: data)
guard payload.status == "success" else { throw URLError(.cannotParseResponse) }
return payload.data.result.compactMap { result in
Double(result.rawValue).map { PrometheusSample(metric: result.metric, value: $0) }
return try payload.data.result.map { result in
guard let value = Double(result.rawValue), value.isFinite, value >= 0 else {
throw LocalObservabilityError.invalidSample
}
return PrometheusSample(metric: result.metric, value: value)
}
}

private func requireSuccess(response: URLResponse, data: Data, source: String) throws {
private func requireSuccess(response: URLResponse, source: String) throws {
guard let http = response as? HTTPURLResponse, (200..<300).contains(http.statusCode) else {
let status = (response as? HTTPURLResponse)?.statusCode ?? -1
let body = String(data: data, encoding: .utf8)?.prefix(180) ?? ""
throw LocalObservabilityError.queryRejected(status: status, source: source, body: String(body))
throw LocalObservabilityError.queryRejected(status: status, source: source)
}
}

// Never forward provider bodies, URLs, decoding debug descriptions or userInfo.
private static func safeError(_ error: Error) -> String {
if let error = error as? LocalObservabilityError { return error.localizedDescription }
if let error = error as? URLError { return "Telemetry transport failed (code \(error.code.rawValue))." }
if error is DecodingError { return "Telemetry response could not be decoded." }
return "Telemetry request failed. Error content was omitted."
}

private func adapter(for runtime: AgentRuntime) -> AgentTelemetry {
let supported = runtime == .codex || runtime == .claude || runtime == .grok
let note = supported
Expand Down Expand Up @@ -378,12 +390,13 @@ struct PrometheusSample: Sendable, Equatable {
}

enum LocalObservabilityError: LocalizedError {
case queryRejected(status: Int, source: String, body: String)
case queryRejected(status: Int, source: String)
case invalidSample

var errorDescription: String? {
switch self {
case let .queryRejected(status, source, body):
"\(source) returned HTTP \(status): \(body)"
case let .queryRejected(status, source): "\(source) returned HTTP \(status). Response content was omitted."
case .invalidSample: "Telemetry returned an invalid numeric sample."
}
}
}
Expand Down
2 changes: 1 addition & 1 deletion Sources/ContextCore/PluginInventory.swift
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ public struct PluginInventoryEntry: Sendable, Equatable, Identifiable {
return preferences[0].enabled ? "Enabled in checked settings" : "Disabled in checked settings"
}
public var reviewReason: String {
if !unreferencedVersions.isEmpty { return "\(unreferencedVersions.count) versions not referenced by the checked registry" }
if !unreferencedVersions.isEmpty { return "\(unreferencedVersions.count) \(unreferencedVersions.count == 1 ? "version" : "versions") not referenced by the checked registry" }
if versions.count > 1 { return "\(versions.count) cached versions; current installation needs verification" }
if !localMatches.isEmpty { return "Matching instructions also exist in your local library" }
if preferences.isEmpty { return "Check whether this plugin is enabled" }
Expand Down
16 changes: 13 additions & 3 deletions Sources/ContextCore/SkillPolicyModels.swift
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,13 @@ public struct SkillExposure: Identifiable, Codable, Sendable, Equatable {
public var id: String { logicalPath }
}

public struct SkillPrecedence: Codable, Sendable, Equatable {
public enum State: String, Codable, Sendable { case preferred, shadowed, coexisting, unverified }
public let state: State
public let preferredDefinitionID: String?
public let source: String
}

public struct SkillRuntimePolicy: Identifiable, Codable, Sendable, Equatable {
public let runtime: AgentRuntime
public let mode: InvocationMode
Expand All @@ -46,19 +53,22 @@ public struct SkillRuntimePolicy: Identifiable, Codable, Sendable, Equatable {
public let invocation: String
public let isExposed: Bool
public let desiredMode: InvocationMode?
/// Resolved only among discovered routes; this never proves installed activation.
public let precedence: SkillPrecedence?

public init(runtime: AgentRuntime, mode: InvocationMode, explicit: Bool, reason: String, invocation: String, isExposed: Bool = true, desiredMode: InvocationMode? = nil) {
public init(runtime: AgentRuntime, mode: InvocationMode, explicit: Bool, reason: String, invocation: String, isExposed: Bool = true, desiredMode: InvocationMode? = nil, precedence: SkillPrecedence? = nil) {
self.runtime = runtime
self.mode = mode
self.explicit = explicit
self.reason = reason
self.invocation = invocation
self.isExposed = isExposed
self.desiredMode = desiredMode
self.precedence = precedence
}

public var id: String { runtime.rawValue }
public var evidence: EvidenceQuality { explicit ? .measured : .derived }
public var evidence: EvidenceQuality { mode == .unverified ? .unavailable : explicit ? .measured : .derived }
}

public struct SkillRecord: Identifiable, Codable, Sendable, Equatable {
Expand All @@ -68,7 +78,7 @@ public struct SkillRecord: Identifiable, Codable, Sendable, Equatable {
public let logicalBytes: Int64
public let modified: Date
public let exposures: [SkillExposure]
public let policies: [SkillRuntimePolicy]
public var policies: [SkillRuntimePolicy]
/// SHA-256 of the complete SKILL.md when it fit inside the bounded reader.
/// A nil value means exact-copy analysis is unavailable, not that content differs.
public let contentFingerprint: String?
Expand Down
59 changes: 58 additions & 1 deletion Sources/ContextCore/SkillPolicyResolver.swift
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ public enum SkillPolicyResolver {
var records = groups.compactMap { physicalPath, items in
makeRecord(physicalPath: physicalPath, items: items)
}
records = resolvePrecedence(records)
let recordsByName = Dictionary(grouping: records, by: { $0.name.lowercased() })
records = records.map { record in
var record = record
Expand All @@ -24,10 +25,66 @@ public enum SkillPolicyResolver {
}
record.definitionConflictCount = max(1, overlapping.count)
return record
}.sorted { $0.name.localizedCaseInsensitiveCompare($1.name) == .orderedAscending }
}.sorted {
let order = $0.name.localizedCaseInsensitiveCompare($1.name)
return order == .orderedSame ? $0.id < $1.id : order == .orderedAscending
}
return SkillCatalogSnapshot(records: records, coverage: report.coverage, generatedAt: Date())
}

private static func resolvePrecedence(_ records: [SkillRecord]) -> [SkillRecord] {
let byName = Dictionary(grouping: records, by: \.name)
return records.map { original in
var record = original
record.policies = original.policies.map { policy in
let peers = byName[record.name, default: []].filter { $0.policy(for: policy.runtime)?.isExposed == true }
guard policy.isExposed, peers.count > 1 else { return policy }
var state = SkillPrecedence.State.unverified
var preferred: String?
var source = "No qualified winner rule for this runtime or these routes."
var reason = "Same-name precedence is unverified; no winner was inferred from scan order."
// Scope is not enough: recognize only direct skill-directory routes,
// excluding nested/custom roots, plugin namespaces and renamed frontmatter.
func routes(_ item: SkillRecord) -> [SkillExposure] {
item.exposures.filter { $0.provider == .claude && $0.applicability != .installedOnly }
}
func ordinary(_ item: SkillRecord) -> Bool {
let exposures = routes(item)
return !exposures.isEmpty && exposures.allSatisfy {
$0.logicalPath.hasSuffix("/.claude/skills/\(item.name)/SKILL.md") &&
(($0.scope == .global && $0.source == "Claude · Personal skills") || $0.scope == .project)
}
}
if policy.runtime == .claude, peers.allSatisfy(ordinary) {
let personal = peers.filter { routes($0).contains { $0.scope == .global } }
if personal.count == 1, let winner = personal.first {
preferred = winner.id
state = record.id == winner.id ? .preferred : .shadowed
source = "https://code.claude.com/docs/en/skills#resolve-skills-that-share-a-name"
reason = state == .shadowed
? "Shadowed among discovered Claude routes by personal definition \(winner.id). Personal skills override project skills."
: "Preferred among discovered Claude routes: personal skills override project skills."
reason += " Enterprise, synced skills and session overrides were not qualified; runtime activation remains unverified."
}
} else if policy.runtime == .codex, peers.allSatisfy({ item in
item.exposures.filter { [.codex, .agents].contains($0.provider) && $0.applicability != .installedOnly }
.allSatisfy { $0.logicalPath.contains("/.agents/skills/") }
}) {
state = .coexisting
source = "https://learn.chatgpt.com/docs/build-skills#where-codex-loads-local-skills"
reason = "Codex can list both same-name skills; these discovered routes have no exclusive winner. Runtime activation remains unverified."
}
return SkillRuntimePolicy(runtime: policy.runtime,
mode: state == .shadowed || state == .unverified ? .unverified : policy.mode,
explicit: state == .shadowed || state == .unverified ? false : policy.explicit,
reason: policy.reason + " " + reason, invocation: policy.invocation,
isExposed: policy.isExposed, desiredMode: policy.desiredMode,
precedence: SkillPrecedence(state: state, preferredDefinitionID: preferred, source: source))
}
return record
}
}

private static func makeRecord(physicalPath: String, items: [AIContextItem]) -> SkillRecord? {
guard let first = items.sorted(by: { $0.path < $1.path }).first else { return nil }
let document = readDocument(at: URL(fileURLWithPath: physicalPath))
Expand Down
Loading
Loading