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