use std::fs::{self, File, OpenOptions};
use std::io::{self, Write};
use std::path::{Path, PathBuf};
use crate::backup::{new_op_id, BackupStore};
use crate::hashline::apply::{
apply_section_ops, commit_registers_if_complete, FileClassification, FileResult, MutationState,
PlannedFile, RegisterStore, SectionPlanInput, StagedRegisters,
};
use crate::hashline::scan::Snapshot;
use crate::hashline::snapshot::{
invalidate_removed_source, publish_edit_response_snapshot, AffectedRegion,
EditResponseSnapshot, SnapshotStore,
};
use crate::hashline::syntax::{
Baseline, HashlineRejection, MvOperation, Operation, ResolvedOperation,
};
#[derive(Clone, Debug)]
pub struct MvDestinationInput<'a> {
pub canonical_path: &'a Path,
pub requested_path: &'a str,
pub baseline_bytes: Option<&'a [u8]>,
}
#[derive(Clone, Debug)]
pub struct TransactionSectionInput<'a> {
pub canonical_path: &'a Path,
pub requested_path: &'a str,
pub baseline: &'a Baseline,
pub snapshot: &'a Snapshot,
pub operations: &'a [Operation],
pub resolved: &'a [ResolvedOperation],
pub mv_destination: Option<MvDestinationInput<'a>>,
}
#[derive(Clone, Debug)]
pub struct TransactionPlan {
pub steps: Vec<PlannedStep>,
pub staged_registers: StagedRegisters,
}
#[derive(Clone, Debug)]
pub enum PlannedStep {
Mutate(PlannedFile),
Mv(PlannedMv),
}
#[derive(Clone, Debug)]
pub struct PlannedMv {
pub source_canonical: PathBuf,
pub source_requested: String,
pub source_baseline_bytes: Vec<u8>,
pub dest_canonical: PathBuf,
pub dest_requested: String,
pub dest_existed: bool,
pub dest_baseline_bytes: Option<Vec<u8>>,
pub final_bytes: Vec<u8>,
pub affected: AffectedRegion,
pub warnings: Vec<String>,
pub repair_layers: Vec<&'static str>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum FileRole {
Primary,
MvDestination,
MvSource,
}
impl FileRole {
pub const fn as_str(self) -> &'static str {
match self {
Self::Primary => "primary",
Self::MvDestination => "mv_destination",
Self::MvSource => "mv_source",
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct FileOutcome {
pub canonical_path: PathBuf,
pub requested_path: String,
pub role: FileRole,
pub classification: FileClassification,
pub mutation_state: MutationState,
pub final_bytes: Option<Vec<u8>>,
pub final_tag: Option<String>,
pub affected: AffectedRegion,
pub warnings: Vec<String>,
pub format_skipped_reason: Option<String>,
pub backup_id: Option<String>,
pub remove_file: bool,
pub tag_notice: Option<String>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TransactionEnvelope {
pub success: bool,
pub complete: bool,
pub files: Vec<FileOutcome>,
pub op_id: Option<String>,
pub stop_reason: Option<&'static str>,
pub registers_committed: bool,
pub preview: bool,
pub summary_text: String,
}
impl TransactionEnvelope {
pub fn to_apply_envelope(&self) -> crate::hashline::apply::ApplyResultEnvelope {
crate::hashline::apply::ApplyResultEnvelope {
success: self.success,
complete: self.complete,
files: self
.files
.iter()
.filter(|file| file.role != FileRole::MvSource || file.remove_file)
.map(|file| FileResult {
canonical_path: file.canonical_path.clone(),
requested_path: file.requested_path.clone(),
classification: file.classification,
mutation_state: file.mutation_state,
final_bytes: file.final_bytes.clone(),
affected: file.affected.clone(),
warnings: file.warnings.clone(),
remove_file: file.remove_file,
})
.collect(),
registers_committed: self.registers_committed,
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum ExecuteFault {
BaselineDrift { step: usize },
Backup { step: usize },
Write { step: usize },
Durability { step: usize },
SourceUnlink { step: usize },
FinalTagUnavailable { step: usize },
ValidationFailure { step: usize },
}
pub struct ExecuteContext<'a> {
pub session: &'a str,
pub backups: &'a mut BackupStore,
pub snapshots: &'a mut SnapshotStore,
pub registers: &'a mut RegisterStore,
pub backups_enabled: bool,
pub fault: Option<ExecuteFault>,
}
pub fn plan_transaction(
sections: &[TransactionSectionInput<'_>],
session_registers: &RegisterStore,
backups_enabled: bool,
) -> Result<TransactionPlan, HashlineRejection> {
let mut staged = session_registers.stage();
let mut steps = Vec::with_capacity(sections.len());
for section in sections {
let (line_ops, line_resolved, mv_op) = split_mv(section.operations, section.resolved)?;
if mv_op.is_some() && section.mv_destination.is_none() {
return Err(HashlineRejection::parse(
"MV section is missing resolved destination coordinates",
));
}
if mv_op.is_none() && section.mv_destination.is_some() {
return Err(HashlineRejection::parse(
"destination coordinates supplied without an MV operation",
));
}
let planned_source = if line_ops.is_empty() {
PlannedFile {
canonical_path: section.canonical_path.to_path_buf(),
requested_path: section.requested_path.to_string(),
baseline_bytes: section.baseline.bytes.clone(),
final_bytes: section.baseline.bytes.clone(),
affected: AffectedRegion::default(),
remove_file: false,
warnings: Vec::new(),
repair_layers: Vec::new(),
}
} else {
let input = SectionPlanInput {
canonical_path: section.canonical_path,
requested_path: section.requested_path,
baseline: section.baseline,
snapshot: section.snapshot,
operations: line_ops,
resolved: line_resolved,
};
let plan = plan_line_section(&input, &mut staged)?;
plan
};
let step = if let Some(dest) = section.mv_destination.as_ref() {
let dest_existed = dest.baseline_bytes.is_some();
PlannedStep::Mv(PlannedMv {
source_canonical: planned_source.canonical_path,
source_requested: planned_source.requested_path,
source_baseline_bytes: planned_source.baseline_bytes,
dest_canonical: dest.canonical_path.to_path_buf(),
dest_requested: dest.requested_path.to_string(),
dest_existed,
dest_baseline_bytes: dest.baseline_bytes.map(|bytes| bytes.to_vec()),
final_bytes: planned_source.final_bytes,
affected: planned_source.affected,
warnings: planned_source.warnings,
repair_layers: planned_source.repair_layers,
})
} else {
PlannedStep::Mutate(planned_source)
};
assert_rollback_available(backups_enabled, &step)?;
steps.push(step);
}
if steps.is_empty() {
return Err(HashlineRejection::parse(
"transaction plan requires at least one section",
));
}
Ok(TransactionPlan {
steps,
staged_registers: staged,
})
}
pub fn preview_transaction(plan: TransactionPlan) -> TransactionEnvelope {
let files = plan
.steps
.iter()
.flat_map(preview_step_files)
.collect::<Vec<_>>();
let summary_text = summary_counts(&files);
RegisterStore::discard(plan.staged_registers);
TransactionEnvelope {
success: true,
complete: true,
files,
op_id: None,
stop_reason: None,
registers_committed: false,
preview: true,
summary_text,
}
}
pub fn execute_transaction(
plan: TransactionPlan,
ctx: &mut ExecuteContext<'_>,
) -> TransactionEnvelope {
let TransactionPlan {
steps,
staged_registers,
} = plan;
let op_id = new_op_id();
let mut journaled = false;
let mut files = Vec::new();
let mut stopped = false;
let mut stop_reason: Option<&'static str> = None;
for (step_index, step) in steps.into_iter().enumerate() {
if stopped {
files.extend(not_attempted_for_step(&step));
continue;
}
match execute_step(step_index, step, &op_id, &mut journaled, ctx) {
StepExec::Applied(mut outcomes) => files.append(&mut outcomes),
StepExec::Stopped {
mut outcomes,
reason,
} => {
files.append(&mut outcomes);
stopped = true;
stop_reason = Some(reason);
}
}
}
let classifications: Vec<FileClassification> = files
.iter()
.filter(|file| counts_toward_completion(file))
.map(|file| file.classification)
.collect();
let registers_committed =
commit_registers_if_complete(ctx.registers, staged_registers, &classifications);
let applied = classifications
.iter()
.filter(|classification| classification.is_applied_star())
.count();
let planned_primary = classifications.len();
let success = applied > 0;
let complete = applied == planned_primary && planned_primary > 0;
let summary_text = summary_counts(&files);
TransactionEnvelope {
success,
complete,
files,
op_id: journaled.then_some(op_id),
stop_reason,
registers_committed,
preview: false,
summary_text,
}
}
pub fn run_transaction(
sections: &[TransactionSectionInput<'_>],
session_registers: &RegisterStore,
ctx: &mut ExecuteContext<'_>,
preview: bool,
) -> Result<TransactionEnvelope, HashlineRejection> {
let plan = plan_transaction(sections, session_registers, ctx.backups_enabled)?;
if preview {
Ok(preview_transaction(plan))
} else {
Ok(execute_transaction(plan, ctx))
}
}
fn split_mv<'a>(
operations: &'a [Operation],
resolved: &'a [ResolvedOperation],
) -> Result<
(
&'a [Operation],
&'a [ResolvedOperation],
Option<&'a MvOperation>,
),
HashlineRejection,
> {
if operations.len() != resolved.len() {
return Err(HashlineRejection::parse(
"resolved operation count does not match the parsed section",
));
}
if let Some(Operation::Mv(mv)) = operations.last() {
let line_len = operations.len() - 1;
if operations[..line_len]
.iter()
.any(|operation| matches!(operation, Operation::Mv(_)))
{
return Err(HashlineRejection::parse(
"MV must occur once and after all line operations",
));
}
return Ok((&operations[..line_len], &resolved[..line_len], Some(mv)));
}
if operations
.iter()
.any(|operation| matches!(operation, Operation::Mv(_)))
{
return Err(HashlineRejection::parse(
"MV must occur once and after all line operations",
));
}
Ok((operations, resolved, None))
}
fn plan_line_section(
section: &SectionPlanInput<'_>,
staged: &mut StagedRegisters,
) -> Result<PlannedFile, HashlineRejection> {
for resolved in section.resolved {
match crate::hashline::syntax::verify_exact(
section.snapshot,
section.baseline,
resolved.address,
) {
crate::hashline::syntax::VerificationOutcome::Exact => {}
crate::hashline::syntax::VerificationOutcome::RecoveryRequired(_) => {
return Err(HashlineRejection::new(
crate::hashline::syntax::HashlineRejectionCode::StaleTag,
crate::hashline::syntax::RejectionStage::Recovery,
"addressed content no longer matches the Phase-1 baseline",
));
}
crate::hashline::syntax::VerificationOutcome::Rejected(rejection) => {
return Err(rejection)
}
crate::hashline::syntax::VerificationOutcome::BlockNeedsResolution { .. } => {
return Err(HashlineRejection::new(
crate::hashline::syntax::HashlineRejectionCode::BoundaryIneligible,
crate::hashline::syntax::RejectionStage::Eligibility,
"block address was not expanded before transaction planning",
));
}
}
}
apply_section_ops(
section.requested_path,
section.canonical_path,
section.baseline,
section.operations,
section.resolved,
staged,
)
}
fn assert_rollback_available(
backups_enabled: bool,
step: &PlannedStep,
) -> Result<(), HashlineRejection> {
if backups_enabled {
return Ok(());
}
match step {
PlannedStep::Mutate(_) => Err(HashlineRejection::backup_unavailable(
"backups are disabled; refusing destructive hashline mutation without a restore record",
)),
PlannedStep::Mv(mv) if mv.dest_existed => Err(HashlineRejection::backup_unavailable(
"backups are disabled; refusing MV onto an existing destination without a restore record",
)),
PlannedStep::Mv(_) => Ok(()),
}
}
fn preview_step_files(step: &PlannedStep) -> Vec<FileOutcome> {
match step {
PlannedStep::Mutate(file) => {
vec![FileOutcome {
canonical_path: file.canonical_path.clone(),
requested_path: file.requested_path.clone(),
role: FileRole::Primary,
classification: FileClassification::Applied,
mutation_state: MutationState::Unmutated,
final_bytes: Some(file.final_bytes.clone()),
final_tag: None,
affected: file.affected.clone(),
warnings: file.warnings.clone(),
format_skipped_reason: None,
backup_id: None,
remove_file: file.remove_file,
tag_notice: Some("preview: no final tag or undo identity".into()),
}]
}
PlannedStep::Mv(mv) => vec![
FileOutcome {
canonical_path: mv.dest_canonical.clone(),
requested_path: mv.dest_requested.clone(),
role: FileRole::MvDestination,
classification: FileClassification::Applied,
mutation_state: MutationState::Unmutated,
final_bytes: Some(mv.final_bytes.clone()),
final_tag: None,
affected: mv.affected.clone(),
warnings: mv.warnings.clone(),
format_skipped_reason: None,
backup_id: None,
remove_file: false,
tag_notice: Some("preview: no final tag or undo identity".into()),
},
FileOutcome {
canonical_path: mv.source_canonical.clone(),
requested_path: mv.source_requested.clone(),
role: FileRole::MvSource,
classification: FileClassification::Applied,
mutation_state: MutationState::Unmutated,
final_bytes: None,
final_tag: None,
affected: AffectedRegion::default(),
warnings: Vec::new(),
format_skipped_reason: None,
backup_id: None,
remove_file: true,
tag_notice: Some("preview: source removal not performed".into()),
},
],
}
}
enum StepExec {
Applied(Vec<FileOutcome>),
Stopped {
outcomes: Vec<FileOutcome>,
reason: &'static str,
},
}
fn execute_step(
step_index: usize,
step: PlannedStep,
op_id: &str,
journaled: &mut bool,
ctx: &mut ExecuteContext<'_>,
) -> StepExec {
match step {
PlannedStep::Mutate(file) => execute_mutate(step_index, file, op_id, journaled, ctx),
PlannedStep::Mv(mv) => execute_mv(step_index, mv, op_id, journaled, ctx),
}
}
fn execute_mutate(
step_index: usize,
file: PlannedFile,
op_id: &str,
journaled: &mut bool,
ctx: &mut ExecuteContext<'_>,
) -> StepExec {
if fault_is(ctx, ExecuteFault::Backup { step: step_index }) {
return StepExec::Stopped {
outcomes: vec![failed_outcome(
&file.canonical_path,
&file.requested_path,
FileRole::Primary,
FileClassification::FailedBackup,
file.remove_file,
file.warnings.clone(),
)],
reason: "failed_backup",
};
}
let backup_id = match journal_existing_or_skip(
ctx,
op_id,
&file.canonical_path,
file.remove_file,
"hashline: pre-mutation backup",
) {
Ok(id) => {
if id.is_some() {
*journaled = true;
}
id
}
Err(_) => {
return StepExec::Stopped {
outcomes: vec![failed_outcome(
&file.canonical_path,
&file.requested_path,
FileRole::Primary,
FileClassification::FailedBackup,
file.remove_file,
file.warnings.clone(),
)],
reason: "failed_backup",
};
}
};
if backup_id.is_none() && path_exists(&file.canonical_path) {
return StepExec::Stopped {
outcomes: vec![failed_outcome(
&file.canonical_path,
&file.requested_path,
FileRole::Primary,
FileClassification::FailedBackup,
file.remove_file,
file.warnings.clone(),
)],
reason: "failed_backup",
};
}
if fault_is(ctx, ExecuteFault::BaselineDrift { step: step_index })
|| !baseline_matches(&file.canonical_path, &file.baseline_bytes)
{
return StepExec::Stopped {
outcomes: vec![failed_outcome(
&file.canonical_path,
&file.requested_path,
FileRole::Primary,
FileClassification::FailedBaselineDrift,
file.remove_file,
file.warnings.clone(),
)],
reason: "hashline_baseline_drift",
};
}
if fault_is(ctx, ExecuteFault::Write { step: step_index }) {
return StepExec::Stopped {
outcomes: vec![failed_outcome(
&file.canonical_path,
&file.requested_path,
FileRole::Primary,
FileClassification::FailedWrite,
file.remove_file,
file.warnings.clone(),
)],
reason: "failed_write",
};
}
if file.remove_file {
if let Err(error) = fs::remove_file(&file.canonical_path) {
if error.kind() != io::ErrorKind::NotFound {
return StepExec::Stopped {
outcomes: vec![failed_outcome(
&file.canonical_path,
&file.requested_path,
FileRole::Primary,
FileClassification::FailedWrite,
true,
file.warnings.clone(),
)],
reason: "failed_write",
};
}
}
invalidate_removed_source(ctx.snapshots, &file.canonical_path);
return StepExec::Applied(vec![FileOutcome {
canonical_path: file.canonical_path,
requested_path: file.requested_path,
role: FileRole::Primary,
classification: FileClassification::Applied,
mutation_state: MutationState::Applied,
final_bytes: None,
final_tag: None,
affected: AffectedRegion::default(),
warnings: file.warnings,
format_skipped_reason: None,
backup_id,
remove_file: true,
tag_notice: Some("source path removed; no final tag".into()),
}]);
}
if let Err(error) = durable_write(&file.canonical_path, &file.final_bytes) {
let classification = if error.to_string().contains("durability") {
FileClassification::FailedDurability
} else {
FileClassification::FailedWrite
};
return StepExec::Stopped {
outcomes: vec![failed_outcome(
&file.canonical_path,
&file.requested_path,
FileRole::Primary,
classification,
false,
file.warnings.clone(),
)],
reason: classification.as_str(),
};
}
if fault_is(ctx, ExecuteFault::Durability { step: step_index }) {
return StepExec::Stopped {
outcomes: vec![failed_outcome(
&file.canonical_path,
&file.requested_path,
FileRole::Primary,
FileClassification::FailedDurability,
false,
file.warnings.clone(),
)],
reason: "failed_durability",
};
}
let on_disk = match fs::read(&file.canonical_path) {
Ok(bytes) => bytes,
Err(_) => {
return StepExec::Applied(vec![FileOutcome {
canonical_path: file.canonical_path,
requested_path: file.requested_path,
role: FileRole::Primary,
classification: FileClassification::AppliedTagUnavailable,
mutation_state: MutationState::Applied,
final_bytes: Some(file.final_bytes),
final_tag: None,
affected: file.affected,
warnings: file.warnings,
format_skipped_reason: None,
backup_id,
remove_file: false,
tag_notice: Some("final bytes could not be re-read for tagging".into()),
}]);
}
};
let mut classification = FileClassification::Applied;
if fault_is(ctx, ExecuteFault::ValidationFailure { step: step_index }) {
classification = FileClassification::AppliedWithValidationFailure;
}
let (final_tag, tag_notice, classification) =
if fault_is(ctx, ExecuteFault::FinalTagUnavailable { step: step_index }) {
(
None,
Some("final tag unavailable; re-read before chaining".into()),
FileClassification::AppliedTagUnavailable,
)
} else {
let published = publish_edit_response_snapshot(
ctx.snapshots,
&file.canonical_path,
file.requested_path.clone(),
&on_disk,
&file.affected,
);
tag_from_publish(published, classification)
};
StepExec::Applied(vec![FileOutcome {
canonical_path: file.canonical_path,
requested_path: file.requested_path,
role: FileRole::Primary,
classification,
mutation_state: classification.mutation_state(),
final_bytes: Some(on_disk),
final_tag,
affected: file.affected,
warnings: file.warnings,
format_skipped_reason: None,
backup_id,
remove_file: false,
tag_notice,
}])
}
fn execute_mv(
step_index: usize,
mv: PlannedMv,
op_id: &str,
journaled: &mut bool,
ctx: &mut ExecuteContext<'_>,
) -> StepExec {
if fault_is(ctx, ExecuteFault::Backup { step: step_index }) {
return StepExec::Stopped {
outcomes: mv_failed_pair(
&mv,
FileClassification::FailedBackup,
FileClassification::NotAttempted,
),
reason: "failed_backup",
};
}
let dest_backup_id = if mv.dest_existed {
match ctx.backups.snapshot_with_op(
ctx.session,
&mv.dest_canonical,
"hashline: MV destination backup",
Some(op_id),
) {
Ok(Some(id)) => {
*journaled = true;
Some(id)
}
Ok(None) => {
return StepExec::Stopped {
outcomes: mv_failed_pair(
&mv,
FileClassification::FailedBackup,
FileClassification::NotAttempted,
),
reason: "failed_backup",
};
}
Err(_) => {
return StepExec::Stopped {
outcomes: mv_failed_pair(
&mv,
FileClassification::FailedBackup,
FileClassification::NotAttempted,
),
reason: "failed_backup",
};
}
}
} else {
match ctx.backups.snapshot_op_tombstone(
ctx.session,
op_id,
&mv.dest_canonical,
"hashline: MV created destination",
) {
Ok(Some(id)) => {
*journaled = true;
Some(id)
}
Ok(None) => {
return StepExec::Stopped {
outcomes: mv_failed_pair(
&mv,
FileClassification::FailedBackup,
FileClassification::NotAttempted,
),
reason: "failed_backup",
};
}
Err(_) => {
return StepExec::Stopped {
outcomes: mv_failed_pair(
&mv,
FileClassification::FailedBackup,
FileClassification::NotAttempted,
),
reason: "failed_backup",
};
}
}
};
let source_backup_id = match ctx.backups.snapshot_with_op(
ctx.session,
&mv.source_canonical,
"hashline: MV source backup",
Some(op_id),
) {
Ok(Some(id)) => {
*journaled = true;
Some(id)
}
Ok(None) | Err(_) => {
return StepExec::Stopped {
outcomes: mv_failed_pair(
&mv,
FileClassification::FailedBackup,
FileClassification::NotAttempted,
),
reason: "failed_backup",
};
}
};
if fault_is(ctx, ExecuteFault::BaselineDrift { step: step_index })
|| !baseline_matches(&mv.source_canonical, &mv.source_baseline_bytes)
|| mv
.dest_baseline_bytes
.as_ref()
.is_some_and(|expected| !baseline_matches(&mv.dest_canonical, expected))
{
return StepExec::Stopped {
outcomes: mv_failed_pair(
&mv,
FileClassification::FailedBaselineDrift,
FileClassification::NotAttempted,
),
reason: "hashline_baseline_drift",
};
}
if fault_is(ctx, ExecuteFault::Write { step: step_index }) {
return StepExec::Stopped {
outcomes: mv_failed_pair(
&mv,
FileClassification::FailedWrite,
FileClassification::NotAttempted,
),
reason: "failed_write",
};
}
if let Err(error) = ensure_parent_dirs(&mv.dest_canonical)
.and_then(|_| durable_write(&mv.dest_canonical, &mv.final_bytes))
{
let classification = if error.to_string().contains("durability") {
FileClassification::FailedDurability
} else {
FileClassification::FailedWrite
};
return StepExec::Stopped {
outcomes: mv_failed_pair(&mv, classification, FileClassification::NotAttempted),
reason: classification.as_str(),
};
}
if fault_is(ctx, ExecuteFault::Durability { step: step_index }) {
return StepExec::Stopped {
outcomes: mv_failed_pair(
&mv,
FileClassification::FailedDurability,
FileClassification::NotAttempted,
),
reason: "failed_durability",
};
}
let dest_on_disk = fs::read(&mv.dest_canonical).unwrap_or_else(|_| mv.final_bytes.clone());
if fault_is(ctx, ExecuteFault::SourceUnlink { step: step_index })
|| fs::remove_file(&mv.source_canonical).is_err()
{
let (final_tag, tag_notice, dest_class) =
observe_dest_tag(ctx, &mv, &dest_on_disk, step_index);
return StepExec::Stopped {
outcomes: vec![
FileOutcome {
canonical_path: mv.dest_canonical,
requested_path: mv.dest_requested,
role: FileRole::MvDestination,
classification: dest_class,
mutation_state: dest_class.mutation_state(),
final_bytes: Some(dest_on_disk),
final_tag,
affected: mv.affected,
warnings: mv.warnings,
format_skipped_reason: None,
backup_id: dest_backup_id,
remove_file: false,
tag_notice,
},
FileOutcome {
canonical_path: mv.source_canonical,
requested_path: mv.source_requested,
role: FileRole::MvSource,
classification: FileClassification::FailedSourceUnlink,
mutation_state: MutationState::PartialMv,
final_bytes: Some(mv.source_baseline_bytes),
final_tag: None,
affected: AffectedRegion::default(),
warnings: Vec::new(),
format_skipped_reason: None,
backup_id: source_backup_id,
remove_file: false,
tag_notice: Some(
"destination written; source unlink failed — partial MV under shared op_id"
.into(),
),
},
],
reason: "failed_source_unlink",
};
}
invalidate_removed_source(ctx.snapshots, &mv.source_canonical);
let (final_tag, tag_notice, dest_class) = observe_dest_tag(ctx, &mv, &dest_on_disk, step_index);
StepExec::Applied(vec![
FileOutcome {
canonical_path: mv.dest_canonical,
requested_path: mv.dest_requested,
role: FileRole::MvDestination,
classification: dest_class,
mutation_state: dest_class.mutation_state(),
final_bytes: Some(dest_on_disk),
final_tag,
affected: mv.affected,
warnings: mv.warnings,
format_skipped_reason: None,
backup_id: dest_backup_id,
remove_file: false,
tag_notice,
},
FileOutcome {
canonical_path: mv.source_canonical,
requested_path: mv.source_requested,
role: FileRole::MvSource,
classification: FileClassification::Applied,
mutation_state: MutationState::Applied,
final_bytes: None,
final_tag: None,
affected: AffectedRegion::default(),
warnings: Vec::new(),
format_skipped_reason: None,
backup_id: source_backup_id,
remove_file: true,
tag_notice: Some("source path removed; no final tag".into()),
},
])
}
fn observe_dest_tag(
ctx: &mut ExecuteContext<'_>,
mv: &PlannedMv,
dest_on_disk: &[u8],
step_index: usize,
) -> (Option<String>, Option<String>, FileClassification) {
if fault_is(ctx, ExecuteFault::FinalTagUnavailable { step: step_index }) {
return (
None,
Some("final tag unavailable; re-read before chaining".into()),
FileClassification::AppliedTagUnavailable,
);
}
let mut classification = FileClassification::Applied;
if fault_is(ctx, ExecuteFault::ValidationFailure { step: step_index }) {
classification = FileClassification::AppliedWithValidationFailure;
}
let published = publish_edit_response_snapshot(
ctx.snapshots,
&mv.dest_canonical,
mv.dest_requested.clone(),
dest_on_disk,
&mv.affected,
);
tag_from_publish(published, classification)
}
fn tag_from_publish(
published: EditResponseSnapshot,
classification: FileClassification,
) -> (Option<String>, Option<String>, FileClassification) {
if let Some(snapshot) = published.snapshot {
(Some(snapshot.tag.clone()), published.notice, classification)
} else {
(
None,
published
.notice
.or_else(|| Some("final tag unavailable; re-read before chaining".into())),
FileClassification::AppliedTagUnavailable,
)
}
}
fn journal_existing_or_skip(
ctx: &mut ExecuteContext<'_>,
op_id: &str,
path: &Path,
_remove_file: bool,
description: &str,
) -> Result<Option<String>, ()> {
if !path_exists(path) {
return Ok(None);
}
match ctx
.backups
.snapshot_with_op(ctx.session, path, description, Some(op_id))
{
Ok(id) => Ok(id),
Err(_) => Err(()),
}
}
fn path_exists(path: &Path) -> bool {
fs::symlink_metadata(path).is_ok()
}
fn baseline_matches(path: &Path, expected: &[u8]) -> bool {
match fs::read(path) {
Ok(bytes) => bytes == expected,
Err(error) if error.kind() == io::ErrorKind::NotFound => expected.is_empty(),
Err(_) => false,
}
}
fn ensure_parent_dirs(path: &Path) -> io::Result<()> {
if let Some(parent) = path.parent() {
if !parent.as_os_str().is_empty() {
fs::create_dir_all(parent)?;
}
}
Ok(())
}
fn durable_write(path: &Path, bytes: &[u8]) -> io::Result<()> {
ensure_parent_dirs(path)?;
let parent = path.parent().unwrap_or_else(|| Path::new("."));
let file_name = path
.file_name()
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "path has no file name"))?;
let temp_name = {
let mut name = std::ffi::OsString::from(".aft-hashline-");
name.push(file_name);
name.push(".tmp");
name
};
let temp_path = parent.join(temp_name);
let write_result = (|| {
let mut file = OpenOptions::new()
.write(true)
.create(true)
.truncate(true)
.open(&temp_path)?;
file.write_all(bytes)?;
file.sync_all()
.map_err(|error| io::Error::new(error.kind(), format!("durability: {error}")))?;
fs::rename(&temp_path, path)?;
if let Ok(dir) = File::open(parent) {
let _ = dir.sync_all();
}
Ok(())
})();
if write_result.is_err() {
let _ = fs::remove_file(&temp_path);
}
write_result
}
fn failed_outcome(
path: &Path,
requested: &str,
role: FileRole,
classification: FileClassification,
remove_file: bool,
warnings: Vec<String>,
) -> FileOutcome {
FileOutcome {
canonical_path: path.to_path_buf(),
requested_path: requested.to_string(),
role,
classification,
mutation_state: classification.mutation_state(),
final_bytes: None,
final_tag: None,
affected: AffectedRegion::default(),
warnings,
format_skipped_reason: None,
backup_id: None,
remove_file,
tag_notice: None,
}
}
fn mv_failed_pair(
mv: &PlannedMv,
dest_class: FileClassification,
source_class: FileClassification,
) -> Vec<FileOutcome> {
vec![
failed_outcome(
&mv.dest_canonical,
&mv.dest_requested,
FileRole::MvDestination,
dest_class,
false,
mv.warnings.clone(),
),
failed_outcome(
&mv.source_canonical,
&mv.source_requested,
FileRole::MvSource,
source_class,
false,
Vec::new(),
),
]
}
fn not_attempted_for_step(step: &PlannedStep) -> Vec<FileOutcome> {
match step {
PlannedStep::Mutate(file) => vec![failed_outcome(
&file.canonical_path,
&file.requested_path,
FileRole::Primary,
FileClassification::NotAttempted,
file.remove_file,
Vec::new(),
)],
PlannedStep::Mv(mv) => mv_failed_pair(
mv,
FileClassification::NotAttempted,
FileClassification::NotAttempted,
),
}
}
fn fault_is(ctx: &ExecuteContext<'_>, want: ExecuteFault) -> bool {
ctx.fault.as_ref() == Some(&want)
}
fn counts_toward_completion(file: &FileOutcome) -> bool {
!(file.role == FileRole::MvSource && file.classification.is_applied_star())
}
fn summary_counts(files: &[FileOutcome]) -> String {
let primary: Vec<_> = files
.iter()
.filter(|file| counts_toward_completion(file))
.collect();
let applied = primary
.iter()
.filter(|file| file.classification.is_applied_star())
.count();
let total = primary.len();
format!("{applied} of {total} files applied")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::backup::BackupPolicy;
use crate::hashline::scan::scan_bytes;
use crate::hashline::snapshot::{capture_taggable_read, ReadPublication, ReadSelection};
use crate::hashline::syntax::{
parse_address, resolve_address, resolve_snapshot, PutOperation, PutSource, RegisterRef,
ResolvedAddress,
};
const SESSION: &str = "hashline-tx-test";
fn whole_snapshot(bytes: &[u8]) -> Snapshot {
scan_bytes(bytes)
}
fn put_text(address: &str, body: &[&str]) -> Operation {
Operation::Put(PutOperation {
address: parse_address(address).unwrap(),
source: PutSource::Text(body.iter().map(|line| (*line).to_string()).collect()),
line: 1,
})
}
fn resolve_one(snapshot: &Snapshot, operation: &Operation) -> ResolvedOperation {
let address = match operation.address() {
Some(address) => resolve_address(address, snapshot).unwrap(),
None => ResolvedAddress::WholeFile,
};
ResolvedOperation {
operation_index: 0,
address,
}
}
fn write_file(path: &Path, bytes: &[u8]) {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent).unwrap();
}
fs::write(path, bytes).unwrap();
}
fn backup_store(dir: &Path) -> BackupStore {
let mut store = BackupStore::new();
store.set_storage_dir(dir.to_path_buf(), 72);
store
}
fn ctx<'a>(
backups: &'a mut BackupStore,
snapshots: &'a mut SnapshotStore,
registers: &'a mut RegisterStore,
backups_enabled: bool,
fault: Option<ExecuteFault>,
) -> ExecuteContext<'a> {
ExecuteContext {
session: SESSION,
backups,
snapshots,
registers,
backups_enabled,
fault,
}
}
fn section_put<'a>(
path: &'a Path,
requested: &'a str,
baseline: &'a Baseline,
snapshot: &'a Snapshot,
ops: &'a [Operation],
resolved: &'a [ResolvedOperation],
) -> TransactionSectionInput<'a> {
TransactionSectionInput {
canonical_path: path,
requested_path: requested,
baseline,
snapshot,
operations: ops,
resolved,
mv_destination: None,
}
}
fn put_after_reads(
selections: impl IntoIterator<Item = ReadSelection>,
) -> Result<Vec<u8>, HashlineRejection> {
let temp = tempfile::tempdir().unwrap();
let path = temp.path().join("reread.txt");
let original = b"one\ntwo\nthree\nfour\n";
write_file(&path, original);
let mut snapshots = SnapshotStore::new();
let mut tag = None;
for selection in selections {
let publication =
capture_taggable_read(&mut snapshots, &path, "reread.txt", selection).unwrap();
let ReadPublication::Tagged { snapshot, .. } = publication else {
panic!("fixture read must publish a tagged snapshot");
};
tag.get_or_insert(snapshot.tag);
}
let snapshot = resolve_snapshot(
&mut snapshots,
&path,
tag.as_deref().expect("at least one read selection"),
)?;
let baseline = Baseline::from_bytes(original.to_vec());
let operations = vec![put_text("2", &["TWO"])];
let resolved = vec![resolve_one(&snapshot, &operations[0])];
let sections = [section_put(
&path,
"reread.txt",
&baseline,
&snapshot,
&operations,
&resolved,
)];
let session_registers = RegisterStore::new();
let plan = plan_transaction(§ions, &session_registers, true)?;
let mut backups = backup_store(&temp.path().join("backups"));
let mut execution_registers = RegisterStore::new();
let mut execution = ctx(
&mut backups,
&mut snapshots,
&mut execution_registers,
true,
None,
);
let envelope = execute_transaction(plan, &mut execution);
assert!(envelope.success);
assert!(envelope.complete);
Ok(fs::read(path).unwrap())
}
#[test]
fn two_ranged_reads_of_one_version_then_put_applies() {
let bytes = put_after_reads([ReadSelection::range(1, 2), ReadSelection::range(3, 4)])
.expect("same-version ranged reads must resolve");
assert_eq!(bytes, b"one\nTWO\nthree\nfour\n");
}
#[test]
fn ranged_then_whole_read_of_one_version_then_put_applies() {
let bytes = put_after_reads([ReadSelection::range(2, 2), ReadSelection::WholeFile])
.expect("same-version ranged and whole reads must resolve");
assert_eq!(bytes, b"one\nTWO\nthree\nfour\n");
}
#[test]
fn second_read_without_intervening_mutation_does_not_enter_refusal_loop() {
let bytes = put_after_reads([ReadSelection::range(1, 1), ReadSelection::range(2, 2)])
.expect("a second read of unchanged content must leave the tag editable");
assert_eq!(bytes, b"one\nTWO\nthree\nfour\n");
}
#[test]
fn a8_phase1_mutation_free_and_phase2_ordered() {
let temp = tempfile::tempdir().unwrap();
let a = temp.path().join("a.txt");
let b = temp.path().join("b.txt");
write_file(&a, b"alpha\n");
write_file(&b, b"beta\n");
let bytes_a = fs::read(&a).unwrap();
let bytes_b = fs::read(&b).unwrap();
let snap_a = whole_snapshot(&bytes_a);
let snap_b = whole_snapshot(&bytes_b);
let base_a = Baseline::from_bytes(bytes_a.clone());
let base_b = Baseline::from_bytes(bytes_b.clone());
let ops_a = vec![put_text("1", &["ALPHA"])];
let ops_b = vec![put_text("1", &["BETA"])];
let res_a = vec![resolve_one(&snap_a, &ops_a[0])];
let res_b = vec![resolve_one(&snap_b, &ops_b[0])];
let sections = [
section_put(&a, "a.txt", &base_a, &snap_a, &ops_a, &res_a),
section_put(&b, "b.txt", &base_b, &snap_b, &ops_b, &res_b),
];
let registers = RegisterStore::new();
let plan = plan_transaction(§ions, ®isters, true).expect("phase1");
assert_eq!(fs::read(&a).unwrap(), b"alpha\n");
assert_eq!(fs::read(&b).unwrap(), b"beta\n");
assert_eq!(plan.steps.len(), 2);
let backup_dir = temp.path().join("backups");
let mut backups = backup_store(&backup_dir);
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
let mut exec = ctx(&mut backups, &mut snapshots, &mut session_regs, true, None);
let envelope = execute_transaction(plan, &mut exec);
assert!(envelope.success);
assert!(envelope.complete);
assert!(envelope.op_id.is_some());
assert_eq!(envelope.files.len(), 2);
assert_eq!(envelope.files[0].requested_path, "a.txt");
assert_eq!(envelope.files[1].requested_path, "b.txt");
assert_eq!(
envelope.files[0].classification,
FileClassification::Applied
);
assert_eq!(envelope.files[0].mutation_state, MutationState::Applied);
assert_eq!(fs::read(&a).unwrap(), b"ALPHA\n");
assert_eq!(fs::read(&b).unwrap(), b"BETA\n");
assert!(envelope.summary_text.contains("2 of 2 files applied"));
let op_id = envelope.op_id.clone().unwrap();
let restored = backups.restore_last_operation(SESSION).unwrap();
assert_eq!(restored.op_id, op_id);
assert_eq!(fs::read(&a).unwrap(), b"alpha\n");
assert_eq!(fs::read(&b).unwrap(), b"beta\n");
}
#[test]
fn a8_all_failed_emits_success_false_with_envelope() {
let temp = tempfile::tempdir().unwrap();
let a = temp.path().join("a.txt");
let b = temp.path().join("b.txt");
write_file(&a, b"a\n");
write_file(&b, b"b\n");
let bytes_a = fs::read(&a).unwrap();
let bytes_b = fs::read(&b).unwrap();
let snap_a = whole_snapshot(&bytes_a);
let snap_b = whole_snapshot(&bytes_b);
let base_a = Baseline::from_bytes(bytes_a);
let base_b = Baseline::from_bytes(bytes_b);
let ops_a = vec![put_text("1", &["A"])];
let ops_b = vec![put_text("1", &["B"])];
let res_a = vec![resolve_one(&snap_a, &ops_a[0])];
let res_b = vec![resolve_one(&snap_b, &ops_b[0])];
let sections = [
section_put(&a, "a.txt", &base_a, &snap_a, &ops_a, &res_a),
section_put(&b, "b.txt", &base_b, &snap_b, &ops_b, &res_b),
];
let registers = RegisterStore::new();
let plan = plan_transaction(§ions, ®isters, true).unwrap();
let mut backups = backup_store(&temp.path().join("backups"));
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
let mut exec = ctx(
&mut backups,
&mut snapshots,
&mut session_regs,
true,
Some(ExecuteFault::BaselineDrift { step: 0 }),
);
let envelope = execute_transaction(plan, &mut exec);
assert!(!envelope.success);
assert!(!envelope.complete);
assert_eq!(
envelope.files[0].classification,
FileClassification::FailedBaselineDrift
);
assert_eq!(envelope.files[0].mutation_state, MutationState::Unmutated);
assert_eq!(
envelope.files[1].classification,
FileClassification::NotAttempted
);
assert_eq!(envelope.files[1].mutation_state, MutationState::Unmutated);
assert_eq!(envelope.stop_reason, Some("hashline_baseline_drift"));
assert!(envelope.summary_text.starts_with("0 of 2 files applied"));
assert!(envelope.op_id.is_some());
assert_eq!(fs::read(&a).unwrap(), b"a\n");
assert_eq!(fs::read(&b).unwrap(), b"b\n");
}
#[test]
fn a8_partial_failure_keeps_prior_under_shared_op_id() {
let temp = tempfile::tempdir().unwrap();
let a = temp.path().join("a.txt");
let b = temp.path().join("b.txt");
write_file(&a, b"a\n");
write_file(&b, b"b\n");
let bytes_a = fs::read(&a).unwrap();
let bytes_b = fs::read(&b).unwrap();
let snap_a = whole_snapshot(&bytes_a);
let snap_b = whole_snapshot(&bytes_b);
let base_a = Baseline::from_bytes(bytes_a);
let base_b = Baseline::from_bytes(bytes_b);
let ops_a = vec![put_text("1", &["A"])];
let ops_b = vec![put_text("1", &["B"])];
let res_a = vec![resolve_one(&snap_a, &ops_a[0])];
let res_b = vec![resolve_one(&snap_b, &ops_b[0])];
let sections = [
section_put(&a, "a.txt", &base_a, &snap_a, &ops_a, &res_a),
section_put(&b, "b.txt", &base_b, &snap_b, &ops_b, &res_b),
];
let registers = RegisterStore::new();
let plan = plan_transaction(§ions, ®isters, true).unwrap();
let mut backups = backup_store(&temp.path().join("backups"));
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
let mut exec = ctx(
&mut backups,
&mut snapshots,
&mut session_regs,
true,
Some(ExecuteFault::Write { step: 1 }),
);
let envelope = execute_transaction(plan, &mut exec);
assert!(envelope.success);
assert!(!envelope.complete);
assert!(envelope.op_id.is_some());
assert_eq!(
envelope.files[0].classification,
FileClassification::Applied
);
assert_eq!(
envelope.files[1].classification,
FileClassification::FailedWrite
);
assert_eq!(
envelope.files[1].mutation_state,
MutationState::UnknownPossiblyMutated
);
assert_eq!(fs::read(&a).unwrap(), b"A\n");
assert_eq!(fs::read(&b).unwrap(), b"b\n");
let op_id = envelope.op_id.unwrap();
let restored = backups.restore_last_operation(SESSION).unwrap();
assert_eq!(restored.op_id, op_id);
assert_eq!(fs::read(&a).unwrap(), b"a\n");
}
#[test]
fn a8_mv_new_and_existing_destination_with_undo() {
let temp = tempfile::tempdir().unwrap();
let src = temp.path().join("src.txt");
let new_dest = temp.path().join("new_dest.txt");
write_file(&src, b"move-me\n");
let bytes = fs::read(&src).unwrap();
let snap = whole_snapshot(&bytes);
let base = Baseline::from_bytes(bytes.clone());
let ops = vec![Operation::Mv(MvOperation {
destination: "new_dest.txt".into(),
line: 1,
})];
let resolved = vec![ResolvedOperation {
operation_index: 0,
address: ResolvedAddress::WholeFile,
}];
let sections = [TransactionSectionInput {
canonical_path: &src,
requested_path: "src.txt",
baseline: &base,
snapshot: &snap,
operations: &ops,
resolved: &resolved,
mv_destination: Some(MvDestinationInput {
canonical_path: &new_dest,
requested_path: "new_dest.txt",
baseline_bytes: None,
}),
}];
let registers = RegisterStore::new();
let plan = plan_transaction(§ions, ®isters, true).unwrap();
let mut backups = backup_store(&temp.path().join("backups"));
let mut snapshots = SnapshotStore::new();
snapshots.publish(&src, snap.clone());
let mut session_regs = RegisterStore::new();
let mut exec = ctx(&mut backups, &mut snapshots, &mut session_regs, true, None);
let envelope = execute_transaction(plan, &mut exec);
assert!(envelope.success && envelope.complete);
assert!(envelope.op_id.is_some());
assert_eq!(envelope.files[0].role, FileRole::MvDestination);
assert_eq!(envelope.files[1].role, FileRole::MvSource);
assert!(envelope.files[0].final_tag.is_some());
assert!(envelope.files[1].remove_file);
assert_eq!(fs::read(&new_dest).unwrap(), b"move-me\n");
assert!(!src.exists());
assert!(snapshots.lookup(&src, &snap.tag).is_err());
let op_id = envelope.op_id.unwrap();
let restored = backups.restore_last_operation(SESSION).unwrap();
assert_eq!(restored.op_id, op_id);
assert_eq!(fs::read(&src).unwrap(), b"move-me\n");
assert!(!new_dest.exists(), "created destination removed on undo");
let src2 = temp.path().join("src2.txt");
let dest2 = temp.path().join("dest2.txt");
write_file(&src2, b"from\n");
write_file(&dest2, b"old-dest\n");
let bytes2 = fs::read(&src2).unwrap();
let snap2 = whole_snapshot(&bytes2);
let base2 = Baseline::from_bytes(bytes2);
let dest_bytes = fs::read(&dest2).unwrap();
let ops2 = vec![Operation::Mv(MvOperation {
destination: "dest2.txt".into(),
line: 1,
})];
let resolved2 = vec![ResolvedOperation {
operation_index: 0,
address: ResolvedAddress::WholeFile,
}];
let sections2 = [TransactionSectionInput {
canonical_path: &src2,
requested_path: "src2.txt",
baseline: &base2,
snapshot: &snap2,
operations: &ops2,
resolved: &resolved2,
mv_destination: Some(MvDestinationInput {
canonical_path: &dest2,
requested_path: "dest2.txt",
baseline_bytes: Some(&dest_bytes),
}),
}];
let plan2 = plan_transaction(§ions2, ®isters, true).unwrap();
let mut exec2 = ctx(&mut backups, &mut snapshots, &mut session_regs, true, None);
let envelope2 = execute_transaction(plan2, &mut exec2);
assert!(envelope2.success);
assert_eq!(fs::read(&dest2).unwrap(), b"from\n");
assert!(!src2.exists());
let restored2 = backups.restore_last_operation(SESSION).unwrap();
assert_eq!(restored2.op_id, envelope2.op_id.unwrap());
assert_eq!(fs::read(&src2).unwrap(), b"from\n");
assert_eq!(fs::read(&dest2).unwrap(), b"old-dest\n");
}
#[test]
fn a8_mv_source_unlink_failure_is_partial_mv() {
let temp = tempfile::tempdir().unwrap();
let src = temp.path().join("src.txt");
let dest = temp.path().join("dest.txt");
write_file(&src, b"body\n");
let bytes = fs::read(&src).unwrap();
let snap = whole_snapshot(&bytes);
let base = Baseline::from_bytes(bytes);
let ops = vec![Operation::Mv(MvOperation {
destination: "dest.txt".into(),
line: 1,
})];
let resolved = vec![ResolvedOperation {
operation_index: 0,
address: ResolvedAddress::WholeFile,
}];
let sections = [TransactionSectionInput {
canonical_path: &src,
requested_path: "src.txt",
baseline: &base,
snapshot: &snap,
operations: &ops,
resolved: &resolved,
mv_destination: Some(MvDestinationInput {
canonical_path: &dest,
requested_path: "dest.txt",
baseline_bytes: None,
}),
}];
let registers = RegisterStore::new();
let plan = plan_transaction(§ions, ®isters, true).unwrap();
let mut backups = backup_store(&temp.path().join("backups"));
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
let mut exec = ctx(
&mut backups,
&mut snapshots,
&mut session_regs,
true,
Some(ExecuteFault::SourceUnlink { step: 0 }),
);
let envelope = execute_transaction(plan, &mut exec);
assert!(envelope.success);
assert!(!envelope.complete);
assert!(envelope.op_id.is_some());
assert_eq!(
envelope.files[0].classification,
FileClassification::Applied
);
assert_eq!(
envelope.files[1].classification,
FileClassification::FailedSourceUnlink
);
assert_eq!(envelope.files[1].mutation_state, MutationState::PartialMv);
assert_eq!(fs::read(&dest).unwrap(), b"body\n");
assert!(src.exists(), "source remains after unlink failure");
}
#[test]
fn a8_register_commit_only_when_all_applied() {
let temp = tempfile::tempdir().unwrap();
let a = temp.path().join("a.txt");
write_file(&a, b"one\ntwo\n");
let bytes = fs::read(&a).unwrap();
let snap = whole_snapshot(&bytes);
let base = Baseline::from_bytes(bytes);
let ops = vec![Operation::Cut(crate::hashline::syntax::CutOperation {
address: parse_address("1").unwrap(),
register: Some(RegisterRef::Named("clip".into())),
line: 1,
})];
let resolved = vec![resolve_one(&snap, &ops[0])];
let sections = [section_put(&a, "a.txt", &base, &snap, &ops, &resolved)];
let registers = RegisterStore::new();
let plan = plan_transaction(§ions, ®isters, true).unwrap();
let mut backups = backup_store(&temp.path().join("backups"));
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
let mut exec = ctx(&mut backups, &mut snapshots, &mut session_regs, true, None);
let envelope = execute_transaction(plan, &mut exec);
assert!(envelope.registers_committed);
assert_eq!(
session_regs.get(&RegisterRef::Named("clip".into())),
Some(["one".to_string()].as_slice())
);
write_file(&a, b"one\ntwo\n");
let plan2 = plan_transaction(§ions, &RegisterStore::new(), true).unwrap();
let mut session_regs2 = RegisterStore::new();
let mut exec2 = ctx(
&mut backups,
&mut snapshots,
&mut session_regs2,
true,
Some(ExecuteFault::Write { step: 0 }),
);
let envelope2 = execute_transaction(plan2, &mut exec2);
assert!(!envelope2.registers_committed);
assert!(session_regs2
.get(&RegisterRef::Named("clip".into()))
.is_none());
}
#[test]
fn a8_applied_star_variants() {
let temp = tempfile::tempdir().unwrap();
let path = temp.path().join("v.txt");
write_file(&path, b"x\n");
let bytes = fs::read(&path).unwrap();
let snap = whole_snapshot(&bytes);
let base = Baseline::from_bytes(bytes);
let ops = vec![put_text("1", &["Y"])];
let resolved = vec![resolve_one(&snap, &ops[0])];
let sections = [section_put(&path, "v.txt", &base, &snap, &ops, &resolved)];
let registers = RegisterStore::new();
let mut backups = backup_store(&temp.path().join("backups"));
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
let plan = plan_transaction(§ions, ®isters, true).unwrap();
let mut exec = ctx(
&mut backups,
&mut snapshots,
&mut session_regs,
true,
Some(ExecuteFault::ValidationFailure { step: 0 }),
);
let envelope = execute_transaction(plan, &mut exec);
assert_eq!(
envelope.files[0].classification,
FileClassification::AppliedWithValidationFailure
);
assert_eq!(fs::read(&path).unwrap(), b"Y\n");
write_file(&path, b"x\n");
let bytes = fs::read(&path).unwrap();
let snap = whole_snapshot(&bytes);
let base = Baseline::from_bytes(bytes);
let sections = [section_put(&path, "v.txt", &base, &snap, &ops, &resolved)];
let plan = plan_transaction(§ions, ®isters, true).unwrap();
let mut exec = ctx(
&mut backups,
&mut snapshots,
&mut session_regs,
true,
Some(ExecuteFault::FinalTagUnavailable { step: 0 }),
);
let envelope = execute_transaction(plan, &mut exec);
assert_eq!(
envelope.files[0].classification,
FileClassification::AppliedTagUnavailable
);
assert!(envelope.files[0].final_tag.is_none());
assert!(envelope.files[0].tag_notice.is_some());
}
#[test]
fn a10_preview_mutates_nothing() {
let temp = tempfile::tempdir().unwrap();
let path = temp.path().join("p.txt");
let dest = temp.path().join("p-dest.txt");
write_file(&path, b"preview\n");
let bytes = fs::read(&path).unwrap();
let snap = whole_snapshot(&bytes);
let base = Baseline::from_bytes(bytes.clone());
let ops = vec![
put_text("1", &["PREVIEWED"]),
Operation::Mv(MvOperation {
destination: "p-dest.txt".into(),
line: 2,
}),
];
let resolved = vec![
resolve_one(&snap, &ops[0]),
ResolvedOperation {
operation_index: 1,
address: ResolvedAddress::WholeFile,
},
];
let sections = [TransactionSectionInput {
canonical_path: &path,
requested_path: "p.txt",
baseline: &base,
snapshot: &snap,
operations: &ops,
resolved: &resolved,
mv_destination: Some(MvDestinationInput {
canonical_path: &dest,
requested_path: "p-dest.txt",
baseline_bytes: None,
}),
}];
let mut registers = RegisterStore::new();
{
let mut staged = registers.stage();
staged
.capture(RegisterRef::Named("keep".into()), vec!["seed".into()])
.unwrap();
registers.commit(staged);
}
let plan = plan_transaction(§ions, ®isters, true).unwrap();
let before_reg = registers
.get(&RegisterRef::Named("keep".into()))
.map(|lines| lines.to_vec());
let backups = backup_store(&temp.path().join("backups"));
let tracked_before = backups.tracked_files(SESSION);
let mut snapshots = SnapshotStore::new();
snapshots.publish(&path, snap.clone());
let envelope = preview_transaction(plan);
assert!(envelope.preview);
assert!(envelope.op_id.is_none());
assert!(!envelope.registers_committed);
assert_eq!(fs::read(&path).unwrap(), b"preview\n");
assert!(!dest.exists());
assert_eq!(backups.tracked_files(SESSION), tracked_before);
assert!(snapshots.lookup(&path, &snap.tag).is_ok());
assert_eq!(
registers
.get(&RegisterRef::Named("keep".into()))
.map(|lines| lines.to_vec()),
before_reg
);
assert!(envelope.files.iter().all(|f| f.final_tag.is_none()));
assert!(envelope
.files
.iter()
.all(|f| f.mutation_state == MutationState::Unmutated));
}
#[test]
fn a12_baseline_drift_stops_later_files_and_keeps_prior_op_id() {
let temp = tempfile::tempdir().unwrap();
let a = temp.path().join("a.txt");
let b = temp.path().join("b.txt");
write_file(&a, b"a0\n");
write_file(&b, b"b0\n");
let bytes_a = fs::read(&a).unwrap();
let bytes_b = fs::read(&b).unwrap();
let snap_a = whole_snapshot(&bytes_a);
let snap_b = whole_snapshot(&bytes_b);
let base_a = Baseline::from_bytes(bytes_a);
let base_b = Baseline::from_bytes(bytes_b);
let ops_a = vec![put_text("1", &["A1"])];
let ops_b = vec![put_text("1", &["B1"])];
let res_a = vec![resolve_one(&snap_a, &ops_a[0])];
let res_b = vec![resolve_one(&snap_b, &ops_b[0])];
let sections = [
section_put(&a, "a.txt", &base_a, &snap_a, &ops_a, &res_a),
section_put(&b, "b.txt", &base_b, &snap_b, &ops_b, &res_b),
];
let registers = RegisterStore::new();
let plan = plan_transaction(§ions, ®isters, true).unwrap();
write_file(&b, b"b-EXTERNAL\n");
let mut backups = backup_store(&temp.path().join("backups"));
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
let mut exec = ctx(&mut backups, &mut snapshots, &mut session_regs, true, None);
let envelope = execute_transaction(plan, &mut exec);
assert!(envelope.success);
assert!(!envelope.complete);
assert_eq!(
envelope.files[0].classification,
FileClassification::Applied
);
assert_eq!(
envelope.files[1].classification,
FileClassification::FailedBaselineDrift
);
assert_eq!(envelope.files[1].mutation_state, MutationState::Unmutated);
assert_eq!(envelope.stop_reason, Some("hashline_baseline_drift"));
assert!(envelope.op_id.is_some());
assert_eq!(fs::read(&a).unwrap(), b"A1\n");
assert_eq!(fs::read(&b).unwrap(), b"b-EXTERNAL\n");
let op_id = envelope.op_id.unwrap();
let restored = backups.restore_last_operation(SESSION).unwrap();
assert_eq!(restored.op_id, op_id);
assert_eq!(fs::read(&a).unwrap(), b"a0\n");
}
#[test]
fn a17_backup_unavailable_refusals_and_new_dest_mv() {
let temp = tempfile::tempdir().unwrap();
let path = temp.path().join("t.txt");
write_file(&path, b"t\n");
let bytes = fs::read(&path).unwrap();
let snap = whole_snapshot(&bytes);
let base = Baseline::from_bytes(bytes);
let ops = vec![put_text("1", &["T"])];
let resolved = vec![resolve_one(&snap, &ops[0])];
let sections = [section_put(&path, "t.txt", &base, &snap, &ops, &resolved)];
let registers = RegisterStore::new();
let err = plan_transaction(§ions, ®isters, false).unwrap_err();
assert_eq!(
err.code,
crate::hashline::syntax::HashlineRejectionCode::BackupUnavailable
);
assert_eq!(err.stage, crate::hashline::syntax::RejectionStage::Baseline);
assert_eq!(fs::read(&path).unwrap(), b"t\n");
let src = temp.path().join("s.txt");
let dest = temp.path().join("d.txt");
write_file(&src, b"s\n");
write_file(&dest, b"d\n");
let s_bytes = fs::read(&src).unwrap();
let d_bytes = fs::read(&dest).unwrap();
let s_snap = whole_snapshot(&s_bytes);
let s_base = Baseline::from_bytes(s_bytes);
let mv_ops = vec![Operation::Mv(MvOperation {
destination: "d.txt".into(),
line: 1,
})];
let mv_resolved = vec![ResolvedOperation {
operation_index: 0,
address: ResolvedAddress::WholeFile,
}];
let mv_sections = [TransactionSectionInput {
canonical_path: &src,
requested_path: "s.txt",
baseline: &s_base,
snapshot: &s_snap,
operations: &mv_ops,
resolved: &mv_resolved,
mv_destination: Some(MvDestinationInput {
canonical_path: &dest,
requested_path: "d.txt",
baseline_bytes: Some(&d_bytes),
}),
}];
let err = plan_transaction(&mv_sections, ®isters, false).unwrap_err();
assert_eq!(
err.code,
crate::hashline::syntax::HashlineRejectionCode::BackupUnavailable
);
assert_eq!(fs::read(&src).unwrap(), b"s\n");
assert_eq!(fs::read(&dest).unwrap(), b"d\n");
let src2 = temp.path().join("s2.txt");
let dest2 = temp.path().join("d2.txt");
write_file(&src2, b"s2\n");
let s2_bytes = fs::read(&src2).unwrap();
let s2_snap = whole_snapshot(&s2_bytes);
let s2_base = Baseline::from_bytes(s2_bytes);
let mv2_ops = vec![Operation::Mv(MvOperation {
destination: "d2.txt".into(),
line: 1,
})];
let mv2_resolved = vec![ResolvedOperation {
operation_index: 0,
address: ResolvedAddress::WholeFile,
}];
let mv2_sections = [TransactionSectionInput {
canonical_path: &src2,
requested_path: "s2.txt",
baseline: &s2_base,
snapshot: &s2_snap,
operations: &mv2_ops,
resolved: &mv2_resolved,
mv_destination: Some(MvDestinationInput {
canonical_path: &dest2,
requested_path: "d2.txt",
baseline_bytes: None,
}),
}];
let plan = plan_transaction(&mv2_sections, ®isters, false).expect("new dest MV plans");
let mut backups = backup_store(&temp.path().join("backups"));
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
let mut exec = ctx(&mut backups, &mut snapshots, &mut session_regs, false, None);
let envelope = execute_transaction(plan, &mut exec);
assert!(envelope.success);
assert!(envelope.op_id.is_some(), "real journaled op_id required");
assert_eq!(fs::read(&dest2).unwrap(), b"s2\n");
assert!(!src2.exists());
let op_id = envelope.op_id.unwrap();
let restored = backups.restore_last_operation(SESSION).unwrap();
assert_eq!(restored.op_id, op_id);
assert_eq!(fs::read(&src2).unwrap(), b"s2\n");
assert!(!dest2.exists());
let mut disabled = BackupStore::new();
disabled.set_policy(BackupPolicy {
enabled: false,
..BackupPolicy::default()
});
write_file(&src2, b"s2\n");
let plan = plan_transaction(&mv2_sections, ®isters, false).unwrap();
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
let mut exec = ctx(
&mut disabled,
&mut snapshots,
&mut session_regs,
false,
None,
);
let envelope = execute_transaction(plan, &mut exec);
assert!(!envelope.success);
assert!(envelope.op_id.is_none());
assert_eq!(
envelope.files[0].classification,
FileClassification::FailedBackup
);
}
#[test]
fn op_id_present_when_journal_entry_exists_before_failure() {
let temp = tempfile::tempdir().unwrap();
let a = temp.path().join("a.txt");
write_file(&a, b"a\n");
let bytes = fs::read(&a).unwrap();
let snap = whole_snapshot(&bytes);
let base = Baseline::from_bytes(bytes);
let ops = vec![put_text("1", &["A"])];
let resolved = vec![resolve_one(&snap, &ops[0])];
let sections = [section_put(&a, "a.txt", &base, &snap, &ops, &resolved)];
let registers = RegisterStore::new();
let plan = plan_transaction(§ions, ®isters, true).unwrap();
let mut backups = backup_store(&temp.path().join("backups"));
let mut snapshots = SnapshotStore::new();
let mut session_regs = RegisterStore::new();
write_file(&a, b"changed\n");
let mut exec = ctx(&mut backups, &mut snapshots, &mut session_regs, true, None);
let envelope = execute_transaction(plan, &mut exec);
assert!(!envelope.success);
assert_eq!(
envelope.files[0].classification,
FileClassification::FailedBaselineDrift
);
assert!(envelope.op_id.is_some());
assert_eq!(fs::read(&a).unwrap(), b"changed\n");
}
}