krz/domain-dig

an ios app for DNS & SSL analysis

clone: git clone https://gitbay.org/krz/domain-dig.git

v4.7.0: DomainDig/IntegrationService.swift · raw

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