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}