prole/prole-app/Sources/ServiceChecker.swift
chrisfu 197c72dd7a Refactor prole-app and establish temporary release process
- Moved prole-tools-app to prole-app at the project root to make it self-contained for transition to its own repository.
- Created prole-tools-app/dist/ directory to host build artifacts.
- Generated distribution artifacts (Prole Tools.app and Prole Tools.zip) using prole-app/build.sh package.
- Checked in the generated artifacts to prole-tools-app/dist/ (bypassing .gitignore for temporary release process).
Changes Summary
•
Renamed directory prole-tools-app/ to prole-app/.
•
Populated prole-tools-app/dist/ with the latest build output from prole-app/build.sh.
•
Staged all changes, including the forced addition of ignored artifacts in prole-tools-app/dist/.
2026-01-18 15:04:23 -08:00

250 lines
11 KiB
Swift

import Foundation
import Network
final class ServiceChecker {
static let statusDidChangeNotification = Notification.Name("ServiceChecker.statusDidChange")
static let forceRefreshNotification = Notification.Name("ServiceChecker.forceRefresh")
private let queue = DispatchQueue(label: "prole.status.checker")
private var timer: DispatchSourceTimer?
private var isRunning: Bool = false
private var lastRefreshAt: Date = .distantPast
private let refreshInterval: TimeInterval = 30
// Public reachability flags
private(set) var k3dReachable = false
private(set) var prometheusReachable = false
private(set) var grafanaReachable = false
private(set) var openbaoReachable = false
private(set) var ollamaReachable = false
private(set) var postgresReachable = false
// Last error messages (for tooltips when red)
private(set) var k3dError: String? = nil
private(set) var prometheusError: String? = nil
private(set) var grafanaError: String? = nil
private(set) var openbaoError: String? = nil
private(set) var ollamaError: String? = nil
private(set) var postgresError: String? = nil
// Generic per-endpoint state so UI can query dynamically from Preferences
struct EndpointState: Equatable {
var reachable: Bool
var latencyMs: Int
var error: String?
var lastChecked: Date
}
// Cached states (by logical keys)
// services: key = service name (Config.ServiceEndpoint.name)
// kubes: key = host:port string
private(set) var serviceStates: [String: EndpointState] = [:]
private(set) var kubeStates: [String: EndpointState] = [:]
init() {
NotificationCenter.default.addObserver(self, selector: #selector(forceRefresh), name: Self.forceRefreshNotification, object: nil)
// When configuration changes, invalidate cache and reload immediately
NotificationCenter.default.addObserver(self, selector: #selector(configDidChange), name: Config.didChangeNotification, object: nil)
}
func start(interval: TimeInterval = 30) {
timer?.cancel()
let t = DispatchSource.makeTimerSource(queue: queue)
t.schedule(deadline: .now(), repeating: interval)
t.setEventHandler { [weak self] in self?.refreshAll() }
dlog("ServiceChecker: starting timer, interval=\(interval)s")
t.resume()
timer = t
}
@objc private func forceRefresh() { queue.async { self.refreshAll() } }
@objc private func configDidChange() {
queue.async {
// Invalidate freshness so the next refresh runs immediately
self.lastRefreshAt = .distantPast
// Clear caches so UI doesn't briefly show stale dynamic entries
self.serviceStates.removeAll()
self.kubeStates.removeAll()
self.refreshAll()
}
}
private func refreshAll() {
// Ensure we don't overlap and we don't rerun faster than every 30s
if isRunning {
dlog("ServiceChecker: skip — previous refresh still running")
return
}
let now = Date()
if now.timeIntervalSince(lastRefreshAt) < refreshInterval {
dlog("ServiceChecker: skip — cache still fresh (< \(Int(refreshInterval))s)")
return
}
isRunning = true
dlog("ServiceChecker: begin refresh round")
let group = DispatchGroup()
// clear previous errors before a new round (legacy fields)
k3dError = nil; prometheusError = nil; grafanaError = nil; openbaoError = nil; ollamaError = nil; postgresError = nil
let cfg = Config.shared
// Iterate Services from Preferences
for svc in cfg.services {
group.enter()
let nameUpper = svc.name.uppercased()
if nameUpper == "OPENBAO" {
// Special check for OpenBAO via kubectl if local
// We use a global utility queue for the script bridge to avoid blocking
DispatchQueue.global(qos: .utility).async {
let res = PFScriptBridge.openbaoStatus()
let isReady = res.out.contains("true")
let errorMsg = isReady ? nil : (res.out.isEmpty ? "No openbao pods found" : "OpenBAO not ready: \(res.out)")
self.queue.async {
let st = EndpointState(reachable: isReady, latencyMs: 0, error: errorMsg, lastChecked: Date())
self.serviceStates[svc.name] = st
self.openbaoReachable = isReady
self.openbaoError = errorMsg
group.leave()
}
}
} else {
tcpPing(host: svc.host, port: UInt16(svc.port)) { [weak self] ok, ms, err in
guard let self = self else { group.leave(); return }
let key = svc.name
let st = EndpointState(reachable: ok, latencyMs: ms, error: err, lastChecked: Date())
self.serviceStates[key] = st
// Update specific flags
switch nameUpper {
case "K3D":
self.k3dReachable = ok
self.k3dError = err
case "PROMETHEUS":
self.prometheusReachable = ok
self.prometheusError = err
case "GRAFANA":
self.grafanaReachable = ok
self.grafanaError = err
case "OLLAMA":
self.ollamaReachable = ok
self.ollamaError = err
case "POSTGRESQL":
self.postgresReachable = ok
self.postgresError = err
default:
break
}
group.leave()
}
}
}
// Iterate Kubernetes endpoints from Preferences
let kubes = cfg.kubernetes
for (idx, k) in kubes.enumerated() {
group.enter()
tcpPing(host: k.host, port: UInt16(k.port)) { [weak self] ok, ms, err in
guard let self = self else { group.leave(); return }
let key = "\(k.host):\(k.port)"
let st = EndpointState(reachable: ok, latencyMs: ms, error: err, lastChecked: Date())
self.kubeStates[key] = st
group.leave()
}
}
group.notify(queue: .main) {
self.lastRefreshAt = Date()
self.isRunning = false
dlog("ServiceChecker: refresh round complete → posting statusDidChangeNotification")
NotificationCenter.default.post(name: Self.statusDidChangeNotification, object: self)
}
}
private func tcpPing(host: String, port: UInt16, timeout: TimeInterval = 2.0, completion: @escaping (Bool, Int, String?) -> Void) {
dlog("tcpPing: attempting \(host):\(port) timeout=\(timeout)s")
let start = DispatchTime.now()
let params = NWParameters.tcp
params.allowLocalEndpointReuse = true
let endpoint = NWEndpoint.hostPort(host: .name(host, nil), port: .init(integerLiteral: port))
let conn = NWConnection(to: endpoint, using: params)
// Ensure completion is invoked exactly once
var finished = false
func finishOnce(_ ok: Bool, _ ms: Int, _ err: String?, reason: String) {
// All state updates and timeout run on `queue` so this is serialized
if finished { return }
finished = true
dlog("tcpPing: finishOnce(\(host):\(port)) reason=\(reason) ok=\(ok) ms=\(ms) err=\(err ?? "nil")")
completion(ok, ms, err)
}
// Prepare timeout work item so we can cancel it on success/failure
let timeoutWork = DispatchWorkItem { [weak conn] in
if finished { return }
dlog("tcpPing: timeout reached for \(host):\(port); cancelling connection")
conn?.cancel()
finishOnce(false, -1, "connection timed out", reason: "timeout")
}
conn.stateUpdateHandler = { state in
switch state {
case .ready:
let elapsed = DispatchTime.now().uptimeNanoseconds - start.uptimeNanoseconds
let ms = Int(Double(elapsed) / 1_000_000.0)
dlog("tcpPing: READY \(host):\(port) in \(ms) ms")
timeoutWork.cancel()
finishOnce(true, ms, nil, reason: "ready")
conn.cancel() // will emit .cancelled; ignored due to finished=true
case .failed(let error):
let msg = Self.describeNWError(error)
dlog("tcpPing: FAILED \(host):\(port)\(msg)")
timeoutWork.cancel()
finishOnce(false, -1, msg, reason: "failed")
conn.cancel()
case .cancelled:
// If we already finished (e.g., due to .ready), this is expected; ignore.
if finished {
dlog("tcpPing: CANCELLED \(host):\(port) after finish — ignoring")
} else {
// Cancel without prior .ready/.failed implies timeout or external cancel
finishOnce(false, -1, "connection cancelled (possible timeout)", reason: "cancelled-before-finish")
}
default:
break
}
}
conn.start(queue: queue)
// Schedule timeout on the same queue
queue.asyncAfter(deadline: .now() + timeout, execute: timeoutWork)
}
// Tooltips
func tooltipFor(name: String) -> String {
let st = serviceStates[name]
let reachable = st?.reachable ?? false
let error = st?.error
let latency = st?.latencyMs ?? -1
let host = Config.shared.services.first(where: { $0.name == name })?.host ?? "unknown"
let port = Config.shared.services.first(where: { $0.name == name })?.port ?? 0
if reachable { return "\(name) (\(host):\(port)) — reachable (\(latency) ms)" }
var s = "\(name) (\(host):\(port)) — unreachable"
if let e = error { s += "\nError: \(e)" }
return s
}
private static func describeNWError(_ error: NWError) -> String {
switch error {
case .posix(let code): return "POSIX \(code.rawValue): \(code)"
case .dns(let code): return "DNS \(code): \(code)"
case .tls(let status): return "TLS/OSStatus \(status)"
case .wifiAware(let reason): return "Wi-Fi Aware: \(reason)"
@unknown default: return "Unknown network error"
}
}
}