use std::path::{Path, PathBuf};
use serde::Deserialize;
use tokio::{fs, io::AsyncWriteExt};
use crate::{
api::RequestContext,
domain::{BusinessGlossary, RepositoryMapType},
project::{
AGENT_CONTRACT_DIR_NAME, KNOWLEDGE_MAP_HISTORY_DIR_NAME, KNOWLEDGE_MAP_RELATIVE_PATH,
KNOWLEDGE_MAP_TOPICS_DIR_NAME, KNOWLEDGE_MAP_V3_RETAINED_BACKUP_FILE_NAME,
KNOWLEDGE_MAP_V3_RETAINED_FILE_NAME, LEGACY_AGENT_CONTRACT_DIR_NAME,
LEGACY_BUSINESS_GLOSSARY_RELATIVE_PATH, LEGACY_KNOWLEDGE_MAP_BACKUP_FILE_NAME,
LEGACY_KNOWLEDGE_MAP_REDIRECT_PREPARED_FILE_NAME,
LEGACY_KNOWLEDGE_MAP_REDIRECT_PREVIOUS_FILE_NAME, LEGACY_KNOWLEDGE_MAP_RELATIVE_PATH,
LEGACY_KNOWLEDGE_MAP_ROLLBACK_PREPARED_FILE_NAME,
LEGACY_KNOWLEDGE_MAP_ROLLBACK_PREVIOUS_FILE_NAME,
},
};
use super::{
ARTIFACT_SCHEMA_VERSION, KnowledgeMapMutationResponse, KnowledgeMapService,
KnowledgeMapServiceError, ensure_owned_directory, read_root_file, temporary_path,
validation::LegacyGlossaryReadPolicy,
};
#[cfg(test)]
use super::{
KnowledgeMapArchiveRef, KnowledgeMapSchemaProbe, KnowledgeMapTopicShard,
LEGACY_ARTIFACT_SCHEMA_VERSION, WRITE_LOCK_TIMEOUT, content_digest, parse_manifest,
parse_v1_map_for_legacy_recovery, publish_immutable_in, read_verified_ref_in, serialize_yaml,
stable_id,
};
#[cfg(test)]
use crate::domain::KnowledgeMap;
const MAX_LEGACY_MIGRATION_ARTIFACT_FILES: usize = 1_024;
const MAX_LEGACY_MIGRATION_ARTIFACT_BYTES: u64 = 64 * 1024 * 1024;
const MAX_LEGACY_MIGRATION_ARTIFACT_FILE_BYTES: u64 = 4 * 1024 * 1024;
#[derive(Deserialize)]
struct LegacyRedirect {
schema_version: u16,
artifact_kind: String,
map_type: RepositoryMapType,
target: String,
}
impl KnowledgeMapService {
pub async fn migrate_to_v4(
&self,
context: &RequestContext,
) -> Result<KnowledgeMapMutationResponse, KnowledgeMapServiceError> {
self.require_knowledge_map("map migrate --to-v4")?;
self.init(context).await
}
#[cfg(test)]
pub(super) async fn rollback_v3(
&self,
context: &RequestContext,
) -> Result<KnowledgeMapMutationResponse, KnowledgeMapServiceError> {
self.require_knowledge_map("map migrate --rollback")?;
let _legacy_lock = self.acquire_legacy_write_lock(WRITE_LOCK_TIMEOUT).await?;
let _lock = self.acquire_write_lock(WRITE_LOCK_TIMEOUT).await?;
self.recover_manifest_backup().await?;
if self.discard_incomplete_forward_migration_staging().await? {
let active_legacy =
read_root_file(&self.repository_root, &self.legacy_map_path()).await?;
let rollback_version = map_version_from_validated_legacy_content(&active_legacy)?;
return Ok(self.mutation_response(
context,
rollback_version,
"discarded incomplete initial migration staging; legacy Knowledge Map remains active"
.to_owned(),
));
}
if self.recover_legacy_rollback_transition().await? {
let restored_legacy =
read_root_file(&self.repository_root, &self.legacy_map_path()).await?;
self.validate_legacy_map_content_in(
LEGACY_AGENT_CONTRACT_DIR_NAME,
&restored_legacy,
LegacyGlossaryReadPolicy::LegacyRootCompatibility,
)
.await?;
let rollback_version = map_version_from_validated_legacy_content(&restored_legacy)?;
return Ok(self.mutation_response(
context,
rollback_version,
"retained committed Knowledge Map v2 rollback and cleaned recovery residue"
.to_owned(),
));
}
if let Some(active_legacy) = self.active_clean_legacy_rollback().await? {
let rollback_version = map_version_from_validated_legacy_content(&active_legacy)?;
return Ok(self.mutation_response(
context,
rollback_version,
"Knowledge Map v2 rollback is already active; retained v3 data remains available"
.to_owned(),
));
}
let legacy_backup = self.validate_legacy_backup().await?;
let rollback_version = map_version_from_validated_legacy_content(&legacy_backup)?;
let legacy = self.legacy_map_path();
let rollback_prepared = self.legacy_rollback_prepared_path();
let rollback_previous = self.legacy_rollback_previous_path();
let current = self.map_path();
let ordinary_backup = self.backup_path();
let retained = self.retained_v3_path();
let retained_backup = self.retained_v3_backup_path();
let legacy_existed = regular_file_exists_or_missing(&legacy).await?;
let current_exists = regular_file_exists_or_missing(¤t).await?;
let ordinary_backup_exists = regular_file_exists_or_missing(&ordinary_backup).await?;
remove_regular_transition_file(&rollback_prepared).await?;
remove_regular_transition_file(&rollback_previous).await?;
if current_exists {
remove_regular_transition_file(&retained).await?;
}
if ordinary_backup_exists {
remove_regular_transition_file(&retained_backup).await?;
}
replace_with_new_synced_file(&rollback_prepared, legacy_backup.as_bytes()).await?;
let current_moved = if current_exists {
match fs::rename(¤t, &retained).await {
Ok(()) => true,
Err(error) => {
let _ = fs::remove_file(&rollback_prepared).await;
return Err(error.into());
}
}
} else {
false
};
let ordinary_backup_moved = if ordinary_backup_exists {
match fs::rename(&ordinary_backup, &retained_backup).await {
Ok(()) => true,
Err(error) => {
restore_visible_rollback_roots(
¤t,
&retained,
current_moved,
&ordinary_backup,
&retained_backup,
false,
)
.await;
let _ = fs::remove_file(&rollback_prepared).await;
return Err(error.into());
}
}
} else {
false
};
let legacy_moved = if legacy_existed {
match fs::rename(&legacy, &rollback_previous).await {
Ok(()) => true,
Err(error) => {
restore_visible_rollback_roots(
¤t,
&retained,
current_moved,
&ordinary_backup,
&retained_backup,
ordinary_backup_moved,
)
.await;
let _ = fs::remove_file(&rollback_prepared).await;
return Err(error.into());
}
}
} else {
false
};
if let Err(error) = fs::rename(&rollback_prepared, &legacy).await {
if legacy_moved {
let _ = fs::rename(&rollback_previous, &legacy).await;
}
restore_visible_rollback_roots(
¤t,
&retained,
current_moved,
&ordinary_backup,
&retained_backup,
ordinary_backup_moved,
)
.await;
let _ = fs::remove_file(&rollback_prepared).await;
return Err(error.into());
}
let _ = remove_regular_transition_file(&rollback_previous).await;
let _ = remove_regular_transition_file(&self.legacy_redirect_prepared_path()).await;
let _ = remove_regular_transition_file(&self.legacy_redirect_previous_path()).await;
Ok(self.mutation_response(
context,
rollback_version,
"restored Knowledge Map v2 root; retained v3 data for forward recovery".to_owned(),
))
}
pub(super) async fn prepare_legacy_migration(&self) -> Result<(), KnowledgeMapServiceError> {
let current = self.map_path();
if self.map_type != RepositoryMapType::Knowledge
|| regular_file_exists_or_missing(¤t).await?
{
return Ok(());
}
let legacy = self.legacy_map_path();
let legacy_content = match read_root_file(&self.repository_root, &legacy).await {
Ok(content) => content,
Err(KnowledgeMapServiceError::Io(error))
if error.kind() == std::io::ErrorKind::NotFound =>
{
return Ok(());
}
Err(error) => return Err(error),
};
self.validate_legacy_map_content_in(
LEGACY_AGENT_CONTRACT_DIR_NAME,
&legacy_content,
LegacyGlossaryReadPolicy::LegacyRootCompatibility,
)
.await?;
let backup = self.legacy_backup_path();
remove_regular_transition_file(&backup).await?;
write_new_synced_file(&backup, legacy_content.as_bytes()).await?;
copy_contract_tree(
&self.repository_root,
&self
.repository_root
.join(LEGACY_AGENT_CONTRACT_DIR_NAME)
.join(KNOWLEDGE_MAP_TOPICS_DIR_NAME),
&self
.repository_root
.join(AGENT_CONTRACT_DIR_NAME)
.join(KNOWLEDGE_MAP_TOPICS_DIR_NAME),
)
.await?;
copy_contract_tree(
&self.repository_root,
&self
.repository_root
.join(LEGACY_AGENT_CONTRACT_DIR_NAME)
.join(KNOWLEDGE_MAP_HISTORY_DIR_NAME),
&self
.repository_root
.join(AGENT_CONTRACT_DIR_NAME)
.join(KNOWLEDGE_MAP_HISTORY_DIR_NAME),
)
.await?;
let legacy_glossary = self
.repository_root
.join(LEGACY_BUSINESS_GLOSSARY_RELATIVE_PATH);
let glossary = self.business_glossary_path();
let legacy_glossary_exists = regular_file_exists_or_missing(&legacy_glossary).await?;
if legacy_glossary_exists {
if let Some(parent) = glossary.parent() {
ensure_owned_directory(&self.repository_root, parent).await?;
}
let content = read_root_file(&self.repository_root, &legacy_glossary).await?;
BusinessGlossary::parse(content.as_bytes())?;
replace_with_new_synced_file(&glossary, content.as_bytes()).await?;
}
if let Some(parent) = current.parent() {
ensure_owned_directory(&self.repository_root, parent).await?;
}
write_new_synced_file(¤t, legacy_content.as_bytes()).await?;
Ok(())
}
pub(super) async fn publish_legacy_redirect(&self) -> Result<(), KnowledgeMapServiceError> {
if self.map_type != RepositoryMapType::Knowledge {
return Ok(());
}
let current = read_root_file(&self.repository_root, &self.map_path()).await?;
if !self.validate_visible_current_map_content(¤t).await? {
return Err(KnowledgeMapServiceError::Integrity(
"legacy redirect publication requires a complete current visible root".to_owned(),
));
}
self.converge_legacy_redirect().await
}
pub(super) async fn legacy_recovery_state_exists(
&self,
) -> Result<bool, KnowledgeMapServiceError> {
if self.map_type != RepositoryMapType::Knowledge {
return Ok(false);
}
for path in [
self.legacy_map_path(),
self.legacy_backup_path(),
self.legacy_redirect_prepared_path(),
self.legacy_redirect_previous_path(),
self.legacy_rollback_prepared_path(),
self.legacy_rollback_previous_path(),
] {
if path_entry_exists(&path).await? {
return Ok(true);
}
}
Ok(false)
}
pub(super) async fn legacy_history_cleanup_is_safe(
&self,
) -> Result<bool, KnowledgeMapServiceError> {
if self.map_type != RepositoryMapType::Knowledge {
return Ok(false);
}
match read_root_file(&self.repository_root, &self.legacy_map_path()).await {
Ok(content) => Ok(is_supported_legacy_redirect(&content)),
Err(KnowledgeMapServiceError::Io(error))
if error.kind() == std::io::ErrorKind::NotFound =>
{
Ok(true)
}
Err(error) => Err(error),
}
}
pub(super) async fn recover_legacy_redirect_transition(
&self,
) -> Result<(), KnowledgeMapServiceError> {
if self.map_type != RepositoryMapType::Knowledge
|| !path_entry_exists(&self.map_path()).await?
|| !path_entry_exists(&self.legacy_backup_path()).await?
{
return Ok(());
}
let current = read_root_file(&self.repository_root, &self.map_path()).await?;
if !self.validate_visible_current_map_content(¤t).await? {
return Ok(());
}
self.converge_legacy_redirect().await
}
pub(super) async fn recover_legacy_rollback_transition(
&self,
) -> Result<bool, KnowledgeMapServiceError> {
if self.map_type != RepositoryMapType::Knowledge {
return Ok(false);
}
let prepared = self.legacy_rollback_prepared_path();
let previous = self.legacy_rollback_previous_path();
let prepared_exists = regular_file_exists_or_missing(&prepared).await?;
let previous_exists = regular_file_exists_or_missing(&previous).await?;
if !prepared_exists && !previous_exists {
return Ok(false);
}
let legacy = self.legacy_map_path();
let legacy_exists = regular_file_exists_or_missing(&legacy).await?;
if !prepared_exists && previous_exists && legacy_exists {
let _ = remove_regular_transition_file(&previous).await;
return Ok(true);
}
let expected_legacy = self.validate_legacy_backup().await?;
if prepared_exists {
let staged = read_root_file(&self.repository_root, &prepared).await?;
if staged.as_bytes() != expected_legacy.as_bytes() {
return Err(KnowledgeMapServiceError::Integrity(
"rollback prepared root differs from the retained legacy backup".to_owned(),
));
}
}
let current = self.map_path();
let retained = self.retained_v3_path();
let ordinary_backup = self.backup_path();
let retained_backup = self.retained_v3_backup_path();
let current_exists = regular_file_exists_or_missing(¤t).await?;
let retained_exists = regular_file_exists_or_missing(&retained).await?;
let ordinary_backup_exists = regular_file_exists_or_missing(&ordinary_backup).await?;
if !current_exists && !retained_exists {
if !regular_file_exists_or_missing(&retained_backup).await? {
return Err(KnowledgeMapServiceError::Integrity(
"rollback recovery has neither a visible nor retained v3 root".to_owned(),
));
}
fs::rename(&retained_backup, &retained).await?;
}
let retained_backup_exists = regular_file_exists_or_missing(&retained_backup).await?;
if !current_exists {
fs::rename(&retained, ¤t).await?;
}
if !ordinary_backup_exists && retained_backup_exists {
fs::rename(&retained_backup, &ordinary_backup).await?;
}
if !legacy_exists && previous_exists {
fs::rename(&previous, &legacy).await?;
}
remove_regular_transition_file(&prepared).await?;
Ok(false)
}
pub(super) async fn readable_retained_v3_root(
&self,
) -> Result<Option<PathBuf>, KnowledgeMapServiceError> {
if self.map_type != RepositoryMapType::Knowledge
|| !regular_file_exists_or_missing(&self.legacy_rollback_prepared_path()).await?
{
return Ok(None);
}
for path in [self.retained_v3_path(), self.retained_v3_backup_path()] {
if regular_file_exists_or_missing(&path).await? {
return Ok(Some(path));
}
}
Ok(None)
}
#[cfg(test)]
async fn active_clean_legacy_rollback(
&self,
) -> Result<Option<String>, KnowledgeMapServiceError> {
if regular_file_exists_or_missing(&self.map_path()).await?
|| regular_file_exists_or_missing(&self.backup_path()).await?
|| !regular_file_exists_or_missing(&self.retained_v3_path()).await?
|| !regular_file_exists_or_missing(&self.legacy_map_path()).await?
{
return Ok(None);
}
let active_legacy = read_root_file(&self.repository_root, &self.legacy_map_path()).await?;
self.validate_legacy_map_content_in(
LEGACY_AGENT_CONTRACT_DIR_NAME,
&active_legacy,
LegacyGlossaryReadPolicy::LegacyRootCompatibility,
)
.await?;
Ok(Some(active_legacy))
}
#[cfg(test)]
async fn discard_incomplete_forward_migration_staging(
&self,
) -> Result<bool, KnowledgeMapServiceError> {
if self.map_type != RepositoryMapType::Knowledge {
return Ok(false);
}
let current = self.map_path();
if !regular_file_exists_or_missing(¤t).await? {
return Ok(false);
}
let current_content = read_root_file(&self.repository_root, ¤t).await?;
if self
.validate_visible_current_map_content(¤t_content)
.await?
{
return Ok(false);
}
let retained = self.retained_v3_path();
if !regular_file_exists_or_missing(&retained).await? {
let active_legacy =
read_root_file(&self.repository_root, &self.legacy_map_path()).await?;
self.validate_legacy_map_content_in(
LEGACY_AGENT_CONTRACT_DIR_NAME,
&active_legacy,
LegacyGlossaryReadPolicy::LegacyRootCompatibility,
)
.await?;
remove_regular_transition_file(¤t).await?;
return Ok(true);
}
let retained_content = read_root_file(&self.repository_root, &retained).await?;
if !self
.validate_visible_current_map_content(&retained_content)
.await?
{
return Err(KnowledgeMapServiceError::Integrity(
"incomplete forward migration retained root is not a complete current map"
.to_owned(),
));
}
remove_regular_transition_file(¤t).await?;
Ok(false)
}
async fn converge_legacy_redirect(&self) -> Result<(), KnowledgeMapServiceError> {
let yaml = format!(
"schema_version: {ARTIFACT_SCHEMA_VERSION}\nartifact_kind: redirect\nmap_type: knowledge\ntarget: {KNOWLEDGE_MAP_RELATIVE_PATH}\n"
);
let legacy = self.legacy_map_path();
let prepared = self.legacy_redirect_prepared_path();
let previous = self.legacy_redirect_previous_path();
let rollback_prepared = self.legacy_rollback_prepared_path();
let rollback_previous = self.legacy_rollback_previous_path();
let live_legacy = match read_root_file(&self.repository_root, &legacy).await {
Ok(content) => Some(content),
Err(KnowledgeMapServiceError::Io(error))
if error.kind() == std::io::ErrorKind::NotFound =>
{
None
}
Err(error) => return Err(error),
};
let live_is_current_redirect = live_legacy
.as_deref()
.is_some_and(|content| content.as_bytes() == yaml.as_bytes());
let live_is_supported_redirect = live_legacy
.as_deref()
.is_some_and(is_supported_legacy_redirect);
let has_residue = path_entry_exists(&prepared).await?
|| path_entry_exists(&previous).await?
|| path_entry_exists(&rollback_prepared).await?
|| path_entry_exists(&rollback_previous).await?;
if live_is_supported_redirect && !has_residue {
if live_is_current_redirect {
return Ok(());
}
return replace_with_new_synced_file(&legacy, yaml.as_bytes()).await;
}
let legacy_backup = self.validate_legacy_backup().await?;
if !live_is_supported_redirect {
self.ensure_legacy_glossary_matches_canonical(&legacy_backup)
.await?;
}
if regular_file_exists_or_missing(&previous).await? {
let previous_legacy = read_root_file(&self.repository_root, &previous).await?;
if previous_legacy != legacy_backup {
return Err(KnowledgeMapServiceError::Integrity(
"legacy redirect recovery preserves a root that diverged from the migration backup"
.to_owned(),
));
}
}
if live_is_supported_redirect {
if !live_is_current_redirect {
replace_with_new_synced_file(&legacy, yaml.as_bytes()).await?;
}
remove_regular_transition_file(&prepared).await?;
remove_regular_transition_file(&previous).await?;
remove_regular_transition_file(&rollback_prepared).await?;
remove_regular_transition_file(&rollback_previous).await?;
return Ok(());
}
if let Some(live_legacy) = &live_legacy
&& live_legacy != &legacy_backup
{
return Err(KnowledgeMapServiceError::Integrity(
"legacy root diverged from the migration backup after current-map publication; refusing redirect to preserve edits"
.to_owned(),
));
}
remove_regular_transition_file(&prepared).await?;
write_new_synced_file(&prepared, yaml.as_bytes()).await?;
let moved_legacy = if fs::try_exists(&legacy).await? {
remove_regular_transition_file(&previous).await?;
fs::rename(&legacy, &previous).await?;
true
} else {
false
};
if let Err(error) = fs::rename(&prepared, &legacy).await {
if moved_legacy {
let _ = fs::rename(&previous, &legacy).await;
}
let _ = fs::remove_file(prepared).await;
return Err(error.into());
}
remove_regular_transition_file(&previous).await?;
remove_regular_transition_file(&rollback_prepared).await?;
remove_regular_transition_file(&rollback_previous).await?;
Ok(())
}
async fn ensure_legacy_glossary_matches_canonical(
&self,
legacy_backup: &str,
) -> Result<(), KnowledgeMapServiceError> {
let Some(_) = self
.routed_business_glossary_source_in(LEGACY_AGENT_CONTRACT_DIR_NAME, legacy_backup)
.await?
else {
return Ok(());
};
let legacy_path = self
.repository_root
.join(LEGACY_BUSINESS_GLOSSARY_RELATIVE_PATH);
if !regular_file_exists_or_missing(&legacy_path).await? {
return Ok(());
}
let legacy_glossary = read_root_file(&self.repository_root, &legacy_path).await?;
BusinessGlossary::parse(legacy_glossary.as_bytes())?;
let canonical_glossary =
read_root_file(&self.repository_root, &self.business_glossary_path()).await?;
if canonical_glossary != legacy_glossary {
return Err(KnowledgeMapServiceError::Integrity(
"legacy business glossary diverged from the migration backup after v3 publication; refusing redirect to preserve edits"
.to_owned(),
));
}
Ok(())
}
async fn validate_legacy_backup(&self) -> Result<String, KnowledgeMapServiceError> {
let path = self.legacy_backup_path();
let content = match read_root_file(&self.repository_root, &path).await {
Ok(content) => content,
Err(KnowledgeMapServiceError::Io(error))
if error.kind() == std::io::ErrorKind::NotFound =>
{
return Err(KnowledgeMapServiceError::InvalidRequest(
"v2 rollback backup is missing".to_owned(),
));
}
Err(error) => return Err(error),
};
self.validate_legacy_map_content_in(
LEGACY_AGENT_CONTRACT_DIR_NAME,
&content,
LegacyGlossaryReadPolicy::ExactRoute,
)
.await?;
Ok(content)
}
pub(super) fn retained_v3_path(&self) -> PathBuf {
self.repository_root
.join(AGENT_CONTRACT_DIR_NAME)
.join(KNOWLEDGE_MAP_V3_RETAINED_FILE_NAME)
}
fn retained_v3_backup_path(&self) -> PathBuf {
self.repository_root
.join(AGENT_CONTRACT_DIR_NAME)
.join(KNOWLEDGE_MAP_V3_RETAINED_BACKUP_FILE_NAME)
}
fn legacy_redirect_prepared_path(&self) -> PathBuf {
self.repository_root
.join(LEGACY_AGENT_CONTRACT_DIR_NAME)
.join(LEGACY_KNOWLEDGE_MAP_REDIRECT_PREPARED_FILE_NAME)
}
fn legacy_redirect_previous_path(&self) -> PathBuf {
self.repository_root
.join(LEGACY_AGENT_CONTRACT_DIR_NAME)
.join(LEGACY_KNOWLEDGE_MAP_REDIRECT_PREVIOUS_FILE_NAME)
}
pub(super) fn legacy_rollback_prepared_path(&self) -> PathBuf {
self.repository_root
.join(LEGACY_AGENT_CONTRACT_DIR_NAME)
.join(LEGACY_KNOWLEDGE_MAP_ROLLBACK_PREPARED_FILE_NAME)
}
fn legacy_rollback_previous_path(&self) -> PathBuf {
self.repository_root
.join(LEGACY_AGENT_CONTRACT_DIR_NAME)
.join(LEGACY_KNOWLEDGE_MAP_ROLLBACK_PREVIOUS_FILE_NAME)
}
pub(super) fn legacy_map_path(&self) -> PathBuf {
self.repository_root
.join(LEGACY_KNOWLEDGE_MAP_RELATIVE_PATH)
}
pub(super) fn legacy_backup_path(&self) -> PathBuf {
self.repository_root
.join(LEGACY_AGENT_CONTRACT_DIR_NAME)
.join(LEGACY_KNOWLEDGE_MAP_BACKUP_FILE_NAME)
}
}
#[cfg(test)]
pub(super) async fn rewrite_contract_schema_for_test(
repository_root: &Path,
contract_dir: &str,
root_path: &Path,
target_schema_version: u16,
history_archive: Option<KnowledgeMapArchiveRef>,
) -> Result<(), KnowledgeMapServiceError> {
if !matches!(
target_schema_version,
LEGACY_ARTIFACT_SCHEMA_VERSION | super::DIRECTORY_ARTIFACT_SCHEMA_VERSION
) {
return Err(KnowledgeMapServiceError::InvalidRequest(
"test fixture target must be a readable legacy artifact schema".to_owned(),
));
}
let content = read_root_file(repository_root, root_path).await?;
let mut manifest = parse_manifest(&content)?;
for topic_ref in &mut manifest.topics {
let content = read_verified_ref_in(
repository_root,
contract_dir,
&topic_ref.r#ref,
&topic_ref.digest,
)
.await?;
let mut shard = serde_norway::from_str::<KnowledgeMapTopicShard>(&content)
.map_err(|error| KnowledgeMapServiceError::Yaml(error.to_string()))?;
shard.schema_version = target_schema_version;
let yaml = serialize_yaml(&shard)?;
let digest = content_digest(yaml.as_bytes());
let relative = format!(
"{KNOWLEDGE_MAP_TOPICS_DIR_NAME}/topic-{}-{digest}.yaml",
stable_id(&topic_ref.id)
);
publish_immutable_in(repository_root, contract_dir, &relative, yaml.as_bytes()).await?;
topic_ref.r#ref = relative;
topic_ref.digest = digest;
}
manifest.schema_version = target_schema_version;
manifest.history.archived_through = manifest.history.omitted_through;
manifest.history.omitted_through = 0;
if (manifest.history.archived_through == 0) != history_archive.is_none() {
return Err(KnowledgeMapServiceError::InvalidRequest(
"test fixture history checkpoint and archive must agree".to_owned(),
));
}
manifest.history.archive = history_archive;
manifest.history.index = None;
replace_with_new_synced_file(root_path, serialize_yaml(&manifest)?.as_bytes()).await
}
fn is_supported_legacy_redirect(content: &str) -> bool {
serde_norway::from_str::<LegacyRedirect>(content).is_ok_and(|redirect| {
matches!(
redirect.schema_version,
super::DIRECTORY_ARTIFACT_SCHEMA_VERSION | ARTIFACT_SCHEMA_VERSION
) && redirect.artifact_kind == "redirect"
&& redirect.map_type == RepositoryMapType::Knowledge
&& redirect.target == KNOWLEDGE_MAP_RELATIVE_PATH
})
}
async fn path_entry_exists(path: &Path) -> Result<bool, KnowledgeMapServiceError> {
match fs::symlink_metadata(path).await {
Ok(_) => Ok(true),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false),
Err(error) => Err(error.into()),
}
}
async fn regular_file_exists_or_missing(path: &Path) -> Result<bool, KnowledgeMapServiceError> {
match fs::symlink_metadata(path).await {
Ok(metadata) if metadata.is_file() && !metadata.file_type().is_symlink() => Ok(true),
Ok(_) => Err(KnowledgeMapServiceError::UnsafePath(
path.display().to_string(),
)),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false),
Err(error) => Err(error.into()),
}
}
#[cfg(test)]
fn map_version_from_validated_legacy_content(
content: &str,
) -> Result<u64, KnowledgeMapServiceError> {
let probe = serde_norway::from_str::<KnowledgeMapSchemaProbe>(content)
.map_err(|error| KnowledgeMapServiceError::Yaml(error.to_string()))?;
if probe.schema_version == KnowledgeMap::SCHEMA_VERSION {
return Ok(parse_v1_map_for_legacy_recovery(content)?.map_version);
}
if matches!(
probe.schema_version,
LEGACY_ARTIFACT_SCHEMA_VERSION
| super::DIRECTORY_ARTIFACT_SCHEMA_VERSION
| ARTIFACT_SCHEMA_VERSION
) {
return Ok(parse_manifest(content)?.map_version);
}
Err(KnowledgeMapServiceError::Yaml(format!(
"unsupported schema_version {}",
probe.schema_version
)))
}
#[cfg(test)]
async fn restore_visible_rollback_roots(
current: &Path,
retained: &Path,
current_moved: bool,
ordinary_backup: &Path,
retained_backup: &Path,
ordinary_backup_moved: bool,
) {
if ordinary_backup_moved {
let _ = fs::rename(retained_backup, ordinary_backup).await;
}
if current_moved {
let _ = fs::rename(retained, current).await;
}
}
async fn write_new_synced_file(
path: &Path,
content: &[u8],
) -> Result<(), KnowledgeMapServiceError> {
let mut file = fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(path)
.await?;
if let Err(error) = file.write_all(content).await {
drop(file);
let _ = fs::remove_file(path).await;
return Err(error.into());
}
if let Err(error) = file.sync_all().await {
drop(file);
let _ = fs::remove_file(path).await;
return Err(error.into());
}
Ok(())
}
async fn replace_with_new_synced_file(
path: &Path,
content: &[u8],
) -> Result<(), KnowledgeMapServiceError> {
regular_file_exists_or_missing(path).await?;
let temporary = temporary_path(path);
write_new_synced_file(&temporary, content).await?;
if let Err(error) = fs::rename(&temporary, path).await {
let _ = fs::remove_file(&temporary).await;
return Err(error.into());
}
Ok(())
}
async fn remove_regular_transition_file(path: &Path) -> Result<(), KnowledgeMapServiceError> {
match fs::symlink_metadata(path).await {
Ok(metadata) if !metadata.is_file() || metadata.file_type().is_symlink() => Err(
KnowledgeMapServiceError::UnsafePath(path.display().to_string()),
),
Ok(_) => {
fs::remove_file(path).await?;
Ok(())
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(error) => Err(error.into()),
}
}
async fn copy_contract_tree(
repository_root: &Path,
source: &Path,
target: &Path,
) -> Result<(), KnowledgeMapServiceError> {
match fs::symlink_metadata(source).await {
Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_dir() => {
return Err(KnowledgeMapServiceError::UnsafePath(
source.display().to_string(),
));
}
Ok(_) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
Err(error) => return Err(error.into()),
}
let repository = fs::canonicalize(repository_root).await?;
let canonical_source = fs::canonicalize(source).await?;
if !canonical_source.starts_with(&repository) {
return Err(KnowledgeMapServiceError::UnsafePath(
source.display().to_string(),
));
}
let target = ensure_owned_directory(repository_root, target).await?;
let mut entries = fs::read_dir(source).await?;
let mut copied_files = 0usize;
let mut copied_bytes = 0u64;
while let Some(entry) = entries.next_entry().await? {
let file_type = entry.file_type().await?;
if file_type.is_symlink() || !file_type.is_file() {
return Err(KnowledgeMapServiceError::UnsafePath(
entry.path().display().to_string(),
));
}
copied_files = copied_files.checked_add(1).ok_or_else(|| {
KnowledgeMapServiceError::Integrity(
"legacy migration artifact file count overflow".to_owned(),
)
})?;
if copied_files > MAX_LEGACY_MIGRATION_ARTIFACT_FILES {
return Err(KnowledgeMapServiceError::Integrity(format!(
"legacy migration artifact count exceeds {MAX_LEGACY_MIGRATION_ARTIFACT_FILES}"
)));
}
let file_bytes = entry.metadata().await?.len();
if file_bytes > MAX_LEGACY_MIGRATION_ARTIFACT_FILE_BYTES {
return Err(KnowledgeMapServiceError::Integrity(format!(
"legacy migration artifact '{}' exceeds {MAX_LEGACY_MIGRATION_ARTIFACT_FILE_BYTES} bytes",
entry.path().display()
)));
}
copied_bytes = copied_bytes.checked_add(file_bytes).ok_or_else(|| {
KnowledgeMapServiceError::Integrity(
"legacy migration artifact byte count overflow".to_owned(),
)
})?;
if copied_bytes > MAX_LEGACY_MIGRATION_ARTIFACT_BYTES {
return Err(KnowledgeMapServiceError::Integrity(format!(
"legacy migration artifacts exceed {MAX_LEGACY_MIGRATION_ARTIFACT_BYTES} bytes"
)));
}
let content = read_root_file(repository_root, &entry.path()).await?;
let destination = target.join(entry.file_name());
if regular_file_exists_or_missing(&destination).await? {
let existing = read_root_file(repository_root, &destination).await?;
if existing.as_bytes() != content.as_bytes() {
return Err(KnowledgeMapServiceError::Integrity(format!(
"migration target '{}' differs from the retained legacy artifact",
destination.display()
)));
}
} else {
write_new_synced_file(&destination, content.as_bytes()).await?;
}
}
Ok(())
}
#[cfg(test)]
#[path = "migration_tests.rs"]
mod tests;
#[cfg(test)]
#[path = "migration_recovery_tests.rs"]
mod recovery_tests;