mod optimize;
mod shard_read;
mod snapshots;
mod update;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use crate::wal::WalOptions;
use crate::common::save_on_disk::SaveOnDisk;
use fs_err as fs;
use parking_lot::Mutex;
use crate::segment::common::operation_error::{OperationError, OperationResult};
use crate::segment::entry::ReadSegmentEntry as _;
use crate::segment::segment_constructor::{load_segment, normalize_segment_dir};
use crate::shard::files::{PAYLOAD_INDEX_CONFIG_FILE, SEGMENTS_PATH, segment_manifest_path};
use crate::shard::operations::CollectionUpdateOperations;
use crate::shard::segment_holder::locked::LockedSegmentHolder;
use crate::shard::segment_holder::{FlushMode, SegmentHolder};
use crate::shard::segment_manifest::SegmentsManifest;
use crate::shard::wal::SerdeWal;
use uuid::Uuid;
use crate::edge::config::optimizers::EdgeOptimizersConfig;
use crate::edge::config::shard::{EDGE_CONFIG_FILE, EdgeConfig};
use crate::edge::read_view::build_segment_pool;
#[derive(Debug)]
pub struct EdgeShard {
path: PathBuf,
config: Arc<SaveOnDisk<EdgeConfig>>,
wal: Mutex<SerdeWal<CollectionUpdateOperations>>,
segments: LockedSegmentHolder,
segment_manifest: Option<SaveOnDisk<SegmentsManifest>>,
search_pool: Arc<rayon::ThreadPool>,
}
const WAL_PATH: &str = "wal";
impl EdgeShard {
pub fn new(path: &Path, config: EdgeConfig) -> OperationResult<Self> {
if has_existing_segments(path) {
return Err(OperationError::service_error(
"cannot create edge shard: path already contains segment data",
));
}
let wal_options = config.wal_options.clone().unwrap_or_default();
let (wal, segments_path) = ensure_dirs_and_open_wal(path, wal_options)?;
config.save(path)?;
let mut segments = SegmentHolder::default();
ensure_appendable_segment(&mut segments, path, &segments_path, &config)?;
let search_pool = build_segment_pool(
"edge-search",
config.search_thread_count(),
config.search_pool_core,
)?;
let config_path = path.join(EDGE_CONFIG_FILE);
let config = Arc::new(
SaveOnDisk::new(&config_path, config)
.map_err(|e| OperationError::service_error(e.to_string()))?,
);
let segment_manifest = init_segment_manifest(path, &segments)?;
Ok(Self {
path: path.into(),
config,
wal: parking_lot::Mutex::new(wal),
segments: LockedSegmentHolder::new(segments),
segment_manifest,
search_pool,
})
}
pub fn load(path: &Path, config: Option<EdgeConfig>) -> OperationResult<Self> {
let resolved = resolve_initial_config(path, config)?;
let wal_options = resolved
.as_ref()
.and_then(|c| c.wal_options.clone())
.unwrap_or_default();
let (wal, segments_path) = ensure_dirs_and_open_wal(path, wal_options)?;
let (mut segments, derived) = load_segments(&segments_path)?;
let config = match (resolved, derived) {
(Some(resolved), Some(derived)) => {
let merged = resolved.fill_unspecified_from(&derived);
merged
.check_compatible_with_segment_config(&derived.plain_segment_config())
.map_err(|err| {
OperationError::service_error(format!(
"config is incompatible with existing segments: {err}"
))
})?;
merged
}
(Some(resolved), None) => resolved,
(None, Some(derived)) => derived,
(None, None) => {
return Err(OperationError::service_error(
"edge config is not provided and no segments were loaded",
));
}
};
ensure_appendable_segment(&mut segments, path, &segments_path, &config)?;
let search_pool = build_segment_pool(
"edge-search",
config.search_thread_count(),
config.search_pool_core,
)?;
let config_path = path.join(EDGE_CONFIG_FILE);
let config = Arc::new(
SaveOnDisk::new(&config_path, config)
.map_err(|e| OperationError::service_error(e.to_string()))?,
);
let segment_manifest = init_segment_manifest(path, &segments)?;
Ok(Self {
path: path.into(),
config,
wal: parking_lot::Mutex::new(wal),
segments: LockedSegmentHolder::new(segments),
segment_manifest,
search_pool,
})
}
pub(crate) fn update_segment_manifest(&self) -> OperationResult<()> {
let Some(manifest) = &self.segment_manifest else {
return Ok(());
};
let rebuilt = {
let holder = self.segments.read();
SegmentsManifest::from_segment_holder(&holder)
};
manifest
.write_optional(|previous| {
let current = rebuilt.preserving(previous);
(*previous != current).then_some(current)
})
.map_err(|err| OperationError::service_error(err.to_string()))?;
Ok(())
}
pub fn config(&self) -> parking_lot::RwLockReadGuard<'_, EdgeConfig> {
self.config.read()
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn set_hnsw_config(&self, hnsw_config: crate::segment::types::HnswConfig) -> OperationResult<()> {
self.config
.write(|cfg| cfg.set_hnsw_config(hnsw_config))
.map_err(|e| OperationError::service_error(e.to_string()))
}
pub fn set_vector_hnsw_config(
&self,
vector_name: &str,
hnsw_config: crate::segment::types::HnswConfig,
) -> OperationResult<()> {
let mut mutation = Ok(());
self.config
.write_optional(|cfg| {
let mut updated = cfg.clone();
match updated.set_vector_hnsw_config(vector_name, hnsw_config) {
Ok(()) => Some(updated),
Err(e) => {
mutation = Err(e);
None
}
}
})
.map_err(|e| OperationError::service_error(e.to_string()))
.and(mutation)
}
pub fn set_optimizers_config(&self, optimizers: EdgeOptimizersConfig) -> OperationResult<()> {
self.config
.write(|cfg| cfg.set_optimizers_config(optimizers))
.map_err(|e| OperationError::service_error(e.to_string()))
}
pub fn flush(&self) -> OperationResult<()> {
self.wal
.lock()
.flush()
.map_err(|e| OperationError::service_error(format!("WAL flush failed: {e}")))?;
self.segments.read().flush_all(FlushMode::Sync, true)?;
Ok(())
}
}
impl Drop for EdgeShard {
fn drop(&mut self) {
if let Err(e) = self.flush() {
log::error!("EdgeShard flush during drop failed: {e}");
}
}
}
fn init_segment_manifest(
path: &Path,
segments: &SegmentHolder,
) -> OperationResult<Option<SaveOnDisk<SegmentsManifest>>> {
if !crate::common::flags::feature_flags().write_segment_manifest {
return Ok(None);
}
let manifest_path = segment_manifest_path(path);
let rebuilt = SegmentsManifest::from_segment_holder(segments);
let manifest = match fs::read(&manifest_path)
.ok()
.and_then(|bytes| serde_json::from_slice::<SegmentsManifest>(&bytes).ok())
{
Some(previous) => rebuilt.preserving(&previous),
None => rebuilt,
};
let manifest = SaveOnDisk::new(manifest_path, manifest)
.map_err(|err| OperationError::service_error(err.to_string()))?;
Ok(Some(manifest))
}
fn has_existing_segments(path: &Path) -> bool {
let segments_path = path.join(SEGMENTS_PATH);
let Ok(entries) = fs::read_dir(&segments_path) else {
return false;
};
for entry in entries.flatten() {
let p = entry.path();
if !p.is_dir() {
continue;
}
if p.file_name()
.and_then(|n| n.to_str())
.is_some_and(|n| n.starts_with('.'))
{
continue;
}
if normalize_segment_dir(&p).ok().flatten().is_some() {
return true;
}
}
false
}
fn ensure_dirs_and_open_wal(
path: &Path,
wal_options: WalOptions,
) -> OperationResult<(SerdeWal<CollectionUpdateOperations>, PathBuf)> {
let wal_path = path.join(WAL_PATH);
if !wal_path.exists() {
fs::create_dir(&wal_path).map_err(|err| {
OperationError::service_error(format!("failed to create WAL directory: {err}"))
})?;
}
let wal = SerdeWal::new(&wal_path, wal_options).map_err(|err| {
OperationError::service_error(format!("failed to open WAL {}: {err}", wal_path.display(),))
})?;
let segments_path = path.join(SEGMENTS_PATH);
if !segments_path.exists() {
fs::create_dir(&segments_path).map_err(|err| {
OperationError::service_error(format!("failed to create segments directory: {err}"))
})?;
}
Ok((wal, segments_path))
}
fn resolve_initial_config(
path: &Path,
config: Option<EdgeConfig>,
) -> OperationResult<Option<EdgeConfig>> {
let persisted = match EdgeConfig::load(path) {
Some(Ok(c)) => Some(c),
Some(Err(e)) => return Err(e),
None => None,
};
Ok(match (config, persisted) {
(Some(provided), Some(persisted)) => Some(provided.fill_unspecified_from(&persisted)),
(Some(provided), None) => Some(provided),
(None, persisted) => persisted,
})
}
pub(crate) fn scan_segment_dirs(segments_path: &Path) -> OperationResult<HashMap<Uuid, PathBuf>> {
let segments_dir = fs::read_dir(segments_path).map_err(|err| {
OperationError::service_error(format!("failed to read segments directory: {err}"))
})?;
let mut result = HashMap::new();
for entry in segments_dir {
let entry = entry.map_err(|err| {
OperationError::service_error(format!(
"failed to read entry in segments directory: {err}",
))
})?;
let segment_path = entry.path();
if !segment_path.is_dir() {
log::warn!(
"Skipping non-directory segment entry {}",
segment_path.display(),
);
continue;
}
if segment_path
.file_name()
.and_then(|n| n.to_str())
.is_some_and(|n| n.starts_with('.'))
{
log::warn!(
"Skipping hidden segment directory {}",
segment_path.display(),
);
continue;
}
let Some((segment_path, segment_uuid)) = normalize_segment_dir(&segment_path)? else {
continue;
};
result.insert(segment_uuid, segment_path);
}
Ok(result)
}
fn load_segments(segments_path: &Path) -> OperationResult<(SegmentHolder, Option<EdgeConfig>)> {
let mut segments = SegmentHolder::default();
let mut derived: Option<EdgeConfig> = None;
let mut segment_dirs: Vec<_> = scan_segment_dirs(segments_path)?.into_iter().collect();
segment_dirs.sort_unstable_by_key(|(segment_uuid, _)| *segment_uuid);
for (segment_uuid, segment_path) in segment_dirs {
let mut segment = load_segment(&segment_path, segment_uuid, None, &AtomicBool::new(false))
.map_err(|err| {
OperationError::service_error(format!(
"failed to load segment {}: {err}",
segment_path.display(),
))
})?;
let segment_cfg = segment.config();
if let Some(acc) = derived.as_ref() {
acc.check_compatible_with_segment_config(segment_cfg)
.map_err(|err| {
OperationError::service_error(format!(
"segment {} is incompatible with previously loaded segments: {err}",
segment_path.display(),
))
})?;
}
derived = Some(EdgeConfig::fold_from_segment_config(derived, segment_cfg));
segment.check_consistency_and_repair().map_err(|err| {
OperationError::service_error(format!(
"failed to repair segment {}: {err}",
segment_path.display(),
))
})?;
segments.add_new(segment);
}
Ok((segments, derived))
}
fn ensure_appendable_segment(
segments: &mut SegmentHolder,
path: &Path,
segments_path: &Path,
config: &EdgeConfig,
) -> OperationResult<()> {
if segments.has_appendable_segment() {
return Ok(());
}
let payload_index_schema_path = path.join(PAYLOAD_INDEX_CONFIG_FILE);
let payload_index_schema = SaveOnDisk::load_or_init_default(&payload_index_schema_path)
.map_err(|err| {
OperationError::service_error(format!(
"failed to initialize payload index schema file {}: {err}",
payload_index_schema_path.display(),
))
})?;
segments.create_appendable_segment(
segments_path,
config.plain_segment_config(),
Arc::new(payload_index_schema),
None,
)?;
debug_assert!(segments.has_appendable_segment());
Ok(())
}