Merge pull request 'WP-1: durable SpoolStore with crash reconciliation and removed/' (#4) from wp1/spool-store-20260831 into feat/shotdeck-20260830
This commit was merged in pull request #4.
This commit is contained in:
@@ -0,0 +1,418 @@
|
||||
import Foundation
|
||||
import ImageIO
|
||||
import CoreGraphics
|
||||
import Darwin
|
||||
|
||||
public actor SpoolStore {
|
||||
private let paths: AppSupportPaths
|
||||
private var openSession: CaptureSession
|
||||
|
||||
public init(paths: AppSupportPaths) throws {
|
||||
self.paths = paths
|
||||
let fm = FileManager.default
|
||||
try Self.supersedeSpoolArchiveOverlaps(paths: paths, fileManager: fm)
|
||||
try Self.finishInterruptedArchives(paths: paths, fileManager: fm)
|
||||
let remainingIDs = try Self.listUUIDDirectories(in: paths.spool, fileManager: fm)
|
||||
if remainingIDs.isEmpty {
|
||||
self.openSession = try Self.createFreshSession(paths: paths)
|
||||
} else {
|
||||
var candidates: [CaptureSession] = []
|
||||
for id in remainingIDs {
|
||||
let dir = paths.sessionDirectory(id)
|
||||
candidates.append(
|
||||
try Self.reconcileSessionDirectory(
|
||||
at: dir, id: id, assumedStateIfRebuilt: .open, fileManager: fm))
|
||||
}
|
||||
self.openSession = Self.pickNewest(first: candidates[0], rest: Array(candidates.dropFirst()))
|
||||
}
|
||||
}
|
||||
|
||||
/// The session currently accepting captures. Cheap accessor — all reconciliation already
|
||||
/// happened once, inside init.
|
||||
public func currentSession() throws -> CaptureSession {
|
||||
openSession
|
||||
}
|
||||
|
||||
/// Writes `pngData` to disk and fsyncs it BEFORE the manifest is touched, then updates and
|
||||
/// durably writes the manifest. Order is non-negotiable: image durable -> manifest durable
|
||||
/// -> return. Uses AtomicFile.write/writeJSON for both writes — never Data.write(to:).
|
||||
public func append(
|
||||
pngData: Data, pixelWidth: Int, pixelHeight: Int,
|
||||
scale: CGFloat, capturedAt: Date
|
||||
) throws -> Capture {
|
||||
let captureID = UUID()
|
||||
let sequence = openSession.nextSequence
|
||||
let fileName = "\(String(format: "%03d", sequence))-\(Self.hexSuffix(captureID)).png"
|
||||
let sessionDir = paths.sessionDirectory(openSession.id)
|
||||
let fileURL = sessionDir.appendingPathComponent(fileName)
|
||||
try AtomicFile.write(pngData, to: fileURL)
|
||||
let capture = Capture(
|
||||
id: captureID, sequence: sequence, fileName: fileName,
|
||||
pixelWidth: pixelWidth, pixelHeight: pixelHeight,
|
||||
scale: scale, capturedAt: capturedAt)
|
||||
let updated = openSession.appending(capture)
|
||||
try AtomicFile.writeJSON(updated, to: sessionDir.appendingPathComponent("session.json"))
|
||||
openSession = updated
|
||||
return capture
|
||||
}
|
||||
|
||||
/// D-11: moves the capture's PNG into `<sessionDir>/removed/` (created lazily) and drops
|
||||
/// its manifest entry. NEVER unlinks/deletes a user PNG. Throws (no filesystem change) if
|
||||
/// `captureID` is not present in the open session.
|
||||
public func remove(captureID: UUID) throws -> CaptureSession {
|
||||
guard let capture = openSession.captures.first(where: { $0.id == captureID }) else {
|
||||
throw ShotdeckError.spoolWriteFailed(
|
||||
path: paths.sessionDirectory(openSession.id).path,
|
||||
underlying: "capture \(captureID) is not in the open session")
|
||||
}
|
||||
let sessionDir = paths.sessionDirectory(openSession.id)
|
||||
let removedDir = sessionDir.appendingPathComponent("removed", isDirectory: true)
|
||||
do {
|
||||
try FileManager.default.createDirectory(at: removedDir, withIntermediateDirectories: true)
|
||||
} catch {
|
||||
throw ShotdeckError.spoolWriteFailed(path: removedDir.path, underlying: error.localizedDescription)
|
||||
}
|
||||
let sourceURL = sessionDir.appendingPathComponent(capture.fileName)
|
||||
let destURL = removedDir.appendingPathComponent(capture.fileName)
|
||||
guard rename(sourceURL.path, destURL.path) == 0 else {
|
||||
throw ShotdeckError.spoolWriteFailed(
|
||||
path: destURL.path,
|
||||
underlying: "could not move the capture into removed/: \(String(cString: strerror(errno)))")
|
||||
}
|
||||
try AtomicFile.fsyncDirectory(at: removedDir)
|
||||
let updated = openSession.removing(captureID: captureID)
|
||||
try AtomicFile.writeJSON(updated, to: sessionDir.appendingPathComponent("session.json"))
|
||||
openSession = updated
|
||||
return updated
|
||||
}
|
||||
|
||||
/// Closes the open session (must be non-empty), moves its directory under archive/, records
|
||||
/// pdfFileName, and starts a fresh empty open session. Returns the archived one.
|
||||
public func archiveCurrent(pdfFileName: String) throws -> CaptureSession {
|
||||
guard !openSession.isEmpty else {
|
||||
throw ShotdeckError.spoolWriteFailed(
|
||||
path: paths.sessionDirectory(openSession.id).path,
|
||||
underlying: "cannot archive an empty session")
|
||||
}
|
||||
let archived = openSession.markArchived(pdfFileName: pdfFileName)
|
||||
let sessionDir = paths.sessionDirectory(openSession.id)
|
||||
try AtomicFile.writeJSON(archived, to: sessionDir.appendingPathComponent("session.json"))
|
||||
let archiveDir = paths.archiveDirectory(openSession.id)
|
||||
guard rename(sessionDir.path, archiveDir.path) == 0 else {
|
||||
throw ShotdeckError.spoolWriteFailed(
|
||||
path: sessionDir.path,
|
||||
underlying: "could not move the session into archive/: \(String(cString: strerror(errno)))")
|
||||
}
|
||||
try AtomicFile.fsyncDirectory(at: paths.archive)
|
||||
let fresh = try Self.createFreshSession(paths: paths)
|
||||
openSession = fresh
|
||||
return archived
|
||||
}
|
||||
|
||||
/// Abandons the open session only if it is empty (mints a new id/dir/manifest); throws,
|
||||
/// with no filesystem change, if the open session has captures.
|
||||
public func startNewSession() throws -> CaptureSession {
|
||||
guard openSession.isEmpty else {
|
||||
throw ShotdeckError.spoolWriteFailed(
|
||||
path: paths.sessionDirectory(openSession.id).path,
|
||||
underlying: "cannot start a new session: \(openSession.captures.count) capture(s) present in the open session")
|
||||
}
|
||||
let fresh = try Self.createFreshSession(paths: paths)
|
||||
openSession = fresh
|
||||
return fresh
|
||||
}
|
||||
|
||||
/// Absolute URL of a capture's PNG, in whichever top-level directory its session lives
|
||||
/// (spool/ if session.state == .open, archive/ if .archived).
|
||||
public func imageURL(for capture: Capture, in session: CaptureSession) -> URL {
|
||||
let dir = session.state == .open ? paths.sessionDirectory(session.id) : paths.archiveDirectory(session.id)
|
||||
return dir.appendingPathComponent(capture.fileName)
|
||||
}
|
||||
|
||||
/// Archived sessions, newest createdAt first. Lazily reconciles each archive/ directory
|
||||
/// the same way init reconciles spool/ candidates (orphan recovery, missing-drop, corrupt
|
||||
/// rebuild) — an archived session's manifest can degrade too and must self-heal without
|
||||
/// ever losing a PNG.
|
||||
public func archivedSessions() throws -> [CaptureSession] {
|
||||
let ids = try Self.listUUIDDirectories(in: paths.archive, fileManager: .default)
|
||||
var sessions: [CaptureSession] = []
|
||||
for id in ids {
|
||||
sessions.append(
|
||||
try Self.reconcileSessionDirectory(
|
||||
at: paths.archiveDirectory(id), id: id, assumedStateIfRebuilt: .archived, fileManager: .default))
|
||||
}
|
||||
return sessions.sorted { ($0.createdAt, $0.id.uuidString) > ($1.createdAt, $1.id.uuidString) }
|
||||
}
|
||||
|
||||
// MARK: - Reconciliation (static so they can run inside init)
|
||||
|
||||
private static func listUUIDDirectories(in parent: URL, fileManager: FileManager) throws -> [UUID] {
|
||||
let entries: [URL]
|
||||
do {
|
||||
entries = try fileManager.contentsOfDirectory(
|
||||
at: parent,
|
||||
includingPropertiesForKeys: [.isDirectoryKey],
|
||||
options: [])
|
||||
} catch {
|
||||
throw ShotdeckError.spoolWriteFailed(path: parent.path, underlying: error.localizedDescription)
|
||||
}
|
||||
var ids: [UUID] = []
|
||||
for url in entries {
|
||||
let isDirectory = (try? url.resourceValues(forKeys: [.isDirectoryKey]).isDirectory) ?? false
|
||||
guard isDirectory else { continue }
|
||||
if let id = UUID(uuidString: url.lastPathComponent) {
|
||||
ids.append(id)
|
||||
}
|
||||
}
|
||||
return ids
|
||||
}
|
||||
|
||||
private static func supersedeSpoolArchiveOverlaps(paths: AppSupportPaths, fileManager: FileManager) throws {
|
||||
let spoolIDs = Set(try listUUIDDirectories(in: paths.spool, fileManager: fileManager))
|
||||
let archiveIDs = Set(try listUUIDDirectories(in: paths.archive, fileManager: fileManager))
|
||||
for id in spoolIDs.intersection(archiveIDs) {
|
||||
let spoolDir = paths.sessionDirectory(id)
|
||||
let supersededDir = paths.spool.appendingPathComponent(
|
||||
"\(id.uuidString).superseded-\(DubaiTime.fileStamp(Date()))", isDirectory: true)
|
||||
guard rename(spoolDir.path, supersededDir.path) == 0 else {
|
||||
throw ShotdeckError.spoolWriteFailed(
|
||||
path: spoolDir.path,
|
||||
underlying: "could not supersede a duplicate spool copy: \(String(cString: strerror(errno)))")
|
||||
}
|
||||
try AtomicFile.fsyncDirectory(at: paths.spool)
|
||||
Log.spool.warning("Found session \(id.uuidString, privacy: .public) in both spool/ and archive/; kept the archive copy and superseded the spool copy — nothing was deleted.")
|
||||
}
|
||||
}
|
||||
|
||||
private static func finishInterruptedArchives(paths: AppSupportPaths, fileManager: FileManager) throws {
|
||||
for id in try listUUIDDirectories(in: paths.spool, fileManager: fileManager) {
|
||||
let spoolDir = paths.sessionDirectory(id)
|
||||
let manifestURL = spoolDir.appendingPathComponent("session.json")
|
||||
guard let data = try? Data(contentsOf: manifestURL),
|
||||
let decoded = try? decodeSession(from: data),
|
||||
decoded.id == id, decoded.state == .archived
|
||||
else { continue }
|
||||
let archiveDir = paths.archiveDirectory(id)
|
||||
guard rename(spoolDir.path, archiveDir.path) == 0 else {
|
||||
throw ShotdeckError.spoolWriteFailed(
|
||||
path: spoolDir.path,
|
||||
underlying: "could not complete an interrupted archive move: \(String(cString: strerror(errno)))")
|
||||
}
|
||||
try AtomicFile.fsyncDirectory(at: paths.archive)
|
||||
Log.spool.warning("Completed an archive move for \(id.uuidString, privacy: .public) that was interrupted before this launch.")
|
||||
}
|
||||
}
|
||||
|
||||
private static func reconcileSessionDirectory(
|
||||
at dir: URL,
|
||||
id: UUID,
|
||||
assumedStateIfRebuilt: SessionState,
|
||||
fileManager: FileManager
|
||||
) throws -> CaptureSession {
|
||||
let manifestURL = dir.appendingPathComponent("session.json")
|
||||
let decoded: CaptureSession?
|
||||
if let data = try? Data(contentsOf: manifestURL),
|
||||
let session = try? decodeSession(from: data),
|
||||
session.id == id {
|
||||
decoded = session
|
||||
} else {
|
||||
decoded = nil
|
||||
}
|
||||
|
||||
guard let session = decoded else {
|
||||
return try rebuildManifest(
|
||||
at: dir,
|
||||
id: id,
|
||||
assumedState: assumedStateIfRebuilt,
|
||||
fileManager: fileManager)
|
||||
}
|
||||
|
||||
var present: [Capture] = []
|
||||
var missingCount = 0
|
||||
for capture in session.captures {
|
||||
let fileURL = dir.appendingPathComponent(capture.fileName)
|
||||
if fileManager.fileExists(atPath: fileURL.path) {
|
||||
present.append(capture)
|
||||
} else {
|
||||
missingCount += 1
|
||||
}
|
||||
}
|
||||
let referenced = Set(session.captures.map(\.fileName))
|
||||
let pngs = try listCapturePNGs(in: dir, fileManager: fileManager)
|
||||
var recoveredOrphans: [Capture] = []
|
||||
for pngURL in pngs where !referenced.contains(pngURL.lastPathComponent) {
|
||||
if let recovered = recoverCapture(from: pngURL, fileManager: fileManager) {
|
||||
recoveredOrphans.append(recovered)
|
||||
} else {
|
||||
Log.spool.error("Could not decode orphan PNG at \(pngURL.path, privacy: .public); leaving it on disk.")
|
||||
}
|
||||
}
|
||||
let deletedTmp = deleteStrayTmpFiles(in: dir, fileManager: fileManager)
|
||||
if missingCount == 0 && recoveredOrphans.isEmpty && !deletedTmp {
|
||||
return session
|
||||
}
|
||||
let finalCaptures = (present + recoveredOrphans).sorted { $0.sequence < $1.sequence }
|
||||
let updated = CaptureSession(
|
||||
id: session.id,
|
||||
createdAt: session.createdAt,
|
||||
state: session.state,
|
||||
captures: finalCaptures,
|
||||
pdfFileName: session.pdfFileName)
|
||||
try AtomicFile.writeJSON(updated, to: manifestURL)
|
||||
Log.spool.warning("Reconciled session \(id.uuidString, privacy: .public): dropped \(missingCount) missing PNG(s), recovered \(recoveredOrphans.count) orphan(s).")
|
||||
return updated
|
||||
}
|
||||
|
||||
private static func rebuildManifest(
|
||||
at dir: URL,
|
||||
id: UUID,
|
||||
assumedState: SessionState,
|
||||
fileManager: FileManager
|
||||
) throws -> CaptureSession {
|
||||
let manifestURL = dir.appendingPathComponent("session.json")
|
||||
if fileManager.fileExists(atPath: manifestURL.path) {
|
||||
let corruptURL = dir.appendingPathComponent(
|
||||
"session.json.corrupt-\(DubaiTime.fileStamp(Date()))")
|
||||
do {
|
||||
try fileManager.moveItem(at: manifestURL, to: corruptURL)
|
||||
} catch {
|
||||
throw ShotdeckError.spoolWriteFailed(
|
||||
path: manifestURL.path,
|
||||
underlying: "could not quarantine a corrupt manifest: \(error.localizedDescription)")
|
||||
}
|
||||
}
|
||||
|
||||
var recovered: [Capture] = []
|
||||
for pngURL in try listCapturePNGs(in: dir, fileManager: fileManager) {
|
||||
if let capture = recoverCapture(from: pngURL, fileManager: fileManager) {
|
||||
recovered.append(capture)
|
||||
} else {
|
||||
Log.spool.error("Could not decode PNG at \(pngURL.path, privacy: .public) while rebuilding the manifest; leaving it on disk.")
|
||||
}
|
||||
}
|
||||
_ = deleteStrayTmpFiles(in: dir, fileManager: fileManager)
|
||||
recovered.sort { $0.sequence < $1.sequence }
|
||||
|
||||
let createdAt: Date
|
||||
if let earliest = recovered.map(\.capturedAt).min() {
|
||||
createdAt = earliest
|
||||
} else {
|
||||
createdAt = (try? dir.resourceValues(forKeys: [.creationDateKey]))?.creationDate ?? Date()
|
||||
}
|
||||
|
||||
let rebuilt = CaptureSession(
|
||||
id: id,
|
||||
createdAt: createdAt,
|
||||
state: assumedState,
|
||||
captures: recovered,
|
||||
pdfFileName: nil)
|
||||
try AtomicFile.writeJSON(rebuilt, to: dir.appendingPathComponent("session.json"))
|
||||
Log.spool.warning("Rebuilt manifest for \(id.uuidString, privacy: .public) from \(recovered.count) recovered PNG(s).")
|
||||
if assumedState == .archived {
|
||||
Log.spool.warning("Rebuilt an archived session \(id.uuidString, privacy: .public) from PNGs; pdfFileName could not be recovered.")
|
||||
}
|
||||
return rebuilt
|
||||
}
|
||||
|
||||
private static func recoverCapture(from fileURL: URL, fileManager: FileManager) -> Capture? {
|
||||
guard fileManager.fileExists(atPath: fileURL.path) else { return nil }
|
||||
guard let source = CGImageSourceCreateWithURL(fileURL as CFURL, nil),
|
||||
let properties = CGImageSourceCopyPropertiesAtIndex(source, 0, nil) as NSDictionary?,
|
||||
let width = (properties[kCGImagePropertyPixelWidth] as? NSNumber)?.intValue,
|
||||
let height = (properties[kCGImagePropertyPixelHeight] as? NSNumber)?.intValue
|
||||
else { return nil }
|
||||
let name = fileURL.lastPathComponent
|
||||
let sequence = Int(name.prefix(3)) ?? 1
|
||||
let capturedAt = (try? fileURL.resourceValues(forKeys: [.creationDateKey]))?.creationDate ?? Date()
|
||||
return Capture(
|
||||
id: UUID(),
|
||||
sequence: sequence,
|
||||
fileName: name,
|
||||
pixelWidth: width,
|
||||
pixelHeight: height,
|
||||
scale: 1.0,
|
||||
capturedAt: capturedAt)
|
||||
}
|
||||
|
||||
private static func pickNewest(first: CaptureSession, rest: [CaptureSession]) -> CaptureSession {
|
||||
rest.reduce(first) { current, candidate in
|
||||
let currentKey = (current.createdAt, current.id.uuidString)
|
||||
let candidateKey = (candidate.createdAt, candidate.id.uuidString)
|
||||
return candidateKey > currentKey ? candidate : current
|
||||
}
|
||||
}
|
||||
|
||||
private static func createFreshSession(paths: AppSupportPaths) throws -> CaptureSession {
|
||||
let id = UUID()
|
||||
let createdAt = Date()
|
||||
let dir = paths.sessionDirectory(id)
|
||||
do {
|
||||
try FileManager.default.createDirectory(at: dir, withIntermediateDirectories: true)
|
||||
} catch {
|
||||
throw ShotdeckError.spoolWriteFailed(path: dir.path, underlying: error.localizedDescription)
|
||||
}
|
||||
let session = CaptureSession(id: id, createdAt: createdAt, state: .open, captures: [], pdfFileName: nil)
|
||||
try AtomicFile.writeJSON(session, to: dir.appendingPathComponent("session.json"))
|
||||
return session
|
||||
}
|
||||
|
||||
private static func hexSuffix(_ id: UUID) -> String {
|
||||
String(id.uuidString.replacingOccurrences(of: "-", with: "").prefix(8)).uppercased()
|
||||
}
|
||||
|
||||
private static func decodeSession(from data: Data) throws -> CaptureSession {
|
||||
let decoder = JSONDecoder()
|
||||
decoder.dateDecodingStrategy = .iso8601
|
||||
return try decoder.decode(CaptureSession.self, from: data)
|
||||
}
|
||||
|
||||
private static func isCapturePNGName(_ name: String) -> Bool {
|
||||
guard name.hasSuffix(".png") else { return false }
|
||||
let stem = String(name.dropLast(4))
|
||||
let parts = stem.split(separator: "-", maxSplits: 1, omittingEmptySubsequences: false)
|
||||
guard parts.count == 2,
|
||||
parts[0].count == 3,
|
||||
parts[0].allSatisfy(\.isNumber),
|
||||
parts[1].count == 8,
|
||||
parts[1].allSatisfy(\.isHexDigit)
|
||||
else { return false }
|
||||
return true
|
||||
}
|
||||
|
||||
private static func listCapturePNGs(in dir: URL, fileManager: FileManager) throws -> [URL] {
|
||||
let entries: [URL]
|
||||
do {
|
||||
entries = try fileManager.contentsOfDirectory(
|
||||
at: dir,
|
||||
includingPropertiesForKeys: [.isDirectoryKey],
|
||||
options: [])
|
||||
} catch {
|
||||
throw ShotdeckError.spoolWriteFailed(path: dir.path, underlying: error.localizedDescription)
|
||||
}
|
||||
return entries.filter { url in
|
||||
let isDirectory = (try? url.resourceValues(forKeys: [.isDirectoryKey]).isDirectory) ?? false
|
||||
guard !isDirectory else { return false }
|
||||
return isCapturePNGName(url.lastPathComponent)
|
||||
}
|
||||
}
|
||||
|
||||
private static func deleteStrayTmpFiles(in dir: URL, fileManager: FileManager) -> Bool {
|
||||
let entries = (try? fileManager.contentsOfDirectory(
|
||||
at: dir,
|
||||
includingPropertiesForKeys: [.isDirectoryKey],
|
||||
options: [])) ?? []
|
||||
var deleted = false
|
||||
for url in entries {
|
||||
guard url.lastPathComponent.hasSuffix(".tmp") else { continue }
|
||||
let isDirectory = (try? url.resourceValues(forKeys: [.isDirectoryKey]).isDirectory) ?? false
|
||||
guard !isDirectory else { continue }
|
||||
do {
|
||||
try fileManager.removeItem(at: url)
|
||||
deleted = true
|
||||
} catch {
|
||||
Log.spool.error("Could not remove stray temp file at \(url.path, privacy: .public): \(error.localizedDescription, privacy: .public)")
|
||||
}
|
||||
}
|
||||
return deleted
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user