use std::path::{Path, PathBuf};
use std::sync::atomic::AtomicBool;
use std::time::Instant;
use crate::common::defaults::log_load_timing;
use crate::common::fs::{safe_delete_with_suffix, sync_parent_dir};
use crate::common::storage_version::StorageVersion;
use crate::common::types::PointOffsetType;
use crate::common::universal_io::MmapFs;
use fs_err as fs;
use log::info;
use uuid::Uuid;
use super::create_segment::create_segment;
use super::legacy_state::{load_segment_state_v3, load_segment_state_v5};
use crate::segment::common::operation_error::{OperationError, OperationResult};
use crate::segment::index::struct_payload_index::IndexLoadMode;
use crate::segment::segment::{Segment, SegmentVersion};
use crate::segment::types::SegmentConfig;
#[must_use = "a newly built segment must be registered in the segment manifest, or explicitly dropped"]
pub struct NewSegmentToken(Uuid);
impl NewSegmentToken {
pub(crate) fn new(uuid: Uuid) -> Self {
NewSegmentToken(uuid)
}
pub fn id(&self) -> Uuid {
self.0
}
}
pub fn normalize_segment_dir(path: &Path) -> OperationResult<Option<(PathBuf, Uuid)>> {
if path
.extension()
.and_then(|ext| ext.to_str())
.map(|ext| ext == "deleted")
.unwrap_or(false)
{
log::warn!("Deleting leftover segment: {}", path.display());
safe_delete_with_suffix(path).map_err(|err| {
OperationError::service_error(format!("failed to delete leftover segment: {err}"))
})?;
return Ok(None);
}
if SegmentVersion::load_universal(&MmapFs, path)?.is_none() {
log::warn!("Deleting segment without version file: {}", path.display());
safe_delete_with_suffix(path).map_err(|err| {
OperationError::service_error(format!("failed to delete leftover segment: {err}"))
})?;
return Ok(None);
}
let file_name = path
.file_name()
.and_then(|fname| fname.to_str())
.ok_or_else(|| {
OperationError::service_error(format!(
"Failed to get segment folder name: {}",
path.display()
))
})?;
match Uuid::try_parse(file_name) {
Ok(uuid) => Ok(Some((path.to_path_buf(), uuid))),
Err(_) => {
let segment_uuid = Uuid::new_v4();
let new_path = path.with_file_name(segment_uuid.to_string());
log::warn!(
"Segment name is not a valid UUID: {}. Renaming to {segment_uuid}",
path.display(),
);
fs::rename(path, &new_path)?;
sync_parent_dir(&new_path)?;
Ok(Some((new_path, segment_uuid)))
}
}
}
pub fn load_segment(
path: &Path,
uuid: Uuid,
deferred_internal_id: Option<PointOffsetType>,
stopped: &AtomicBool,
) -> OperationResult<Segment> {
let total_started = Instant::now();
let stored_version = SegmentVersion::load_universal(&MmapFs, path)?.ok_or_else(|| {
OperationError::service_error(format!(
"Segment version file not found in segment: {}",
path.display()
))
})?;
let app_version = SegmentVersion::current();
if stored_version != app_version {
info!("Migrating segment {stored_version} -> {app_version}");
if stored_version > app_version {
return Err(OperationError::service_error(format!(
"Data version {stored_version} is newer than application version {app_version}. \
Please upgrade the application. Compatibility is not guaranteed."
)));
}
if stored_version.major == 0 && stored_version.minor < 3 {
return Err(OperationError::service_error(format!(
"Segment version({stored_version}) is not compatible with current version({app_version})"
)));
}
if stored_version.major == 0 && stored_version.minor == 3 {
let segment_state = load_segment_state_v3(path)?;
Segment::save_state(&segment_state, path)?;
} else if stored_version.major == 0 && stored_version.minor <= 5 {
let segment_state = load_segment_state_v5(path)?;
Segment::save_state(&segment_state, path)?;
}
SegmentVersion::save(path)?
}
let started = Instant::now();
let segment_state = Segment::load_state(path)?;
log_load_timing(path, "load_state", started);
let segment = create_segment(
segment_state.initial_version,
segment_state.version,
path,
uuid,
deferred_internal_id,
&segment_state.config,
stopped,
IndexLoadMode::LoadExisting,
)?;
log_load_timing(path, "total", total_started);
Ok(segment)
}
pub fn build_segment(
segments_path: &Path,
config: &SegmentConfig,
deferred_internal_id: Option<PointOffsetType>,
ready: bool,
) -> OperationResult<(Segment, NewSegmentToken)> {
let uuid = Uuid::new_v4();
let token = NewSegmentToken::new(uuid);
let segment_path = segments_path.join(uuid.to_string());
let stopped = AtomicBool::new(false);
fs::create_dir_all(&segment_path)?;
let segment = create_segment(
None,
None,
&segment_path,
uuid,
deferred_internal_id,
config,
&stopped,
IndexLoadMode::CreateIfMissing,
)?;
segment.save_current_state()?;
if ready {
SegmentVersion::save(&segment_path)?;
}
Ok((segment, token))
}