bugger/Sources/Services/BitableEventService.swift

228 lines
8.0 KiB
Swift

import Foundation
/// Connects to the feishu-app change-notification service via Server-Sent
/// Events (SSE). When the assignee of a tracked Bitable record changes, the
/// server pushes a `change` event down this stream; Bugger responds by
/// triggering `PollerService.fetchNow()` so the bug list refreshes instantly.
///
/// The stream URL is built from `AppConfig.feishuAppBaseURL`:
/// `{feishuAppBaseURL}/api/v1/bitable/events?file_token=...&assignee_field=...&assignee_name=...`
///
/// When `feishuAppBaseURL` is empty the service stays disabled and no
/// network requests are made polling remains the only refresh mechanism.
final class BitableEventService {
static let shared = BitableEventService()
private var session: URLSession?
private var task: URLSessionDataTask?
private var delegate: SSESessionDelegate?
private var reconnectWorkItem: DispatchWorkItem?
/// The user's intent to be connected. Distinguishes an unexpected stream
/// drop (should auto-reconnect) from an explicit `disconnect()` (should not).
private var isEnabled = false
private let reconnectDelay: TimeInterval = 30
private init() {}
// MARK: - Public
/// Open the SSE connection. Safe to call repeatedly a no-op if already
/// enabled, and a no-op if `feishuAppBaseURL` is not configured.
func connect() {
guard !isEnabled else { return }
guard let url = buildURL() else {
BuggerLog.info("BitableEventService: subscribe base URL not configured, skipping SSE")
return
}
isEnabled = true
openStream(at: url)
}
/// Close the SSE connection. No auto-reconnect will be attempted.
func disconnect() {
isEnabled = false
closeStream()
}
/// Apply the current config: tear down any existing stream, then connect
/// if the subscribe base URL is set. Call after the user changes the URL.
func reconnect() {
disconnect()
connect()
}
// MARK: - Private
private func openStream(at url: URL) {
closeStream()
let config = URLSessionConfiguration.ephemeral
// The server sends a heartbeat every ~30s, so 5 min of silence is a
// safe "something went wrong" threshold.
config.timeoutIntervalForRequest = 300
config.timeoutIntervalForResource = .infinity
config.waitsForConnectivity = true
let delegate = SSESessionDelegate(
onEvent: { [weak self] event, data in
self?.handleEvent(event, data: data)
},
onCompletion: { [weak self] error in
self?.handleStreamEnd(error: error)
}
)
self.delegate = delegate
let session = URLSession(configuration: config, delegate: delegate, delegateQueue: nil)
self.session = session
var request = URLRequest(url: url)
request.setValue("text/event-stream", forHTTPHeaderField: "Accept")
request.timeoutInterval = 300
let task = session.dataTask(with: request)
self.task = task
task.resume()
BuggerLog.info("BitableEventService: connecting to \(url.absoluteString)")
}
private func closeStream() {
reconnectWorkItem?.cancel()
reconnectWorkItem = nil
task?.cancel()
task = nil
session?.invalidateAndCancel()
session = nil
delegate = nil
}
private func handleStreamEnd(error: Error?) {
task = nil
session = nil
delegate = nil
if let error {
BuggerLog.error("BitableEventService: stream ended (\(error.localizedDescription))")
} else {
BuggerLog.info("BitableEventService: stream ended")
}
// Only auto-reconnect if the user still wants to be connected.
guard isEnabled else { return }
let work = DispatchWorkItem { [weak self] in
guard let self, self.isEnabled else { return }
guard let url = self.buildURL() else { return }
BuggerLog.info("BitableEventService: reconnecting...")
self.openStream(at: url)
}
reconnectWorkItem = work
DispatchQueue.main.asyncAfter(deadline: .now() + reconnectDelay, execute: work)
}
private func handleEvent(_ event: String, data: String) {
switch event {
case "change":
BuggerLog.info("BitableEventService: change push received (\(data)), refreshing now")
Task { await PollerService.shared.fetchNow() }
case "error":
BuggerLog.error("BitableEventService: server error \(data)")
case "connected", "heartbeat":
break
default:
BuggerLog.debug("BitableEventService: unknown event \(event)")
}
}
private func buildURL() -> URL? {
guard let config = AppStateService.shared.config,
config.isConfigured,
!config.feishuAppBaseURL.isEmpty else {
return nil
}
let assigneeName = config.assigneeName
let assigneeField = config.fieldMappings.assigneeField
guard !assigneeName.isEmpty, !assigneeField.isEmpty else {
BuggerLog.info("BitableEventService: assignee name or field not configured")
return nil
}
let base = config.feishuAppBaseURL.trimmingCharacters(in: CharacterSet(charactersIn: "/"))
guard var comps = URLComponents(string: "\(base)/api/v1/bitable/events") else {
return nil
}
comps.queryItems = [
URLQueryItem(name: "file_token", value: config.appToken),
URLQueryItem(name: "assignee_field", value: assigneeField),
URLQueryItem(name: "assignee_name", value: assigneeName),
]
return comps.url
}
}
// MARK: - SSE parsing
private final class SSESessionDelegate: NSObject, URLSessionDataDelegate {
private let onEvent: (String, String) -> Void
private let onCompletion: (Error?) -> Void
private var buffer = ""
init(onEvent: @escaping (String, String) -> Void, onCompletion: @escaping (Error?) -> Void) {
self.onEvent = onEvent
self.onCompletion = onCompletion
}
func urlSession(_ session: URLSession, dataTask: URLSessionDataTask, didReceive data: Data) {
guard let chunk = String(data: data, encoding: .utf8) else { return }
// Normalize line endings: SSE allows \n, \r, or \r\n.
buffer.append(chunk.replacingOccurrences(of: "\r\n", with: "\n")
.replacingOccurrences(of: "\r", with: "\n"))
// SSE frames are delimited by a blank line (\n\n).
while let range = buffer.range(of: "\n\n") {
let frame = String(buffer[..<range.lowerBound])
buffer.removeSubrange(..<range.upperBound)
processFrame(frame)
}
}
func urlSession(_ session: URLSession, task: URLSessionTask, didCompleteWithError error: Error?) {
// Flush any trailing frame that wasn't followed by a blank line.
if !buffer.isEmpty {
processFrame(buffer)
buffer = ""
}
// Route back to the main queue so BitableEventService state is only
// ever touched from one thread.
DispatchQueue.main.async { [self] in
self.onCompletion(error)
}
}
private func processFrame(_ frame: String) {
var eventType = ""
var data = ""
for line in frame.split(separator: "\n", omittingEmptySubsequences: false) {
if line.hasPrefix("event:") {
eventType = String(line.dropFirst("event:".count)).trimmingCharacters(in: .whitespaces)
} else if line.hasPrefix("data:") {
data = String(line.dropFirst("data:".count)).trimmingCharacters(in: .whitespaces)
}
// Lines starting with ":" are comments (heartbeats) ignored.
}
// Frames without an event type (e.g. heartbeat-only frames) are skipped.
guard !eventType.isEmpty else { return }
DispatchQueue.main.async { [self] in
self.onEvent(eventType, data)
}
}
}