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}