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