260 lines
9.6 KiB
Swift
260 lines
9.6 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
|
|
|
|
// -- Exponential backoff with jitter -----------------------------------
|
|
|
|
/// Consecutive reconnect attempts since the last successful `connected`
|
|
/// event. Reset to 0 when the server confirms the stream is alive.
|
|
private var reconnectAttempt = 0
|
|
|
|
/// Base delay (seconds) for the first reconnect attempt.
|
|
private let reconnectBaseDelay: TimeInterval = 1
|
|
/// Multiplier applied per attempt.
|
|
private let reconnectBackoffFactor: Double = 2
|
|
/// Maximum delay (seconds) — preserves the previous steady-state cadence.
|
|
private let reconnectMaxDelay: TimeInterval = 30
|
|
/// Jitter range applied to the computed delay (uniform in [1-range, 1+range]).
|
|
private let reconnectJitter: Double = 0.25
|
|
|
|
private init() {}
|
|
|
|
/// Compute the reconnect delay for the current attempt count using
|
|
/// exponential backoff capped at `reconnectMaxDelay`, then apply ±jitter.
|
|
/// Formula: delay = min(cap, base * factor^(attempt-1)) * random(1±jitter)
|
|
private func computeReconnectDelay() -> TimeInterval {
|
|
let exponent = Double(max(reconnectAttempt, 1) - 1)
|
|
let raw = min(reconnectMaxDelay, reconnectBaseDelay * pow(reconnectBackoffFactor, exponent))
|
|
let lo = 1 - reconnectJitter
|
|
let hi = 1 + reconnectJitter
|
|
let jitterFactor = Double.random(in: lo...hi)
|
|
return raw * jitterFactor
|
|
}
|
|
|
|
// 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
|
|
reconnectAttempt = 0
|
|
}
|
|
|
|
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 }
|
|
|
|
reconnectAttempt += 1
|
|
let delay = computeReconnectDelay()
|
|
BuggerLog.info("BitableEventService: reconnecting in \(String(format: "%.1f", delay))s (attempt \(reconnectAttempt))")
|
|
|
|
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() + delay, 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 "connected":
|
|
reconnectAttempt = 0
|
|
case "error":
|
|
BuggerLog.error("BitableEventService: server error \(data)")
|
|
case "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)
|
|
}
|
|
}
|
|
}
|