krz/domain-dig
an ios app for DNS & SSL analysis
clone: git clone https://gitbay.org/krz/domain-dig.git
v4.4.1: DomainDig/IntegrationService.swift · raw
1import Foundation
2import Network
3import Observation
4import Security
5
6@MainActor
7@Observable
8final class IntegrationService {
9 static let shared = IntegrationService()
10
11 var targets: [IntegrationTarget]
12 var deliveryRecords: [DeliveryRecord]
13 var queue: [QueuedDelivery]
14 var statusMessage: String?
15
16 private let defaults: UserDefaults
17 private var processingTask: Task<Void, Never>?
18
19 private init(defaults: UserDefaults = .standard) {
20 self.defaults = defaults
21 self.targets = Self.loadTargets(defaults: defaults)
22 self.deliveryRecords = Self.loadRecords(defaults: defaults)
23 self.queue = Self.loadQueue(defaults: defaults)
24 }
25
26 func refresh() {
27 targets = Self.loadTargets(defaults: defaults)
28 deliveryRecords = Self.loadRecords(defaults: defaults)
29 queue = Self.loadQueue(defaults: defaults)
30 }
31
32 func upsert(
33 target: IntegrationTarget,
34 webhookURL: String? = nil,
35 slackWebhookURL: String? = nil,
36 emailPassword: String? = nil
37 ) throws {
38 var updatedTarget = target
39
40 switch updatedTarget.configuration {
41 case .webhook(var configuration):
42 if let webhookURL {
43 let reference = configuration.credentialReference ?? Self.secretReference(for: updatedTarget.id, suffix: "webhook")
44 try IntegrationSecretStore.save(secret: webhookURL, reference: reference)
45 configuration.credentialReference = reference
46 configuration.endpointDisplayHost = Self.hostLabel(from: webhookURL)
47 updatedTarget.configuration = .webhook(configuration)
48 }
49 case .slack(var configuration):
50 if let slackWebhookURL {
51 let reference = configuration.credentialReference ?? Self.secretReference(for: updatedTarget.id, suffix: "slack")
52 try IntegrationSecretStore.save(secret: slackWebhookURL, reference: reference)
53 configuration.credentialReference = reference
54 configuration.destinationLabel = Self.hostLabel(from: slackWebhookURL)
55 updatedTarget.configuration = .slack(configuration)
56 }
57 case .email(var configuration):
58 if let emailPassword {
59 let reference = configuration.credentialReference ?? Self.secretReference(for: updatedTarget.id, suffix: "smtp")
60 try IntegrationSecretStore.save(secret: emailPassword, reference: reference)
61 configuration.credentialReference = reference
62 updatedTarget.configuration = .email(configuration)
63 }
64 }
65
66 if let index = targets.firstIndex(where: { $0.id == updatedTarget.id }) {
67 targets[index] = updatedTarget
68 } else {
69 targets.append(updatedTarget)
70 }
71 targets.sort { $0.name.localizedCaseInsensitiveCompare($1.name) == .orderedAscending }
72 persistTargets()
73 statusMessage = "Saved integration settings."
74 }
75
76 func delete(targetID: UUID) {
77 guard let target = targets.first(where: { $0.id == targetID }) else { return }
78 deleteSecrets(for: target)
79 targets.removeAll { $0.id == targetID }
80 queue.removeAll { $0.integrationID == targetID }
81 deliveryRecords.removeAll { $0.integrationID == targetID }
82 persistTargets()
83 persistQueue()
84 persistRecords()
85 }
86
87 func setEnabled(_ isEnabled: Bool, for targetID: UUID) {
88 guard let index = targets.firstIndex(where: { $0.id == targetID }) else { return }
89 targets[index].isEnabled = isEnabled
90 persistTargets()
91 }
92
93 func deliveryRecords(for targetID: UUID) -> [DeliveryRecord] {
94 deliveryRecords
95 .filter { $0.integrationID == targetID }
96 .sorted { $0.timestamp > $1.timestamp }
97 }
98
99 func enqueue(events: [MonitoringEvent]) {
100 guard !events.isEmpty else { return }
101 let eligibleTargets = targets.filter(\.isEnabled)
102 for event in events {
103 for target in eligibleTargets {
104 if let reason = filterMismatchReason(for: event, target: target) {
105 appendRecord(
106 DeliveryRecord(
107 integrationID: target.id,
108 eventID: event.id,
109 status: .skipped,
110 destination: destinationLabel(for: target),
111 summary: event.summary,
112 failureReason: reason
113 )
114 )
115 continue
116 }
117
118 queue.append(QueuedDelivery(integrationID: target.id, event: event))
119 appendRecord(
120 DeliveryRecord(
121 integrationID: target.id,
122 eventID: event.id,
123 status: .pending,
124 destination: destinationLabel(for: target),
125 summary: event.summary
126 )
127 )
128 }
129 }
130 persistQueue()
131 scheduleProcessing()
132 }
133
134 func recordNoOutboundEvents(for runSummary: String) {
135 let eligibleTargets = targets.filter(\.isEnabled)
136 for target in eligibleTargets {
137 appendRecord(
138 DeliveryRecord(
139 integrationID: target.id,
140 eventID: UUID(),
141 status: .skipped,
142 destination: destinationLabel(for: target),
143 summary: runSummary,
144 failureReason: "Monitoring run produced no outbound events."
145 )
146 )
147 }
148 }
149
150 func sendTest(for targetID: UUID) {
151 guard targets.contains(where: { $0.id == targetID }) else { return }
152 let event = MonitoringEvent(
153 type: .test,
154 severity: .info,
155 domain: "example.com",
156 summary: "DomainDig integration test",
157 details: [
158 "source": "manual test",
159 "environment": "local-first"
160 ]
161 )
162 queue.append(QueuedDelivery(integrationID: targetID, event: event))
163 appendRecord(
164 DeliveryRecord(
165 integrationID: targetID,
166 eventID: event.id,
167 status: .pending,
168 destination: targets.first(where: { $0.id == targetID }).map(destinationLabel(for:)) ?? "Unknown",
169 summary: event.summary
170 )
171 )
172 persistQueue()
173 scheduleProcessing()
174 }
175
176 func processQueueNow() {
177 scheduleProcessing(force: true)
178 }
179
180 func localSecretReferences() -> [String] {
181 targets.compactMap { target in
182 switch target.configuration {
183 case .webhook(let configuration):
184 configuration.credentialReference
185 case .slack(let configuration):
186 configuration.credentialReference
187 case .email(let configuration):
188 configuration.credentialReference
189 }
190 }
191 }
192
193 func resetAfterLocalWipe() {
194 processingTask?.cancel()
195 processingTask = nil
196 targets = []
197 deliveryRecords = []
198 queue = []
199 statusMessage = nil
200 }
201
202 private func scheduleProcessing(force: Bool = false) {
203 if force {
204 processingTask?.cancel()
205 processingTask = nil
206 }
207 guard processingTask == nil else { return }
208 processingTask = Task { [weak self] in
209 guard let self else { return }
210 await self.processQueueLoop()
211 }
212 }
213
214 private func processQueueLoop() async {
215 defer { processingTask = nil }
216
217 while true {
218 let dueItems = queue
219 .enumerated()
220 .filter { $0.element.nextAttemptAt <= Date() }
221
222 if dueItems.isEmpty {
223 guard let nextAttemptAt = queue.map(\.nextAttemptAt).min() else {
224 break
225 }
226
227 let delay = max(0.25, nextAttemptAt.timeIntervalSinceNow)
228 do {
229 try await Task.sleep(nanoseconds: UInt64(delay * 1_000_000_000))
230 continue
231 } catch {
232 break
233 }
234 }
235
236 for entry in dueItems.reversed() {
237 guard entry.offset < queue.count else { continue }
238 let item = queue[entry.offset]
239 await process(item: item, at: entry.offset)
240 }
241 }
242 }
243
244 private func process(item: QueuedDelivery, at index: Int) async {
245 guard let target = targets.first(where: { $0.id == item.integrationID }) else {
246 queue.remove(at: index)
247 persistQueue()
248 return
249 }
250
251 guard item.expiresAt > Date() else {
252 queue.remove(at: index)
253 persistQueue()
254 appendRecord(
255 DeliveryRecord(
256 integrationID: target.id,
257 eventID: item.event.id,
258 status: .expired,
259 destination: destinationLabel(for: target),
260 summary: item.event.summary,
261 failureReason: "Delivery expired before succeeding.",
262 attemptCount: item.attemptCount
263 )
264 )
265 return
266 }
267
268 do {
269 try await deliver(item.event, to: target)
270 queue.remove(at: index)
271 persistQueue()
272 appendRecord(
273 DeliveryRecord(
274 integrationID: target.id,
275 eventID: item.event.id,
276 status: .delivered,
277 destination: destinationLabel(for: target),
278 summary: item.event.summary,
279 attemptCount: item.attemptCount + 1
280 )
281 )
282 statusMessage = "Delivered \(item.event.summary)."
283 } catch {
284 var updated = item
285 updated.attemptCount += 1
286 updated.lastError = error.localizedDescription
287
288 if updated.attemptCount >= 5 {
289 queue.remove(at: index)
290 appendRecord(
291 DeliveryRecord(
292 integrationID: target.id,
293 eventID: item.event.id,
294 status: .failed,
295 destination: destinationLabel(for: target),
296 summary: item.event.summary,
297 failureReason: error.localizedDescription,
298 attemptCount: updated.attemptCount
299 )
300 )
301 } else {
302 let backoff = min(pow(2, Double(updated.attemptCount)) * 30, 3600)
303 updated.nextAttemptAt = Date().addingTimeInterval(backoff)
304 queue[index] = updated
305 appendRecord(
306 DeliveryRecord(
307 integrationID: target.id,
308 eventID: item.event.id,
309 status: .retrying,
310 destination: destinationLabel(for: target),
311 summary: item.event.summary,
312 failureReason: error.localizedDescription,
313 attemptCount: updated.attemptCount
314 )
315 )
316 }
317
318 persistQueue()
319 statusMessage = error.localizedDescription
320 }
321 }
322
323 private func deliver(_ event: MonitoringEvent, to target: IntegrationTarget) async throws {
324 switch target.configuration {
325 case .webhook(let configuration):
326 guard let reference = configuration.credentialReference else {
327 throw IntegrationError.missingSecret
328 }
329 let webhookURLString = try IntegrationSecretStore.secret(reference: reference)
330 try await HTTPIntegrationClient.sendJSON(
331 payload: IntegrationEventPayload(event: event),
332 to: webhookURLString,
333 headers: configuration.additionalHeaders,
334 timeoutSeconds: configuration.timeoutSeconds
335 )
336 case .slack(let configuration):
337 guard let reference = configuration.credentialReference else {
338 throw IntegrationError.missingSecret
339 }
340 let webhookURLString = try IntegrationSecretStore.secret(reference: reference)
341 try await HTTPIntegrationClient.sendJSON(
342 payload: SlackPayload(event: event),
343 to: webhookURLString,
344 headers: [:],
345 timeoutSeconds: 15
346 )
347 case .email(let configuration):
348 guard let reference = configuration.credentialReference else {
349 throw IntegrationError.missingSecret
350 }
351 let password = try IntegrationSecretStore.secret(reference: reference)
352 try await SMTPClient.send(
353 event: event,
354 configuration: configuration,
355 password: password
356 )
357 }
358 }
359
360 private func appendRecord(_ record: DeliveryRecord) {
361 deliveryRecords.insert(record, at: 0)
362 deliveryRecords = Array(deliveryRecords.prefix(250))
363 persistRecords()
364 }
365
366 private func persistTargets() {
367 Self.save(targets, key: StorageKey.targets, defaults: defaults)
368 }
369
370 private func persistRecords() {
371 Self.save(deliveryRecords, key: StorageKey.records, defaults: defaults)
372 }
373
374 private func persistQueue() {
375 Self.save(queue, key: StorageKey.queue, defaults: defaults)
376 }
377
378 private func deleteSecrets(for target: IntegrationTarget) {
379 switch target.configuration {
380 case .webhook(let configuration):
381 if let reference = configuration.credentialReference {
382 try? IntegrationSecretStore.delete(reference: reference)
383 }
384 case .slack(let configuration):
385 if let reference = configuration.credentialReference {
386 try? IntegrationSecretStore.delete(reference: reference)
387 }
388 case .email(let configuration):
389 if let reference = configuration.credentialReference {
390 try? IntegrationSecretStore.delete(reference: reference)
391 }
392 }
393 }
394
395 private func destinationLabel(for target: IntegrationTarget) -> String {
396 switch target.configuration {
397 case .webhook(let configuration):
398 return configuration.endpointDisplayHost.isEmpty ? target.name : configuration.endpointDisplayHost
399 case .slack(let configuration):
400 return configuration.destinationLabel
401 case .email(let configuration):
402 return configuration.recipientAddresses.joined(separator: ", ")
403 }
404 }
405
406 private func filterMismatchReason(for event: MonitoringEvent, target: IntegrationTarget) -> String? {
407 let filters = target.filters
408
409 if event.severity < filters.minimumSeverity {
410 return "Filtered by severity. Event was \(event.severity.title), target requires \(filters.minimumSeverity.title)."
411 }
412
413 if !filters.eventTypes.isEmpty, !filters.eventTypes.contains(event.type) {
414 return "Filtered by event type. Event was \(event.type.title)."
415 }
416
417 if !filters.domains.isEmpty, !filters.domains.map({ $0.lowercased() }).contains(event.domain.lowercased()) {
418 return "Filtered by domain. Event was for \(event.domain)."
419 }
420
421 return nil
422 }
423
424 private static func hostLabel(from string: String) -> String {
425 URL(string: string)?.host ?? "Configured"
426 }
427
428 private static func secretReference(for integrationID: UUID, suffix: String) -> String {
429 "integration.\(integrationID.uuidString).\(suffix)"
430 }
431
432 private static func loadTargets(defaults: UserDefaults) -> [IntegrationTarget] {
433 load([IntegrationTarget].self, key: StorageKey.targets, defaults: defaults) ?? []
434 }
435
436 private static func loadRecords(defaults: UserDefaults) -> [DeliveryRecord] {
437 load([DeliveryRecord].self, key: StorageKey.records, defaults: defaults) ?? []
438 }
439
440 private static func loadQueue(defaults: UserDefaults) -> [QueuedDelivery] {
441 load([QueuedDelivery].self, key: StorageKey.queue, defaults: defaults) ?? []
442 }
443
444 private static func load<T: Decodable>(_ type: T.Type, key: String, defaults: UserDefaults) -> T? {
445 guard let data = defaults.data(forKey: key) else {
446 return nil
447 }
448 return try? JSONDecoder().decode(type, from: data)
449 }
450
451 private static func save<T: Encodable>(_ value: T, key: String, defaults: UserDefaults) {
452 if let data = try? JSONEncoder().encode(value) {
453 defaults.set(data, forKey: key)
454 }
455 }
456
457 private enum StorageKey {
458 static let targets = "integrations.targets"
459 static let records = "integrations.records"
460 static let queue = "integrations.queue"
461 }
462}
463
464private struct IntegrationEventPayload: Encodable {
465 let eventType: String
466 let domain: String
467 let timestamp: Date
468 let severity: String
469 let summary: String
470 let details: [String: String]
471
472 init(event: MonitoringEvent) {
473 self.eventType = event.type.rawValue
474 self.domain = event.domain
475 self.timestamp = event.timestamp
476 self.severity = event.severity.rawValue
477 self.summary = event.summary
478 self.details = event.details
479 }
480}
481
482private struct SlackPayload: Encodable {
483 let text: String
484 let blocks: [SlackBlock]
485
486 init(event: MonitoringEvent) {
487 let title = "\(event.severity.title.uppercased()) • \(event.domain)"
488 let detailLines = event.details
489 .sorted { $0.key < $1.key }
490 .prefix(6)
491 .map { "\($0.key): \($0.value)" }
492 .joined(separator: "\n")
493
494 self.text = "\(title) — \(event.summary)"
495 self.blocks = [
496 SlackBlock(
497 type: "section",
498 text: .init(type: "mrkdwn", text: "*\(title)*\n\(event.summary)")
499 ),
500 SlackBlock(
501 type: "section",
502 text: .init(
503 type: "mrkdwn",
504 text: "*Event*: \(event.type.title)\n*Timestamp*: \(event.timestamp.formatted(date: .abbreviated, time: .shortened))"
505 )
506 ),
507 SlackBlock(
508 type: "section",
509 text: .init(type: "mrkdwn", text: detailLines.isEmpty ? "_No extra details_" : detailLines)
510 )
511 ]
512 }
513}
514
515private struct SlackBlock: Encodable {
516 let type: String
517 let text: SlackText
518}
519
520private struct SlackText: Encodable {
521 let type: String
522 let text: String
523}
524
525private enum IntegrationError: LocalizedError {
526 case invalidURL
527 case missingSecret
528 case invalidResponse(Int)
529 case invalidSMTPPort
530 case smtp(String)
531 case streamClosed
532
533 var errorDescription: String? {
534 switch self {
535 case .invalidURL:
536 return "The integration URL is invalid."
537 case .missingSecret:
538 return "This integration is missing a saved secret."
539 case .invalidResponse(let statusCode):
540 return "The remote endpoint returned \(statusCode)."
541 case .invalidSMTPPort:
542 return "The SMTP port is invalid."
543 case .smtp(let message):
544 return message
545 case .streamClosed:
546 return "The SMTP connection closed unexpectedly."
547 }
548 }
549}
550
551private enum HTTPIntegrationClient {
552 static func sendJSON<T: Encodable>(
553 payload: T,
554 to urlString: String,
555 headers: [String: String],
556 timeoutSeconds: Double
557 ) async throws {
558 guard let url = URL(string: urlString) else {
559 throw IntegrationError.invalidURL
560 }
561
562 var request = URLRequest(url: url, timeoutInterval: timeoutSeconds)
563 request.httpMethod = "POST"
564 request.setValue("application/json", forHTTPHeaderField: "Content-Type")
565 for (key, value) in headers {
566 request.setValue(value, forHTTPHeaderField: key)
567 }
568
569 let encoder = JSONEncoder()
570 encoder.dateEncodingStrategy = .iso8601
571 request.httpBody = try encoder.encode(payload)
572
573 let (_, response) = try await URLSession.shared.data(for: request)
574 guard let httpResponse = response as? HTTPURLResponse else {
575 throw IntegrationError.invalidResponse(-1)
576 }
577 guard (200..<300).contains(httpResponse.statusCode) else {
578 throw IntegrationError.invalidResponse(httpResponse.statusCode)
579 }
580 }
581}
582
583private enum IntegrationSecretStore {
584 static func save(secret: String, reference: String) throws {
585 let data = Data(secret.utf8)
586 try? delete(reference: reference)
587
588 let query: [String: Any] = [
589 kSecClass as String: kSecClassGenericPassword,
590 kSecAttrAccount as String: reference,
591 kSecValueData as String: data,
592 kSecAttrAccessible as String: kSecAttrAccessibleAfterFirstUnlock
593 ]
594
595 let status = SecItemAdd(query as CFDictionary, nil)
596 guard status == errSecSuccess else {
597 throw IntegrationError.smtp("Could not save integration secret.")
598 }
599 }
600
601 static func secret(reference: String) throws -> String {
602 let query: [String: Any] = [
603 kSecClass as String: kSecClassGenericPassword,
604 kSecAttrAccount as String: reference,
605 kSecReturnData as String: true,
606 kSecMatchLimit as String: kSecMatchLimitOne
607 ]
608
609 var result: CFTypeRef?
610 let status = SecItemCopyMatching(query as CFDictionary, &result)
611 guard status == errSecSuccess,
612 let data = result as? Data,
613 let secret = String(data: data, encoding: .utf8) else {
614 throw IntegrationError.missingSecret
615 }
616
617 return secret
618 }
619
620 static func delete(reference: String) throws {
621 let query: [String: Any] = [
622 kSecClass as String: kSecClassGenericPassword,
623 kSecAttrAccount as String: reference
624 ]
625 SecItemDelete(query as CFDictionary)
626 }
627}
628
629private enum SMTPClient {
630 static func send(
631 event: MonitoringEvent,
632 configuration: EmailIntegrationConfiguration,
633 password: String
634 ) async throws {
635 guard let port = NWEndpoint.Port(rawValue: UInt16(configuration.port)) else {
636 throw IntegrationError.invalidSMTPPort
637 }
638
639 let parameters: NWParameters = {
640 switch configuration.securityMode {
641 case .plain:
642 return .tcp
643 case .directTLS:
644 let tls = NWProtocolTLS.Options()
645 return NWParameters(tls: tls, tcp: NWProtocolTCP.Options())
646 }
647 }()
648
649 let channel = SMTPChannel(host: configuration.smtpHost, port: port, parameters: parameters)
650 try await channel.start()
651 _ = try await channel.readResponse(expecting: [220])
652 _ = try await channel.sendCommand("EHLO domaindig.local", expecting: [250])
653
654 if !configuration.username.isEmpty {
655 _ = try await channel.sendCommand("AUTH LOGIN", expecting: [334])
656 _ = try await channel.sendCommand(Data(configuration.username.utf8).base64EncodedString(), expecting: [334])
657 _ = try await channel.sendCommand(Data(password.utf8).base64EncodedString(), expecting: [235])
658 }
659
660 _ = try await channel.sendCommand("MAIL FROM:<\(configuration.senderAddress)>", expecting: [250])
661 for recipient in configuration.recipientAddresses {
662 _ = try await channel.sendCommand("RCPT TO:<\(recipient)>", expecting: [250, 251])
663 }
664 _ = try await channel.sendCommand("DATA", expecting: [354])
665
666 let detailLines = event.details
667 .sorted { $0.key < $1.key }
668 .map { "\($0.key): \($0.value)" }
669 .joined(separator: "\r\n")
670 let body = [
671 "From: DomainDig <\(configuration.senderAddress)>",
672 "To: \(configuration.recipientAddresses.joined(separator: ", "))",
673 "Subject: [DomainDig] \(event.severity.title) \(event.domain) \(event.type.title)",
674 "Date: \(DateFormatter.rfc2822.string(from: Date()))",
675 "",
676 event.summary,
677 "",
678 "Domain: \(event.domain)",
679 "Severity: \(event.severity.title)",
680 "Event: \(event.type.title)",
681 "Timestamp: \(event.timestamp.formatted(date: .abbreviated, time: .shortened))",
682 detailLines
683 ]
684 .joined(separator: "\r\n")
685
686 try await channel.sendRaw(body + "\r\n.\r\n")
687 _ = try await channel.readResponse(expecting: [250])
688 _ = try await channel.sendCommand("QUIT", expecting: [221])
689 channel.cancel()
690 }
691}
692
693private final class SMTPChannel {
694 private let connection: NWConnection
695 private var parsedLines: [String] = []
696 private var lineWaiters: [CheckedContinuation<String, Error>] = []
697 private var receiveBuffer = Data()
698
699 init(host: String, port: NWEndpoint.Port, parameters: NWParameters) {
700 connection = NWConnection(host: NWEndpoint.Host(host), port: port, using: parameters)
701 }
702
703 func start() async throws {
704 try await withCheckedThrowingContinuation { (continuation: CheckedContinuation<Void, Error>) in
705 connection.stateUpdateHandler = { [weak self] state in
706 switch state {
707 case .ready:
708 DispatchQueue.global(qos: .utility).async {
709 self?.startReceiveLoop()
710 }
711 continuation.resume()
712 case .failed(let error):
713 continuation.resume(throwing: error)
714 default:
715 break
716 }
717 }
718 connection.start(queue: .global(qos: .utility))
719 }
720 }
721
722 func cancel() {
723 connection.cancel()
724 }
725
726 func sendCommand(_ command: String, expecting codes: Set<Int>) async throws -> String {
727 try await sendRaw(command + "\r\n")
728 return try await readResponse(expecting: codes)
729 }
730
731 func sendCommand(_ command: String, expecting codes: [Int]) async throws -> String {
732 try await sendCommand(command, expecting: Set(codes))
733 }
734
735 func sendRaw(_ string: String) async throws {
736 let data = Data(string.utf8)
737 try await withCheckedThrowingContinuation { (continuation: CheckedContinuation<Void, Error>) in
738 connection.send(content: data, completion: .contentProcessed { error in
739 if let error {
740 continuation.resume(throwing: error)
741 } else {
742 continuation.resume()
743 }
744 })
745 }
746 }
747
748 func readResponse(expecting codes: Set<Int>) async throws -> String {
749 var lines: [String] = []
750
751 while true {
752 let line = try await readLine()
753 lines.append(line)
754
755 guard line.count >= 4,
756 let code = Int(line.prefix(3)) else {
757 continue
758 }
759
760 let delimiterIndex = line.index(line.startIndex, offsetBy: 3)
761 if line[delimiterIndex] == " " {
762 guard codes.contains(code) else {
763 throw IntegrationError.smtp(line)
764 }
765 return lines.joined(separator: "\n")
766 }
767 }
768 }
769
770 private func readLine() async throws -> String {
771 if !parsedLines.isEmpty {
772 return parsedLines.removeFirst()
773 }
774
775 return try await withCheckedThrowingContinuation { continuation in
776 lineWaiters.append(continuation)
777 }
778 }
779
780 private func startReceiveLoop() {
781 connection.receive(minimumIncompleteLength: 1, maximumLength: 4096) { [weak self] data, _, isComplete, error in
782 guard let self else { return }
783
784 if let error {
785 self.failWaiters(with: error)
786 return
787 }
788
789 if let data, !data.isEmpty {
790 self.receiveBuffer.append(data)
791 self.flushBuffer()
792 }
793
794 if isComplete {
795 self.failWaiters(with: IntegrationError.streamClosed)
796 return
797 }
798
799 self.startReceiveLoop()
800 }
801 }
802
803 private func flushBuffer() {
804 let delimiter = Data("\r\n".utf8)
805 while let range = receiveBuffer.range(of: delimiter) {
806 let lineData = receiveBuffer.subdata(in: receiveBuffer.startIndex..<range.lowerBound)
807 receiveBuffer.removeSubrange(receiveBuffer.startIndex..<range.upperBound)
808 let line = String(data: lineData, encoding: .utf8) ?? ""
809 if !lineWaiters.isEmpty {
810 let continuation = lineWaiters.removeFirst()
811 continuation.resume(returning: line)
812 } else {
813 parsedLines.append(line)
814 }
815 }
816 }
817
818 private func failWaiters(with error: Error) {
819 let waiters = lineWaiters
820 lineWaiters.removeAll()
821 for waiter in waiters {
822 waiter.resume(throwing: error)
823 }
824 }
825}
826
827private extension DateFormatter {
828 static let rfc2822: DateFormatter = {
829 let formatter = DateFormatter()
830 formatter.locale = Locale(identifier: "en_US_POSIX")
831 formatter.timeZone = TimeZone(secondsFromGMT: 0)
832 formatter.dateFormat = "EEE, dd MMM yyyy HH:mm:ss Z"
833 return formatter
834 }()
835}