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 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 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; 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 "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" } } }