Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
186 changes: 166 additions & 20 deletions src/crates/assembly/core/src/infrastructure/app_paths/path_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,15 @@ use std::env;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};

const LEGACY_REMOTE_RUNTIME_DIRECTORIES: &[&str] = &[
"sessions",
"request-traces",
"snapshots",
"locks",
"config",
"plans",
];

/// Storage level
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum StorageLevel {
Expand Down Expand Up @@ -469,10 +478,10 @@ impl PathManager {
remote_root_norm: &str,
runtime_root: &Path,
) -> PathBuf {
if runtime_root.exists() {
return runtime_root.to_path_buf();
}

static MIGRATION_LOCK: Mutex<()> = Mutex::new(());
let _guard = MIGRATION_LOCK
.lock()
.expect("remote runtime migration lock poisoned");
let legacy_root = remote_workspace_runtime_root(
self.remote_ssh_mirror_root_dir(),
ssh_host,
Expand All @@ -482,26 +491,61 @@ impl PathManager {
return runtime_root.to_path_buf();
}

if let Err(error) = std::fs::create_dir_all(self.projects_root()) {
warn!(
"Failed to prepare workspace projects root for remote runtime migration: root={}, error={}",
self.projects_root().display(),
error
);
return legacy_root;
let directories: Vec<_> = LEGACY_REMOTE_RUNTIME_DIRECTORIES
.iter()
.map(|directory| (legacy_root.join(directory), runtime_root.join(directory)))
.filter(|(source, _)| source.is_dir())
.collect();
for (source, destination) in &directories {
if destination.exists() {
warn!(
"Cannot migrate legacy remote workspace runtime into an existing directory: source={}, destination={}",
source.display(),
destination.display()
);
return legacy_root;
}
}

if let Err(error) = std::fs::rename(&legacy_root, runtime_root) {
if let Err(error) = std::fs::create_dir_all(runtime_root) {
warn!(
"Failed to migrate legacy remote workspace runtime: legacy_root={}, runtime_root={}, error={}",
legacy_root.display(),
"Failed to prepare new remote workspace runtime: runtime_root={}, error={}",
runtime_root.display(),
error
);
if !runtime_root.exists() {
return legacy_root;
}

for (index, (source, destination)) in directories.iter().enumerate() {
if let Err(error) = std::fs::rename(source, destination) {
warn!(
"Failed to migrate legacy remote workspace runtime directory: source={}, destination={}, error={}",
source.display(),
destination.display(),
error
);
// Preserve the legacy-root fallback if a later directory fails.
for (moved_source, moved_destination) in directories[..index].iter().rev() {
if let Err(error) = std::fs::rename(moved_destination, moved_source) {
error!(
"Failed to roll back remote runtime migration: source={}, destination={}, error={}",
moved_destination.display(),
moved_source.display(),
error
);
}
}
return legacy_root;
}
}

// Keep nested legacy workspace runtimes in place.
if let Ok(mut entries) = std::fs::read_dir(&legacy_root) {
if entries.next().is_none() {
let _ = std::fs::remove_dir(&legacy_root);
}
}

runtime_root.to_path_buf()
}

Expand Down Expand Up @@ -779,7 +823,10 @@ pub fn try_get_path_manager_arc() -> OpenBitFunResult<Arc<PathManager>> {

#[cfg(test)]
mod tests {
use super::{GlobalPathManagerState, PathManager};
use super::{GlobalPathManagerState, PathManager, LEGACY_REMOTE_RUNTIME_DIRECTORIES};
use openbitfun_services_core::workspace_identity::{
remote_workspace_runtime_key, remote_workspace_runtime_root,
};
use std::ffi::OsString;
use std::path::Path;
use std::sync::{Arc, Mutex};
Expand Down Expand Up @@ -988,9 +1035,11 @@ mod tests {
host,
remote_root,
);
std::fs::create_dir_all(legacy.join("sessions")).expect("legacy runtime should exist");
std::fs::write(legacy.join("sessions").join("marker"), b"legacy")
.expect("legacy marker should be written");
for directory in LEGACY_REMOTE_RUNTIME_DIRECTORIES {
std::fs::create_dir_all(legacy.join(directory)).expect("legacy directory should exist");
std::fs::write(legacy.join(directory).join("marker"), b"legacy")
.expect("legacy marker should be written");
}

let runtime = pm.remote_workspace_runtime_root(host, remote_root);

Expand All @@ -1001,8 +1050,105 @@ mod tests {
.map(str::len),
Some(24)
);
assert!(runtime.join("sessions").join("marker").exists());
for directory in LEGACY_REMOTE_RUNTIME_DIRECTORIES {
assert_eq!(
std::fs::read(runtime.join(directory).join("marker")).unwrap(),
b"legacy"
);
}
assert!(!legacy.exists());
std::fs::remove_dir_all(base).expect("test root should be removed");
}

#[test]
fn remote_runtime_migration_does_not_move_nested_workspace_runtime() {
let base = std::env::temp_dir().join(format!(
"openbitfun-remote-nested-runtime-migration-{}",
uuid::Uuid::new_v4()
));
let pm = PathManager::with_user_root_for_tests(base.join("user"));
let host = "Example.COM";
let parent_root = "/home/project-A";
let child_root = "/home/project-A/sub-project-B";
let parent_legacy =
openbitfun_services_core::workspace_identity::remote_workspace_runtime_root(
pm.remote_ssh_mirror_root_dir(),
host,
parent_root,
);
let child_legacy =
openbitfun_services_core::workspace_identity::remote_workspace_runtime_root(
pm.remote_ssh_mirror_root_dir(),
host,
child_root,
);
std::fs::create_dir_all(parent_legacy.join("sessions"))
.expect("parent legacy runtime should exist");
std::fs::write(parent_legacy.join("sessions").join("parent"), b"parent")
.expect("parent marker should be written");
std::fs::create_dir_all(child_legacy.join("sessions"))
.expect("child legacy runtime should exist");
std::fs::write(child_legacy.join("sessions").join("child"), b"child")
.expect("child marker should be written");

let parent_runtime = pm.remote_workspace_runtime_root(host, parent_root);

assert!(parent_runtime.join("sessions").join("parent").exists());
assert!(child_legacy.join("sessions").join("child").exists());
assert!(child_legacy.exists());
assert!(!parent_runtime.join("sub-project-B").exists());
assert_eq!(
pm.remote_workspace_runtime_root(host, parent_root),
parent_runtime
);

let child_runtime = pm.remote_workspace_runtime_root(host, child_root);
assert_ne!(child_runtime, parent_runtime);
assert_eq!(
std::fs::read(child_runtime.join("sessions").join("child")).unwrap(),
b"child"
);
assert!(!child_legacy.exists());
std::fs::remove_dir_all(base).expect("test root should be removed");
}

#[test]
fn remote_runtime_migration_preserves_legacy_data_on_destination_conflict() {
let base = std::env::temp_dir().join(format!(
"openbitfun-remote-runtime-conflict-{}",
uuid::Uuid::new_v4()
));
let pm = PathManager::with_user_root_for_tests(base.join("user"));
let host = "example.com";
let remote_root = "/home/project-A";
let legacy =
remote_workspace_runtime_root(pm.remote_ssh_mirror_root_dir(), host, remote_root);
let runtime = pm
.projects_root()
.join(remote_workspace_runtime_key(host, remote_root));
std::fs::create_dir_all(legacy.join("sessions")).unwrap();
std::fs::write(legacy.join("sessions/marker"), b"legacy").unwrap();
std::fs::create_dir_all(legacy.join("snapshots")).unwrap();
std::fs::create_dir_all(runtime.join("snapshots")).unwrap();
std::fs::write(runtime.join("snapshots/marker"), b"new").unwrap();

assert_eq!(pm.remote_workspace_runtime_root(host, remote_root), legacy);
assert_eq!(
std::fs::read(legacy.join("sessions/marker")).unwrap(),
b"legacy"
);
assert_eq!(
std::fs::read(runtime.join("snapshots/marker")).unwrap(),
b"new"
);
assert!(!runtime.join("sessions").exists());

std::fs::remove_file(runtime.join("snapshots/marker")).unwrap();
std::fs::remove_dir(runtime.join("snapshots")).unwrap();
assert_eq!(pm.remote_workspace_runtime_root(host, remote_root), runtime);
assert!(runtime.join("sessions/marker").exists());
assert!(!legacy.exists());
std::fs::remove_dir_all(base).expect("test root should be removed");
}

#[test]
Expand Down
Loading