diff --git a/Cargo.lock b/Cargo.lock index a417117..f073f82 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1873,6 +1873,10 @@ dependencies = [ "inkfinite-crdt", "inkfinite-engine", "inkfinite-model", + "inkfinite-protocol", + "serde", + "serde_json", + "thiserror 2.0.18", ] [[package]] diff --git a/ROADMAP.md b/ROADMAP.md index f08821e..33ad4d5 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -166,6 +166,14 @@ hit test for the frozen 10,000-shape board. The V1 medians are 0.61 ms and 0.22 ms. V2-11 may add a spatial index only if its linear query path misses the 1 ms budget. +V2-07 completed on July 17, 2026. inkfinite-file imports the frozen v1 desktop +and web envelopes into normalized pages, default layers, scene containers, +bindings, styles, and deterministic draw order. It writes compact canonical +Automerge files through same-directory temporary files, flushes before +replacement, holds advisory locks, and retains bounded recovery snapshots plus +encoded change journals across failed writes. JSON export is deterministic and +history-free, as documented in [docs/v2-file-format.md](docs/v2-file-format.md). + One Inkfinite transaction maps to one Automerge change. Causal heads, rather than a scalar revision, are the concurrency token. A local sequence number may be displayed, but callers use inspected heads and operation preconditions. diff --git a/TODO.md b/TODO.md index 6bb153e..f98bc4b 100644 --- a/TODO.md +++ b/TODO.md @@ -198,14 +198,14 @@ Blocked by: V2-04, V2-05 Acceptance criteria: -- [ ] Every V2-01 valid fixture imports with the same page, shape, binding, +- [x] Every V2-01 valid fixture imports with the same page, shape, binding, group, style, and draw order; each page gains one default layer. -- [ ] Invalid or newer formats produce typed errors and never overwrite input. -- [ ] Saves use a same-directory temporary file, flush, atomic replacement where +- [x] Invalid or newer formats produce typed errors and never overwrite input. +- [x] Saves use a same-directory temporary file, flush, atomic replacement where supported, recovery copies, and advisory locking. -- [ ] Recovery stores a compact snapshot plus bounded change journal and can +- [x] Recovery stores a compact snapshot plus bounded change journal and can restore after failures at each write step. -- [ ] JSON export is deterministic and documented as a snapshot that cannot +- [x] JSON export is deterministic and documented as a snapshot that cannot preserve CRDT history. Verification: @@ -603,5 +603,5 @@ pnpm --filter inkfinite-web lint ## Frontier -V2-07 is the current frontier. Work it in a fresh implementation context before -starting tickets that depend on safe v2 import and persistence. +V2-08 is the current frontier. Work it in a fresh implementation context before +starting tickets that depend on Rust-owned document sessions. diff --git a/crates/inkfinite-file/Cargo.toml b/crates/inkfinite-file/Cargo.toml index a8babdf..cce9085 100644 --- a/crates/inkfinite-file/Cargo.toml +++ b/crates/inkfinite-file/Cargo.toml @@ -10,7 +10,12 @@ version.workspace = true inkfinite-crdt.workspace = true inkfinite-engine.workspace = true inkfinite-model.workspace = true +serde.workspace = true +serde_json.workspace = true +thiserror.workspace = true + +[dev-dependencies] +inkfinite-protocol.workspace = true [lints] workspace = true - diff --git a/crates/inkfinite-file/src/lib.rs b/crates/inkfinite-file/src/lib.rs index 84baf26..3dc5312 100644 --- a/crates/inkfinite-file/src/lib.rs +++ b/crates/inkfinite-file/src/lib.rs @@ -1,5 +1,77 @@ -//! Import, migration, and durable file boundary. +#![forbid(unsafe_code)] + +//! Import, migration, and durable file boundary for Inkfinite documents. +//! +//! The file boundary accepts the frozen v1 JSON envelope, creates a normalized +//! v2 document, and persists the Rust-owned CRDT as compact Automerge bytes. +//! Snapshot JSON is an inspection/export format; it does not contain the CRDT +//! history needed to reproduce a canonical `.inkfinite` file. + +use std::path::PathBuf; + +use inkfinite_engine::EngineError; +use thiserror::Error; + +/// Recoverable failure at the import, persistence, or recovery boundary. +#[derive(Debug, Error)] +pub enum FileError { + /// The input was valid JSON but not a valid frozen v1 envelope. + #[error("invalid v1 document: {0}")] + InvalidV1(String), + /// JSON could not be parsed or serialized. + #[error("JSON error: {0}")] + Json(#[from] serde_json::Error), + /// A recognized format version is newer than this implementation supports. + #[error("unsupported document format {format:?} version {version}")] + UnsupportedFormat { format: String, version: u32 }, + /// A v1 shape kind has no v2 registry entry. + #[error("unsupported v1 shape kind {kind:?} for shape {shape_id}")] + UnsupportedShapeKind { kind: String, shape_id: String }, + /// A source and destination path were the same file. + #[error("refusing to import a file over itself: {path}")] + SamePath { path: PathBuf }, + /// Another cooperating writer owns the document lock. + #[error("document is locked by another writer: {path}")] + Locked { path: PathBuf }, + /// The requested recovery record does not exist. + #[error("no recovery record exists for {path}")] + RecoveryNotFound { path: PathBuf }, + /// A recovery record is malformed or does not match its document. + #[error("invalid recovery record: {0}")] + InvalidRecovery(String), + /// A recovery record contains newer state than the currently opened file. + #[error("recovery state is ahead of the opened document: {path}")] + RecoveryAhead { path: PathBuf }, + /// The transaction engine rejected a document or CRDT operation. + #[error(transparent)] + Engine(#[from] EngineError), + /// A filesystem operation failed. + #[error("{operation} {path}: {source}")] + Io { + /// Operation attempted when the error occurred. + operation: &'static str, + /// Path involved in the operation. + path: PathBuf, + /// Underlying operating-system error. + #[source] + source: std::io::Error, + }, + /// A canonical destination already exists when creating a new document. + #[error("document already exists: {path}")] + AlreadyExists { path: PathBuf }, +} + +mod migration; +mod persistence; pub use inkfinite_crdt::CrdtDocument; -pub use inkfinite_engine::{CommitResult, TransactionDraft}; +pub use inkfinite_engine::{CommitResult, TransactionDraft, TransactionEngine}; pub use inkfinite_model::DocumentSnapshot; +pub use migration::{ImportedV1, import_v1_json, parse_v1_json}; +pub use persistence::{ + DocumentFile, PersistenceOptions, SaveResult, export_snapshot_json, import_v1_file, + import_v1_file_with_options, read_v1_file, recovery_path_for, write_snapshot_json, +}; + +#[cfg(test)] +mod tests; diff --git a/crates/inkfinite-file/src/migration.rs b/crates/inkfinite-file/src/migration.rs new file mode 100644 index 0000000..80d5256 --- /dev/null +++ b/crates/inkfinite-file/src/migration.rs @@ -0,0 +1,733 @@ +//! Migration from the frozen v1 desktop/web JSON envelope to the v2 model. + +use std::collections::{BTreeMap, BTreeSet}; + +use inkfinite_engine::{TransactionEngine, validate_document}; +use inkfinite_model::{ + ActorId, BindingAnchor, BindingId, BindingKind, BindingRecord, Document, DocumentId, LayerId, + LayerRecord, Opacity, Origin, PageId, PageRecord, Provenance, RecordVersion, SemanticMetadata, + ShapeId, ShapeKind, ShapeParent, ShapeProperties, ShapeRecord, ShapeStyle, Timestamp, + Transform, Vec2, builtin_shape_kinds, validate_shape_properties, +}; +use serde::Deserialize; +use serde_json::{Map, Value}; + +use crate::FileError; + +/// A normalized v2 document produced by importing a v1 JSON file. +#[derive(Clone, Debug)] +pub struct ImportedV1 { + /// The v1 board ID used as the v2 document ID. + pub document_id: DocumentId, + /// Board name retained for callers that display migration information. + pub board_name: String, + /// Original v1 creation timestamp. + pub created_at: Timestamp, + /// Original v1 update timestamp. + pub updated_at: Timestamp, + /// Normalized v2 records. + pub document: Document, +} + +impl ImportedV1 { + /// Consumes the import result and returns its normalized document. + #[must_use] + pub fn into_document(self) -> Document { + self.document + } + + /// Creates a Rust-owned transaction engine from the imported document. + /// + /// # Errors + /// + /// Returns an error when the actor ID is empty or the normalized document + /// cannot be encoded by the CRDT adapter. + pub fn into_engine(self, actor_id: ActorId) -> Result { + if actor_id.as_str().trim().is_empty() { + return Err(invalid_v1("import actor ID must not be empty")); + } + Ok(TransactionEngine::create( + self.document_id, + actor_id, + self.document, + )?) + } +} + +/// Parses and migrates a frozen v1 JSON envelope. +/// +/// The importer does not repair malformed input. It validates ownership and +/// ordering, creates one stable default layer per page, and returns a typed +/// error before any destination file can be touched. +/// +/// # Errors +/// +/// Returns [`FileError::Json`] for malformed JSON, [`FileError::UnsupportedFormat`] +/// for a recognized newer envelope, or [`FileError::InvalidV1`] for invalid v1 +/// records and references. +#[allow(clippy::needless_pass_by_value)] +pub fn import_v1_json(input: &str, actor_id: ActorId) -> Result { + if actor_id.as_str().trim().is_empty() { + return Err(invalid_v1("import actor ID must not be empty")); + } + let value: Value = serde_json::from_str(input)?; + reject_newer_format(&value)?; + let envelope: LegacyEnvelope = serde_json::from_value(value) + .map_err(|error| invalid_v1(format!("invalid envelope: {error}")))?; + migrate(envelope, &actor_id) +} + +/// Alias for [`import_v1_json`] for callers that prefer parser terminology. +/// +/// # Errors +/// +/// Returns the same typed import, format, and JSON errors as +/// [`import_v1_json`]. +pub fn parse_v1_json(input: &str, actor_id: ActorId) -> Result { + import_v1_json(input, actor_id) +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct LegacyEnvelope { + board: LegacyBoard, + doc: LegacyDocument, + order: LegacyOrder, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct LegacyBoard { + id: String, + name: String, + created_at: i64, + updated_at: i64, +} + +#[derive(Debug, Deserialize)] +struct LegacyDocument { + pages: BTreeMap, + shapes: BTreeMap, + bindings: BTreeMap, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct LegacyPage { + id: String, + name: String, + shape_ids: Vec, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct LegacyShape { + id: String, + #[serde(rename = "type")] + shape_type: String, + page_id: String, + x: f64, + y: f64, + rot: f64, + #[serde(default)] + group_id: Option, + props: Map, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct LegacyBinding { + id: String, + #[serde(rename = "type")] + binding_type: String, + from_shape_id: String, + to_shape_id: String, + handle: String, + anchor: LegacyAnchor, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "lowercase", tag = "kind")] +enum LegacyAnchor { + Center, + Edge { nx: f64, ny: f64 }, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct LegacyOrder { + page_ids: Vec, + #[serde(default)] + shape_order: Option>>, +} + +#[derive(Clone)] +struct GroupInfo { + page_id: PageId, + shape_ids: Vec, + positions: Vec, +} + +#[allow(clippy::too_many_lines)] +fn migrate(envelope: LegacyEnvelope, actor_id: &ActorId) -> Result { + let LegacyEnvelope { board, doc, order } = envelope; + if board.id.trim().is_empty() { + return Err(invalid_v1("board.id must not be empty")); + } + if doc.pages.is_empty() { + return Err(invalid_v1("doc.pages must contain at least one page")); + } + + let document_id = DocumentId::from(board.id.clone()); + let page_ids = validate_page_order(&doc.pages, &order)?; + let ordered_shapes = validate_shape_order(&doc, &order, &page_ids)?; + validate_shape_keys(&doc.shapes)?; + validate_binding_keys(&doc.bindings)?; + + let timestamp = Timestamp(board.updated_at); + let mut groups = BTreeMap::::new(); + let mut seen_shapes = BTreeSet::new(); + for page_id in &page_ids { + let shape_ids = ordered_shapes + .get(page_id) + .ok_or_else(|| invalid_v1(format!("missing order for page {page_id}")))?; + for (position, shape_id) in shape_ids.iter().enumerate() { + if !seen_shapes.insert(shape_id.clone()) { + return Err(invalid_v1(format!( + "shape {shape_id} appears more than once in persisted order" + ))); + } + let legacy_shape = doc.shapes.get(shape_id).ok_or_else(|| { + invalid_v1(format!("page {page_id} refers to missing shape {shape_id}")) + })?; + if legacy_shape.page_id != page_id.as_str() { + return Err(invalid_v1(format!( + "shape {shape_id} belongs to page {}, not {page_id}", + legacy_shape.page_id + ))); + } + if let Some(group_id) = legacy_shape.group_id.as_deref() { + if group_id.trim().is_empty() { + return Err(invalid_v1(format!("shape {shape_id} has an empty groupId"))); + } + let group = groups + .entry(group_id.to_owned()) + .or_insert_with(|| GroupInfo { + page_id: page_id.clone(), + shape_ids: Vec::new(), + positions: Vec::new(), + }); + if group.page_id != *page_id { + return Err(invalid_v1(format!("group {group_id} spans multiple pages"))); + } + group.shape_ids.push(ShapeId::from(shape_id.as_str())); + group.positions.push(position); + } + } + } + if seen_shapes.len() != doc.shapes.len() { + let missing = doc + .shapes + .keys() + .find(|shape_id| !seen_shapes.contains(*shape_id)) + .cloned() + .unwrap_or_else(|| "".into()); + return Err(invalid_v1(format!( + "shape {missing} is not present in any persisted page order" + ))); + } + for group_id in groups.keys() { + if doc.shapes.contains_key(group_id) { + return Err(invalid_v1(format!( + "group ID {group_id} collides with a shape ID" + ))); + } + } + + let container_groups: BTreeSet = groups + .iter() + .filter(|(_, group)| group.shape_ids.len() > 1 && is_contiguous(&group.positions)) + .map(|(group_id, _)| group_id.clone()) + .collect(); + + let mut pages = BTreeMap::new(); + let mut layers = BTreeMap::new(); + for page_id in &page_ids { + let legacy_page = doc + .pages + .get(page_id.as_str()) + .ok_or_else(|| invalid_v1(format!("page order refers to missing page {page_id}")))?; + let layer_id = default_layer_id(page_id); + let shape_ids = ordered_shapes + .get(page_id) + .ok_or_else(|| invalid_v1(format!("missing order for page {page_id}")))?; + let mut layer_shape_ids = Vec::new(); + let mut added_groups = BTreeSet::new(); + for shape_id in shape_ids { + let group_id = doc + .shapes + .get(shape_id) + .and_then(|shape| shape.group_id.as_deref()); + if let Some(group_id) = group_id + && container_groups.contains(group_id) + { + if added_groups.insert(group_id.to_owned()) { + layer_shape_ids.push(ShapeId::from(group_id)); + } + } else { + layer_shape_ids.push(ShapeId::from(shape_id.as_str())); + } + } + pages.insert( + page_id.clone(), + PageRecord { + id: page_id.clone(), + name: checked_name(&legacy_page.name, "page", page_id.as_str())?, + layer_ids: vec![layer_id.clone()], + version: RecordVersion(1), + }, + ); + layers.insert( + layer_id.clone(), + LayerRecord { + id: layer_id, + page_id: page_id.clone(), + name: "Default".into(), + shape_ids: layer_shape_ids, + visible: true, + locked: false, + opacity: Opacity::OPAQUE, + version: RecordVersion(1), + }, + ); + } + + let mut shapes = BTreeMap::new(); + for page_id in &page_ids { + let layer_id = default_layer_id(page_id); + for shape_id in ordered_shapes + .get(page_id) + .ok_or_else(|| invalid_v1(format!("missing order for page {page_id}")))? + { + let legacy_shape = doc.shapes.get(shape_id).ok_or_else(|| { + invalid_v1(format!("page {page_id} refers to missing shape {shape_id}")) + })?; + let kind = checked_shape_kind(legacy_shape)?; + let group_parent = legacy_shape + .group_id + .as_deref() + .filter(|group_id| container_groups.contains(*group_id)) + .map_or_else( + || ShapeParent::Layer(layer_id.clone()), + |group_id| ShapeParent::Shape(ShapeId::from(group_id)), + ); + let mut properties = migrate_properties(&kind, &legacy_shape.props)?; + if let Some(group_id) = legacy_shape + .group_id + .as_deref() + .filter(|group_id| !container_groups.contains(*group_id)) + { + properties.insert("legacy_group_id".into(), Value::String(group_id.into())); + } + let style = migrate_style(&kind, &legacy_shape.props, shape_id)?; + let shape = ShapeRecord { + id: ShapeId::from(shape_id.as_str()), + kind: ShapeKind::from(kind), + parent: group_parent, + transform: Transform { + translation: Vec2 { + x: legacy_shape.x, + y: legacy_shape.y, + }, + rotation: legacy_shape.rot, + scale_x: 1.0, + scale_y: 1.0, + }, + child_ids: Vec::new(), + layout: None, + properties, + metadata: imported_metadata(actor_id, timestamp, None), + style, + version: RecordVersion(1), + }; + shapes.insert(ShapeId::from(shape_id.as_str()), shape); + } + } + + for (group_id, group) in &groups { + if !container_groups.contains(group_id) { + continue; + } + let layer_id = default_layer_id(&group.page_id); + let container_id = ShapeId::from(group_id.as_str()); + shapes.insert( + container_id.clone(), + ShapeRecord { + id: container_id, + kind: ShapeKind::from("container"), + parent: ShapeParent::Layer(layer_id), + transform: identity_transform(), + child_ids: group.shape_ids.clone(), + layout: Some(inkfinite_model::ContainerLayout::Free), + properties: ShapeProperties::new(), + metadata: imported_metadata(actor_id, timestamp, Some(group_id.clone())), + style: ShapeStyle { + opacity: Opacity::OPAQUE, + fill_opacity: None, + stroke_opacity: None, + }, + version: RecordVersion(1), + }, + ); + } + + let mut bindings = BTreeMap::new(); + for (binding_key, legacy_binding) in &doc.bindings { + if legacy_binding.id != *binding_key { + return Err(invalid_v1(format!( + "binding map key {binding_key} does not match id {}", + legacy_binding.id + ))); + } + if legacy_binding.id.trim().is_empty() + || legacy_binding.binding_type.trim().is_empty() + || legacy_binding.from_shape_id.trim().is_empty() + || legacy_binding.to_shape_id.trim().is_empty() + || legacy_binding.handle.trim().is_empty() + { + return Err(invalid_v1(format!( + "binding {binding_key} has an empty field" + ))); + } + let anchor = match &legacy_binding.anchor { + LegacyAnchor::Center => BindingAnchor::Center, + LegacyAnchor::Edge { nx, ny } => { + if !nx.is_finite() + || !ny.is_finite() + || !(-1.0..=1.0).contains(nx) + || !(-1.0..=1.0).contains(ny) + { + return Err(invalid_v1(format!( + "binding {binding_key} has an invalid edge anchor" + ))); + } + BindingAnchor::Edge { x: *nx, y: *ny } + } + }; + let source_shape_id = ShapeId::from(legacy_binding.from_shape_id.as_str()); + let target_shape_id = ShapeId::from(legacy_binding.to_shape_id.as_str()); + if !shapes.contains_key(&source_shape_id) { + return Err(invalid_v1(format!( + "binding {binding_key} refers to missing source shape {}", + legacy_binding.from_shape_id + ))); + } + if shapes[&source_shape_id].kind.as_str() != inkfinite_model::ARROW_KIND { + return Err(invalid_v1(format!( + "binding {binding_key} source shape {} is not an arrow", + legacy_binding.from_shape_id + ))); + } + if !shapes.contains_key(&target_shape_id) { + return Err(invalid_v1(format!( + "binding {binding_key} refers to missing target shape {}", + legacy_binding.to_shape_id + ))); + } + bindings.insert( + BindingId::from(binding_key.as_str()), + BindingRecord { + id: BindingId::from(binding_key.as_str()), + kind: BindingKind::from(legacy_binding.binding_type.clone()), + source_shape_id, + target_shape_id, + source_handle: legacy_binding.handle.clone(), + anchor, + version: RecordVersion(1), + }, + ); + } + + let document = Document { + pages, + page_ids, + layers, + shapes, + bindings, + assets: BTreeMap::new(), + }; + validate_document(&document)?; + Ok(ImportedV1 { + document_id, + board_name: board.name, + created_at: Timestamp(board.created_at), + updated_at: timestamp, + document, + }) +} + +fn validate_page_order( + pages: &BTreeMap, + order: &LegacyOrder, +) -> Result, FileError> { + let page_ids: Vec = order + .page_ids + .iter() + .map(|id| PageId::from(id.clone())) + .collect(); + ensure_unique(&page_ids, "page")?; + if page_ids.len() != pages.len() + || page_ids + .iter() + .any(|page_id| !pages.contains_key(page_id.as_str())) + { + return Err(invalid_v1( + "order.pageIds must contain every page exactly once", + )); + } + for (key, page) in pages { + if key != &page.id { + return Err(invalid_v1(format!( + "page map key {key} does not match id {}", + page.id + ))); + } + if page.id.trim().is_empty() { + return Err(invalid_v1("page IDs must not be empty")); + } + } + if let Some(shape_order) = &order.shape_order { + for page_id in shape_order.keys() { + if !pages.contains_key(page_id) { + return Err(invalid_v1(format!( + "shapeOrder contains unknown page {page_id}" + ))); + } + } + } + Ok(page_ids) +} + +fn validate_shape_order( + document: &LegacyDocument, + order: &LegacyOrder, + page_ids: &[PageId], +) -> Result>, FileError> { + let mut result = BTreeMap::new(); + for page_id in page_ids { + let page = document + .pages + .get(page_id.as_str()) + .ok_or_else(|| invalid_v1(format!("page order refers to missing page {page_id}")))?; + let page_shape_ids = page.shape_ids.clone(); + ensure_unique_strings(&page_shape_ids, "shape")?; + let shape_ids = order + .shape_order + .as_ref() + .and_then(|shape_order| shape_order.get(page_id.as_str())) + .cloned() + .unwrap_or(page_shape_ids.clone()); + ensure_unique_strings(&shape_ids, "shape")?; + let expected: BTreeSet<_> = page_shape_ids.iter().collect(); + let actual: BTreeSet<_> = shape_ids.iter().collect(); + if expected != actual { + return Err(invalid_v1(format!( + "shape order for page {page_id} does not match page.shapeIds" + ))); + } + result.insert(page_id.clone(), shape_ids); + } + Ok(result) +} + +fn validate_shape_keys(shapes: &BTreeMap) -> Result<(), FileError> { + for (key, shape) in shapes { + if key != &shape.id { + return Err(invalid_v1(format!( + "shape map key {key} does not match id {}", + shape.id + ))); + } + if shape.id.trim().is_empty() || shape.page_id.trim().is_empty() { + return Err(invalid_v1(format!("shape {key} has an empty ID or pageId"))); + } + if !shape.x.is_finite() || !shape.y.is_finite() || !shape.rot.is_finite() { + return Err(invalid_v1(format!( + "shape {key} has a non-finite transform" + ))); + } + } + Ok(()) +} + +fn validate_binding_keys(bindings: &BTreeMap) -> Result<(), FileError> { + for (key, binding) in bindings { + if key != &binding.id { + return Err(invalid_v1(format!( + "binding map key {key} does not match id {}", + binding.id + ))); + } + } + Ok(()) +} + +fn checked_shape_kind(shape: &LegacyShape) -> Result { + if !builtin_shape_kinds().contains(&shape.shape_type.as_str()) { + return Err(FileError::UnsupportedShapeKind { + kind: shape.shape_type.clone(), + shape_id: shape.id.clone(), + }); + } + Ok(shape.shape_type.clone()) +} + +fn migrate_properties( + kind: &str, + properties: &Map, +) -> Result { + let mut result: ShapeProperties = properties.clone().into_iter().collect(); + for (legacy_name, v2_name) in [("w", "width"), ("h", "height")] { + if let Some(value) = result.remove(legacy_name) { + if result.contains_key(v2_name) { + return Err(invalid_v1(format!( + "shape properties contain both {legacy_name} and {v2_name}" + ))); + } + result.insert(v2_name.into(), value); + } + } + validate_shape_properties(kind, &result) + .map_err(|error| invalid_v1(format!("shape kind {kind} properties: {error}")))?; + Ok(result) +} + +fn migrate_style( + kind: &str, + properties: &Map, + shape_id: &str, +) -> Result { + let mut style = ShapeStyle { + opacity: Opacity::OPAQUE, + fill_opacity: None, + stroke_opacity: None, + }; + if kind == "stroke" + && let Some(opacity) = properties + .get("style") + .and_then(Value::as_object) + .and_then(|style| style.get("opacity")) + { + let value = opacity.as_f64().ok_or_else(|| { + invalid_v1(format!("stroke {shape_id} style.opacity must be a number")) + })?; + let value = value + .to_string() + .parse::() + .map_err(|_| invalid_v1(format!("stroke {shape_id} style.opacity is out of range")))?; + style.stroke_opacity = + Some(Opacity::new(value).map_err(|error| { + invalid_v1(format!("stroke {shape_id} style.opacity: {error}")) + })?); + } + Ok(style) +} + +fn imported_metadata( + actor_id: &ActorId, + timestamp: Timestamp, + name: Option, +) -> SemanticMetadata { + SemanticMetadata { + name, + role: None, + description: None, + tags: Vec::new(), + locked: false, + agent_editable: true, + provenance: Provenance { + actor_id: actor_id.clone(), + origin: Origin::Import, + timestamp, + source: Some("v1-import".into()), + }, + } +} + +fn identity_transform() -> Transform { + Transform { + translation: Vec2 { x: 0.0, y: 0.0 }, + rotation: 0.0, + scale_x: 1.0, + scale_y: 1.0, + } +} + +fn default_layer_id(page_id: &PageId) -> LayerId { + LayerId::new(format!("layer:{}:default", page_id.as_str())) +} + +fn is_contiguous(positions: &[usize]) -> bool { + positions + .first() + .zip(positions.last()) + .is_some_and(|(first, last)| last - first + 1 == positions.len()) +} + +fn checked_name(name: &str, kind: &str, id: &str) -> Result { + if name.trim().is_empty() { + return Err(invalid_v1(format!("{kind} {id} has an empty name"))); + } + Ok(name.to_owned()) +} + +fn ensure_unique(values: &[T], kind: &str) -> Result<(), FileError> +where + T: Ord + std::fmt::Display, +{ + let mut seen = BTreeSet::new(); + for value in values { + if !seen.insert(value) { + return Err(invalid_v1(format!("{kind} {value} appears more than once"))); + } + } + Ok(()) +} + +fn ensure_unique_strings(values: &[String], kind: &str) -> Result<(), FileError> { + let mut seen = BTreeSet::new(); + for value in values { + if !seen.insert(value) { + return Err(invalid_v1(format!("{kind} {value} appears more than once"))); + } + } + Ok(()) +} + +fn reject_newer_format(value: &Value) -> Result<(), FileError> { + let Some(object) = value.as_object() else { + return Err(invalid_v1("v1 envelope must be a JSON object")); + }; + let Some(format) = object.get("format") else { + return Ok(()); + }; + let format = format + .as_str() + .ok_or_else(|| invalid_v1("format must be a string"))?; + let version_value = object + .get("format_version") + .or_else(|| object.get("formatVersion")) + .or_else(|| object.get("version")); + let version = version_value + .and_then(Value::as_u64) + .and_then(|value| u32::try_from(value).ok()) + .unwrap_or(0); + Err(FileError::UnsupportedFormat { + format: format.to_owned(), + version, + }) +} + +fn invalid_v1(message: impl Into) -> FileError { + FileError::InvalidV1(message.into()) +} diff --git a/crates/inkfinite-file/src/persistence.rs b/crates/inkfinite-file/src/persistence.rs new file mode 100644 index 0000000..d9eaa59 --- /dev/null +++ b/crates/inkfinite-file/src/persistence.rs @@ -0,0 +1,897 @@ +//! Canonical file persistence, advisory locks, atomic replacement, and recovery. + +use std::fmt::Write as _; +use std::fs::{self, File, OpenOptions}; +use std::io::Write; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicU64, Ordering}; + +use inkfinite_crdt::EncodedChange; +use inkfinite_engine::{TransactionDraft, TransactionEngine}; +use inkfinite_model::{ActorId, ChangeHash, Document, DocumentId, DocumentSnapshot}; +use serde::{Deserialize, Serialize}; + +use crate::{FileError, ImportedV1, import_v1_json}; + +const RECOVERY_FORMAT: &str = "inkfinite.recovery"; +const RECOVERY_VERSION: u32 = 1; +const DEFAULT_MAX_JOURNAL_ENTRIES: usize = 32; +const DEFAULT_MAX_JOURNAL_BYTES: usize = 8 * 1024 * 1024; +static TEMP_COUNTER: AtomicU64 = AtomicU64::new(0); + +/// Options controlling where recovery records live and how large their +/// incremental journal may become. +#[derive(Clone, Debug)] +pub struct PersistenceOptions { + /// Optional app-data directory for recovery records. When absent, a + /// .inkfinite-recovery directory is created beside the canonical file. + pub recovery_directory: Option, + /// Maximum number of encoded changes retained in one recovery journal. + pub max_journal_entries: usize, + /// Maximum total encoded byte length retained in one recovery journal. + pub max_journal_bytes: usize, +} + +impl Default for PersistenceOptions { + fn default() -> Self { + Self { + recovery_directory: None, + max_journal_entries: DEFAULT_MAX_JOURNAL_ENTRIES, + max_journal_bytes: DEFAULT_MAX_JOURNAL_BYTES, + } + } +} + +impl PersistenceOptions { + /// Returns options using the supplied directory for recovery records. + #[must_use] + pub fn with_recovery_directory(directory: impl Into) -> Self { + Self { + recovery_directory: Some(directory.into()), + ..Self::default() + } + } + + fn recovery_directory_for(&self, document_path: &Path) -> PathBuf { + self.recovery_directory.clone().unwrap_or_else(|| { + document_path + .parent() + .unwrap_or_else(|| Path::new(".")) + .join(".inkfinite-recovery") + }) + } + + fn max_journal_entries(&self) -> usize { + self.max_journal_entries.max(1) + } + + fn max_journal_bytes(&self) -> usize { + self.max_journal_bytes.max(1) + } +} + +/// Result of a successful canonical save. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct SaveResult { + /// Canonical file that was replaced. + pub path: PathBuf, + /// Causal heads written to the file. + pub heads: Vec, + /// Number of canonical bytes written. + pub bytes_written: usize, + /// Recovery sidecar associated with the document. + pub recovery_path: PathBuf, + /// Whether cleanup of the recovery sidecar could not be confirmed. + pub recovery_retained: bool, +} + +#[derive(Clone, Debug, Serialize, Deserialize)] +struct RecoveryFile { + format: String, + version: u32, + source_path: String, + document_id: DocumentId, + base_heads: Vec, + current_heads: Vec, + /// Compact Automerge bytes at the beginning of the recovery window. + snapshot: Vec, + /// Changes after `base_heads`, retained in causal order. + journal: Vec, +} + +impl RecoveryFile { + fn new( + source_path: &Path, + document_id: DocumentId, + snapshot: Vec, + heads: Vec, + ) -> Self { + Self { + format: RECOVERY_FORMAT.into(), + version: RECOVERY_VERSION, + source_path: path_string(source_path), + document_id, + base_heads: heads.clone(), + current_heads: heads, + snapshot, + journal: Vec::new(), + } + } +} + +/// A lock-held Rust document session. +pub struct DocumentFile { + path: PathBuf, + actor_id: ActorId, + engine: TransactionEngine, + options: PersistenceOptions, + baseline_bytes: Vec, + baseline_heads: Vec, + pending_recovery: Option, + _lock: AdvisoryLock, +} + +impl DocumentFile { + /// Opens a canonical .inkfinite file and holds its advisory lock for the + /// lifetime of the session. + /// + /// # Errors + /// + /// Returns a typed filesystem, CRDT, or validation error. An existing + /// recovery record is left untouched; call recover explicitly to adopt it. + pub fn open(path: impl AsRef, actor_id: ActorId) -> Result { + Self::open_with_options(path, actor_id, PersistenceOptions::default()) + } + + /// Opens a canonical file with explicit recovery settings. + /// + /// # Errors + /// + /// Returns [`FileError`] when the actor, lock, canonical bytes, or + /// materialized document is invalid. + pub fn open_with_options( + path: impl AsRef, + actor_id: ActorId, + options: PersistenceOptions, + ) -> Result { + ensure_actor(&actor_id)?; + let path = absolute_path(path.as_ref())?; + let lock = AdvisoryLock::acquire(&path)?; + let bytes = read_bytes(&path, "read canonical document")?; + let mut engine = load_canonical_bytes(&bytes, actor_id.clone())?; + let heads = engine.snapshot()?.heads; + Ok(Self { + path, + actor_id, + engine, + options, + baseline_bytes: bytes, + baseline_heads: heads, + pending_recovery: None, + _lock: lock, + }) + } + + /// Creates and safely persists a new canonical document. + /// + /// # Errors + /// + /// Returns [`FileError`] when the document is invalid, the destination + /// exists, or a safe write cannot complete. + pub fn create( + path: impl AsRef, + document_id: DocumentId, + actor_id: ActorId, + document: Document, + ) -> Result { + Self::create_with_options( + path, + document_id, + actor_id, + document, + PersistenceOptions::default(), + ) + } + + /// Creates and safely persists a new document with explicit recovery + /// settings. + /// + /// # Errors + /// + /// Returns [`FileError`] when the document is invalid, the destination + /// exists, or a safe write cannot complete. + pub fn create_with_options( + path: impl AsRef, + document_id: DocumentId, + actor_id: ActorId, + document: Document, + options: PersistenceOptions, + ) -> Result { + ensure_actor(&actor_id)?; + let path = absolute_path(path.as_ref())?; + let lock = AdvisoryLock::acquire(&path)?; + if path.exists() { + return Err(FileError::AlreadyExists { path }); + } + let mut engine = TransactionEngine::create(document_id, actor_id.clone(), document)?; + let baseline_bytes = engine.save()?; + let baseline_heads = engine.snapshot()?.heads; + let mut session = Self { + path, + actor_id, + engine, + options, + baseline_bytes, + baseline_heads, + pending_recovery: None, + _lock: lock, + }; + session.save()?; + Ok(session) + } + + /// Imports a v1 file into a newly persisted canonical destination. + /// + /// # Errors + /// + /// Returns [`FileError`] when migration or canonical persistence fails. + pub fn import_v1( + source: impl AsRef, + destination: impl AsRef, + actor_id: ActorId, + ) -> Result { + import_v1_file(source, destination, actor_id) + } + + /// Imports a v1 file with explicit recovery settings. + /// + /// # Errors + /// + /// Returns [`FileError`] when migration or canonical persistence fails. + pub fn import_v1_with_options( + source: impl AsRef, + destination: impl AsRef, + actor_id: ActorId, + options: PersistenceOptions, + ) -> Result { + import_v1_file_with_options(source, destination, actor_id, options) + } + + /// Recovers the newest interrupted save associated with the path. + /// + /// Recovery loads the compact base snapshot, applies its bounded change + /// journal, validates the result through the transaction engine, and keeps + /// the recovery record until the caller successfully saves the recovered + /// document. + /// + /// # Errors + /// + /// Returns [`FileError`] when the recovery record is absent, malformed, or + /// cannot be validated and adopted. + pub fn recover( + path: impl AsRef, + actor_id: ActorId, + options: PersistenceOptions, + ) -> Result { + ensure_actor(&actor_id)?; + let path = absolute_path(path.as_ref())?; + let recovery_path = find_recovery_path(&path, &options)?; + let lock = AdvisoryLock::acquire(&path)?; + let recovery = read_recovery(&recovery_path, &path, &options)?; + let mut engine = TransactionEngine::load(&recovery.snapshot, actor_id.clone())?; + let base_snapshot = engine.snapshot()?; + if base_snapshot.document_id != recovery.document_id + || !same_heads(&base_snapshot.heads, &recovery.base_heads) + { + return Err(FileError::InvalidRecovery( + "base snapshot identity does not match recovery metadata".into(), + )); + } + if !recovery.journal.is_empty() { + engine.merge_changes(&recovery.journal)?; + } + let recovered_snapshot = engine.snapshot()?; + if recovered_snapshot.document_id != recovery.document_id + || !same_heads(&recovered_snapshot.heads, &recovery.current_heads) + { + return Err(FileError::InvalidRecovery( + "recovery journal did not produce the recorded heads".into(), + )); + } + Ok(Self { + path, + actor_id, + engine, + options, + baseline_bytes: recovery.snapshot.clone(), + baseline_heads: recovery.base_heads.clone(), + pending_recovery: Some(recovery), + _lock: lock, + }) + } + + /// Returns the canonical path held by this session. + #[must_use] + pub fn path(&self) -> &Path { + &self.path + } + + /// Returns the actor used for future local changes. + #[must_use] + pub fn actor_id(&self) -> &ActorId { + &self.actor_id + } + + /// Borrows the transaction engine for read-only inspection. + #[must_use] + pub fn engine(&self) -> &TransactionEngine { + &self.engine + } + + /// Borrows the transaction engine for commits, undo, redo, and queries. + pub fn engine_mut(&mut self) -> &mut TransactionEngine { + &mut self.engine + } + + /// Commits one validated transaction through the held document session. + /// + /// # Errors + /// + /// Returns [`FileError`] when the transaction engine rejects the draft. + pub fn commit( + &mut self, + transaction: TransactionDraft, + ) -> Result { + Ok(self.engine.commit(transaction)?) + } + + /// Materializes the current v2 snapshot. + /// + /// # Errors + /// + /// Returns [`FileError`] when the CRDT snapshot cannot be materialized. + pub fn snapshot(&mut self) -> Result { + Ok(self.engine.snapshot()?) + } + + /// Returns deterministic inspection JSON for the current snapshot. + /// + /// This projection contains records and causal heads only. It cannot + /// preserve the Automerge history of the canonical file. + /// + /// # Errors + /// + /// Returns [`FileError`] when the snapshot cannot be materialized or + /// serialized. + pub fn export_json(&mut self) -> Result { + let snapshot = self.snapshot()?; + export_snapshot_json(&snapshot) + } + + /// Writes deterministic snapshot JSON to a separate file atomically. + /// + /// # Errors + /// + /// Returns [`FileError`] when the snapshot cannot be materialized or the + /// destination cannot be written safely. + pub fn export_json_to(&mut self, path: impl AsRef) -> Result<(), FileError> { + let destination = absolute_path(path.as_ref())?; + if paths_equivalent(&self.path, &destination) { + return Err(FileError::SamePath { + path: self.path.clone(), + }); + } + let snapshot = self.snapshot()?; + write_snapshot_json(destination, &snapshot) + } + + /// Returns the expected recovery sidecar path for this document. + /// + /// # Errors + /// + /// Returns [`FileError`] when the current snapshot cannot be materialized. + pub fn recovery_path(&mut self) -> Result { + let document_id = self.snapshot()?.document_id; + Ok(recovery_path_for(&self.path, &document_id, &self.options)) + } + + /// Reports whether a recovery sidecar is present for this session. + /// + /// # Errors + /// + /// Returns [`FileError`] when the current snapshot cannot be materialized. + pub fn recovery_available(&mut self) -> Result { + Ok(self.recovery_path()?.exists()) + } + + /// Safely persists compact Automerge bytes and retains recovery until the + /// canonical replacement succeeds. + /// + /// # Errors + /// + /// Returns [`FileError`] when recovery preparation, flushing, replacement, + /// or validation fails. A failed canonical replacement leaves recovery data + /// available for [`Self::recover`]. + pub fn save(&mut self) -> Result { + let snapshot = self.engine.snapshot()?; + let document_id = snapshot.document_id.clone(); + let heads = snapshot.heads.clone(); + let bytes = self.engine.save()?; + let recovery_path = recovery_path_for(&self.path, &document_id, &self.options); + + let mut recovery = if let Some(recovery) = self.pending_recovery.clone() { + recovery + } else if recovery_path.exists() { + read_recovery(&recovery_path, &self.path, &self.options)? + } else { + RecoveryFile::new( + &self.path, + document_id.clone(), + self.baseline_bytes.clone(), + self.baseline_heads.clone(), + ) + }; + if recovery.document_id != document_id { + return Err(FileError::RecoveryAhead { + path: recovery_path, + }); + } + if !same_heads(&recovery.current_heads, &heads) { + let changes = self + .engine + .changes_since(&recovery.current_heads) + .map_err(|_| FileError::RecoveryAhead { + path: recovery_path.clone(), + })?; + recovery.journal.extend(changes); + } + recovery.current_heads.clone_from(&heads); + if journal_bytes(&recovery.journal) > self.options.max_journal_bytes() + || recovery.journal.len() > self.options.max_journal_entries() + { + recovery.snapshot.clone_from(&bytes); + recovery.base_heads.clone_from(&heads); + recovery.journal.clear(); + } + self.pending_recovery = Some(recovery.clone()); + let recovery_directory = recovery_path + .parent() + .ok_or_else(|| FileError::InvalidRecovery("recovery path has no parent".into()))?; + fs::create_dir_all(recovery_directory).map_err(|error| { + io_error( + "create recovery directory", + recovery_directory.to_owned(), + error, + ) + })?; + write_recovery(&recovery_path, &recovery)?; + + if let Err(error) = atomic_write(&self.path, &bytes) { + self.pending_recovery = Some(recovery); + return Err(error); + } + self.baseline_bytes.clone_from(&bytes); + self.baseline_heads.clone_from(&heads); + self.pending_recovery = None; + + let recovery_retained = match fs::remove_file(&recovery_path) { + Ok(()) => false, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => false, + Err(_) => true, + }; + Ok(SaveResult { + path: self.path.clone(), + heads, + bytes_written: bytes.len(), + recovery_path, + recovery_retained, + }) + } +} + +/// Reads and migrates a v1 JSON file without writing any destination. +/// +/// # Errors +/// +/// Returns [`FileError`] when the source cannot be read, locked, parsed, or +/// migrated. +pub fn read_v1_file(path: impl AsRef, actor_id: ActorId) -> Result { + ensure_actor(&actor_id)?; + let path = absolute_path(path.as_ref())?; + let _lock = AdvisoryLock::acquire(&path)?; + let input = String::from_utf8(read_bytes(&path, "read v1 document")?) + .map_err(|error| FileError::InvalidV1(format!("v1 document is not UTF-8: {error}")))?; + import_v1_json(&input, actor_id) +} + +/// Imports a v1 file and safely writes its canonical v2 representation. +/// +/// # Errors +/// +/// Returns [`FileError`] when the source cannot be migrated or the destination +/// cannot be written safely. +pub fn import_v1_file( + source: impl AsRef, + destination: impl AsRef, + actor_id: ActorId, +) -> Result { + import_v1_file_with_options(source, destination, actor_id, PersistenceOptions::default()) +} + +/// Imports a v1 file and safely writes its canonical v2 representation with +/// explicit recovery settings. +/// +/// # Errors +/// +/// Returns [`FileError`] when the source cannot be migrated or the destination +/// cannot be written safely. +pub fn import_v1_file_with_options( + source: impl AsRef, + destination: impl AsRef, + actor_id: ActorId, + options: PersistenceOptions, +) -> Result { + ensure_actor(&actor_id)?; + let source = absolute_path(source.as_ref())?; + let destination = absolute_path(destination.as_ref())?; + if paths_equivalent(&source, &destination) { + return Err(FileError::SamePath { path: source }); + } + let imported = read_v1_file(&source, actor_id.clone())?; + let engine = imported.into_engine(actor_id.clone())?; + create_session_from_engine(destination, actor_id, engine, options) +} + +/// Loads canonical Automerge bytes into a validated transaction engine. +/// +/// # Errors +/// +/// Returns [`FileError`] when the actor, CRDT bytes, or materialized document is +/// invalid. +pub fn load_canonical_bytes( + bytes: &[u8], + actor_id: ActorId, +) -> Result { + ensure_actor(&actor_id)?; + Ok(TransactionEngine::load(bytes, actor_id)?) +} + +/// Serializes a v2 materialized snapshot in deterministic, human-readable JSON. +/// +/// Map keys are ordered by the v2 model's `BTreeMap` fields and causal heads are +/// sorted for stable output. The result is a snapshot projection: applying it +/// to a new CRDT would create new history rather than preserve the original +/// Automerge changes. +/// +/// # Errors +/// +/// Returns [`FileError::Json`] when the snapshot cannot be serialized. +pub fn export_snapshot_json(snapshot: &DocumentSnapshot) -> Result { + let mut snapshot = snapshot.clone(); + snapshot.heads.sort(); + Ok(format!("{}\n", serde_json::to_string_pretty(&snapshot)?)) +} + +/// Writes a deterministic snapshot JSON file with a same-directory temporary +/// file and atomic replacement. +/// +/// # Errors +/// +/// Returns [`FileError`] when the snapshot cannot be serialized or the +/// destination cannot be written safely. +pub fn write_snapshot_json( + path: impl AsRef, + snapshot: &DocumentSnapshot, +) -> Result<(), FileError> { + let path = absolute_path(path.as_ref())?; + let _lock = AdvisoryLock::acquire(&path)?; + let contents = export_snapshot_json(snapshot)?; + atomic_write(&path, contents.as_bytes()).map(|_| ()) +} + +/// Returns the recovery sidecar path for a document ID and persistence policy. +pub fn recovery_path_for( + document_path: impl AsRef, + document_id: &DocumentId, + options: &PersistenceOptions, +) -> PathBuf { + let document_path = document_path.as_ref(); + options.recovery_directory_for(document_path).join(format!( + "{}.recovery", + encode_path_component(document_id.as_str()) + )) +} + +fn create_session_from_engine( + path: PathBuf, + actor_id: ActorId, + mut engine: TransactionEngine, + options: PersistenceOptions, +) -> Result { + let lock = AdvisoryLock::acquire(&path)?; + let baseline_bytes = engine.save()?; + let baseline_heads = engine.snapshot()?.heads; + let mut session = DocumentFile { + path, + actor_id, + engine, + options, + baseline_bytes, + baseline_heads, + pending_recovery: None, + _lock: lock, + }; + session.save()?; + Ok(session) +} + +fn write_recovery(path: &Path, recovery: &RecoveryFile) -> Result<(), FileError> { + let bytes = serde_json::to_vec(recovery)?; + atomic_write(path, &bytes).map(|_| ()) +} + +fn read_recovery( + path: &Path, + document_path: &Path, + options: &PersistenceOptions, +) -> Result { + let bytes = read_bytes(path, "read recovery record")?; + let recovery: RecoveryFile = serde_json::from_slice(&bytes) + .map_err(|error| FileError::InvalidRecovery(format!("{}: {error}", path.display())))?; + if recovery.format != RECOVERY_FORMAT || recovery.version != RECOVERY_VERSION { + return Err(FileError::InvalidRecovery(format!( + "unsupported recovery format {:?} version {}", + recovery.format, recovery.version + ))); + } + if recovery.source_path != path_string(document_path) { + return Err(FileError::InvalidRecovery( + "recovery record belongs to another document".into(), + )); + } + if recovery.snapshot.is_empty() + || recovery.base_heads.is_empty() + || recovery.current_heads.is_empty() + || recovery.journal.len() > options.max_journal_entries() + || journal_bytes(&recovery.journal) > options.max_journal_bytes() + { + return Err(FileError::InvalidRecovery( + "recovery snapshot or bounded journal is invalid".into(), + )); + } + Ok(recovery) +} + +fn find_recovery_path( + document_path: &Path, + options: &PersistenceOptions, +) -> Result { + let directory = options.recovery_directory_for(document_path); + let entries = fs::read_dir(&directory).map_err(|error| { + if error.kind() == std::io::ErrorKind::NotFound { + FileError::RecoveryNotFound { + path: document_path.to_owned(), + } + } else { + io_error("list recovery directory", directory.clone(), error) + } + })?; + let expected_source = path_string(document_path); + let mut candidates = Vec::new(); + for entry in entries { + let entry = + entry.map_err(|error| io_error("read recovery entry", directory.clone(), error))?; + let path = entry.path(); + if path.extension().and_then(|value| value.to_str()) != Some("recovery") { + continue; + } + let Ok(bytes) = fs::read(&path) else { + continue; + }; + let Ok(recovery) = serde_json::from_slice::(&bytes) else { + continue; + }; + if recovery.source_path == expected_source { + candidates.push(path); + } + } + candidates.sort(); + candidates + .into_iter() + .next() + .ok_or_else(|| FileError::RecoveryNotFound { + path: document_path.to_owned(), + }) +} + +fn ensure_actor(actor_id: &ActorId) -> Result<(), FileError> { + if actor_id.as_str().trim().is_empty() { + Err(FileError::InvalidV1("actor ID must not be empty".into())) + } else { + Ok(()) + } +} + +fn journal_bytes(journal: &[EncodedChange]) -> usize { + journal.iter().map(|change| change.as_bytes().len()).sum() +} + +fn same_heads(left: &[ChangeHash], right: &[ChangeHash]) -> bool { + let mut left = left.to_vec(); + let mut right = right.to_vec(); + left.sort(); + right.sort(); + left == right +} + +fn read_bytes(path: &Path, operation: &'static str) -> Result, FileError> { + fs::read(path).map_err(|error| io_error(operation, path.to_owned(), error)) +} + +fn absolute_path(path: &Path) -> Result { + if path.is_absolute() { + return Ok(path.to_owned()); + } + let current = std::env::current_dir() + .map_err(|error| io_error("resolve current directory", PathBuf::from("."), error))?; + Ok(current.join(path)) +} + +fn paths_equivalent(left: &Path, right: &Path) -> bool { + let left = fs::canonicalize(left).unwrap_or_else(|_| left.to_owned()); + let right = fs::canonicalize(right).unwrap_or_else(|_| right.to_owned()); + left == right +} + +fn path_string(path: &Path) -> String { + path.to_string_lossy().into_owned() +} + +fn encode_path_component(value: &str) -> String { + let mut encoded = String::new(); + for byte in value.bytes() { + if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.') { + encoded.push(char::from(byte)); + } else { + let _ = write!(&mut encoded, "%{byte:02X}"); + } + } + if encoded.is_empty() { + "document".into() + } else { + encoded + } +} + +fn lock_path(path: &Path) -> PathBuf { + let file_name = path + .file_name() + .and_then(|value| value.to_str()) + .unwrap_or("document"); + path.parent() + .unwrap_or_else(|| Path::new(".")) + .join(format!(".{file_name}.lock")) +} + +struct AdvisoryLock { + path: PathBuf, + _file: File, +} + +impl AdvisoryLock { + fn acquire(document_path: &Path) -> Result { + let path = lock_path(document_path); + let mut file = match OpenOptions::new().write(true).create_new(true).open(&path) { + Ok(file) => file, + Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => { + return Err(FileError::Locked { + path: document_path.to_owned(), + }); + } + Err(error) => return Err(io_error("create document lock", path, error)), + }; + let owner = format!("pid={}\n", std::process::id()); + if let Err(error) = file + .write_all(owner.as_bytes()) + .and_then(|()| file.sync_all()) + { + let _ = fs::remove_file(&path); + return Err(io_error("write document lock", path, error)); + } + Ok(Self { path, _file: file }) + } +} + +impl Drop for AdvisoryLock { + fn drop(&mut self) { + let _ = fs::remove_file(&self.path); + } +} + +fn atomic_write(path: &Path, bytes: &[u8]) -> Result { + let parent = path.parent().unwrap_or_else(|| Path::new(".")); + let file_name = path + .file_name() + .and_then(|value| value.to_str()) + .unwrap_or("document"); + let temporary = parent.join(format!( + ".{file_name}.tmp-{}-{}", + std::process::id(), + TEMP_COUNTER.fetch_add(1, Ordering::Relaxed) + )); + let mut file = match OpenOptions::new() + .write(true) + .create_new(true) + .open(&temporary) + { + Ok(file) => file, + Err(error) => return Err(io_error("create temporary document", temporary, error)), + }; + let result = file + .write_all(bytes) + .and_then(|()| file.flush()) + .and_then(|()| file.sync_all()); + if let Err(error) = result { + let _ = fs::remove_file(&temporary); + return Err(io_error("flush temporary document", temporary, error)); + } + drop(file); + if let Err(error) = replace_file(&temporary, path) { + let _ = fs::remove_file(&temporary); + return Err(io_error( + "replace canonical document", + path.to_owned(), + error, + )); + } + if let Err(error) = sync_directory(parent) { + return Err(io_error( + "flush document directory", + parent.to_owned(), + error, + )); + } + Ok(bytes.len()) +} + +#[cfg(not(windows))] +fn replace_file(temporary: &Path, destination: &Path) -> std::io::Result<()> { + fs::rename(temporary, destination) +} + +#[cfg(windows)] +fn replace_file(temporary: &Path, destination: &Path) -> std::io::Result<()> { + if !destination.exists() { + return fs::rename(temporary, destination); + } + let backup = destination.with_extension("inkfinite-replace-backup"); + fs::rename(destination, &backup)?; + match fs::rename(temporary, destination) { + Ok(()) => { + let _ = fs::remove_file(backup); + Ok(()) + } + Err(error) => { + let _ = fs::rename(&backup, destination); + Err(error) + } + } +} + +#[cfg(unix)] +fn sync_directory(path: &Path) -> std::io::Result<()> { + File::open(path)?.sync_all() +} + +#[cfg(not(unix))] +fn sync_directory(_path: &Path) -> std::io::Result<()> { + Ok(()) +} + +fn io_error(operation: &'static str, path: PathBuf, source: std::io::Error) -> FileError { + FileError::Io { + operation, + path, + source, + } +} diff --git a/crates/inkfinite-file/src/tests.rs b/crates/inkfinite-file/src/tests.rs new file mode 100644 index 0000000..3f7794d --- /dev/null +++ b/crates/inkfinite-file/src/tests.rs @@ -0,0 +1,365 @@ +use std::collections::BTreeMap; +use std::fs; +use std::path::PathBuf; +use std::sync::atomic::{AtomicU64, Ordering}; + +use inkfinite_engine::TransactionDraft; +use inkfinite_model::{ + ActorId, Document, DocumentId, LayerId, LayerRecord, Opacity, Origin, PageId, PageRecord, + RecordVersion, Timestamp, +}; +use inkfinite_protocol::{Operation, TransactionId}; +use serde_json::Value; + +use super::{ + DocumentFile, FileError, PersistenceOptions, export_snapshot_json, import_v1_file, + import_v1_json, +}; + +static TEST_COUNTER: AtomicU64 = AtomicU64::new(0); + +#[test] +fn imports_all_valid_v1_features_into_normalized_v2_records() { + let input = include_str!("../../../fixtures/v1/desktop/all-features.inkfinite.json"); + let imported = import_v1_json(input, ActorId::from("actor:test")).expect("fixture imports"); + let source: Value = serde_json::from_str(input).expect("fixture JSON"); + let document = &imported.document; + + let expected_page_ids = source["order"]["pageIds"] + .as_array() + .expect("page order") + .iter() + .map(|id| id.as_str().expect("page ID")) + .collect::>(); + assert_eq!( + document + .page_ids + .iter() + .map(inkfinite_model::PageId::as_str) + .collect::>(), + expected_page_ids + ); + assert_eq!(document.pages.len(), 2); + assert_eq!(document.layers.len(), 2); + for page in document.pages.values() { + assert_eq!(page.layer_ids.len(), 1); + let layer = &document.layers[&page.layer_ids[0]]; + let flattened = flatten_shape_order(document, &layer.shape_ids); + let expected = source["order"]["shapeOrder"][page.id.as_str()] + .as_array() + .expect("shape order") + .iter() + .map(|id| id.as_str().expect("shape ID")) + .map(str::to_owned) + .collect::>(); + assert_eq!(flattened, expected); + } + + assert!( + document + .shapes + .contains_key(&"shape:stencil-process".into()) + ); + assert!(document.shapes.contains_key(&"shape:arrow".into())); + assert!(document.shapes.contains_key(&"shape:markdown".into())); + assert_eq!(document.shapes.len(), 16); + assert_eq!( + document.shapes[&"shape:stencil-process".into()].properties["width"].as_f64(), + Some(120.0) + ); + assert_eq!( + document.shapes[&"shape:stencil-process".into()].properties["height"].as_f64(), + Some(80.0) + ); + let stroke_opacity = document.shapes[&"shape:stroke".into()] + .style + .stroke_opacity + .expect("stroke opacity") + .get(); + assert!((stroke_opacity - 0.75).abs() < f32::EPSILON); + assert_eq!( + document.shapes[&"shape:text".into()].properties["legacy_group_id"], + Value::from("group:content") + ); + let card_group = &document.shapes[&"group:card".into()]; + assert_eq!( + card_group.child_ids, + vec![ + "shape:stencil-card".into(), + "shape:stencil-card-divider".into() + ] + ); + assert!(matches!( + document.bindings[&"binding:arrow-start".into()].anchor, + inkfinite_model::BindingAnchor::Edge { x: 1.0, y: 0.0 } + )); + assert!(matches!( + document.bindings[&"binding:arrow-end".into()].anchor, + inkfinite_model::BindingAnchor::Center + )); +} + +#[test] +fn imports_web_and_performance_v1_fixtures() { + let web = include_str!("../../../fixtures/v1/web/all-features.web.json"); + let imported = import_v1_json(web, ActorId::from("actor:web")).expect("web fixture imports"); + assert_eq!(imported.document.page_ids.len(), 2); + assert_eq!(imported.document.layers.len(), 2); + assert_eq!(imported.document.shapes.len(), 16); + + let performance = include_str!("../../../fixtures/v1/performance/board-10000.inkfinite.json"); + let imported = import_v1_json(performance, ActorId::from("actor:performance")) + .expect("large fixture imports"); + assert_eq!(imported.document.page_ids, vec![PageId::from("page:10000")]); + assert_eq!(imported.document.shapes.len(), 10_000); + let page = &imported.document.pages[&PageId::from("page:10000")]; + let layer = &imported.document.layers[&page.layer_ids[0]]; + assert_eq!(layer.shape_ids.len(), 10_000); + assert_eq!( + flatten_shape_order(&imported.document, &layer.shape_ids).len(), + 10_000 + ); +} + +#[test] +fn rejects_invalid_and_newer_inputs_before_persistence() { + let actor = ActorId::from("actor:test"); + let invalid_inputs = [ + include_str!("../../../fixtures/v1/invalid/malformed-json.inkfinite.json"), + include_str!("../../../fixtures/v1/invalid/missing-envelope-fields.json"), + include_str!("../../../fixtures/v1/invalid/dangling-references.inkfinite.json"), + include_str!("../../../fixtures/v1/invalid/duplicate-order.inkfinite.json"), + ]; + for input in invalid_inputs { + assert!( + matches!( + import_v1_json(input, actor.clone()), + Err(FileError::Json(_) | FileError::InvalidV1(_)) + ), + "input should be rejected" + ); + } + assert!(matches!( + import_v1_json( + r#"{"format":"inkfinite.document","format_version":3}"#, + actor.clone() + ), + Err(FileError::UnsupportedFormat { version: 3, .. }) + )); + + let temporary = TestDirectory::new(); + let source = temporary.path.join("invalid.inkfinite.json"); + let destination = temporary.path.join("existing.inkfinite"); + fs::write( + &source, + include_str!("../../../fixtures/v1/invalid/duplicate-order.inkfinite.json"), + ) + .expect("write source"); + fs::write(&destination, b"keep this file").expect("write destination"); + let result = import_v1_file(&source, &destination, actor); + assert!(matches!(result, Err(FileError::InvalidV1(_)))); + assert_eq!( + fs::read(&destination).expect("read destination"), + b"keep this file" + ); +} + +#[test] +fn canonical_sessions_lock_save_reopen_and_export_deterministically() { + let input = include_str!("../../../fixtures/v1/desktop/all-features.inkfinite.json"); + let imported = import_v1_json(input, ActorId::from("actor:test")).expect("fixture imports"); + let expected_document = imported.document.clone(); + let temporary = TestDirectory::new(); + let canonical = temporary.path.join("board.inkfinite"); + let snapshot_json = temporary.path.join("board.inkfinite.json"); + let options = PersistenceOptions::with_recovery_directory(temporary.path.join("recovery")); + + let mut session = DocumentFile::create_with_options( + &canonical, + imported.document_id.clone(), + ActorId::from("actor:test"), + imported.document, + options.clone(), + ) + .expect("create canonical document"); + let first_json = session.export_json().expect("export JSON"); + assert_eq!(first_json, session.export_json().expect("repeat export")); + assert!(first_json.contains("\"format\": \"inkfinite.document\"")); + session + .export_json_to(&snapshot_json) + .expect("write snapshot JSON"); + assert!(matches!( + session.export_json_to(&canonical), + Err(FileError::SamePath { .. }) + )); + assert!(matches!( + DocumentFile::open_with_options(&canonical, ActorId::from("actor:other"), options.clone()), + Err(FileError::Locked { .. }) + )); + drop(session); + + let canonical_bytes = fs::read(&canonical).expect("canonical bytes"); + assert!(serde_json::from_slice::(&canonical_bytes).is_err()); + let exported: inkfinite_model::DocumentSnapshot = + serde_json::from_str(&fs::read_to_string(&snapshot_json).expect("snapshot JSON")) + .expect("snapshot parses"); + assert_eq!(exported.document, expected_document); + + let mut reopened = + DocumentFile::open_with_options(&canonical, ActorId::from("actor:reopen"), options) + .expect("reopen canonical document"); + assert_eq!( + reopened.snapshot().expect("snapshot").document, + expected_document + ); + assert_eq!( + export_snapshot_json(&reopened.snapshot().expect("snapshot")).expect("export"), + first_json + ); +} + +#[test] +fn recovery_restores_journal_after_canonical_replace_failure() { + let temporary = TestDirectory::new(); + let canonical = temporary.path.join("board.inkfinite"); + let recovery_directory = temporary.path.join("recovery"); + let options = PersistenceOptions { + recovery_directory: Some(recovery_directory.clone()), + max_journal_entries: 1, + max_journal_bytes: 1024 * 1024, + }; + let mut session = DocumentFile::create_with_options( + &canonical, + DocumentId::from("document:recovery"), + ActorId::from("actor:writer"), + simple_document(), + options.clone(), + ) + .expect("create document"); + let base = session.snapshot().expect("base snapshot"); + session + .commit(TransactionDraft { + id: TransactionId("transaction:rename".into()), + actor_id: ActorId::from("actor:writer"), + origin: Origin::Human, + base_heads: base.heads, + description: "rename page".into(), + operations: vec![Operation::RenamePage { + page_id: PageId::from("page:one"), + name: "Recovered".into(), + expected_version: Some(RecordVersion(1)), + }], + timestamp: Timestamp(1), + }) + .expect("commit edit"); + let recovery_path = session.recovery_path().expect("recovery path"); + let backup = temporary.path.join("old.inkfinite"); + fs::rename(&canonical, &backup).expect("move canonical aside"); + fs::create_dir(&canonical).expect("make replacement fail"); + + let result = session.save(); + assert!(matches!(result, Err(FileError::Io { .. }))); + assert!(recovery_path.exists()); + let recovery: Value = + serde_json::from_slice(&fs::read(&recovery_path).expect("recovery bytes")) + .expect("recovery JSON"); + assert!( + recovery["snapshot"] + .as_array() + .is_some_and(|bytes| !bytes.is_empty()) + ); + assert_eq!(recovery["journal"].as_array().expect("journal").len(), 1); + drop(session); + + fs::remove_dir(&canonical).expect("remove failure directory"); + fs::rename(&backup, &canonical).expect("restore old canonical"); + let mut recovered = DocumentFile::recover(&canonical, ActorId::from("actor:recovery"), options) + .expect("recover interrupted save"); + assert_eq!( + recovered + .snapshot() + .expect("recovered snapshot") + .document + .pages[&PageId::from("page:one")] + .name, + "Recovered" + ); + let save = recovered.save().expect("save recovered document"); + assert!(!save.recovery_retained); + assert!(!recovery_path.exists()); + drop(recovered); + + let mut reopened = + DocumentFile::open(&canonical, ActorId::from("actor:verify")).expect("reopen recovered"); + assert_eq!( + reopened.snapshot().expect("snapshot").document.pages[&PageId::from("page:one")].name, + "Recovered" + ); +} + +fn flatten_shape_order(document: &Document, shape_ids: &[inkfinite_model::ShapeId]) -> Vec { + let mut flattened = Vec::new(); + for shape_id in shape_ids { + if let Some(shape) = document.shapes.get(shape_id) { + if shape.kind.as_str() != inkfinite_model::CONTAINER_KIND { + flattened.push(shape_id.as_str().to_owned()); + } + flattened.extend(flatten_shape_order(document, &shape.child_ids)); + } else { + flattened.push(shape_id.as_str().to_owned()); + } + } + flattened +} + +fn simple_document() -> Document { + let page_id = PageId::from("page:one"); + let layer_id = LayerId::from("layer:one"); + Document { + pages: BTreeMap::from([( + page_id.clone(), + PageRecord { + id: page_id.clone(), + name: "Page".into(), + layer_ids: vec![layer_id.clone()], + version: RecordVersion(1), + }, + )]), + page_ids: vec![page_id.clone()], + layers: BTreeMap::from([( + layer_id.clone(), + LayerRecord { + id: layer_id, + page_id, + name: "Default".into(), + shape_ids: Vec::new(), + visible: true, + locked: false, + opacity: Opacity::OPAQUE, + version: RecordVersion(1), + }, + )]), + shapes: BTreeMap::new(), + bindings: BTreeMap::new(), + assets: BTreeMap::new(), + } +} + +struct TestDirectory { + path: PathBuf, +} + +impl TestDirectory { + fn new() -> Self { + let id = TEST_COUNTER.fetch_add(1, Ordering::Relaxed); + let path = std::env::temp_dir().join(format!("inkfinite-file-test-{id}")); + fs::create_dir_all(&path).expect("create test directory"); + Self { path } + } +} + +impl Drop for TestDirectory { + fn drop(&mut self) { + let _ = fs::remove_dir_all(&self.path); + } +} diff --git a/docs/v2-file-format.md b/docs/v2-file-format.md new file mode 100644 index 0000000..0774155 --- /dev/null +++ b/docs/v2-file-format.md @@ -0,0 +1,45 @@ +# V2 file format + +Inkfinite stores the canonical document as an Automerge file with the +.inkfinite extension. The .inkfinite.json format is used for v1 imports and +stable JSON snapshots; it is not a CRDT round-trip format. + +## V1 migration + +The file crate validates the v1 envelope, page order, shape ownership, draw +order, bindings, and references before it creates a v2 document. + +- board.id becomes the v2 document ID. +- order.pageIds becomes document page order. +- Each page receives one stable default layer named Default. +- A page's shapeOrder is used when present; otherwise its shapeIds are used. +- Shape properties are retained. The v1 w and h properties become v2 width and + height. +- Contiguous legacy groups become v2 containers with their original group ID. + Their children retain the legacy draw order. +- Non-contiguous or singleton groups remain flat so the exact draw order is + preserved, with legacy_group_id retained as compatibility metadata. +- Freehand stroke opacity becomes the v2 stroke opacity value. + +Invalid and newer inputs return typed errors before a destination file is +created or replaced. + +## Safe writes and recovery + +DocumentFile holds an advisory sidecar lock for the canonical path. A save +serializes a compact CRDT snapshot, records a bounded journal of changes since +the last durable snapshot, writes recovery data atomically, and then replaces +the canonical file through a same-directory temporary file after flushing and +syncing it. Recovery data is removed only after the canonical replacement +succeeds. + +Recovery combines the saved compact snapshot with its encoded change journal, +validates the resulting document and causal heads, and can then be saved as a +new canonical baseline. The recovery journal is compacted when its configured +entry or byte bound is reached. + +## JSON snapshots + +The JSON export is deterministic: CRDT heads are sorted and the materialized +model uses ordered maps and lists. It represents one materialized snapshot and +cannot preserve Automerge history, causal heads, or an undo/redo journal.