319cdc8b18
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
268 lines
11 KiB
Swift
268 lines
11 KiB
Swift
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 "элементов"
|
||
}
|