Files
brov-macbook/NotchBuddy/Sources/App/N8nPoller.swift
T

268 lines
11 KiB
Swift
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import Foundation
// MARK: - N8nPoller
// Polls n8n for the latest workflow execution every 15s.
// Tries /api/v1/executions first, falls back to /rest/executions.
// On a new terminal execution: fetches full details, stores [name, detail] in steps.
final class N8nPoller: @unchecked Sendable {
static let shared = N8nPoller()
private var timer: DispatchSourceTimer?
private var lastExecutionId: String = ""
private init() {}
func start() {
guard timer == nil else { return }
let t = DispatchSource.makeTimerSource(queue: .global(qos: .background))
t.schedule(deadline: .now() + 3, repeating: 15)
t.setEventHandler { [weak self] in self?.poll() }
t.resume()
timer = t
}
// MARK: - Poll list endpoint
private func poll() {
guard let apiKey = KeychainStore.shared.get("n8n-api-key"),
let rawBase = KeychainStore.shared.get("n8n-url") else {
n8nLog("No API key or URL configured")
return
}
let base = rawBase.trimmingCharacters(in: CharacterSet(charactersIn: "/"))
let endpoints = [
"\(base)/api/v1/executions?limit=1&includeData=false",
"\(base)/rest/executions?limit=1&includeData=false",
]
tryList(endpoints, apiKey: apiKey, base: base, idx: 0)
}
private func tryList(_ urls: [String], apiKey: String, base: String, idx: Int) {
guard idx < urls.count, let url = URL(string: urls[idx]) else {
n8nLog("All list endpoints failed")
return
}
var req = URLRequest(url: url, timeoutInterval: 10)
req.setValue(apiKey, forHTTPHeaderField: "X-N8N-API-KEY")
req.setValue("application/json", forHTTPHeaderField: "Accept")
n8nLog("Polling \(url.host ?? "?")\(url.path)")
URLSession.shared.dataTask(with: req) { [weak self] data, response, error in
guard let self else { return }
let code = (response as? HTTPURLResponse)?.statusCode ?? 0
if let error {
self.n8nLog("Network error: \(error.localizedDescription)")
self.tryList(urls, apiKey: apiKey, base: base, idx: idx + 1)
return
}
guard let data else {
self.tryList(urls, apiKey: apiKey, base: base, idx: idx + 1)
return
}
self.n8nLog("HTTP \(code) · \(data.count) bytes")
guard code == 200 else {
self.tryList(urls, apiKey: apiKey, base: base, idx: idx + 1)
return
}
guard let json = try? JSONSerialization.jsonObject(with: data) else { return }
// Response is either { "data": [...] } or [...]
let items: [[String: Any]]
if let obj = json as? [String: Any], let arr = obj["data"] as? [[String: Any]] {
items = arr
} else if let arr = json as? [[String: Any]] {
items = arr
} else {
self.n8nLog("Unexpected response shape")
return
}
guard let first = items.first else { self.n8nLog("No executions found"); return }
let id: String
if let s = first["id"] as? String { id = s }
else if let n = first["id"] as? Int { id = "\(n)" }
else { self.n8nLog("No id in execution"); return }
guard id != self.lastExecutionId else {
self.n8nLog("Same id=\(id) — no change")
return
}
// Terminal check: use status field — more reliable than the `finished` bool
// (published workflows often have finished=false on error)
let status = first["status"] as? String ?? ""
let isTerminal = ["success", "error", "crashed", "canceled", "failed"].contains(status)
guard isTerminal else {
self.n8nLog("id=\(id) status=\(status.isEmpty ? "?" : status) — not terminal")
return
}
self.lastExecutionId = id
let success = status == "success"
self.n8nLog("New execution id=\(id) status=\(status)")
// Fetch full detail (includeData=true required in some n8n versions)
let detailUrls = [
"\(base)/api/v1/executions/\(id)?includeData=true",
"\(base)/api/v1/executions/\(id)",
"\(base)/rest/executions/\(id)?includeData=true",
"\(base)/rest/executions/\(id)",
]
self.fetchDetail(detailUrls, apiKey: apiKey, success: success, idx: 0)
}.resume()
}
// MARK: - Fetch full execution detail
private func fetchDetail(_ urls: [String], apiKey: String, success: Bool, idx: Int) {
guard idx < urls.count, let url = URL(string: urls[idx]) else {
dispatch(success: success, name: "Workflow", detail: nil)
return
}
var req = URLRequest(url: url, timeoutInterval: 10)
req.setValue(apiKey, forHTTPHeaderField: "X-N8N-API-KEY")
req.setValue("application/json", forHTTPHeaderField: "Accept")
URLSession.shared.dataTask(with: req) { [weak self] data, response, _ in
guard let self else { return }
let code = (response as? HTTPURLResponse)?.statusCode ?? 0
guard let data, code == 200 else {
self.n8nLog("Detail HTTP \(code) \(url.host ?? "?")\(url.path)")
self.fetchDetail(urls, apiKey: apiKey, success: success, idx: idx + 1)
return
}
guard let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any] else {
self.fetchDetail(urls, apiKey: apiKey, success: success, idx: idx + 1)
return
}
let name = self.extractWorkflowName(from: json)
let detail = self.parseDetail(from: json, success: success)
self.n8nLog("Parsed: \(name)")
self.dispatch(success: success, name: name, detail: detail)
}.resume()
}
private func extractWorkflowName(from json: [String: Any]) -> String {
if let wd = json["workflowData"] as? [String: Any], let name = wd["name"] as? String { return name }
if let name = json["name"] as? String { return name }
return "Workflow"
}
// MARK: - Parse execution output / error message
private func parseDetail(from json: [String: Any], success: Bool) -> String? {
guard let execData = json["data"] as? [String: Any],
let resultData = execData["resultData"] as? [String: Any] else { return nil }
if success {
return parseSuccessDetail(resultData: resultData)
} else {
return parseErrorDetail(resultData: resultData)
}
}
private func parseErrorDetail(resultData: [String: Any]) -> String? {
// Top-level error
if let error = resultData["error"] as? [String: Any] {
let msg = error["message"] as? String ?? ""
if let node = (error["node"] as? [String: Any])?["name"] as? String, !node.isEmpty {
return "\(node)\n\(msg)"
}
return msg
}
// Scan runData for first node error
if let runData = resultData["runData"] as? [String: Any] {
for (nodeName, runs) in runData {
if let runs = runs as? [[String: Any]],
let run = runs.first,
let err = run["error"] as? [String: Any],
let msg = err["message"] as? String {
return "\(nodeName)\n\(msg)"
}
}
}
return nil
}
private func parseSuccessDetail(resultData: [String: Any]) -> String? {
guard let lastNode = resultData["lastNodeExecuted"] as? String,
let runData = resultData["runData"] as? [String: Any],
let nodeRuns = runData[lastNode] as? [[String: Any]],
let run = nodeRuns.first,
let data = run["data"] as? [String: Any],
let main = data["main"] as? [[[String: Any]]],
let items = main.first else { return nil }
let count = items.count
let header = "→ \(lastNode) · \(count) \(ruItems(count))"
// Preview first item's JSON keys (up to 4)
if let firstItem = items.first,
let jsonObj = firstItem["json"] as? [String: Any], !jsonObj.isEmpty {
let lines = jsonObj.prefix(4).map { "\($0.key): \(fmtValue($0.value))" }
return "\(header)\n\(lines.joined(separator: "\n"))"
}
return header
}
private func fmtValue(_ v: Any) -> String {
if let s = v as? String { return String(s.prefix(50)) }
if let n = v as? NSNumber { return n.stringValue }
if let a = v as? [Any] { return "[\(a.count)]" }
if v is [String: Any] { return "{…}" }
return "\(v)"
}
// MARK: - Dispatch to main
private func dispatch(success: Bool, name: String, detail: String?) {
DispatchQueue.main.async { self.handleExecution(success: success, name: name, detail: detail) }
}
@MainActor
private func handleExecution(success: Bool, name: String, detail: String?) {
let state = AppState.shared
// Apply workflow filter (empty = all workflows)
if !state.n8nWorkflowFilter.isEmpty && !state.n8nWorkflowFilter.contains(name) { return }
state.n8nRuns = Array(([N8nRun(workflow: name, detail: detail, success: success, date: Date())]
+ state.n8nRuns).prefix(10))
guard let idx = state.tasks.firstIndex(where: { $0.id == "integration_n8n" }) else { return }
let focused = state.focusId == "integration_n8n"
state.tasks[idx].state = success ? .finished : .error
state.tasks[idx].steps = detail != nil ? [name, detail!] : [name]
if !focused {
state.tasks[idx].pillBadge = success ? .finished : .error
}
SoundEngine.shared.play(success ? "finish" : "error")
// Auto-clear after 60s (user needs time to read detail)
DispatchQueue.main.asyncAfter(deadline: .now() + 60) {
guard let i = state.tasks.firstIndex(where: { $0.id == "integration_n8n" }) else { return }
guard state.tasks[i].state == .finished || state.tasks[i].state == .error else { return }
state.tasks[i].state = .idle
state.tasks[i].steps = []
state.tasks[i].pillBadge = nil
}
}
// MARK: - Logging
private func n8nLog(_ message: String) {
appendAppLog("n8n.log", message, timestampFormat: "HH:mm:ss")
}
}
/// Russian plural for "элемент": 1 элемент, 2 элемента, 5 элементов.
fileprivate func ruItems(_ n: Int) -> String {
let mod10 = n % 10, mod100 = n % 100
if mod10 == 1 && mod100 != 11 { return "элемент" }
if (2...4).contains(mod10) && !(12...14).contains(mod100) { return "элемента" }
return "элементов"
}