Make private-media deletion transactional (#1468)

Receiver-side private-media deletion (per-bubble delete and /clear) becomes a write-ahead transaction: a deletion journal in BLEPrivateMediaReceiptStore is the atomic commit point, materialization (tombstones + payload unlinks) is idempotent and retried on recovery, path reservations in BLEIncomingFileStore prevent delete racing an in-flight arrival, and overlapping /clear operations are serialized with panic-generation invalidation. Integrates with the receipt quarantine (a pending journal entry outranks quarantine; materialization never resurrects a quarantined ID).

Review fixes included: (1) /clear no longer deletes outgoing media mirrored into another conversation (alias protection now mirrors the incoming path, with a regression test that fails pre-fix); (2) explicit delete of legacy incoming media actually unlinks the decrypted payload when unreferenced — gated on pending-delivery/reservation state, restoring main's delete semantics safely instead of leaving plaintext for quota cleanup; (3) refused deletions surface a localized system message in the affected chat instead of failing silently (30-locale key). Full local suite 1876 green.
This commit is contained in:
jack
2026-07-26 14:33:12 +02:00
committed by GitHub
parent d6bd4f0681
commit c72bb4ca2e
14 changed files with 3954 additions and 194 deletions
+36
View File
@@ -1,6 +1,42 @@
{
"sourceLanguage" : "en",
"strings" : {
"content.system.media_delete_refused" : {
"comment" : "System message shown in the affected chat when an explicit media delete or /clear was refused and bubbles/files were kept",
"extractionState" : "manual",
"localizations" : {
"ar" : { "stringUnit" : { "state" : "needs_review", "value" : "تعذّر حذف بعض الوسائط. جرّب حذف الوسائط الأقدم أولاً." } },
"bn" : { "stringUnit" : { "state" : "needs_review", "value" : "কিছু মিডিয়া মুছে ফেলা যায়নি। আগে পুরোনো মিডিয়া মুছে ফেলার চেষ্টা করুন।" } },
"de" : { "stringUnit" : { "state" : "needs_review", "value" : "einige medien konnten nicht gelöscht werden. versuche zuerst, ältere medien zu löschen." } },
"en" : { "stringUnit" : { "state" : "translated", "value" : "some media could not be deleted. try deleting older media first." } },
"es" : { "stringUnit" : { "state" : "needs_review", "value" : "no se pudieron eliminar algunos archivos multimedia. prueba a eliminar primero los más antiguos." } },
"fa" : { "stringUnit" : { "state" : "needs_review", "value" : "برخی رسانه‌ها حذف نشدند. ابتدا رسانه‌های قدیمی‌تر را حذف کنید." } },
"fil" : { "stringUnit" : { "state" : "needs_review", "value" : "hindi ma-delete ang ilang media. subukang i-delete muna ang mas lumang media." } },
"fr" : { "stringUnit" : { "state" : "needs_review", "value" : "impossible de supprimer certains médias. essaie d'abord de supprimer les médias plus anciens." } },
"he" : { "stringUnit" : { "state" : "needs_review", "value" : "לא ניתן למחוק חלק מהמדיה. נסה קודם למחוק מדיה ישנה יותר." } },
"hi" : { "stringUnit" : { "state" : "needs_review", "value" : "कुछ मीडिया हटाई नहीं जा सकी। पहले पुरानी मीडिया हटाने का प्रयास करें।" } },
"id" : { "stringUnit" : { "state" : "needs_review", "value" : "sebagian media tidak dapat dihapus. coba hapus media yang lebih lama dulu." } },
"it" : { "stringUnit" : { "state" : "needs_review", "value" : "impossibile eliminare alcuni contenuti multimediali. prova prima a eliminare quelli più vecchi." } },
"ja" : { "stringUnit" : { "state" : "needs_review", "value" : "一部のメディアを削除できませんでした。先に古いメディアを削除してみてください。" } },
"ko" : { "stringUnit" : { "state" : "needs_review", "value" : "일부 미디어를 삭제하지 못했습니다. 먼저 오래된 미디어를 삭제해 보세요." } },
"ms" : { "stringUnit" : { "state" : "needs_review", "value" : "sesetengah media tidak dapat dipadamkan. cuba padamkan media yang lebih lama dahulu." } },
"ne" : { "stringUnit" : { "state" : "needs_review", "value" : "केही मिडिया मेटाउन सकिएन। पहिले पुराना मिडिया मेटाउने प्रयास गर।" } },
"nl" : { "stringUnit" : { "state" : "needs_review", "value" : "sommige media konden niet worden verwijderd. probeer eerst oudere media te verwijderen." } },
"pl" : { "stringUnit" : { "state" : "needs_review", "value" : "nie udało się usunąć części multimediów. spróbuj najpierw usunąć starsze multimedia." } },
"pt" : { "stringUnit" : { "state" : "needs_review", "value" : "não foi possível eliminar alguns ficheiros multimédia. tenta eliminar primeiro os mais antigos." } },
"pt-BR" : { "stringUnit" : { "state" : "needs_review", "value" : "não foi possível excluir algumas mídias. tente excluir primeiro as mídias mais antigas." } },
"ru" : { "stringUnit" : { "state" : "needs_review", "value" : "не удалось удалить часть медиафайлов. попробуй сначала удалить более старые." } },
"sv" : { "stringUnit" : { "state" : "needs_review", "value" : "vissa medier kunde inte raderas. prova att radera äldre medier först." } },
"ta" : { "stringUnit" : { "state" : "needs_review", "value" : "சில ஊடகங்களை நீக்க முடியவில்லை. முதலில் பழைய ஊடகங்களை நீக்க முயற்சிக்கவும்." } },
"th" : { "stringUnit" : { "state" : "needs_review", "value" : "ไม่สามารถลบสื่อบางรายการได้ ลองลบสื่อที่เก่ากว่าก่อน" } },
"tr" : { "stringUnit" : { "state" : "needs_review", "value" : "bazı medya silinemedi. önce daha eski medyayı silmeyi dene." } },
"uk" : { "stringUnit" : { "state" : "needs_review", "value" : "не вдалося видалити частину медіафайлів. спробуй спочатку видалити старіші." } },
"ur" : { "stringUnit" : { "state" : "needs_review", "value" : "کچھ میڈیا حذف نہیں ہو سکا۔ پہلے پرانا میڈیا حذف کرنے کی کوشش کریں۔" } },
"vi" : { "stringUnit" : { "state" : "needs_review", "value" : "không thể xóa một số phương tiện. hãy thử xóa phương tiện cũ hơn trước." } },
"zh-Hans" : { "stringUnit" : { "state" : "needs_review", "value" : "部分媒体无法删除。请先尝试删除较早的媒体。" } },
"zh-Hant" : { "stringUnit" : { "state" : "needs_review", "value" : "部分媒體無法刪除。請先嘗試刪除較舊的媒體。" } }
}
},
"notification.action.wave" : {
"comment" : "Title of the notification action button that sends a friendly wave back to a nearby person",
"extractionState" : "manual",
@@ -40,6 +40,9 @@ struct BLEFileTransferHandlerEnvironment {
let commitPrivateMediaFile: (_ messageID: String, _ storedURL: URL) -> Bool
/// Rolls back a saved payload when its durable receipt commit fails.
let removeIncomingFile: (_ storedURL: URL) -> Void
/// Releases the allocator's save-to-UI ownership guard after synchronous
/// conversation insertion has completed.
let finishIncomingFileDelivery: (_ storedURL: URL) -> Void
/// Checks the authenticated sender before any private-media disk work.
let isPrivateMediaSenderBlocked: (PeerID) -> Bool
/// Updates the registry last-seen timestamp for the peer (async barrier write).
@@ -49,11 +52,14 @@ struct BLEFileTransferHandlerEnvironment {
let acknowledgePrivateMedia: (_ messageID: String, _ peerID: PeerID) -> Void
/// Delivers `.messageReceived` as one main-actor hop while
/// `shouldDeliver` remains true before and after the synchronous sink.
/// The completion authorizes the stable-media ACK.
/// The completion authorizes the stable-media ACK. Finalization runs after
/// every delivery attempt, including rejection, so allocator ownership
/// cannot leak indefinitely.
let deliverMessage: (
_ message: BitchatMessage,
_ shouldDeliver: @escaping () -> Bool,
_ completion: @escaping () -> Void
_ completion: @escaping () -> Void,
_ finalization: @escaping (TransportEventDeliveryOutcome) -> Void
) -> Void
}
@@ -359,7 +365,24 @@ final class BLEFileTransferHandler {
env: env
)
} else {
env.deliverMessage(message, { true }, {})
env.deliverMessage(
message,
{ true },
{},
{ outcome in
if outcome == .rejected {
// Raw media has no durable receipt that can redeliver
// it later. Do not leave a newly saved, UI-unowned file
// available for a stale fallback path to misidentify.
env.removeIncomingFile(destination)
} else {
// Plain delegates are invoked without synchronous
// insertion confirmation. Preserve the payload for
// that supported delivery path.
env.finishIncomingFileDelivery(destination)
}
}
)
}
return true
}
@@ -383,6 +406,9 @@ final class BLEFileTransferHandler {
},
{
env.acknowledgePrivateMedia(messageID, peerID)
},
{ _ in
env.finishIncomingFileDelivery(expectedURL)
}
)
}
+237 -15
View File
@@ -93,7 +93,7 @@ struct PanicRecoveryOperations {
}
}
struct BLEIncomingFileStore {
struct BLEIncomingFileStore: @unchecked Sendable {
enum PanicRecoveryError: Error {
case externalMarkerCommitFailed
case markerWriteFailed(Error)
@@ -103,7 +103,17 @@ struct BLEIncomingFileStore {
)
}
private static let quotaBytes: Int64 = 100 * 1024 * 1024
struct PrivateMediaDeletionReservation: Sendable {
fileprivate let id: UUID
}
private final class PayloadCoordination: @unchecked Sendable {
let lock = NSLock()
var pendingDeliveryPaths: Set<String> = []
var deletionReservations: [UUID: Set<String>] = [:]
}
private static let defaultQuotaBytes: Int64 = 100 * 1024 * 1024
/// Kept outside `files/` so deleting the media tree cannot erase the
/// fail-closed startup decision before the full panic has committed.
private static let panicRecoveryPendingMarkerFileName =
@@ -120,7 +130,6 @@ struct BLEIncomingFileStore {
"files/incoming",
"files/outgoing"
]
/// Name prefix of in-flight live voice captures (progressively written by
/// `ChatLiveVoiceCoordinator`). Quota eviction skips them by pattern
/// deleting one mid-stream unlinks the inode under an open `FileHandle`
@@ -134,7 +143,9 @@ struct BLEIncomingFileStore {
private let baseDirectory: URL?
private let dateProvider: () -> Date
private let panicMarkerWriter: (Data, URL) throws -> Void
private let quotaBytes: Int64
private let privateMediaReceipts: BLEPrivateMediaReceiptStore
private let payloadCoordination: PayloadCoordination
init(
fileManager: FileManager = .default,
@@ -142,17 +153,20 @@ struct BLEIncomingFileStore {
dateProvider: @escaping () -> Date = Date.init,
panicMarkerWriter: @escaping (Data, URL) throws -> Void = {
try $0.write(to: $1, options: .atomic)
}
},
quotaBytes: Int64 = Self.defaultQuotaBytes
) {
self.fileManager = fileManager
self.baseDirectory = baseDirectory
self.dateProvider = dateProvider
self.panicMarkerWriter = panicMarkerWriter
self.quotaBytes = max(0, quotaBytes)
self.privateMediaReceipts = BLEPrivateMediaReceiptStore(
fileManager: fileManager,
baseDirectory: baseDirectory,
now: dateProvider
)
self.payloadCoordination = PayloadCoordination()
}
/// Panic-wipe every managed incoming and outgoing media artifact before
@@ -165,10 +179,21 @@ struct BLEIncomingFileStore {
func panicWipe(
hasDurablePendingMarker: Bool = false
) throws {
// The receipt index caches tombstones as well as accepted payloads.
// Always invalidate it on return, including partial-failure paths, so
// no pre-panic receiver decision survives after identity reset.
defer { privateMediaReceipts.resetForPanic() }
// The receipt index caches tombstones as well as accepted payloads,
// while payload coordination retains save/delete reservations. Always
// invalidate both on return, including partial-failure paths, so no
// pre-panic receiver decision survives after identity reset.
defer {
privateMediaReceipts.resetForPanic()
payloadCoordination.lock.lock()
payloadCoordination.pendingDeliveryPaths.removeAll(
keepingCapacity: false
)
payloadCoordination.deletionReservations.removeAll(
keepingCapacity: false
)
payloadCoordination.lock.unlock()
}
let markerError: Error?
do {
@@ -251,6 +276,9 @@ struct BLEIncomingFileStore {
fallbackExtension: String?,
defaultPrefix: String
) -> URL? {
payloadCoordination.lock.lock()
defer { payloadCoordination.lock.unlock() }
do {
let base = try filesDirectory().appendingPathComponent(subdirectory, isDirectory: true)
try fileManager.createDirectory(at: base, withIntermediateDirectories: true, attributes: nil)
@@ -259,8 +287,26 @@ struct BLEIncomingFileStore {
defaultName: "\(defaultPrefix)_\(Self.timestampString(from: dateProvider()))",
fallbackExtension: fallbackExtension
)
let destination = uniqueFileURL(in: base, fileName: sanitized)
let reservedPaths = privateMediaReceipts.reservedPayloadPaths()
let deletionPaths = payloadCoordination
.deletionReservations.values.reduce(into: Set<String>()) {
$0.formUnion($1)
}
let allocationReservations = deletionPaths.union(
payloadCoordination.pendingDeliveryPaths
)
let destination = uniqueFileURL(
in: base,
fileName: sanitized,
reservedPaths: (reservedPaths ?? []).union(
allocationReservations
),
forceRandomizedName: reservedPaths == nil
)
try data.write(to: destination, options: .atomic)
payloadCoordination.pendingDeliveryPaths.insert(
destination.standardizedFileURL.path
)
return destination
} catch {
SecureLogger.error("❌ Failed to persist incoming media: \(error)", category: .session)
@@ -295,8 +341,152 @@ struct BLEIncomingFileStore {
)
}
/// Reserves every receipt/UI path before the asynchronous deletion
/// barrier. Allocation and reservation share one lock, so either an
/// in-flight raw arrival is observed and deletion fails closed, or the
/// arrival is forced onto a different filename.
func reservePrivateMediaDeletion(
messageIDs: [String],
payloadRelativePaths: [String: String]
) -> PrivateMediaDeletionReservation? {
payloadCoordination.lock.lock()
defer { payloadCoordination.lock.unlock() }
guard let paths = privateMediaReceipts
.prospectiveDeletionPayloadPaths(
messageIDs: messageIDs,
payloadRelativePaths: payloadRelativePaths
),
paths.isDisjoint(
with: payloadCoordination.pendingDeliveryPaths
) else {
return nil
}
let reservation = PrivateMediaDeletionReservation(id: UUID())
payloadCoordination.deletionReservations[reservation.id] = paths
return reservation
}
func commitPrivateMediaDeletion(
reservation: PrivateMediaDeletionReservation,
messageIDs: [String],
payloadRelativePaths: [String: String],
protectedPayloadRelativePaths: Set<String>
) -> Bool {
payloadCoordination.lock.lock()
defer {
payloadCoordination.deletionReservations.removeValue(
forKey: reservation.id
)
payloadCoordination.lock.unlock()
}
guard payloadCoordination.deletionReservations[reservation.id] != nil
else {
return false
}
return privateMediaReceipts.recordDeleted(
messageIDs: messageIDs,
payloadRelativePaths: payloadRelativePaths,
protectedPayloadRelativePaths: protectedPayloadRelativePaths
)
}
/// Explicit deletion of a LEGACY (non-stable-ID) incoming payload.
///
/// Legacy media has no durable receipt, so the only safe unlink is one
/// that can prove no other owner may hold the basename: the path must
/// not be pending delivery, must not belong to an in-flight deletion
/// reservation, and must not be owned by a stable receipt or journal
/// entry. When any of those hold or receipt state cannot be read
/// the file stays for bounded quota cleanup (the fail-safe fallback).
/// Returns true only when the payload was verifiably unlinked.
@discardableResult
func removeLegacyIncomingFile(relativePath: String) -> Bool {
payloadCoordination.lock.lock()
defer { payloadCoordination.lock.unlock() }
guard let payload = incomingPayloadURL(
relativePath: relativePath
) else {
return false
}
let standardizedPath = payload.standardizedFileURL.path
let reservedByDeletion = payloadCoordination.deletionReservations
.values.contains { $0.contains(standardizedPath) }
guard !reservedByDeletion,
!payloadCoordination.pendingDeliveryPaths.contains(
standardizedPath
),
let receiptOwnedPaths =
privateMediaReceipts.reservedPayloadPaths(),
!receiptOwnedPaths.contains(standardizedPath) else {
return false
}
guard fileManager.fileExists(atPath: payload.path),
(try? payload.resourceValues(
forKeys: [.isRegularFileKey]
).isRegularFile) == true else {
return false
}
do {
try fileManager.removeItem(at: payload)
return !fileManager.fileExists(atPath: payload.path)
} catch {
SecureLogger.warning(
"⚠️ Failed to remove explicitly deleted legacy media: \(error)",
category: .session
)
return false
}
}
/// Resolves a `files/`-relative path iff it lands directly inside one of
/// the incoming media directories. Anything else is not a deletable
/// incoming payload.
private func incomingPayloadURL(relativePath: String) -> URL? {
guard !relativePath.isEmpty,
let base = try? filesDirectory().standardizedFileURL else {
return nil
}
let candidate = base
.appendingPathComponent(relativePath, isDirectory: false)
.standardizedFileURL
let parentPath = candidate.deletingLastPathComponent().path
let incomingDirectories = [
"voicenotes/incoming",
"images/incoming",
"files/incoming"
]
guard incomingDirectories.contains(where: { relativeDirectory in
base.appendingPathComponent(
relativeDirectory,
isDirectory: true
).standardizedFileURL.path == parentPath
}) else {
return nil
}
return candidate
}
/// Releases the short window between disk save and synchronous
/// conversation insertion. Before this callback, a deletion transaction
/// may not infer ownership from a stale bubble that names the same path.
func finishIncomingFileDelivery(at storedURL: URL) {
payloadCoordination.lock.lock()
defer { payloadCoordination.lock.unlock() }
payloadCoordination.pendingDeliveryPaths.remove(
storedURL.standardizedFileURL.path
)
}
/// Best-effort rollback for a payload whose durable receipt commit failed.
func removeIncomingFile(at storedURL: URL) {
payloadCoordination.lock.lock()
defer { payloadCoordination.lock.unlock() }
payloadCoordination.pendingDeliveryPaths.remove(
storedURL.standardizedFileURL.path
)
guard isURLInsideFilesDirectory(storedURL) else { return }
do {
try fileManager.removeItem(at: storedURL)
@@ -314,6 +504,9 @@ struct BLEIncomingFileStore {
/// a finalized transfer can arrive at quota while a burst is still
/// streaming but they still count toward usage.
func enforceQuota(reservingBytes: Int) {
payloadCoordination.lock.lock()
defer { payloadCoordination.lock.unlock() }
do {
let base = try filesDirectory()
let incomingDirs = [
@@ -339,14 +532,26 @@ struct BLEIncomingFileStore {
}
let currentUsage = allFiles.reduce(0) { $0 + $1.size }
let targetUsage = Self.quotaBytes - Int64(reservingBytes)
let targetUsage = quotaBytes - Int64(reservingBytes)
guard currentUsage > targetUsage else { return }
let needToFree = currentUsage - targetUsage
let activeDeletionPaths = payloadCoordination
.deletionReservations.values.reduce(into: Set<String>()) {
$0.formUnion($1)
}
let protectedPaths = activeDeletionPaths.union(
payloadCoordination.pendingDeliveryPaths
)
var freedSpace: Int64 = 0
for file in allFiles.sorted(by: { $0.modified < $1.modified }) {
guard freedSpace < needToFree else { break }
guard !file.url.lastPathComponent.hasPrefix(Self.liveCapturePrefix) else { continue }
guard !protectedPaths.contains(
file.url.standardizedFileURL.path
) else {
continue
}
do {
try fileManager.removeItem(at: file.url)
freedSpace += file.size
@@ -434,11 +639,20 @@ struct BLEIncomingFileStore {
return candidate.isEmpty ? defaultName : candidate
}
private func uniqueFileURL(in directory: URL, fileName: String) -> URL {
private func uniqueFileURL(
in directory: URL,
fileName: String,
reservedPaths: Set<String>,
forceRandomizedName: Bool
) -> URL {
let directoryPath = directory.standardizedFileURL.path
func isInsideDirectory(_ url: URL) -> Bool {
url.standardizedFileURL.path.hasPrefix(directoryPath + "/")
}
func isAvailable(_ url: URL) -> Bool {
!reservedPaths.contains(url.standardizedFileURL.path)
&& !fileManager.fileExists(atPath: url.path)
}
var candidate = directory.appendingPathComponent(fileName)
guard isInsideDirectory(candidate) else {
@@ -446,19 +660,27 @@ struct BLEIncomingFileStore {
return directory.appendingPathComponent("blocked_\(UUID().uuidString)")
}
if !fileManager.fileExists(atPath: candidate.path) {
let baseName = (fileName as NSString).deletingPathExtension
let ext = (fileName as NSString).pathExtension
if forceRandomizedName {
let suffix = UUID().uuidString
let randomizedName = ext.isEmpty
? "\(baseName)_\(suffix)"
: "\(baseName)_\(suffix).\(ext)"
return directory.appendingPathComponent(randomizedName)
}
if isAvailable(candidate) {
return candidate
}
let baseName = (fileName as NSString).deletingPathExtension
let ext = (fileName as NSString).pathExtension
for counter in 1..<100 {
let newName = ext.isEmpty ? "\(baseName) (\(counter))" : "\(baseName) (\(counter)).\(ext)"
candidate = directory.appendingPathComponent(newName)
guard isInsideDirectory(candidate) else {
return directory.appendingPathComponent("blocked_\(UUID().uuidString)")
}
if !fileManager.fileExists(atPath: candidate.path) {
if isAvailable(candidate) {
return candidate
}
}
@@ -18,17 +18,25 @@ enum BLEPrivateMediaReceiptState: Equatable {
/// Each ID has its own atomic record so one hot lookup never rewrites or
/// decodes the entire ledger. The process-lifetime index is installed only
/// after a complete directory scan. A structural failure of the directory
/// itself (create/enumerate) remains globally fail-closed and retryable.
/// An individual record that cannot be read, decoded, or validated is
/// quarantined instead: the file is moved aside with a `.corrupt` suffix,
/// excluded from future scans, and only that ID stays fail-closed the rest
/// of the ledger keeps working, so one damaged record can never make every
/// inbound private media payload vanish.
/// itself (create/enumerate) or of the single batch deletion journal
/// remains globally fail-closed and retryable. An individual record that
/// cannot be read, decoded, or validated is quarantined instead: the file is
/// moved aside with a `.corrupt` suffix, excluded from future scans, and only
/// that ID stays fail-closed the rest of the ledger keeps working, so one
/// damaged record can never make every inbound private media payload vanish.
final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
typealias DirectoryReader = (_ directory: URL) throws -> [URL]
typealias DataReader = (_ url: URL) throws -> Data
typealias DataWriter = (
_ data: Data,
_ url: URL,
_ options: Data.WritingOptions
) throws -> Void
typealias PayloadRemover = (_ url: URL) throws -> Void
private static let receiptDirectoryName = ".private-media-receipts"
private static let quarantinePathExtension = "corrupt"
private static let deletionJournalFileName = ".deletion-journal.json"
private static let maximumDeletionPathsPerMessage = 2
private struct ReceiptRecord: Codable, Equatable {
enum Kind: String, Codable {
@@ -43,6 +51,19 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
let recordedAt: Date
}
/// One atomic write of this journal is the commit point for an entire
/// explicit-deletion batch. Per-ID records and payload unlinks are
/// idempotent materialization performed only after that commit.
private struct DeletionJournalEntry: Codable, Equatable {
let relativePaths: [String]
let recordedAt: Date
}
private struct DeletionJournal: Codable {
let version: Int
let entries: [String: DeletionJournalEntry]
}
private final class Runtime: @unchecked Sendable {
let lock = NSLock()
var records: [String: ReceiptRecord]?
@@ -50,7 +71,7 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
/// together with `records`; these IDs stay fail-closed while every
/// other record keeps serving.
var quarantined: Set<String> = []
var volatileTombstones: [String: Date] = [:]
var deletionJournal: [String: DeletionJournalEntry]?
}
private let fileManager: FileManager
@@ -60,6 +81,8 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
private let now: () -> Date
private let directoryReader: DirectoryReader?
private let dataReader: DataReader?
private let dataWriter: DataWriter
private let payloadRemover: PayloadRemover
private let runtime = Runtime()
init(
@@ -69,7 +92,11 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
ttl: TimeInterval = TransportConfig.privateMediaReceivedLedgerTTLSeconds,
now: @escaping () -> Date = Date.init,
directoryReader: DirectoryReader? = nil,
dataReader: DataReader? = nil
dataReader: DataReader? = nil,
dataWriter: @escaping DataWriter = {
try $0.write(to: $1, options: $2)
},
payloadRemover: PayloadRemover? = nil
) {
self.fileManager = fileManager
self.baseDirectory = baseDirectory
@@ -78,6 +105,10 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
self.now = now
self.directoryReader = directoryReader
self.dataReader = dataReader
self.dataWriter = dataWriter
self.payloadRemover = payloadRemover ?? {
try fileManager.removeItem(at: $0)
}
}
/// Drops process-lifetime decisions after the enclosing media directory
@@ -88,7 +119,7 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
runtime.lock.lock()
runtime.records = nil
runtime.quarantined.removeAll(keepingCapacity: false)
runtime.volatileTombstones.removeAll(keepingCapacity: false)
runtime.deletionJournal = nil
runtime.lock.unlock()
}
@@ -101,35 +132,61 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
defer { runtime.lock.unlock() }
let date = now()
if let tombstonedAt = runtime.volatileTombstones[messageID] {
if !isExpired(tombstonedAt, at: date) {
return .tombstoned
}
runtime.volatileTombstones.removeValue(forKey: messageID)
}
guard let directory = resolvedReceiptDirectory(),
var records = loadIndexIfNeeded(from: directory, at: date) else {
loadDurableStateIfNeeded(from: directory, at: date) else {
return .unavailable
}
_ = recoverDeletionJournal(in: directory)
// Committed deletion intent outranks quarantine: a pending journal
// entry proves the payload must stay deleted no matter what state
// the per-ID record file is in.
if runtime.deletionJournal?[messageID] != nil {
return .tombstoned
}
// A quarantined record could have been an acceptance or a tombstone;
// only this ID fails closed, so a retry can neither resurrect deleted
// media nor double-deliver, while every other payload keeps working.
if runtime.quarantined.contains(messageID) {
return .unavailable
}
guard let record = records[messageID] else { return .absent }
guard var records = runtime.records else { return .unavailable }
guard var record = records[messageID] else { return .absent }
if isExpired(record.recordedAt, at: date) {
if record.kind == .tombstone, record.relativePath != nil {
guard scrubLegacyTombstone(
record,
messageID: messageID,
in: directory,
records: &records
) else {
return .tombstoned
}
guard let scrubbed = records[messageID] else {
return .unavailable
}
record = scrubbed
}
let retainsAcceptedPath =
record.kind == .accepted
&& record.relativePath.flatMap(existingPayload) != nil
if isExpired(record.recordedAt, at: date),
!retainsAcceptedPath {
guard removeRecord(
messageID: messageID,
from: directory
) else {
return record.kind == .tombstone
? .tombstoned
: .unavailable
}
records.removeValue(forKey: messageID)
runtime.records = records
removeRecord(messageID: messageID, from: directory)
return .absent
}
switch record.kind {
case .tombstone:
removePayloadRecordedByTombstone(record)
return .tombstoned
case .accepted:
@@ -137,9 +194,14 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
let existingURL = existingPayload(relativePath: relativePath) else {
// Quota cleanup is not explicit deletion. Remove the stale
// receipt so a sender retry can restore the payload and bubble.
guard removeRecord(
messageID: messageID,
from: directory
) else {
return .unavailable
}
records.removeValue(forKey: messageID)
runtime.records = records
removeRecord(messageID: messageID, from: directory)
return .absent
}
return .accepted(existingURL)
@@ -160,14 +222,25 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
defer { runtime.lock.unlock() }
let date = now()
if let tombstonedAt = runtime.volatileTombstones[messageID],
!isExpired(tombstonedAt, at: date) {
guard let directory = resolvedReceiptDirectory(),
loadDurableStateIfNeeded(from: directory, at: date) else {
return false
}
guard let directory = resolvedReceiptDirectory(),
var records = loadIndexIfNeeded(from: directory, at: date),
!runtime.quarantined.contains(messageID) else {
_ = recoverDeletionJournal(in: directory)
guard runtime.deletionJournal?[messageID] == nil,
!runtime.quarantined.contains(messageID),
var records = runtime.records else {
return false
}
if runtime.deletionJournal?.values.contains(where: {
$0.relativePaths.contains(relativePath)
}) == true {
return false
}
if records.contains(where: { existingMessageID, record in
existingMessageID != messageID
&& record.relativePath == relativePath
}) {
return false
}
if let existing = records[messageID],
@@ -205,74 +278,238 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
return true
}
/// Foundation for explicit media deletion. This branch does not wire the
/// chat-clear UI; it only makes a tombstone durable and fail closed.
func recordDeleted(messageID: String) -> Bool {
guard PrivateMediaMessageIdentity.isStableID(messageID) else {
return false
}
/// Atomically commits explicit deletion for every stable ID in `messageIDs`.
///
/// The journal is the single batch commit point: a failed write changes
/// neither durable nor in-memory receiver state, so callers must preserve
/// their bubbles. Once the write succeeds, every entry is tombstoned even
/// if the process exits before per-ID materialization or payload unlink.
/// Recovery retries both operations on the next lookup or launch.
func recordDeleted(
messageIDs: [String],
payloadRelativePaths: [String: String] = [:],
protectedPayloadRelativePaths: Set<String> = []
) -> Bool {
let stableIDs = Array(
Set(messageIDs.filter(PrivateMediaMessageIdentity.isStableID))
).sorted()
guard !stableIDs.isEmpty else { return true }
guard stableIDs.count <= capacity else { return false }
runtime.lock.lock()
defer { runtime.lock.unlock() }
let date = now()
addVolatileTombstone(messageID, at: date)
guard let directory = resolvedReceiptDirectory(),
var records = loadIndexIfNeeded(from: directory, at: date),
!runtime.quarantined.contains(messageID) else {
runtime.volatileTombstones.removeValue(forKey: messageID)
loadDurableStateIfNeeded(from: directory, at: date) else {
return false
}
if let existing = records[messageID],
existing.kind == .tombstone,
!isExpired(existing.recordedAt, at: date) {
runtime.volatileTombstones.removeValue(forKey: messageID)
removePayloadRecordedByTombstone(existing)
_ = recoverDeletionJournal(in: directory)
guard let records = runtime.records,
var journal = runtime.deletionJournal else {
return false
}
// Quarantined IDs need no special case here: their record is never
// indexed, so deleting one requires a caller-supplied payload path
// (the general pathless-deletion refusal below rejects it otherwise),
// and materialization never rewrites a quarantined ID's record.
let pendingIDs = Set(journal.keys)
let newIDs = stableIDs.filter { messageID in
if pendingIDs.contains(messageID) { return false }
if let current = records[messageID],
current.kind == .tombstone,
!isExpired(current.recordedAt, at: date) {
return false
}
return true
}
let victim = capacityVictim(
for: .tombstone,
replacing: messageID,
in: records
)
if records[messageID]?.kind != .tombstone,
records.values.lazy.filter({ $0.kind == .tombstone }).count >= capacity,
victim == nil {
runtime.volatileTombstones.removeValue(forKey: messageID)
guard !newIDs.isEmpty else {
return true
}
guard Set(journal.keys).union(newIDs).count <= capacity else {
return false
}
let tombstone = ReceiptRecord(
kind: .tombstone,
// Retain the accepted path so a crash between the atomic record
// write and payload unlink can finish cleanup after relaunch.
relativePath: records[messageID]?.relativePath,
recordedAt: date
)
guard persist(tombstone, messageID: messageID, to: directory) else {
runtime.volatileTombstones.removeValue(forKey: messageID)
var newEntries: [String: DeletionJournalEntry] = [:]
for messageID in newIDs {
// Accepted receipts normally supply the exact stored path. The UI
// fallback is required when that receipt aged or was capacity
// evicted while its bubble and payload remain. They may differ
// after a retry selected a suffixed filename, so journal both.
// Never commit a new pathless tombstone: it could retire
// successfully while leaving an untracked payload behind.
let relativePaths = Array(Set([
records[messageID]?.relativePath,
payloadRelativePaths[messageID]
].compactMap { $0 })).sorted()
guard !relativePaths.isEmpty,
relativePaths.count <=
Self.maximumDeletionPathsPerMessage,
relativePaths.allSatisfy({
isSafeDeletionTarget(relativePath: $0)
}),
protectedPayloadRelativePaths.isDisjoint(
with: relativePaths
),
!records.contains(where: { otherMessageID, record in
otherMessageID != messageID
&& !stableIDs.contains(otherMessageID)
&& record.relativePath.map(
relativePaths.contains
) == true
}),
!journal.contains(where: { otherMessageID, entry in
otherMessageID != messageID
&& !stableIDs.contains(otherMessageID)
&& !Set(entry.relativePaths).isDisjoint(
with: relativePaths
)
}) else {
return false
}
newEntries[messageID] = DeletionJournalEntry(
// The journal retains every owned path so payload deletion
// remains recoverable across a crash or unlink failure.
relativePaths: relativePaths,
recordedAt: date
)
}
journal.merge(newEntries) { _, new in new }
// Do not install any process-local tombstone before this succeeds.
// A failed delete must continue to resolve to its prior accepted
// state, otherwise a retry could be falsely ACKed while UI remains.
guard persistDeletionJournal(journal, in: directory) else {
return false
}
records[messageID] = tombstone
if let victim, victim != messageID {
records.removeValue(forKey: victim)
removeRecord(messageID: victim, from: directory)
}
runtime.records = records
runtime.volatileTombstones.removeValue(forKey: messageID)
removePayloadRecordedByTombstone(tombstone)
runtime.deletionJournal = journal
_ = recoverDeletionJournal(in: directory)
return true
}
private func loadIndexIfNeeded(
func recordDeleted(messageID: String) -> Bool {
guard PrivateMediaMessageIdentity.isStableID(messageID) else {
return false
}
return recordDeleted(messageIDs: [messageID])
}
/// Resolves every path a deletion transaction may target. The incoming
/// allocator reserves this set at the journal barrier so a concurrent raw
/// arrival cannot reuse a missing UI fallback.
func prospectiveDeletionPayloadPaths(
messageIDs: [String],
payloadRelativePaths: [String: String]
) -> Set<String>? {
let stableIDs = Set(
messageIDs.filter(PrivateMediaMessageIdentity.isStableID)
)
guard !stableIDs.isEmpty else { return [] }
runtime.lock.lock()
defer { runtime.lock.unlock() }
let date = now()
guard let directory = resolvedReceiptDirectory(),
loadDurableStateIfNeeded(from: directory, at: date) else {
return nil
}
_ = recoverDeletionJournal(in: directory)
guard let journal = runtime.deletionJournal,
let records = runtime.records else {
return nil
}
var paths: Set<String> = []
for messageID in stableIDs {
if let entry = journal[messageID] {
paths.formUnion(entry.relativePaths)
continue
}
if let relativePath = records[messageID]?.relativePath {
paths.insert(relativePath)
}
if let relativePath = payloadRelativePaths[messageID] {
paths.insert(relativePath)
}
}
guard paths.allSatisfy({
candidatePayload(relativePath: $0) != nil
}) else {
return nil
}
return Set(paths.compactMap {
candidatePayload(relativePath: $0)?
.standardizedFileURL.path
})
}
/// Paths owned by accepted receipts, legacy pathful tombstones, or the
/// deletion journal. Incoming allocation must not reuse any of them for a
/// different ID. Quarantined records are unreadable, so any path they may
/// have owned cannot be reserved; their bytes remain preserved in the
/// `.corrupt` file for offline inspection.
func reservedPayloadPaths() -> Set<String>? {
runtime.lock.lock()
defer { runtime.lock.unlock() }
let date = now()
guard let directory = resolvedReceiptDirectory(),
loadDurableStateIfNeeded(from: directory, at: date) else {
return nil
}
_ = recoverDeletionJournal(in: directory)
guard let journal = runtime.deletionJournal,
let records = runtime.records else {
return nil
}
let recordPaths = records.values.compactMap(\.relativePath)
let journalPaths = journal.values.flatMap(\.relativePaths)
return Set((recordPaths + journalPaths).compactMap { relativePath in
candidatePayload(relativePath: relativePath)?
.standardizedFileURL.path
})
}
private func scrubLegacyTombstone(
_ tombstone: ReceiptRecord,
messageID: String,
in directory: URL,
records: inout [String: ReceiptRecord]
) -> Bool {
guard tombstone.kind == .tombstone,
tombstone.relativePath != nil else {
return true
}
// Pre-journal tombstones cannot prove that the current file is still
// the payload they originally described. Older allocators did not
// reserve these paths, so a raw/public arrival may have reused the
// basename. Shed the ambiguous path without unlinking anything.
let pathless = ReceiptRecord(
kind: .tombstone,
relativePath: nil,
recordedAt: tombstone.recordedAt
)
guard persist(
pathless,
messageID: messageID,
to: directory
) else {
return false
}
records[messageID] = pathless
runtime.records = records
return true
}
private func loadDurableStateIfNeeded(
from directory: URL,
at date: Date
) -> [String: ReceiptRecord]? {
if let records = runtime.records {
return records
) -> Bool {
if runtime.records != nil, runtime.deletionJournal != nil {
return true
}
do {
@@ -286,7 +523,7 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
"❌ Failed to create private-media receipt directory: \(error)",
category: .session
)
return nil
return false
}
let urls: [URL]
@@ -305,13 +542,14 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
"❌ Failed to enumerate private-media receipts: \(error)",
category: .session
)
return nil
return false
}
var records: [String: ReceiptRecord] = [:]
var scannedRecords: [String: ReceiptRecord] = [:]
var quarantined: Set<String> = []
var expired: [String] = []
var tombstones: [ReceiptRecord] = []
var tombstones: [(messageID: String, record: ReceiptRecord)] = []
for url in urls {
// Records quarantined by an earlier scan stay fail-closed on
// every launch without being re-read: only their ID matters.
@@ -354,39 +592,271 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
quarantined.insert(messageID)
continue
}
if isExpired(record.recordedAt, at: date) {
scannedRecords[messageID] = record
if record.kind == .tombstone {
tombstones.append((messageID, record))
}
let retainsAcceptedPath =
record.kind == .accepted
&& record.relativePath.flatMap(existingPayload) != nil
if isExpired(record.recordedAt, at: date),
!retainsAcceptedPath {
expired.append(messageID)
continue
}
records[messageID] = record
if record.kind == .tombstone {
tombstones.append(record)
}
}
// A readable duplicate of a quarantined ID must not override the
// fail-closed decision.
// fail-closed decision, not even through the failed-prune restore
// or legacy-tombstone scrub paths below.
for messageID in quarantined {
records.removeValue(forKey: messageID)
scannedRecords.removeValue(forKey: messageID)
}
tombstones.removeAll { quarantined.contains($0.messageID) }
let overflow = overflowVictims(in: records)
for messageID in overflow {
records.removeValue(forKey: messageID)
}
// Install the index only after the whole directory was scanned.
// Cleanup cannot influence a failed scan.
runtime.records = records
runtime.quarantined = quarantined
let journal: [String: DeletionJournalEntry]
let journalURL = deletionJournalURL(in: directory)
if fileManager.fileExists(atPath: journalURL.path) {
do {
let data = try dataReader?(journalURL)
?? Data(contentsOf: journalURL)
let snapshot = try JSONDecoder().decode(
DeletionJournal.self,
from: data
)
guard snapshot.version == 1,
snapshot.entries.count <= capacity,
snapshot.entries.allSatisfy({ messageID, entry in
PrivateMediaMessageIdentity.isStableID(messageID)
&& !entry.relativePaths.isEmpty
&& entry.relativePaths.count <=
Self.maximumDeletionPathsPerMessage
&& Set(entry.relativePaths).count
== entry.relativePaths.count
&& entry.relativePaths.allSatisfy {
candidatePayload(relativePath: $0) != nil
}
}) else {
// The journal is a single batch commit file: like a
// directory-level failure it stays globally fail-closed
// and retryable, never quarantined per-ID.
SecureLogger.error(
"❌ Invalid private-media deletion journal",
category: .session
)
return false
}
journal = snapshot.entries
} catch {
// The journal is the all-ID commit record. It may never be
// skipped or treated as empty when unreadable.
SecureLogger.error(
"❌ Failed to read private-media deletion journal: \(error)",
category: .session
)
return false
}
} else {
journal = [:]
}
var protectedFromPruning: Set<String> = []
// Old per-ID tombstones may still carry a payload path from before the
// deletion journal existed. That path is inherently ambiguous because
// old allocators did not reserve it. Convert it to a pathless
// tombstone without unlinking any current file.
for (messageID, tombstone) in tombstones
where journal[messageID] == nil
&& tombstone.relativePath != nil {
let pathless = ReceiptRecord(
kind: .tombstone,
relativePath: nil,
recordedAt: tombstone.recordedAt
)
guard persist(
pathless,
messageID: messageID,
to: directory
) else {
records[messageID] = tombstone
protectedFromPruning.insert(messageID)
continue
}
scannedRecords[messageID] = pathless
if records[messageID] != nil {
records[messageID] = pathless
}
}
for messageID in expired + overflow {
removeRecord(messageID: messageID, from: directory)
guard !protectedFromPruning.contains(messageID) else { continue }
if !removeRecord(messageID: messageID, from: directory),
let record = scannedRecords[messageID] {
records[messageID] = record
}
}
for tombstone in tombstones {
removePayloadRecordedByTombstone(tombstone)
// Install caches only after every per-ID record and the batch journal
// have been read and validated. Failed legacy cleanup remains indexed
// and path-reserved instead of blocking unrelated media.
runtime.records = records
runtime.quarantined = quarantined
runtime.deletionJournal = journal
return true
}
/// Idempotently materializes the write-ahead journal. An entry leaves the
/// journal only after its per-ID tombstone is durable and every recorded
/// payload is absent. Any failure keeps the journal authoritative for a
/// later lookup or process restart.
@discardableResult
private func recoverDeletionJournal(in directory: URL) -> Bool {
guard let journal = runtime.deletionJournal,
!journal.isEmpty,
var records = runtime.records else {
return true
}
return records
let preservedMessageIDs = Set(journal.keys)
var remaining = journal
let orderedEntries = journal.sorted { lhs, rhs in
if lhs.value.recordedAt == rhs.value.recordedAt {
return lhs.key < rhs.key
}
return lhs.value.recordedAt < rhs.value.recordedAt
}
for (messageID, journalEntry) in orderedEntries {
// Never materialize a per-ID record for a quarantined ID: the
// scanner ignores readable duplicates of a quarantined record,
// and quarantine is already permanently fail-closed
// (state == .unavailable, commitAccepted refuses it), which
// subsumes the tombstone's no-resurrection guarantee. Only the
// payload unlinks below still need to run before the entry can
// retire.
if !runtime.quarantined.contains(messageID) {
let pathlessTombstone = ReceiptRecord(
kind: .tombstone,
relativePath: nil,
recordedAt: journalEntry.recordedAt
)
if records[messageID] != pathlessTombstone {
let victim = capacityVictim(
for: .tombstone,
replacing: messageID,
in: records,
preserving: preservedMessageIDs
)
if records[messageID]?.kind != .tombstone,
records.values.lazy.filter({
$0.kind == .tombstone
}).count >= capacity,
victim == nil {
continue
}
guard persist(
pathlessTombstone,
messageID: messageID,
to: directory
) else {
continue
}
records[messageID] = pathlessTombstone
if let victim, victim != messageID {
records.removeValue(forKey: victim)
removeRecord(messageID: victim, from: directory)
}
}
}
// Only the journal retains the paths. Once every unlink succeeds,
// the durable per-ID tombstone is pathless and cannot later
// delete a different payload that reused the basename.
let removedEveryPayload = journalEntry.relativePaths.allSatisfy {
removePayloadRecordedByTombstone(ReceiptRecord(
kind: .tombstone,
relativePath: $0,
recordedAt: journalEntry.recordedAt
))
}
guard removedEveryPayload else {
continue
}
remaining.removeValue(forKey: messageID)
}
runtime.records = records
guard remaining != journal else { return false }
if remaining.isEmpty {
guard removeDeletionJournal(in: directory) else { return false }
} else {
guard persistDeletionJournal(remaining, in: directory) else {
return false
}
}
runtime.deletionJournal = remaining
return remaining.isEmpty
}
private func persistDeletionJournal(
_ entries: [String: DeletionJournalEntry],
in directory: URL
) -> Bool {
do {
try fileManager.createDirectory(
at: directory,
withIntermediateDirectories: true,
attributes: nil
)
let data = try JSONEncoder().encode(
DeletionJournal(version: 1, entries: entries)
)
var options: Data.WritingOptions = [.atomic]
#if os(iOS)
options.insert(
.completeFileProtectionUntilFirstUserAuthentication
)
#endif
try dataWriter(data, deletionJournalURL(in: directory), options)
return true
} catch {
SecureLogger.error(
"❌ Failed to persist private-media deletion journal: \(error)",
category: .session
)
return false
}
}
private func removeDeletionJournal(in directory: URL) -> Bool {
let url = deletionJournalURL(in: directory)
guard fileManager.fileExists(atPath: url.path) else { return true }
do {
try fileManager.removeItem(at: url)
return true
} catch {
SecureLogger.warning(
"⚠️ Failed to retire private-media deletion journal: \(error)",
category: .session
)
return false
}
}
private func deletionJournalURL(in directory: URL) -> URL {
directory.appendingPathComponent(
Self.deletionJournalFileName,
isDirectory: false
)
}
/// Moves an unreadable record aside so future scans skip it while its ID
@@ -437,10 +907,14 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
) -> [String] {
var victims: [String] = []
for kind in [ReceiptRecord.Kind.accepted, .tombstone] {
let matching = records.filter { $0.value.kind == kind }
let overflow = matching.count - capacity
let allMatching = records.filter { $0.value.kind == kind }
let overflow = allMatching.count - capacity
guard overflow > 0 else { continue }
victims.append(contentsOf: matching.sorted { lhs, rhs in
let eligible = allMatching.filter { _, record in
guard kind == .accepted else { return true }
return record.relativePath.flatMap(existingPayload) == nil
}
victims.append(contentsOf: eligible.sorted { lhs, rhs in
if lhs.value.recordedAt == rhs.value.recordedAt {
return lhs.key < rhs.key
}
@@ -457,13 +931,24 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
private func capacityVictim(
for incomingKind: ReceiptRecord.Kind,
replacing messageID: String,
in records: [String: ReceiptRecord]
in records: [String: ReceiptRecord],
preserving preservedMessageIDs: Set<String> = []
) -> String? {
guard records[messageID]?.kind != incomingKind else { return nil }
// During a live session an accepted receipt is the only durable owner
// of its filename, even after quota removes the payload. Evicting it
// here could let another ID reuse the path while the old bubble still
// exists. The large capacity therefore acts as admission control.
guard incomingKind != .accepted else { return nil }
let matching = records.filter {
$0.key != messageID && $0.value.kind == incomingKind
$0.key != messageID
&& !preservedMessageIDs.contains($0.key)
&& $0.value.kind == incomingKind
}
guard matching.count >= capacity else { return nil }
let currentCount = records.values.lazy.filter {
$0.kind == incomingKind
}.count
guard currentCount >= capacity else { return nil }
return matching.min { lhs, rhs in
if lhs.value.recordedAt == rhs.value.recordedAt {
return lhs.key < rhs.key
@@ -489,7 +974,7 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
options.insert(.completeFileProtectionUntilFirstUserAuthentication)
#endif
let url = recordURL(messageID: messageID, in: directory)
try data.write(to: url, options: options)
try dataWriter(data, url, options)
return true
} catch {
SecureLogger.error(
@@ -500,16 +985,22 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
}
}
private func removeRecord(messageID: String, from directory: URL) {
@discardableResult
private func removeRecord(
messageID: String,
from directory: URL
) -> Bool {
let url = recordURL(messageID: messageID, in: directory)
guard fileManager.fileExists(atPath: url.path) else { return }
guard fileManager.fileExists(atPath: url.path) else { return true }
do {
try fileManager.removeItem(at: url)
return true
} catch {
SecureLogger.warning(
"⚠️ Failed to prune private-media receipt \(messageID.prefix(12))…: \(error)",
category: .session
)
return false
}
}
@@ -519,39 +1010,55 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
.appendingPathExtension("json")
}
private func removePayloadRecordedByTombstone(_ record: ReceiptRecord) {
@discardableResult
private func removePayloadRecordedByTombstone(
_ record: ReceiptRecord
) -> Bool {
guard record.kind == .tombstone,
let relativePath = record.relativePath,
let payload = candidatePayload(relativePath: relativePath),
fileManager.fileExists(atPath: payload.path) else {
return
return true
}
guard let values = try? payload.resourceValues(
forKeys: [.isRegularFileKey]
),
values.isRegularFile == true else {
SecureLogger.warning(
"⚠️ Refusing to remove non-file private-media payload",
category: .session
)
return false
}
do {
try fileManager.removeItem(at: payload)
try payloadRemover(payload)
return !fileManager.fileExists(atPath: payload.path)
} catch {
SecureLogger.warning(
"⚠️ Failed to remove explicitly deleted private media: \(error)",
category: .session
)
return false
}
}
private func addVolatileTombstone(_ messageID: String, at date: Date) {
runtime.volatileTombstones[messageID] = date
let overflow = runtime.volatileTombstones.count - capacity
guard overflow > 0 else { return }
let oldest = runtime.volatileTombstones.sorted {
if $0.value == $1.value { return $0.key < $1.key }
return $0.value < $1.value
private func isSafeDeletionTarget(relativePath: String) -> Bool {
guard let payload = candidatePayload(relativePath: relativePath) else {
return false
}
for (oldMessageID, _) in oldest.prefix(overflow) {
runtime.volatileTombstones.removeValue(forKey: oldMessageID)
guard fileManager.fileExists(atPath: payload.path) else {
return true
}
return (try? payload.resourceValues(
forKeys: [.isRegularFileKey]
).isRegularFile) == true
}
private func validExistingPayload(_ url: URL) -> URL? {
let standardized = url.standardizedFileURL
guard isInsideFilesDirectory(standardized) else { return nil }
guard isInsideIncomingMediaDirectory(standardized) else {
return nil
}
var isDirectory: ObjCBool = false
guard fileManager.fileExists(
atPath: standardized.path,
@@ -588,15 +1095,27 @@ final class BLEPrivateMediaReceiptStore: @unchecked Sendable {
let candidate = filesRoot
.appendingPathComponent(relativePath, isDirectory: false)
.standardizedFileURL
guard candidate.path.hasPrefix(filesRoot.path + "/") else { return nil }
guard isInsideIncomingMediaDirectory(candidate) else { return nil }
return candidate
}
private func isInsideFilesDirectory(_ url: URL) -> Bool {
private func isInsideIncomingMediaDirectory(_ url: URL) -> Bool {
guard let filesRoot = try? filesDirectory().standardizedFileURL else {
return false
}
return url.standardizedFileURL.path.hasPrefix(filesRoot.path + "/")
let parentPath = url.standardizedFileURL
.deletingLastPathComponent().path
return [
"voicenotes/incoming",
"images/incoming",
"files/incoming"
].contains { relativeDirectory in
filesRoot.appendingPathComponent(
relativeDirectory,
isDirectory: true
)
.standardizedFileURL.path == parentPath
}
}
private func resolvedReceiptDirectory() -> URL? {
+148 -17
View File
@@ -2697,6 +2697,18 @@ final class BLEService: NSObject {
removeIncomingFile: { [weak self] storedURL in
self?.incomingFileStore.removeIncomingFile(at: storedURL)
},
finishIncomingFileDelivery: { [weak self] storedURL in
// Serialize pending-owner release behind deletion barriers.
// If /clear snapshots before this UI insertion, its already
// queued barrier must still observe the path as pending. If
// insertion wins first, the next MainActor snapshot sees the
// new bubble and protects the path explicitly.
self?.messageQueue.async(flags: .barrier) {
self?.incomingFileStore.finishIncomingFileDelivery(
at: storedURL
)
}
},
isPrivateMediaSenderBlocked: { [weak self] peerID in
guard let self else { return false }
let senderStaticKey = self.noiseService.getPeerPublicKeyData(peerID)
@@ -2721,11 +2733,12 @@ final class BLEService: NSObject {
}
self.sendDeliveryAck(for: messageID, to: peerID)
},
deliverMessage: { [weak self] message, shouldDeliver, completion in
deliverMessage: { [weak self] message, shouldDeliver, completion, finalization in
self?.emitTransportEvent(
.messageReceived(message),
shouldDeliver: shouldDeliver,
completion: completion
completion: completion,
finalization: finalization
)
}
)
@@ -3532,6 +3545,18 @@ extension BLEService {
}
}
func _test_emitTransportEvent(
_ event: TransportEvent,
completion: @escaping () -> Void,
finalization: @escaping (TransportEventDeliveryOutcome) -> Void
) {
emitTransportEvent(
event,
completion: completion,
finalization: finalization
)
}
func _test_handlePacket(_ packet: BitchatPacket, fromPeerID: PeerID, preseedPeer: Bool = true, signingPublicKey: Data? = nil) {
if preseedPeer {
// Ensure the synthetic peer is known and marked verified for public-message tests
@@ -4451,8 +4476,95 @@ extension BLEService {
// No alias rotation or advertising restarts required.
}
// MARK: - Private Media Deletion
extension BLEService: PrivateMediaDeletionPersisting {
@MainActor
func persistDeletedPrivateMedia(
messageIDs: [String],
payloadRelativePaths: [String: String],
protectedPayloadRelativePaths: Set<String>,
completion: @escaping @MainActor (Bool) -> Void
) {
let fileStore = incomingFileStore
messageQueue.async(flags: .barrier) {
guard let reservation = fileStore
.reservePrivateMediaDeletion(
messageIDs: messageIDs,
payloadRelativePaths: payloadRelativePaths
) else {
Task { @MainActor in
completion(false)
}
return
}
let persisted = fileStore
.commitPrivateMediaDeletion(
reservation: reservation,
messageIDs: messageIDs,
payloadRelativePaths: payloadRelativePaths,
protectedPayloadRelativePaths:
protectedPayloadRelativePaths
)
Task { @MainActor in
completion(persisted)
}
}
}
@MainActor
func removeLegacyPrivateMediaPayload(relativePath: String) {
let fileStore = incomingFileStore
messageQueue.async(flags: .barrier) {
fileStore.removeLegacyIncomingFile(relativePath: relativePath)
}
}
}
// MARK: - Private Helpers
enum TransportEventDeliveryOutcome: Equatable {
/// A synchronous sink inserted the message and revalidation succeeded.
case accepted
/// A supported plain delegate was invoked, but insertion cannot be
/// confirmed synchronously.
case invokedUnconfirmed
/// No sink accepted the event, or receipt revalidation rejected it.
case rejected
}
enum TransportEventDeliveryGate {
/// Runs finalization exactly once for every attempted main-actor delivery,
/// including pre-insertion rejection, a missing/rejecting sink, and
/// post-insertion revalidation failure. Only a fully accepted delivery
/// runs `completion` (for example, a stable-media ACK).
@MainActor
static func attempt(
shouldDeliver: () -> Bool,
deliver: () -> TransportEventDeliveryOutcome,
completion: () -> Void,
finalization: (TransportEventDeliveryOutcome) -> Void
) {
var outcome = TransportEventDeliveryOutcome.rejected
defer { finalization(outcome) }
guard shouldDeliver() else {
return
}
switch deliver() {
case .rejected:
return
case .invokedUnconfirmed:
outcome = .invokedUnconfirmed
return
case .accepted:
break
}
guard shouldDeliver() else { return }
outcome = .accepted
completion()
}
}
extension BLEService {
/// Notify UI on the MainActor to satisfy Swift concurrency isolation
@@ -4476,19 +4588,32 @@ extension BLEService {
private func emitTransportEvent(
_ event: TransportEvent,
shouldDeliver: (() -> Bool)? = nil,
completion: (() -> Void)? = nil
completion: (() -> Void)? = nil,
finalization: ((TransportEventDeliveryOutcome) -> Void)? = nil
) {
notifyUI { [weak self] in
guard let generation = capturePanicLifecycleGeneration() else {
Task { @MainActor in
finalization?(.rejected)
}
return
}
Task { @MainActor [weak self] in
guard let self,
shouldDeliver?() ?? true,
self.deliverTransportEvent(event),
// Quota cleanup can race the asynchronous main-actor hop or
// the synchronous ConversationStore upsert. ACK only while
// the exact durable mapping and file still resolve.
shouldDeliver?() ?? true else {
self.isCurrentPanicLifecycleGeneration(generation) else {
finalization?(.rejected)
return
}
completion?()
TransportEventDeliveryGate.attempt(
shouldDeliver: {
self.isCurrentPanicLifecycleGeneration(generation)
&& (shouldDeliver?() ?? true)
},
deliver: {
return self.deliverTransportEvent(event)
},
completion: { completion?() },
finalization: { finalization?($0) }
)
}
}
@@ -4508,34 +4633,40 @@ extension BLEService {
/// event and `false` when no delegate is installed.
@MainActor
@discardableResult
private func deliverTransportEvent(_ event: TransportEvent) -> Bool {
private func deliverTransportEvent(
_ event: TransportEvent
) -> TransportEventDeliveryOutcome {
if case .messageReceived(let message) = event {
if let synchronousDelegate =
eventDelegate as? SynchronousMessageTransportEventDelegate {
return synchronousDelegate
.didReceiveTransportMessageSynchronously(message)
? .accepted
: .rejected
}
if let eventDelegate {
eventDelegate.didReceiveTransportEvent(event)
return false
return .invokedUnconfirmed
}
if let synchronousDelegate =
delegate as? SynchronousMessageTransportEventDelegate {
return synchronousDelegate
.didReceiveTransportMessageSynchronously(message)
? .accepted
: .rejected
}
}
if let eventDelegate {
eventDelegate.didReceiveTransportEvent(event)
return true
return .accepted
} else {
guard let delegate else { return false }
guard let delegate else { return .rejected }
delegate.receiveTransportEvent(event)
if case .messageReceived = event {
return false
return .invokedUnconfirmed
}
return true
return .accepted
}
}
+21
View File
@@ -97,6 +97,27 @@ enum PrivateMediaSendPolicy: Equatable {
case blockedDowngrade
}
/// Receiver-only persistence surface for explicit private-media deletion.
/// Kept separate from `Transport` so sender retry branches can rebase without
/// inheriting or implementing receiver storage concerns.
protocol PrivateMediaDeletionPersisting: AnyObject {
@MainActor
func persistDeletedPrivateMedia(
messageIDs: [String],
payloadRelativePaths: [String: String],
protectedPayloadRelativePaths: Set<String>,
completion: @escaping @MainActor (Bool) -> Void
)
/// Gated unlink for a LEGACY (non-stable-ID) incoming payload whose
/// bubble was explicitly removed. Implementations delete the file only
/// when its path is not pending delivery and not reserved by any receipt
/// or in-flight deletion transaction; otherwise the file stays for
/// bounded quota cleanup.
@MainActor
func removeLegacyPrivateMediaPayload(relativePath: String)
}
protocol TransportEventDelegate: AnyObject {
@MainActor func didReceiveTransportEvent(_ event: TransportEvent)
}
@@ -86,7 +86,14 @@ protocol ChatMediaTransferContext: AnyObject {
@discardableResult
func appendPublicMessage(_ message: BitchatMessage, to conversationID: ConversationID) -> Bool
func removeMessage(withID messageID: String, cleanupFile: Bool)
/// Removes a media bubble with direction-scoped cleanup instead of the
/// broad compatibility cleanup path.
func removeUntombstonedMediaMessage(withID messageID: String)
func removeOutgoingMediaMessage(withID messageID: String)
func addSystemMessage(_ content: String)
/// Surfaces a refused explicit media deletion in the affected chat so a
/// wedged delete never looks like success.
func notifyMediaDeletionRefused(messageID: String)
/// Signals that message state changed so observers refresh (e.g. `objectWillChange.send()`).
func notifyUIChanged()
@@ -122,6 +129,15 @@ protocol ChatMediaTransferContext: AnyObject {
)
func sendFileBroadcast(_ packet: BitchatFilePacket, transferId: String)
func cancelTransfer(_ transferId: String)
/// Receiver-side stable-ID deletion commit. Implementations must invoke
/// completion only after the entire batch is durably tombstoned.
func persistDeletedPrivateMedia(
messageIDs: [String],
completion: @escaping @MainActor (Bool) -> Void
)
/// Whether any current private-chat copy of this stable ID came from a
/// remote peer and therefore requires a receiver tombstone.
func requiresPrivateMediaTombstone(messageID: String) -> Bool
}
extension ChatViewModel: ChatMediaTransferContext {
@@ -208,6 +224,185 @@ extension ChatViewModel: ChatMediaTransferContext {
func cancelTransfer(_ transferId: String) {
meshService.cancelTransfer(transferId)
}
func removeUntombstonedMediaMessage(withID messageID: String) {
let message = conversations.conversationsByID.values.lazy
.flatMap(\.messages)
.first { $0.id == messageID }
if let message, !isIncomingPrivateMessage(message) {
mediaTransferCoordinator.cleanupOutgoingLocalFile(
forMessage: message
)
}
removeMessage(withID: messageID, cleanupFile: false)
if let message {
cleanupLegacyIncomingMediaPayloads(for: [message])
}
}
/// Explicitly deleted LEGACY (non-stable-ID) incoming media has no
/// durable ID-to-file ownership, so the actual unlink is delegated to
/// the transport's gated cleanup: a basename that is pending delivery or
/// reserved by a receipt/deletion transaction stays on disk for bounded
/// quota cleanup instead. Must run after the bubbles were removed; a
/// surviving reference in any conversation keeps the payload.
func cleanupLegacyIncomingMediaPayloads(for messages: [BitchatMessage]) {
guard let cleanup =
meshService as? any PrivateMediaDeletionPersisting else {
return
}
let legacyPaths = Set(messages.compactMap { message -> String? in
guard !PrivateMediaMessageIdentity.isStableID(message.id),
isIncomingPrivateMessage(message) else {
return nil
}
return incomingMediaRelativePath(for: message)
})
guard !legacyPaths.isEmpty else { return }
let survivingPaths = Set(
conversations.conversationsByID.values.lazy
.flatMap(\.messages)
.compactMap { message -> String? in
guard self.isIncomingPrivateMessage(message) else {
return nil
}
return self.incomingMediaRelativePath(for: message)
}
)
for relativePath in legacyPaths.subtracting(survivingPaths).sorted() {
cleanup.removeLegacyPrivateMediaPayload(
relativePath: relativePath
)
}
}
func removeOutgoingMediaMessage(withID messageID: String) {
let message = conversations.conversationsByID.values.lazy
.flatMap(\.messages)
.first { $0.id == messageID }
if let message {
mediaTransferCoordinator.cleanupOutgoingLocalFile(
forMessage: message
)
}
removeMessage(withID: messageID, cleanupFile: false)
}
func persistDeletedPrivateMedia(
messageIDs: [String],
completion: @escaping @MainActor (Bool) -> Void
) {
guard !messageIDs.isEmpty else {
completion(true)
return
}
guard let persistence =
meshService as? any PrivateMediaDeletionPersisting else {
completion(false)
return
}
let requestedIDs = Set(messageIDs)
let incomingPathReferences = Array(
conversations.conversationsByID.values
.lazy
.flatMap(\.messages)
.compactMap { message -> (
messageID: String,
path: String
)? in
guard self.isIncomingPrivateMessage(message),
let path = self.incomingMediaRelativePath(
for: message
) else {
return nil
}
return (message.id, path)
}
)
let ownerIDsByPath = Dictionary(
grouping: incomingPathReferences,
by: { $0.path }
).mapValues { Set($0.map(\.messageID)) }
let protectedPayloadRelativePaths = Set(
ownerIDsByPath.compactMap { path, ownerIDs in
ownerIDs.isSubset(of: requestedIDs) ? nil : path
}
)
var payloadRelativePaths: [String: String] = [:]
for reference in incomingPathReferences
where requestedIDs.contains(reference.messageID)
&& ownerIDsByPath[reference.path, default: []]
.isSubset(of: requestedIDs) {
payloadRelativePaths[reference.messageID] = reference.path
}
persistence.persistDeletedPrivateMedia(
messageIDs: messageIDs,
payloadRelativePaths: payloadRelativePaths,
protectedPayloadRelativePaths:
protectedPayloadRelativePaths,
completion: completion
)
}
func requiresPrivateMediaTombstone(messageID: String) -> Bool {
guard PrivateMediaMessageIdentity.isStableID(messageID) else {
return false
}
return privateChats.values.lazy.flatMap { $0 }.contains { message in
message.id == messageID && isIncomingPrivateMessage(message)
}
}
func notifyMediaDeletionRefused(messageID: String) {
let owningPeerID = privateChats.first { _, messages in
messages.contains { $0.id == messageID }
}?.key
notifyPrivateMediaDeletionRefused(peerID: owningPeerID)
}
/// A refused deletion/clear previously surfaced only in SecureLogger, so
/// a wedged /clear looked like success. Tell the affected chat that its
/// bubbles and payloads were intentionally kept.
func notifyPrivateMediaDeletionRefused(peerID: PeerID?) {
let copy = String(
localized: "content.system.media_delete_refused",
comment: "System message when an explicit media delete or /clear was refused and bubbles/files were kept"
)
if let peerID = peerID ?? selectedPrivateChatPeer {
addLocalPrivateSystemMessage(copy, to: peerID)
} else {
addSystemMessage(copy)
}
}
private func isIncomingPrivateMessage(
_ message: BitchatMessage
) -> Bool {
if let senderPeerID = message.senderPeerID {
return senderPeerID.toShort() != meshService.myPeerID.toShort()
}
return message.sender != nickname
&& !message.sender.hasPrefix(nickname + "#")
}
private func incomingMediaRelativePath(
for message: BitchatMessage
) -> String? {
let categories: [MimeType.Category] = [.audio, .image, .file]
guard let category = categories.first(where: {
message.content.hasPrefix($0.messagePrefix)
}),
let rawFilename = String(
message.content.dropFirst(category.messagePrefix.count)
).trimmedOrNilIfEmpty,
let safeFilename =
(rawFilename as NSString).lastPathComponent.nilIfEmpty,
safeFilename != ".",
safeFilename != ".." else {
return nil
}
return "\(category.mediaDir)/incoming/\(safeFilename)"
}
}
/// Synchronous boundary between detached image writers and panic deletion.
@@ -297,6 +492,7 @@ final class ChatMediaTransferCoordinator {
private(set) var transferIdToMessageIDs: [String: [String]] = [:]
private(set) var messageIDToTransferId: [String: String] = [:]
private var deletionGeneration: UInt64 = 0
private var reconnectRetryRecords: [
String: PrivateMediaReconnectRetryRecord
] = [:]
@@ -871,8 +1067,7 @@ final class ChatMediaTransferCoordinator {
cancelActiveTransfer: false
)
clearTransferMapping(for: messageID)
context.removeMessage(withID: messageID, cleanupFile: true)
context.removeOutgoingMediaMessage(withID: messageID)
case .rejected(let id, let reason):
guard let messageID = currentMessageID(forTransferID: id) else {
return
@@ -891,6 +1086,40 @@ final class ChatMediaTransferCoordinator {
}
func cleanupLocalFile(forMessage message: BitchatMessage) {
cleanupLocalFile(
forMessage: message,
directions: ["outgoing", "incoming"],
searchAllCategories: true
)
}
/// `/clear` may cancel an outgoing message before receiver tombstones are
/// committed. Restrict cleanup to that message's outgoing directory so a
/// same-name incoming payload cannot be removed prematurely.
func cleanupOutgoingLocalFile(forMessage message: BitchatMessage) {
cleanupLocalFile(
forMessage: message,
directions: ["outgoing"],
searchAllCategories: false
)
}
/// Receiver cleanup runs only after any required tombstone commit. Keep it
/// scoped to the parsed media category and incoming directory so unrelated
/// outgoing or cross-category payloads with the same basename survive.
func cleanupIncomingLocalFile(forMessage message: BitchatMessage) {
cleanupLocalFile(
forMessage: message,
directions: ["incoming"],
searchAllCategories: false
)
}
private func cleanupLocalFile(
forMessage message: BitchatMessage,
directions: [String],
searchAllCategories: Bool
) {
let categories: [MimeType.Category] = [.audio, .image, .file]
guard let category = categories.first(where: { message.content.hasPrefix($0.messagePrefix) }),
let rawFilename = String(message.content.dropFirst(category.messagePrefix.count)).trimmedOrNilIfEmpty,
@@ -901,11 +1130,27 @@ final class ChatMediaTransferCoordinator {
return
}
let subdirs = categories.flatMap { ["\($0.mediaDir)/outgoing", "\($0.mediaDir)/incoming"] }
let targetCategories = searchAllCategories ? categories : [category]
let subdirs = targetCategories.flatMap { category in
directions.map { "\(category.mediaDir)/\($0)" }
}
for subdir in subdirs {
let target = base.appendingPathComponent(subdir, isDirectory: true).appendingPathComponent(safeFilename)
guard target.path.hasPrefix(base.path) else { continue }
guard FileManager.default.fileExists(atPath: target.path) else {
continue
}
guard let values = try? target.resourceValues(
forKeys: [.isRegularFileKey]
),
values.isRegularFile == true else {
SecureLogger.warning(
"Refusing to cleanup non-file media target \(safeFilename)",
category: .session
)
continue
}
do {
try FileManager.default.removeItem(at: target)
} catch CocoaError.fileNoSuchFile {
@@ -917,33 +1162,86 @@ final class ChatMediaTransferCoordinator {
}
func cancelMediaSend(messageID: String) {
discardReconnectRetry(
messageID: messageID,
cancelActiveTransfer: false
)
if let transferId = messageIDToTransferId[messageID],
let active = transferIdToMessageIDs[transferId]?.first,
active == messageID {
context.cancelTransfer(transferId)
}
clearTransferMapping(for: messageID)
context.removeMessage(withID: messageID, cleanupFile: true)
cancelAllMediaSendOwners(messageID: messageID)
context.removeOutgoingMediaMessage(withID: messageID)
}
/// Lets `/clear` cancel send ownership without implicitly deciding which
/// bubbles/files its deletion transaction may remove.
func cancelMediaTransferForConversationClear(messageID: String) {
cancelAllMediaSendOwners(messageID: messageID)
}
func deleteMediaMessage(messageID: String) {
// Stop every exact sender owner before the durable receiver commit.
// Otherwise a retained retry or admitted legacy send could transmit
// after the bubble and payload have been deleted.
cancelAllMediaSendOwners(messageID: messageID)
guard context.requiresPrivateMediaTombstone(
messageID: messageID
) else {
finishMediaDeletion(
messageID: messageID,
receiverJournalOwnsPayload: false
)
return
}
let generation = deletionGeneration
context.persistDeletedPrivateMedia(
messageIDs: [messageID]
) { [weak self] persisted in
guard let self,
self.deletionGeneration == generation else {
return
}
guard persisted else {
SecureLogger.error(
"Refusing to delete private media without a durable tombstone id=\(messageID.prefix(12))",
category: .session
)
self.context.notifyMediaDeletionRefused(
messageID: messageID
)
return
}
self.finishMediaDeletion(
messageID: messageID,
receiverJournalOwnsPayload: true
)
}
}
private func finishMediaDeletion(
messageID: String,
receiverJournalOwnsPayload: Bool
) {
if receiverJournalOwnsPayload {
// The journal already owns the exact path. A basename cleanup here
// could delete a different arrival that reused it after unlink.
context.removeMessage(withID: messageID, cleanupFile: false)
} else {
context.removeUntombstonedMediaMessage(withID: messageID)
}
}
private func cancelAllMediaSendOwners(messageID: String) {
// This releases the retained packet and expiry/retry owner. When its
// exact active transfer still owns the mapping, it cancels that owner
// before any deletion persistence or UI mutation can proceed.
discardReconnectRetry(
messageID: messageID,
cancelActiveTransfer: false
cancelActiveTransfer: true
)
// Delete is also a send cancellation. In particular, an approved
// legacy-clear send may still be waiting on BLEService.messageQueue;
// removing only the UI mapping would let that deferred work transmit.
// In particular, an approved legacy send may still be waiting on
// BLEService.messageQueue. Its admission must be canceled before the
// mapping/consent owner is released.
if let transferId = messageIDToTransferId[messageID],
transferIdToMessageIDs[transferId]?.first == messageID {
context.cancelTransfer(transferId)
}
clearTransferMapping(for: messageID)
context.removeMessage(withID: messageID, cleanupFile: true)
}
/// A raw link callback can arrive before the replacement Noise session
@@ -999,6 +1297,7 @@ final class ChatMediaTransferCoordinator {
/// is the last filesystem mutation before the transaction can complete.
func resetForPanic() {
imagePreparationBarrier.invalidateAndWait()
deletionGeneration &+= 1
peersResolvingReconnectRetry.removeAll(keepingCapacity: false)
for task in reconnectRetryExpiryTasks.values {
task.cancel()
+297 -1
View File
@@ -109,6 +109,16 @@ struct PanicNetworkLifecycle {
}
}
private struct PendingPrivateChatClear {
let peerID: PeerID
let sourceConversationID: ConversationID
let messages: [BitchatMessage]
let otherMessageIDs: Set<String>
let localPeerID: PeerID
let nickname: String
let outgoingMedia: [BitchatMessage]
}
/// Manages the application state and business logic for BitChat.
/// Acts as the primary coordinator between UI components and backend services,
/// implementing the BitchatDelegate protocol to handle network events.
@@ -376,6 +386,11 @@ final class ChatViewModel: ObservableObject, BitchatDelegate, SynchronousMessage
@Published var bluetoothAlertMessage = ""
@Published var bluetoothState: CBManagerState = .unknown
@Published private(set) var legacyPrivateMediaConsentRequest: LegacyPrivateMediaConsentRequest?
@MainActor private var queuedPrivateChatClears: [
PendingPrivateChatClear
] = []
@MainActor private var privateChatClearInFlight = false
@MainActor private var privateChatClearGeneration: UInt64 = 0
private var pendingLegacyPrivateMediaConsents: [PendingLegacyPrivateMediaConsent] = []
private func performDeliveryUpdate(_ update: @escaping @MainActor (ChatDeliveryCoordinator) -> Void) {
@@ -643,7 +658,285 @@ final class ChatViewModel: ObservableObject, BitchatDelegate, SynchronousMessage
/// Empties the peer's chat but keeps the conversation alive (`/clear`).
@MainActor
func clearPrivateChat(_ peerID: PeerID) {
conversations.clear(.directPeer(peerID))
let sourceConversationID = ConversationID.directPeer(peerID)
// An active live-voice row owns an open FileHandle and may be
// republished as frames/final media arrive. Treat it like an in-flight
// arrival rather than unlinking its capture or removing its bubble.
let messages = privateMessages(for: peerID).filter {
!liveVoiceCoordinator.isLiveVoiceMessage($0)
}
let localPeerID = meshService.myPeerID.toShort()
let currentNickname = nickname
let mediaPrefixes = [
MimeType.Category.audio.messagePrefix,
MimeType.Category.image.messagePrefix,
MimeType.Category.file.messagePrefix
]
let outgoingMedia = messages.filter { message in
guard mediaPrefixes.contains(where: {
message.content.hasPrefix($0)
}) else {
return false
}
if let senderPeerID = message.senderPeerID {
return senderPeerID.toShort() == localPeerID
}
return message.sender == currentNickname
|| message.sender.hasPrefix(currentNickname + "#")
}
// Send ownership is canceled at command time even when another clear
// transaction is ahead in the queue. UI and files remain untouched
// until this request's receiver journal commit succeeds.
for message in outgoingMedia {
mediaTransferCoordinator
.cancelMediaTransferForConversationClear(
messageID: message.id
)
}
queuedPrivateChatClears.append(PendingPrivateChatClear(
peerID: peerID,
sourceConversationID: sourceConversationID,
messages: messages,
otherMessageIDs: Set(
privateChats
.filter { $0.key != peerID }
.flatMap { $0.value.map(\.id) }
),
localPeerID: localPeerID,
nickname: currentNickname,
outgoingMedia: outgoingMedia
))
startNextPrivateChatClearIfNeeded()
}
@MainActor
private func startNextPrivateChatClearIfNeeded() {
guard !privateChatClearInFlight,
!queuedPrivateChatClears.isEmpty else {
return
}
privateChatClearInFlight = true
let request = queuedPrivateChatClears.removeFirst()
let generation = privateChatClearGeneration
performPrivateChatClear(
request,
generation: generation
) { [weak self] in
guard let self,
self.privateChatClearGeneration == generation else {
return
}
self.privateChatClearInFlight = false
self.startNextPrivateChatClearIfNeeded()
}
}
@MainActor
private func performPrivateChatClear(
_ request: PendingPrivateChatClear,
generation: UInt64,
completion: @escaping @MainActor () -> Void
) {
guard privateChatClearGeneration == generation else {
completion()
return
}
let peerID = request.peerID
let selectedConversationID = request.sourceConversationID
let messagesToClear = request.messages
guard !messagesToClear.isEmpty else {
completion()
return
}
// Capture the transaction's exact UI set before any off-main receipt
// I/O. Messages arriving while the journal is written are not part of
// this command and must remain visible.
let capturedMessageIDs = Set(messagesToClear.map(\.id))
let survivingMessageIDs = request.otherMessageIDs
let mediaPrefixes = [
MimeType.Category.audio.messagePrefix,
MimeType.Category.image.messagePrefix,
MimeType.Category.file.messagePrefix
]
let localPeerID = request.localPeerID
let isMedia: (BitchatMessage) -> Bool = { message in
mediaPrefixes.contains(where: message.content.hasPrefix)
}
let isFromMe: (BitchatMessage) -> Bool = { [nickname = request.nickname] message in
if let senderPeerID = message.senderPeerID {
return senderPeerID.toShort() == localPeerID
}
return message.sender == nickname
|| message.sender.hasPrefix(nickname + "#")
}
let outgoingMedia = request.outgoingMedia
let capturedExclusiveIDs =
capturedMessageIDs.subtracting(survivingMessageIDs)
let capturedIncomingMedia = messagesToClear.filter {
isMedia($0) && !isFromMe($0)
}
let capturedStableMediaIDs = Set(
capturedIncomingMedia.compactMap { message in
PrivateMediaMessageIdentity.isStableID(message.id)
? message.id
: nil
}
)
func currentRemovalPlan() -> [ConversationID: Set<String>] {
// Identity handoff removes the source conversation and inserts its
// rows elsewhere. The old source may then be recreated by a new
// arrival before journal I/O finishes, so always scan all direct
// conversations. Only IDs exclusive at command time may follow a
// migration; shared aliases remain outside the source.
var plan: [ConversationID: Set<String>] = [:]
for (conversationID, conversation) in
conversations.conversationsByID {
guard case .direct = conversationID else { continue }
let eligibleIDs = conversationID == selectedConversationID
? capturedMessageIDs
: capturedExclusiveIDs
let matchingIDs = Set(conversation.messages.map(\.id))
.intersection(eligibleIDs)
if !matchingIDs.isEmpty {
plan[conversationID] = matchingIDs
}
}
return plan
}
func hasRemainingCopy(
of messageID: String,
after plan: [ConversationID: Set<String>]
) -> Bool {
conversations.conversationsByID.contains { conversationID, conversation in
guard case .direct = conversationID else { return false }
return conversation.messages.contains { message in
message.id == messageID
&& plan[conversationID]?.contains(messageID) != true
}
}
}
@MainActor
func continueClear(
persisted: Bool,
durableStableIDs: Set<String>
) {
guard privateChatClearGeneration == generation else {
completion()
return
}
guard persisted else {
SecureLogger.error(
"Refusing to clear private chat without durable media tombstones peer=\(peerID.id.prefix(8))",
category: .session
)
notifyPrivateMediaDeletionRefused(peerID: peerID)
completion()
return
}
let plan = currentRemovalPlan()
let newlyLastStableIDs = Set(
capturedStableMediaIDs.filter {
!durableStableIDs.contains($0)
&& !hasRemainingCopy(of: $0, after: plan)
}
)
if !newlyLastStableIDs.isEmpty {
persistDeletedPrivateMedia(
messageIDs: Array(newlyLastStableIDs).sorted()
) { persisted in
continueClear(
persisted: persisted,
durableStableIDs:
durableStableIDs.union(newlyLastStableIDs)
)
}
return
}
// A stable receiver tombstone is global for that message ID.
// Remove any alias that arrived while journal I/O was in flight.
if !durableStableIDs.isEmpty {
let directConversationIDs = conversations
.conversationsByID.keys.filter {
if case .direct = $0 { return true }
return false
}
for conversationID in directConversationIDs {
conversations.removeMessages(from: conversationID) {
durableStableIDs.contains($0.id)
}
}
}
// Outgoing media mirrors the incoming alias protection: an ID
// whose copy survives in a conversation this clear does not
// touch (identity-alias handoff) keeps that bubble and its local
// file. Only IDs with no surviving copy are removed from every
// direct conversation and have their payload unlinked.
let outgoingPlan = currentRemovalPlan()
let removableOutgoingMedia = outgoingMedia.filter {
!hasRemainingCopy(of: $0.id, after: outgoingPlan)
}
for message in removableOutgoingMedia {
mediaTransferCoordinator.cleanupOutgoingLocalFile(
forMessage: message
)
}
let removableOutgoingIDs = Set(
removableOutgoingMedia.map(\.id)
)
if !removableOutgoingIDs.isEmpty {
let directConversationIDs = conversations
.conversationsByID.keys.filter {
if case .direct = $0 { return true }
return false
}
for conversationID in directConversationIDs {
conversations.removeMessages(from: conversationID) {
removableOutgoingIDs.contains($0.id)
}
}
}
// Stable payload cleanup belongs entirely to the durable receiver
// journal. Legacy/raw incoming payloads have no durable identity,
// so once their bubbles are gone the transport's gated cleanup
// decides per basename: unlink when unreferenced, or leave any
// pending/reserved path for bounded quota cleanup.
let finalPlan = currentRemovalPlan()
for (conversationID, messageIDs) in finalPlan {
conversations.removeMessages(from: conversationID) {
messageIDs.contains($0.id)
}
}
cleanupLegacyIncomingMediaPayloads(for: capturedIncomingMedia)
completion()
}
let initialPlan = currentRemovalPlan()
let initialStableIDs = Set(
capturedStableMediaIDs.filter {
!hasRemainingCopy(of: $0, after: initialPlan)
}
)
persistDeletedPrivateMedia(
messageIDs: Array(initialStableIDs).sorted()
) { persisted in
continueClear(
persisted: persisted,
durableStableIDs: initialStableIDs
)
}
}
/// Removes the peer's chat entirely, including unread state.
@@ -1275,6 +1568,9 @@ final class ChatViewModel: ObservableObject, BitchatDelegate, SynchronousMessage
// handles before clearing state or removing the media directory.
mediaTransferCoordinator.resetForPanic()
liveVoiceCoordinator.resetForPanic()
privateChatClearGeneration &+= 1
queuedPrivateChatClears.removeAll(keepingCapacity: false)
privateChatClearInFlight = false
// Deny and release any clear-media confirmations before identities,
// message state, and local files are wiped.