Skip to main content

alopex_embedded/
lib.rs

1//! A user-friendly embedded API for the AlopexDB key-value store.
2
3#![deny(missing_docs)]
4
5/// Catalog metadata API (in-memory, primarily for Python bindings).
6pub mod catalog;
7pub mod catalog_api;
8mod cluster_state;
9pub mod columnar_api;
10mod dataframe_api;
11pub mod options;
12pub mod owned_session;
13pub mod owned_sql;
14mod sql_api;
15mod txn_manager;
16
17pub use crate::catalog::{CachedTableInfo, Catalog};
18pub use crate::catalog_api::{
19    CatalogInfo, CatalogManifestCatalog, CatalogManifestColumn, CatalogManifestCompression,
20    CatalogManifestDataSourceFormat, CatalogManifestDataType, CatalogManifestDelta,
21    CatalogManifestIndex, CatalogManifestIndexMethod, CatalogManifestNamespace,
22    CatalogManifestRowIdMode, CatalogManifestStorage, CatalogManifestStorageType,
23    CatalogManifestTable, CatalogManifestTableType, CatalogManifestVectorMetric, ColumnDefinition,
24    ColumnInfo, CreateCatalogRequest, CreateNamespaceRequest, CreateTableRequest, IndexInfo,
25    NamespaceInfo, StorageInfo, TableInfo, CATALOG_MANIFEST_DELTA_FORMAT,
26};
27pub use crate::columnar_api::{
28    ColumnarIndexInfo, ColumnarIndexType, ColumnarRowIterator, EmbeddedConfig, StorageMode,
29};
30pub use crate::options::DatabaseOptions;
31pub use crate::owned_session::{EmbeddedOwnedSessionFactory, OwnedEmbeddedTransaction};
32pub use crate::owned_sql::{OwnedSqlRowOutcome, OwnedSqlStreamPlan};
33pub use crate::sql_api::{SqlStreamingResult, StreamingQueryResult, StreamingRows};
34pub use crate::txn_manager::{TransactionInfo, TransactionManager};
35pub use alopex_dataframe::{DataFrame, JoinKeys, JoinType, SortOptions};
36pub use alopex_sql::{DataSourceFormat, TableType};
37/// `Database::execute_sql()` / `Transaction::execute_sql()` の返却型。
38pub type SqlResult = alopex_sql::SqlResult;
39pub use alopex_core::kv::{OwnedReadOptions, OwnedReadSession, OwnedTransactionSession};
40use alopex_core::vector::hnsw::{HnswTransactionState, SearchStats as HnswSearchStats};
41use alopex_core::{
42    columnar::{
43        kvs_bridge::ColumnarKvsBridge, memory::InMemorySegmentStore, segment_v2::SegmentConfigV2,
44    },
45    kv::any::AnyKVTransaction,
46    kv::memory::MemoryKV,
47    kv::{AnyKV, RangeChangeJournalCapability},
48    validate_dimensions, HnswIndex, KVStore, KVTransaction, Key, LargeValueKind, LargeValueMeta,
49    LargeValueReader, LargeValueWriter, StorageFactory, VectorType, DEFAULT_CHUNK_SIZE,
50};
51pub use alopex_core::{HnswConfig, HnswSearchResult, HnswStats, MemoryStats, Metric, TxnMode};
52/// Streaming query row iterator for FR-7 compliance.
53pub use alopex_sql::executor::QueryRowIterator;
54use alopex_sql::storage::LocalRangeChangeJournal;
55use std::collections::HashMap;
56use std::convert::TryInto;
57use std::fs;
58use std::path::{Path, PathBuf};
59use std::result;
60use std::sync::atomic::{AtomicU64, Ordering};
61use std::sync::{Arc, RwLock};
62
63/// A convenience `Result` type for database operations.
64pub type Result<T> = result::Result<T, Error>;
65
66/// The error type for embedded database operations.
67#[derive(Debug, thiserror::Error)]
68pub enum Error {
69    /// An error from the underlying core storage engine.
70    #[error("core error: {0}")]
71    Core(#[from] alopex_core::Error),
72    /// A requested fenced storage read point is unavailable, expired, or not
73    /// yet readable on this backend.
74    #[error("{0}")]
75    ReadAt(#[from] alopex_core::ReadAtError),
76    /// An error from the SQL execution pipeline.
77    #[error("{0}")]
78    Sql(#[from] alopex_sql::SqlError),
79    /// An error from DataFrame operations.
80    #[error("{0}")]
81    DataFrame(#[from] alopex_dataframe::DataFrameError),
82    /// The transaction has already been completed and cannot be used.
83    #[error("transaction is completed")]
84    TxnCompleted,
85    /// Catalog が見つかりません。
86    #[error("カタログが見つかりません: {0}")]
87    CatalogNotFound(String),
88    /// Catalog が既に存在します。
89    #[error("カタログは既に存在します: {0}")]
90    CatalogAlreadyExists(String),
91    /// Catalog が空ではありません。
92    #[error("カタログが空ではありません: {0}")]
93    CatalogNotEmpty(String),
94    /// Namespace が見つかりません。
95    #[error("ネームスペースが見つかりません: {0}.{1}")]
96    NamespaceNotFound(String, String),
97    /// Namespace が既に存在します。
98    #[error("ネームスペースは既に存在します: {0}.{1}")]
99    NamespaceAlreadyExists(String, String),
100    /// Namespace が空ではありません。
101    #[error("ネームスペースが空ではありません: {0}.{1}")]
102    NamespaceNotEmpty(String, String),
103    /// The requested table was not found or is invalid.
104    #[error("table not found: {0}")]
105    TableNotFound(String),
106    /// Table が既に存在します。
107    #[error("テーブルは既に存在します: {0}")]
108    TableAlreadyExists(String),
109    /// Index が見つかりません。
110    #[error("インデックスが見つかりません: {0}")]
111    IndexNotFound(String),
112    /// default オブジェクトは削除できません。
113    #[error("default オブジェクトは削除できません: {0}")]
114    CannotDeleteDefault(String),
115    /// Managed テーブルにはスキーマが必要です。
116    #[error("managed テーブルにはスキーマが必要です")]
117    SchemaRequired,
118    /// External テーブルには storage_root が必要です。
119    #[error("external テーブルには storage_root が必要です")]
120    StorageRootRequired,
121    /// トランザクションは read-only です。
122    #[error("トランザクションは読み取り専用です")]
123    TxnReadOnly,
124    /// Transaction ID is invalid or missing.
125    #[error("invalid transaction id: {0}")]
126    InvalidTransactionId(String),
127    /// The operation requires in-memory columnar mode.
128    #[error("not in in-memory columnar mode")]
129    NotInMemoryMode,
130    /// The requested data source format is not supported.
131    #[error("unsupported data source format: {0}")]
132    UnsupportedDataSourceFormat(String),
133    /// The catalog store lock was poisoned.
134    #[error("catalog lock poisoned")]
135    CatalogLockPoisoned,
136    /// The embedded cluster state lock was poisoned.
137    #[error("cluster state lock poisoned")]
138    ClusterStateLockPoisoned,
139}
140
141impl Error {
142    /// SQL エラーの場合はエラーコード(例: `ALOPEX-S003`)を返す。
143    pub fn sql_error_code(&self) -> Option<&'static str> {
144        match self {
145            Self::Sql(err) => Some(err.code()),
146            _ => None,
147        }
148    }
149}
150
151/// The main database object.
152pub struct Database {
153    /// The underlying key-value store.
154    pub(crate) store: Arc<AnyKV>,
155    pub(crate) sql_catalog: Arc<RwLock<alopex_sql::catalog::PersistentCatalog<AnyKV>>>,
156    pub(crate) hnsw_cache: RwLock<HashMap<String, Arc<HnswIndex>>>,
157    pub(crate) vector_cache: RwLock<Option<HashMap<Key, CachedVector>>>,
158    /// Table info cache for scan/write operations.
159    pub(crate) table_info_cache: RwLock<HashMap<String, CachedTableInfo>>,
160    /// Cache epoch for invalidation (incremented on DDL operations).
161    pub(crate) table_info_cache_epoch: AtomicU64,
162    pub(crate) columnar_mode: StorageMode,
163    pub(crate) columnar_bridge: ColumnarKvsBridge,
164    pub(crate) columnar_memory: Option<InMemorySegmentStore>,
165    pub(crate) segment_config: SegmentConfigV2,
166    pub(crate) cluster_state: RwLock<cluster_state::EmbeddedClusterState>,
167}
168
169pub(crate) fn disk_data_dir_path(path: &Path) -> std::path::PathBuf {
170    if path.extension().is_some_and(|e| e == "alopex") {
171        // v0.1 file-mode はディレクトリに WAL/SSTable を持つため、`.alopex` の横に
172        // sidecar ディレクトリを作ってそこへ格納する。
173        path.with_extension("alopex.d")
174    } else {
175        path.to_path_buf()
176    }
177}
178
179#[cfg(not(target_arch = "wasm32"))]
180fn read_file_version_from_storage(path: &Path) -> alopex_core::storage::format::FileVersion {
181    use alopex_core::storage::format::{FileHeader, FileVersion, HEADER_SIZE};
182    use std::io::Read;
183
184    let Some(file_path) = resolve_format_file_path(path) else {
185        return FileVersion::CURRENT;
186    };
187
188    let mut header_bytes = [0u8; HEADER_SIZE];
189    let Ok(mut file) = fs::File::open(file_path) else {
190        return FileVersion::CURRENT;
191    };
192    if file.read_exact(&mut header_bytes).is_err() {
193        return FileVersion::CURRENT;
194    }
195    match FileHeader::from_bytes(&header_bytes) {
196        Ok(header) => header.version,
197        Err(_) => FileVersion::CURRENT,
198    }
199}
200
201#[cfg(target_arch = "wasm32")]
202fn read_file_version_from_storage(_path: &Path) -> alopex_core::storage::format::FileVersion {
203    alopex_core::storage::format::FileVersion::CURRENT
204}
205
206#[cfg(not(target_arch = "wasm32"))]
207fn resolve_format_file_path(path: &Path) -> Option<PathBuf> {
208    if path.is_file() {
209        return Some(path.to_path_buf());
210    }
211
212    if path.is_dir() {
213        if let Some(ext) = path.extension() {
214            if ext == "d" {
215                let candidate = path.with_extension("");
216                if candidate.is_file() {
217                    return Some(candidate);
218                }
219            }
220        }
221
222        if let Ok(entries) = fs::read_dir(path) {
223            for entry in entries.flatten() {
224                let entry_path = entry.path();
225                if entry_path.extension().is_some_and(|ext| ext == "alopex") && entry_path.is_file()
226                {
227                    return Some(entry_path);
228                }
229            }
230        }
231    }
232
233    None
234}
235
236impl Database {
237    /// Opens a database at the specified path.
238    pub fn open(path: &Path) -> Result<Self> {
239        let data_dir = disk_data_dir_path(path);
240        let store = StorageFactory::create(alopex_core::StorageMode::Disk {
241            path: data_dir,
242            config: None,
243        })
244        .map_err(Error::Core)?;
245        let mut db = Self::init(store, StorageMode::Disk, None, SegmentConfigV2::default());
246        db.load_sql_catalog()?;
247        Ok(db)
248    }
249
250    /// Creates a new, purely in-memory (transient) database.
251    pub fn new() -> Self {
252        let store = AnyKV::Memory(MemoryKV::new());
253        Self::init(
254            store,
255            StorageMode::InMemory,
256            None,
257            SegmentConfigV2::default(),
258        )
259    }
260
261    /// Opens a database in in-memory mode with default options.
262    pub fn open_in_memory() -> Result<Self> {
263        Self::open_in_memory_with_options(DatabaseOptions::in_memory())
264    }
265
266    /// Opens a database in in-memory mode with the given options.
267    pub fn open_in_memory_with_options(opts: DatabaseOptions) -> Result<Self> {
268        if !opts.memory_mode() {
269            return Err(Error::Core(alopex_core::Error::InvalidFormat(
270                "memory_mode must be enabled for in-memory open".into(),
271            )));
272        }
273        let store = StorageFactory::create(opts.to_storage_mode(None)).map_err(Error::Core)?;
274        let mut db = Self::init(
275            store,
276            StorageMode::InMemory,
277            opts.memory_limit(),
278            SegmentConfigV2::default(),
279        );
280        db.load_sql_catalog()?;
281        Ok(db)
282    }
283
284    /// Opens a database from a URI string.
285    ///
286    /// Supported URI schemes:
287    /// - `file://path` or bare path: Local filesystem
288    /// - `s3://bucket/prefix`: S3-compatible storage (requires `s3` feature)
289    ///
290    /// # Example
291    ///
292    /// ```ignore
293    /// // Local path
294    /// let db = Database::open_with_uri("/path/to/db")?;
295    ///
296    /// // S3 URI (requires s3 feature and credentials)
297    /// let db = Database::open_with_uri("s3://my-bucket/data")?;
298    /// ```
299    pub fn open_with_uri(uri: &str) -> Result<Self> {
300        // Check for S3 URI
301        if uri.starts_with("s3://") {
302            #[cfg(feature = "s3")]
303            {
304                return Self::open_s3(uri);
305            }
306            #[cfg(not(feature = "s3"))]
307            {
308                return Err(Error::Core(alopex_core::Error::InvalidFormat(
309                    "S3 support requires the 's3' feature".into(),
310                )));
311            }
312        }
313
314        // Strip file:// prefix if present
315        let path = if let Some(stripped) = uri.strip_prefix("file://") {
316            stripped
317        } else {
318            uri
319        };
320
321        Self::open(Path::new(path))
322    }
323
324    /// Opens a database backed by S3 storage.
325    ///
326    /// This method downloads data from S3 to a local cache, operates on it
327    /// using LsmKV, and syncs changes back to S3 on flush/close.
328    ///
329    /// # Arguments
330    ///
331    /// * `uri` - S3 URI in the format `s3://bucket/prefix`
332    ///
333    /// # Environment Variables
334    ///
335    /// Required:
336    /// * `AWS_ACCESS_KEY_ID` - AWS access key
337    /// * `AWS_SECRET_ACCESS_KEY` - AWS secret key
338    ///
339    /// Optional:
340    /// * `AWS_REGION` - AWS region (default: us-east-1)
341    /// * `AWS_ENDPOINT_URL` - Custom endpoint for S3-compatible services
342    #[cfg(feature = "s3")]
343    pub fn open_s3(uri: &str) -> Result<Self> {
344        let s3_config = alopex_core::S3Config::from_uri(uri).map_err(Error::Core)?;
345        let store = StorageFactory::create(alopex_core::StorageMode::S3 { config: s3_config })
346            .map_err(Error::Core)?;
347        let mut db = Self::init(store, StorageMode::Disk, None, SegmentConfigV2::default());
348        db.load_sql_catalog()?;
349        Ok(db)
350    }
351
352    pub(crate) fn init(
353        store: AnyKV,
354        columnar_mode: StorageMode,
355        memory_limit: Option<usize>,
356        segment_config: SegmentConfigV2,
357    ) -> Self {
358        let store = Arc::new(store);
359        let sql_catalog = Arc::new(RwLock::new(alopex_sql::catalog::PersistentCatalog::new(
360            store.clone(),
361        )));
362        let columnar_bridge = ColumnarKvsBridge::new(store.clone());
363        let columnar_memory = if matches!(columnar_mode, StorageMode::InMemory) {
364            Some(InMemorySegmentStore::new(memory_limit.map(|v| v as u64)))
365        } else {
366            None
367        };
368
369        Self {
370            store,
371            sql_catalog,
372            hnsw_cache: RwLock::new(HashMap::new()),
373            vector_cache: RwLock::new(None),
374            table_info_cache: RwLock::new(HashMap::new()),
375            table_info_cache_epoch: AtomicU64::new(0),
376            columnar_mode,
377            columnar_bridge,
378            columnar_memory,
379            segment_config,
380            cluster_state: RwLock::new(cluster_state::EmbeddedClusterState::default()),
381        }
382    }
383
384    /// Returns the current cluster status owned by this database instance.
385    pub fn cluster_status_snapshot(&self) -> Result<alopex_cluster::ClusterStatusSnapshot> {
386        let state = self
387            .cluster_state
388            .read()
389            .map_err(|_| Error::ClusterStateLockPoisoned)?;
390        Ok(state.status_snapshot(self.table_info_cache_epoch()))
391    }
392
393    /// Returns the latest routing decision produced by this database instance.
394    pub fn routing_diagnostics(&self) -> Result<alopex_cluster::RoutingDiagnostics> {
395        let state = self
396            .cluster_state
397            .read()
398            .map_err(|_| Error::ClusterStateLockPoisoned)?;
399        Ok(state.routing_diagnostics(self.table_info_cache_epoch()))
400    }
401
402    pub(crate) fn record_routing<C: alopex_sql::Catalog + ?Sized>(
403        &self,
404        catalog: &C,
405        statement: &alopex_sql::Statement,
406        statement_index: usize,
407    ) {
408        let Ok(mut state) = self.cluster_state.write() else {
409            return;
410        };
411        state.record_routing(
412            catalog,
413            statement,
414            statement_index,
415            self.table_info_cache_epoch(),
416        );
417    }
418
419    fn load_sql_catalog(&mut self) -> Result<()> {
420        use alopex_sql::catalog::CatalogError;
421
422        let loaded = match alopex_sql::catalog::PersistentCatalog::load(self.store.clone()) {
423            Ok(catalog) => catalog,
424            Err(CatalogError::Kv(alopex_core::Error::NotFound)) => {
425                alopex_sql::catalog::PersistentCatalog::new(self.store.clone())
426            }
427            Err(other) => return Err(Error::Sql(other.into())),
428        };
429
430        self.sql_catalog = Arc::new(RwLock::new(loaded));
431        Ok(())
432    }
433
434    fn hnsw_cache_get(&self, name: &str) -> Option<Arc<HnswIndex>> {
435        let cache = self.hnsw_cache.read().expect("hnsw cache lock poisoned");
436        cache.get(name).cloned()
437    }
438
439    fn hnsw_cache_insert(&self, name: &str, index: HnswIndex) -> Arc<HnswIndex> {
440        let index = Arc::new(index);
441        let mut cache = self.hnsw_cache.write().expect("hnsw cache lock poisoned");
442        cache.insert(name.to_string(), Arc::clone(&index));
443        index
444    }
445
446    fn hnsw_cache_remove(&self, name: &str) {
447        let mut cache = self.hnsw_cache.write().expect("hnsw cache lock poisoned");
448        cache.remove(name);
449    }
450
451    /// Returns the current table info cache epoch.
452    pub fn table_info_cache_epoch(&self) -> u64 {
453        self.table_info_cache_epoch.load(Ordering::Relaxed)
454    }
455
456    /// Retrieves cached table info if available and epoch matches.
457    pub fn get_cached_table_info(
458        &self,
459        catalog_name: &str,
460        namespace_name: &str,
461        table_name: &str,
462    ) -> Option<CachedTableInfo> {
463        let cache = self
464            .table_info_cache
465            .read()
466            .expect("table info cache lock poisoned");
467        let key = format!("{}.{}.{}", catalog_name, namespace_name, table_name);
468        cache.get(&key).cloned()
469    }
470
471    /// Stores table info in the cache.
472    pub fn cache_table_info(
473        &self,
474        catalog_name: &str,
475        namespace_name: &str,
476        table_name: &str,
477        info: CachedTableInfo,
478    ) {
479        let mut cache = self
480            .table_info_cache
481            .write()
482            .expect("table info cache lock poisoned");
483        let key = format!("{}.{}.{}", catalog_name, namespace_name, table_name);
484        cache.insert(key, info);
485    }
486
487    /// Invalidates the table info cache (called on DDL operations).
488    pub fn invalidate_table_info_cache(&self) {
489        self.table_info_cache_epoch.fetch_add(1, Ordering::Relaxed);
490        let mut cache = self
491            .table_info_cache
492            .write()
493            .expect("table info cache lock poisoned");
494        cache.clear();
495    }
496
497    /// Flushes the current in-memory data to an SSTable on disk (beta).
498    pub fn flush(&self) -> Result<()> {
499        self.store.flush().map_err(Error::Core)
500    }
501
502    /// Returns the file format version supported by the embedded engine.
503    pub fn file_format_version(&self) -> alopex_core::storage::format::FileVersion {
504        use alopex_core::storage::format::FileVersion;
505
506        self.store
507            .file_format_storage_dir()
508            .map_or(FileVersion::CURRENT, read_file_version_from_storage)
509    }
510
511    /// Returns current memory usage statistics (in-memory KV only).
512    pub fn memory_usage(&self) -> Option<MemoryStats> {
513        match self.store.as_ref() {
514            AnyKV::Memory(kv) => Some(kv.memory_stats()),
515            _ => None,
516        }
517    }
518
519    /// Persists the current in-memory database to disk atomically.
520    ///
521    /// `wal_path` は「データディレクトリ」として扱う(file-mode)。
522    pub fn persist_to_disk(&self, wal_path: &Path) -> Result<()> {
523        if !matches!(self.store.as_ref(), AnyKV::Memory(_)) {
524            return Err(Error::NotInMemoryMode);
525        }
526        let data_dir = disk_data_dir_path(wal_path);
527        if wal_path.exists() || data_dir.exists() {
528            return Err(Error::Core(alopex_core::Error::PathExists(
529                wal_path.to_path_buf(),
530            )));
531        }
532
533        let tmp_dir = data_dir.with_extension("tmp");
534        if tmp_dir.exists() {
535            return Err(Error::Core(alopex_core::Error::PathExists(tmp_dir)));
536        }
537
538        let snapshot = self.snapshot_pairs()?;
539        let write_result = (|| -> Result<()> {
540            let store = StorageFactory::create(alopex_core::StorageMode::Disk {
541                path: tmp_dir.clone(),
542                config: None,
543            })
544            .map_err(Error::Core)?;
545
546            let mut txn = store.begin(TxnMode::ReadWrite).map_err(Error::Core)?;
547            for (key, value) in snapshot {
548                txn.put(key, value).map_err(Error::Core)?;
549            }
550            txn.commit_self().map_err(Error::Core)?;
551
552            Ok(())
553        })();
554
555        if let Err(e) = write_result {
556            let _ = fs::remove_dir_all(&tmp_dir);
557            return Err(e);
558        }
559
560        fs::rename(&tmp_dir, &data_dir).map_err(|e| Error::Core(e.into()))?;
561        if wal_path.extension().is_some_and(|e| e == "alopex") {
562            // `.alopex` パスを渡した場合は、存在確認用のマーカーを作る。
563            let _ = fs::OpenOptions::new()
564                .create_new(true)
565                .write(true)
566                .open(wal_path);
567        }
568        Ok(())
569    }
570
571    /// Creates a fully in-memory clone of the current database.
572    pub fn clone_to_memory(&self) -> Result<Self> {
573        let snapshot = self.snapshot_pairs()?;
574        let cloned = Database::open_in_memory()?;
575        if snapshot.is_empty() {
576            return Ok(cloned);
577        }
578
579        let mut txn = cloned.begin(TxnMode::ReadWrite)?;
580        for (key, value) in snapshot {
581            txn.put(&key, &value)?;
582        }
583        txn.commit()?;
584        Ok(cloned)
585    }
586
587    /// Clears all data while keeping the database usable.
588    pub fn clear(&self) -> Result<()> {
589        let keys: Vec<Key> = self.snapshot_pairs()?.into_iter().map(|(k, _)| k).collect();
590        if keys.is_empty() {
591            return Ok(());
592        }
593        let mut txn = self.begin(TxnMode::ReadWrite)?;
594        for key in keys {
595            txn.delete(&key)?;
596        }
597        txn.commit()
598    }
599
600    /// Updates the memory limit in bytes for the underlying in-memory store.
601    pub fn set_memory_limit(&self, bytes: Option<usize>) {
602        if let AnyKV::Memory(kv) = self.store.as_ref() {
603            kv.txn_manager().set_memory_limit(bytes);
604        }
605    }
606
607    /// Returns a read-only snapshot of all key-value pairs.
608    pub fn snapshot(&self) -> Vec<(Key, Vec<u8>)> {
609        self.snapshot_pairs().unwrap_or_default()
610    }
611
612    fn snapshot_pairs(&self) -> Result<Vec<(Key, Vec<u8>)>> {
613        let mut txn = self.store.begin(TxnMode::ReadOnly).map_err(Error::Core)?;
614        let pairs: Vec<(Key, Vec<u8>)> = txn.scan_prefix(b"").map_err(Error::Core)?.collect();
615        txn.commit_self().map_err(Error::Core)?;
616        Ok(pairs)
617    }
618
619    /// HNSW インデックスを作成し、永続化する。
620    pub fn create_hnsw_index(&self, name: &str, config: HnswConfig) -> Result<()> {
621        let mut txn = self.store.begin(TxnMode::ReadWrite).map_err(Error::Core)?;
622        let index = HnswIndex::create(name, config).map_err(Error::Core)?;
623        index.save(&mut txn).map_err(Error::Core)?;
624        txn.commit_self().map_err(Error::Core)?;
625        self.hnsw_cache_insert(name, index);
626        Ok(())
627    }
628
629    /// HNSW インデックスを削除する。
630    pub fn drop_hnsw_index(&self, name: &str) -> Result<()> {
631        let mut txn = self.store.begin(TxnMode::ReadWrite).map_err(Error::Core)?;
632        let index = HnswIndex::load(name, &mut txn).map_err(Error::Core)?;
633        index.drop(&mut txn).map_err(Error::Core)?;
634        txn.commit_self().map_err(Error::Core)?;
635        self.hnsw_cache_remove(name);
636        Ok(())
637    }
638
639    /// HNSW 統計情報を取得する。
640    pub fn get_hnsw_stats(&self, name: &str) -> Result<HnswStats> {
641        if let Some(index) = self.hnsw_cache_get(name) {
642            return Ok(index.stats());
643        }
644        let mut txn = self.store.begin(TxnMode::ReadOnly).map_err(Error::Core)?;
645        let index = HnswIndex::load(name, &mut txn).map_err(Error::Core)?;
646        let stats = index.stats();
647        self.hnsw_cache_insert(name, index);
648        Ok(stats)
649    }
650
651    /// HNSW インデックスをコンパクションし、結果を返す。
652    pub fn compact_hnsw_index(&self, name: &str) -> Result<alopex_core::vector::CompactionResult> {
653        let mut txn = self.store.begin(TxnMode::ReadWrite).map_err(Error::Core)?;
654        let mut index = HnswIndex::load(name, &mut txn).map_err(Error::Core)?;
655        let result = index.compact().map_err(Error::Core)?;
656        index.save(&mut txn).map_err(Error::Core)?;
657        txn.commit_self().map_err(Error::Core)?;
658        self.hnsw_cache_insert(name, index);
659        Ok(result)
660    }
661
662    /// HNSW インデックスに検索を行う。
663    pub fn search_hnsw(
664        &self,
665        name: &str,
666        query: &[f32],
667        k: usize,
668        ef_search: Option<usize>,
669    ) -> Result<(Vec<HnswSearchResult>, HnswSearchStats)> {
670        let profile = std::env::var_os("ALOPEX_PROFILE_HNSW").is_some();
671        let total_start = if profile {
672            Some(std::time::Instant::now())
673        } else {
674            None
675        };
676        if let Some(index) = self.hnsw_cache_get(name) {
677            let search_start = if profile {
678                Some(std::time::Instant::now())
679            } else {
680                None
681            };
682            let result = index.search(query, k, ef_search).map_err(Error::Core)?;
683            if let (true, Some(total_start), Some(search_start)) =
684                (profile, total_start, search_start)
685            {
686                let search_time = search_start.elapsed();
687                let total_time = total_start.elapsed();
688                eprintln!(
689                    "alopex.hnsw_search cache=hit name={} k={} ef_search={:?} search_ms={:.2} total_ms={:.2}",
690                    name,
691                    k,
692                    ef_search,
693                    search_time.as_secs_f64() * 1000.0,
694                    total_time.as_secs_f64() * 1000.0
695                );
696            }
697            return Ok(result);
698        }
699
700        let load_start = if profile {
701            Some(std::time::Instant::now())
702        } else {
703            None
704        };
705        let mut txn = self.store.begin(TxnMode::ReadOnly).map_err(Error::Core)?;
706        let index = HnswIndex::load(name, &mut txn).map_err(Error::Core)?;
707        let load_time = load_start.map(|start| start.elapsed());
708        let index = self.hnsw_cache_insert(name, index);
709        let search_start = if profile {
710            Some(std::time::Instant::now())
711        } else {
712            None
713        };
714        let result = index.search(query, k, ef_search).map_err(Error::Core)?;
715        if let (true, Some(total_start), Some(search_start)) = (profile, total_start, search_start)
716        {
717            let search_time = search_start.elapsed();
718            let total_time = total_start.elapsed();
719            let load_time_ms = load_time
720                .map(|elapsed| elapsed.as_secs_f64() * 1000.0)
721                .unwrap_or(0.0);
722            eprintln!(
723                "alopex.hnsw_search cache=miss name={} k={} ef_search={:?} load_ms={:.2} search_ms={:.2} total_ms={:.2}",
724                name,
725                k,
726                ef_search,
727                load_time_ms,
728                search_time.as_secs_f64() * 1000.0,
729                total_time.as_secs_f64() * 1000.0
730            );
731        }
732        Ok(result)
733    }
734
735    /// Creates a chunked large value writer for opaque blobs (beta).
736    pub fn create_blob_writer(
737        &self,
738        path: &Path,
739        total_len: u64,
740        chunk_size: Option<u32>,
741    ) -> Result<LargeValueWriter> {
742        let meta = LargeValueMeta {
743            kind: LargeValueKind::Blob,
744            total_len,
745            chunk_size: chunk_size.unwrap_or(DEFAULT_CHUNK_SIZE),
746        };
747        LargeValueWriter::create(path, meta).map_err(Error::Core)
748    }
749
750    /// Creates a chunked large value writer for typed payloads (beta).
751    pub fn create_typed_writer(
752        &self,
753        path: &Path,
754        type_id: u16,
755        total_len: u64,
756        chunk_size: Option<u32>,
757    ) -> Result<LargeValueWriter> {
758        let meta = LargeValueMeta {
759            kind: LargeValueKind::Typed(type_id),
760            total_len,
761            chunk_size: chunk_size.unwrap_or(DEFAULT_CHUNK_SIZE),
762        };
763        LargeValueWriter::create(path, meta).map_err(Error::Core)
764    }
765
766    /// Opens a chunked large value reader (beta). Kind/type is read from the file header.
767    pub fn open_large_value(&self, path: &Path) -> Result<LargeValueReader> {
768        LargeValueReader::open(path).map_err(Error::Core)
769    }
770
771    /// Begins a new transaction.
772    pub fn begin(&self, mode: TxnMode) -> Result<Transaction<'_>> {
773        let mut txn = self.store.begin(mode).map_err(Error::Core)?;
774        let journal = if mode == TxnMode::ReadWrite
775            && self.store.range_change_journal_capability()
776                == RangeChangeJournalCapability::Supported
777        {
778            let scope = {
779                let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
780                sql_api::local_journal_scope(&*catalog)
781            };
782            Some(LocalRangeChangeJournal::capture(&mut txn, scope).map_err(Error::Core)?)
783        } else {
784            None
785        };
786        Ok(Transaction {
787            inner: Some(txn),
788            db: self,
789            hnsw_indices: HashMap::new(),
790            overlay: alopex_sql::catalog::CatalogOverlay::new(),
791            vector_cache_updates: HashMap::new(),
792            vector_cache_deletes: Vec::new(),
793            vector_cache_invalidated: false,
794            catalog_modified: false,
795            journal,
796        })
797    }
798}
799
800impl Default for Database {
801    fn default() -> Self {
802        Self::new()
803    }
804}
805
806/// A database transaction.
807pub struct Transaction<'a> {
808    inner: Option<AnyKVTransaction<'a>>,
809    db: &'a Database,
810    hnsw_indices: HashMap<String, (HnswIndex, alopex_core::vector::hnsw::HnswTransactionState)>,
811    overlay: alopex_sql::catalog::CatalogOverlay,
812    vector_cache_updates: HashMap<Key, CachedVector>,
813    vector_cache_deletes: Vec<Key>,
814    vector_cache_invalidated: bool,
815    /// Whether DDL operations were performed in this transaction.
816    pub(crate) catalog_modified: bool,
817    /// SQL row/index state captured before a read-write local transaction.
818    journal: Option<LocalRangeChangeJournal>,
819}
820
821/// A search result row containing key, metadata, and similarity score.
822#[derive(Debug, Clone, PartialEq)]
823pub struct SearchResult {
824    /// User key associated with the vector.
825    pub key: Key,
826    /// Opaque metadata payload stored alongside the vector.
827    pub metadata: Vec<u8>,
828    /// Similarity score for the query/vector pair.
829    pub score: f32,
830}
831
832const VECTOR_INDEX_KEY: &[u8] = b"__alopex_vector_index";
833
834impl<'a> Transaction<'a> {
835    pub(crate) fn catalog_overlay(&self) -> &alopex_sql::catalog::CatalogOverlay {
836        &self.overlay
837    }
838
839    pub(crate) fn catalog_overlay_mut(&mut self) -> &mut alopex_sql::catalog::CatalogOverlay {
840        &mut self.overlay
841    }
842
843    pub(crate) fn txn_mode(&self) -> Result<TxnMode> {
844        let txn = self.inner.as_ref().ok_or(Error::TxnCompleted)?;
845        Ok(txn.mode())
846    }
847    /// Retrieves the value for a given key.
848    pub fn get(&mut self, key: &[u8]) -> Result<Option<Vec<u8>>> {
849        self.inner_mut()?.get(&key.to_vec()).map_err(Error::Core)
850    }
851
852    /// Sets a value for a given key.
853    pub fn put(&mut self, key: &[u8], value: &[u8]) -> Result<()> {
854        self.vector_cache_invalidated = true;
855        self.vector_cache_updates.clear();
856        self.vector_cache_deletes.clear();
857        self.inner_mut()?
858            .put(key.to_vec(), value.to_vec())
859            .map_err(Error::Core)
860    }
861
862    /// Deletes a key-value pair.
863    pub fn delete(&mut self, key: &[u8]) -> Result<()> {
864        self.vector_cache_deletes.push(key.to_vec());
865        self.inner_mut()?.delete(key.to_vec()).map_err(Error::Core)
866    }
867
868    /// Scans all key-value pairs whose keys start with the given prefix.
869    ///
870    /// Returns an iterator over (key, value) pairs.
871    pub fn scan_prefix(
872        &mut self,
873        prefix: &[u8],
874    ) -> Result<Box<dyn Iterator<Item = (Key, Vec<u8>)> + '_>> {
875        self.inner_mut()?.scan_prefix(prefix).map_err(Error::Core)
876    }
877
878    /// HNSW にベクトルをステージング挿入/更新する。
879    pub fn upsert_to_hnsw(
880        &mut self,
881        index_name: &str,
882        key: &[u8],
883        vector: &[f32],
884        metadata: &[u8],
885    ) -> Result<()> {
886        self.ensure_write_txn()?;
887        let (index, state) = self.hnsw_entry_mut(index_name)?;
888        index
889            .upsert_staged(key, vector, metadata, state)
890            .map_err(Error::Core)
891    }
892
893    /// HNSW からキーをステージング削除する。
894    pub fn delete_from_hnsw(&mut self, index_name: &str, key: &[u8]) -> Result<bool> {
895        self.ensure_write_txn()?;
896        let (index, state) = self.hnsw_entry_mut(index_name)?;
897        index.delete_staged(key, state).map_err(Error::Core)
898    }
899
900    /// Upserts a vector and metadata under the provided key after validating dimensions and metric.
901    ///
902    /// A small internal index is maintained to enable scanning for similarity search.
903    pub fn upsert_vector(
904        &mut self,
905        key: &[u8],
906        metadata: &[u8],
907        vector: &[f32],
908        metric: Metric,
909    ) -> Result<()> {
910        if vector.is_empty() {
911            return Err(Error::Core(alopex_core::Error::InvalidFormat(
912                "vector cannot be empty".into(),
913            )));
914        }
915        let vt = VectorType::new(vector.len(), metric);
916        vt.validate(vector).map_err(Error::Core)?;
917
918        let payload = encode_vector_entry(vt, metadata, vector);
919        let txn = self.inner_mut()?;
920        txn.put(key.to_vec(), payload).map_err(Error::Core)?;
921
922        let mut keys = self.load_vector_index()?;
923        if !keys.iter().any(|k| k == key) {
924            keys.push(key.to_vec());
925            self.persist_vector_index(&keys)?;
926        }
927
928        let cached = cached_vector_from_entry(metric, metadata.to_vec(), vector.to_vec());
929        self.vector_cache_updates.insert(key.to_vec(), cached);
930        self.vector_cache_deletes.retain(|k| k != key);
931        Ok(())
932    }
933
934    /// Retrieves a vector stored under the given key.
935    ///
936    /// Returns `None` if the key does not exist. If the key exists but has a different
937    /// metric than specified, returns an error.
938    pub fn get_vector(&mut self, key: &[u8], metric: Metric) -> Result<Option<Vec<f32>>> {
939        let txn = self.inner_mut()?;
940        let key_vec = key.to_vec();
941        let Some(raw) = txn.get(&key_vec).map_err(Error::Core)? else {
942            return Ok(None);
943        };
944        let decoded = decode_vector_entry(&raw).map_err(Error::Core)?;
945        if decoded.metric != metric {
946            return Err(Error::Core(alopex_core::Error::UnsupportedMetric {
947                metric: metric.as_str().to_string(),
948            }));
949        }
950        Ok(Some(decoded.vector))
951    }
952
953    /// Retrieves vectors stored under the given keys in a single transaction call.
954    ///
955    /// Returned list order matches the input `keys` order. Each entry is:
956    /// - `None` if the key does not exist
957    /// - `Some(Vec<f32>)` if the key exists and metric matches
958    ///
959    /// If any existing entry has a different metric than specified, returns an error.
960    pub fn get_vectors(&mut self, keys: &[Key], metric: Metric) -> Result<Vec<Option<Vec<f32>>>> {
961        let txn = self.inner_mut()?;
962        let mut out = Vec::with_capacity(keys.len());
963        for key in keys {
964            let Some(raw) = txn.get(key).map_err(Error::Core)? else {
965                out.push(None);
966                continue;
967            };
968            let decoded = decode_vector_entry(&raw).map_err(Error::Core)?;
969            if decoded.metric != metric {
970                return Err(Error::Core(alopex_core::Error::UnsupportedMetric {
971                    metric: metric.as_str().to_string(),
972                }));
973            }
974            out.push(Some(decoded.vector));
975        }
976        Ok(out)
977    }
978
979    /// Executes a flat similarity search over stored vectors using the provided metric and query.
980    ///
981    /// The optional `filter_keys` restricts the scan to the given keys; otherwise the full
982    /// vector index is scanned. Results are sorted by descending score and truncated to `top_k`.
983    pub fn search_similar(
984        &mut self,
985        query_vector: &[f32],
986        metric: Metric,
987        top_k: usize,
988        filter_keys: Option<&[Key]>,
989    ) -> Result<Vec<SearchResult>> {
990        if top_k == 0 {
991            return Ok(Vec::new());
992        }
993
994        let profile = std::env::var_os("ALOPEX_PROFILE_SEARCH_SIMILAR").is_some();
995        let total_start = if profile {
996            Some(std::time::Instant::now())
997        } else {
998            None
999        };
1000        let query_norm_sq = query_vector.iter().map(|v| v * v).sum::<f32>();
1001        let query_norm = if matches!(metric, Metric::Cosine) {
1002            query_norm_sq.sqrt()
1003        } else {
1004            0.0
1005        };
1006        let inv_query_norm = if query_norm == 0.0 {
1007            0.0
1008        } else {
1009            1.0 / query_norm
1010        };
1011
1012        if filter_keys.is_none() && self.txn_mode()? == TxnMode::ReadOnly {
1013            let cache = self
1014                .db
1015                .vector_cache
1016                .read()
1017                .expect("vector cache lock poisoned");
1018            if let Some(cache) = cache.as_ref() {
1019                if cache.is_empty() {
1020                    return Ok(Vec::new());
1021                }
1022                let keys_len = cache.len();
1023                let mut rows = Vec::with_capacity(keys_len);
1024                let mut score_time = std::time::Duration::ZERO;
1025                for (key, cached) in cache.iter() {
1026                    if cached.metric != metric {
1027                        return Err(Error::Core(alopex_core::Error::UnsupportedMetric {
1028                            metric: metric.as_str().to_string(),
1029                        }));
1030                    }
1031                    validate_dimensions(cached.vector.len(), query_vector.len())
1032                        .map_err(Error::Core)?;
1033                    let score_start = if profile {
1034                        Some(std::time::Instant::now())
1035                    } else {
1036                        None
1037                    };
1038                    let dot = dot_product(query_vector, &cached.vector);
1039                    let score = match metric {
1040                        Metric::Cosine => {
1041                            if cached.inv_norm == 0.0 || inv_query_norm == 0.0 {
1042                                0.0
1043                            } else {
1044                                dot * cached.inv_norm * inv_query_norm
1045                            }
1046                        }
1047                        Metric::L2 => {
1048                            let dist_sq = query_norm_sq + cached.norm_sq - 2.0 * dot;
1049                            -dist_sq.sqrt()
1050                        }
1051                        Metric::InnerProduct => dot,
1052                    };
1053                    if let Some(score_start) = score_start {
1054                        score_time += score_start.elapsed();
1055                    }
1056                    rows.push(SearchResult {
1057                        key: key.clone(),
1058                        metadata: cached.metadata.clone(),
1059                        score,
1060                    });
1061                }
1062
1063                let rows_total = rows.len();
1064                let sort_start = if profile {
1065                    Some(std::time::Instant::now())
1066                } else {
1067                    None
1068                };
1069                if rows.len() > top_k {
1070                    rows.select_nth_unstable_by(top_k - 1, |a, b| {
1071                        b.score.total_cmp(&a.score).then_with(|| a.key.cmp(&b.key))
1072                    });
1073                    rows.truncate(top_k);
1074                }
1075                rows.sort_by(|a, b| b.score.total_cmp(&a.score).then_with(|| a.key.cmp(&b.key)));
1076                if let (true, Some(total_start), Some(sort_start)) =
1077                    (profile, total_start, sort_start)
1078                {
1079                    let sort_time = sort_start.elapsed();
1080                    let total_time = total_start.elapsed();
1081                    eprintln!(
1082                        "alopex.search_similar keys={} results={} top_k={} load_keys_ms={:.2} get_ms={:.2} decode_ms={:.2} score_ms={:.2} sort_ms={:.2} total_ms={:.2}",
1083                        keys_len,
1084                        rows_total,
1085                        top_k,
1086                        0.0,
1087                        0.0,
1088                        0.0,
1089                        score_time.as_secs_f64() * 1000.0,
1090                        sort_time.as_secs_f64() * 1000.0,
1091                        total_time.as_secs_f64() * 1000.0
1092                    );
1093                }
1094                return Ok(rows);
1095            }
1096        }
1097
1098        let (keys, load_keys_time) = if profile {
1099            let start = std::time::Instant::now();
1100            let keys = match filter_keys {
1101                Some(keys) => keys.to_vec(),
1102                None => self.load_vector_index()?,
1103            };
1104            (keys, start.elapsed())
1105        } else {
1106            let keys = match filter_keys {
1107                Some(keys) => keys.to_vec(),
1108                None => self.load_vector_index()?,
1109            };
1110            (keys, std::time::Duration::ZERO)
1111        };
1112        if keys.is_empty() {
1113            return Ok(Vec::new());
1114        }
1115
1116        let keys_len = keys.len();
1117        let mut rows = Vec::with_capacity(keys.len());
1118        let txn = self.inner_mut()?;
1119        let mut get_time = std::time::Duration::ZERO;
1120        let mut decode_time = std::time::Duration::ZERO;
1121        let mut score_time = std::time::Duration::ZERO;
1122        if profile {
1123            for key in keys {
1124                let get_start = std::time::Instant::now();
1125                let Some(raw) = txn.get(&key).map_err(Error::Core)? else {
1126                    get_time += get_start.elapsed();
1127                    continue;
1128                };
1129                get_time += get_start.elapsed();
1130                let decode_start = std::time::Instant::now();
1131                let decoded = decode_vector_entry_view(&raw).map_err(Error::Core)?;
1132                decode_time += decode_start.elapsed();
1133                if decoded.metric != metric {
1134                    return Err(Error::Core(alopex_core::Error::UnsupportedMetric {
1135                        metric: metric.as_str().to_string(),
1136                    }));
1137                }
1138                validate_dimensions(decoded.dim, query_vector.len()).map_err(Error::Core)?;
1139                let score_start = std::time::Instant::now();
1140                let score =
1141                    score_from_bytes(metric, query_vector, query_norm, decoded.vector_bytes)?;
1142                score_time += score_start.elapsed();
1143                rows.push(SearchResult {
1144                    key,
1145                    metadata: decoded.metadata,
1146                    score,
1147                });
1148            }
1149        } else {
1150            for key in keys {
1151                let Some(raw) = txn.get(&key).map_err(Error::Core)? else {
1152                    continue;
1153                };
1154                let decoded = decode_vector_entry_view(&raw).map_err(Error::Core)?;
1155                if decoded.metric != metric {
1156                    return Err(Error::Core(alopex_core::Error::UnsupportedMetric {
1157                        metric: metric.as_str().to_string(),
1158                    }));
1159                }
1160                validate_dimensions(decoded.dim, query_vector.len()).map_err(Error::Core)?;
1161                let score =
1162                    score_from_bytes(metric, query_vector, query_norm, decoded.vector_bytes)?;
1163                rows.push(SearchResult {
1164                    key,
1165                    metadata: decoded.metadata,
1166                    score,
1167                });
1168            }
1169        }
1170
1171        let rows_total = rows.len();
1172        let sort_start = if profile {
1173            Some(std::time::Instant::now())
1174        } else {
1175            None
1176        };
1177        if rows.len() > top_k {
1178            rows.select_nth_unstable_by(top_k - 1, |a, b| {
1179                b.score.total_cmp(&a.score).then_with(|| a.key.cmp(&b.key))
1180            });
1181            rows.truncate(top_k);
1182        }
1183        rows.sort_by(|a, b| b.score.total_cmp(&a.score).then_with(|| a.key.cmp(&b.key)));
1184        if let (true, Some(total_start), Some(sort_start)) = (profile, total_start, sort_start) {
1185            let sort_time = sort_start.elapsed();
1186            let total_time = total_start.elapsed();
1187            eprintln!(
1188                "alopex.search_similar keys={} results={} top_k={} load_keys_ms={:.2} get_ms={:.2} decode_ms={:.2} score_ms={:.2} sort_ms={:.2} total_ms={:.2}",
1189                keys_len,
1190                rows_total,
1191                top_k,
1192                load_keys_time.as_secs_f64() * 1000.0,
1193                get_time.as_secs_f64() * 1000.0,
1194                decode_time.as_secs_f64() * 1000.0,
1195                score_time.as_secs_f64() * 1000.0,
1196                sort_time.as_secs_f64() * 1000.0,
1197                total_time.as_secs_f64() * 1000.0
1198            );
1199        }
1200        Ok(rows)
1201    }
1202
1203    fn load_vector_index(&mut self) -> Result<Vec<Key>> {
1204        let txn = self.inner_mut()?;
1205        let Some(raw) = txn.get(&VECTOR_INDEX_KEY.to_vec()).map_err(Error::Core)? else {
1206            return Ok(Vec::new());
1207        };
1208        decode_index(&raw).map_err(Error::Core)
1209    }
1210
1211    fn persist_vector_index(&mut self, keys: &[Key]) -> Result<()> {
1212        let txn = self.inner_mut()?;
1213        let encoded = encode_index(keys)?;
1214        txn.put(VECTOR_INDEX_KEY.to_vec(), encoded)
1215            .map_err(Error::Core)
1216    }
1217
1218    /// Commits the transaction, applying all changes.
1219    pub fn commit(mut self) -> Result<()> {
1220        {
1221            let txn = self.inner.as_mut().ok_or(Error::TxnCompleted)?;
1222            for (index, state) in self.hnsw_indices.values_mut() {
1223                index.commit_staged(txn, state).map_err(Error::Core)?;
1224            }
1225            let mut catalog = self.db.sql_catalog.write().expect("catalog lock poisoned");
1226            catalog
1227                .persist_overlay(txn, &self.overlay)
1228                .map_err(|err| Error::Sql(err.into()))?;
1229            if let Some(journal) = self.journal.take() {
1230                journal.stage(txn).map_err(Error::Core)?;
1231            }
1232
1233            let vector_cache_invalidated = self.vector_cache_invalidated;
1234            let vector_cache_updates = std::mem::take(&mut self.vector_cache_updates);
1235            let vector_cache_deletes = std::mem::take(&mut self.vector_cache_deletes);
1236            if vector_cache_invalidated {
1237                let mut cache = self
1238                    .db
1239                    .vector_cache
1240                    .write()
1241                    .expect("vector cache lock poisoned");
1242                *cache = None;
1243            } else if !vector_cache_updates.is_empty() || !vector_cache_deletes.is_empty() {
1244                let needs_rebuild = {
1245                    let cache = self
1246                        .db
1247                        .vector_cache
1248                        .read()
1249                        .expect("vector cache lock poisoned");
1250                    cache.is_none()
1251                };
1252                if needs_rebuild {
1253                    let rebuilt = build_vector_cache_from_txn(txn).map_err(Error::Core)?;
1254                    let mut cache = self
1255                        .db
1256                        .vector_cache
1257                        .write()
1258                        .expect("vector cache lock poisoned");
1259                    *cache = Some(rebuilt);
1260                } else {
1261                    let mut cache = self
1262                        .db
1263                        .vector_cache
1264                        .write()
1265                        .expect("vector cache lock poisoned");
1266                    if let Some(cache) = cache.as_mut() {
1267                        for key in vector_cache_deletes {
1268                            cache.remove(&key);
1269                        }
1270                        for (key, cached) in vector_cache_updates {
1271                            cache.insert(key, cached);
1272                        }
1273                    }
1274                }
1275            }
1276        }
1277        let txn = self.inner.take().ok_or(Error::TxnCompleted)?;
1278        let hnsw_indices = std::mem::take(&mut self.hnsw_indices);
1279        txn.commit_self().map_err(Error::Core)?;
1280        if !hnsw_indices.is_empty() {
1281            let mut cache = self
1282                .db
1283                .hnsw_cache
1284                .write()
1285                .expect("hnsw cache lock poisoned");
1286            for (name, (index, _state)) in hnsw_indices {
1287                cache.insert(name, Arc::new(index));
1288            }
1289        }
1290
1291        // KV commit 成功後のみ、カタログにオーバーレイを適用する。
1292        let overlay = std::mem::take(&mut self.overlay);
1293        let catalog_modified = self.catalog_modified;
1294        let mut catalog = self.db.sql_catalog.write().expect("catalog lock poisoned");
1295        catalog.apply_overlay(overlay);
1296        drop(catalog); // Release lock before invalidating cache
1297                       // Invalidate table info cache only if DDL operations were performed
1298        if catalog_modified {
1299            self.db.invalidate_table_info_cache();
1300        }
1301        Ok(())
1302    }
1303
1304    /// トランザクションを消費せずにロールバックする(失敗時の再試行を可能にする)。
1305    pub fn rollback_in_place(&mut self) -> Result<()> {
1306        let txn = self.inner.as_mut().ok_or(Error::TxnCompleted)?;
1307        txn.rollback_in_place().map_err(Error::Core)?;
1308        for (index, state) in self.hnsw_indices.values_mut() {
1309            let _ = index.rollback(state);
1310        }
1311        self.hnsw_indices.clear();
1312        self.overlay = alopex_sql::catalog::CatalogOverlay::default();
1313        self.inner = None;
1314        Ok(())
1315    }
1316
1317    /// Rolls back the transaction, discarding all changes.
1318    pub fn rollback(mut self) -> Result<()> {
1319        if let Some(txn) = self.inner.take() {
1320            for (index, state) in self.hnsw_indices.values_mut() {
1321                let _ = index.rollback(state);
1322            }
1323            self.hnsw_indices.clear();
1324            txn.rollback_self().map_err(Error::Core)
1325        } else {
1326            Err(Error::TxnCompleted)
1327        }
1328    }
1329
1330    fn inner_mut(&mut self) -> Result<&mut AnyKVTransaction<'a>> {
1331        self.inner.as_mut().ok_or(Error::TxnCompleted)
1332    }
1333
1334    fn hnsw_entry_mut(&mut self, name: &str) -> Result<&mut (HnswIndex, HnswTransactionState)> {
1335        if !self.hnsw_indices.contains_key(name) {
1336            let index = {
1337                let txn = self.inner_mut()?;
1338                HnswIndex::load(name, txn).map_err(Error::Core)?
1339            };
1340            self.hnsw_indices
1341                .insert(name.to_string(), (index, HnswTransactionState::default()));
1342        }
1343        Ok(self.hnsw_indices.get_mut(name).unwrap())
1344    }
1345
1346    fn ensure_write_txn(&self) -> Result<()> {
1347        let txn = self.inner.as_ref().ok_or(Error::TxnCompleted)?;
1348        if txn.mode() != TxnMode::ReadWrite {
1349            return Err(Error::Core(alopex_core::Error::TxnReadOnly));
1350        }
1351        Ok(())
1352    }
1353}
1354
1355impl OwnedEmbeddedTransaction {
1356    /// Stage a vector and its metadata in an owned transaction.
1357    pub fn upsert_vector(
1358        &mut self,
1359        key: &[u8],
1360        metadata: &[u8],
1361        vector: &[f32],
1362        metric: Metric,
1363    ) -> Result<()> {
1364        if vector.is_empty() {
1365            return Err(Error::Core(alopex_core::Error::InvalidFormat(
1366                "vector cannot be empty".into(),
1367            )));
1368        }
1369        let vector_type = VectorType::new(vector.len(), metric);
1370        vector_type.validate(vector).map_err(Error::Core)?;
1371        let key = key.to_vec();
1372        let payload = encode_vector_entry(vector_type, metadata, vector);
1373
1374        self.session
1375            .with_transaction(|transaction| {
1376                transaction.put(key.clone(), payload)?;
1377                let mut keys = match transaction.get(&VECTOR_INDEX_KEY.to_vec())? {
1378                    Some(raw) => decode_index(&raw)?,
1379                    None => Vec::new(),
1380                };
1381                if !keys.iter().any(|entry| entry == &key) {
1382                    keys.push(key.clone());
1383                    transaction.put(VECTOR_INDEX_KEY.to_vec(), encode_index(&keys)?)?;
1384                }
1385                Ok(())
1386            })
1387            .map_err(Error::Core)?;
1388        self.vector_cache_invalidated = true;
1389        Ok(())
1390    }
1391
1392    /// Retrieve one vector from an owned transaction.
1393    pub fn get_vector(&mut self, key: &[u8], metric: Metric) -> Result<Option<Vec<f32>>> {
1394        self.session
1395            .with_transaction(|transaction| {
1396                let Some(raw) = transaction.get(&key.to_vec())? else {
1397                    return Ok(None);
1398                };
1399                let decoded = decode_vector_entry(&raw)?;
1400                if decoded.metric != metric {
1401                    return Err(alopex_core::Error::UnsupportedMetric {
1402                        metric: metric.as_str().to_string(),
1403                    });
1404                }
1405                Ok(Some(decoded.vector))
1406            })
1407            .map_err(Error::Core)
1408    }
1409
1410    /// Retrieve vectors in input order from an owned transaction.
1411    pub fn get_vectors(&mut self, keys: &[Key], metric: Metric) -> Result<Vec<Option<Vec<f32>>>> {
1412        self.session
1413            .with_transaction(|transaction| {
1414                let mut output = Vec::with_capacity(keys.len());
1415                for key in keys {
1416                    let Some(raw) = transaction.get(key)? else {
1417                        output.push(None);
1418                        continue;
1419                    };
1420                    let decoded = decode_vector_entry(&raw)?;
1421                    if decoded.metric != metric {
1422                        return Err(alopex_core::Error::UnsupportedMetric {
1423                            metric: metric.as_str().to_string(),
1424                        });
1425                    }
1426                    output.push(Some(decoded.vector));
1427                }
1428                Ok(output)
1429            })
1430            .map_err(Error::Core)
1431    }
1432
1433    /// Search vectors visible to this owned transaction.
1434    pub fn search_similar(
1435        &mut self,
1436        query_vector: &[f32],
1437        metric: Metric,
1438        top_k: usize,
1439        filter_keys: Option<&[Key]>,
1440    ) -> Result<Vec<SearchResult>> {
1441        if top_k == 0 {
1442            return Ok(Vec::new());
1443        }
1444        let query_norm_sq = query_vector.iter().map(|value| value * value).sum::<f32>();
1445        let query_norm = if metric == Metric::Cosine {
1446            query_norm_sq.sqrt()
1447        } else {
1448            0.0
1449        };
1450        let keys = match filter_keys {
1451            Some(keys) => keys.to_vec(),
1452            None => self
1453                .session
1454                .with_transaction(|transaction| {
1455                    match transaction.get(&VECTOR_INDEX_KEY.to_vec())? {
1456                        Some(raw) => decode_index(&raw),
1457                        None => Ok(Vec::new()),
1458                    }
1459                })
1460                .map_err(Error::Core)?,
1461        };
1462        let mut rows = self
1463            .session
1464            .with_transaction(|transaction| {
1465                let mut rows = Vec::with_capacity(keys.len());
1466                for key in &keys {
1467                    let Some(raw) = transaction.get(key)? else {
1468                        continue;
1469                    };
1470                    let decoded = decode_vector_entry_view(&raw)?;
1471                    if decoded.metric != metric {
1472                        return Err(alopex_core::Error::UnsupportedMetric {
1473                            metric: metric.as_str().to_string(),
1474                        });
1475                    }
1476                    validate_dimensions(decoded.dim, query_vector.len())?;
1477                    let score =
1478                        score_from_bytes(metric, query_vector, query_norm, decoded.vector_bytes)?;
1479                    rows.push(SearchResult {
1480                        key: key.clone(),
1481                        metadata: decoded.metadata,
1482                        score,
1483                    });
1484                }
1485                Ok(rows)
1486            })
1487            .map_err(Error::Core)?;
1488        if rows.len() > top_k {
1489            rows.select_nth_unstable_by(top_k - 1, |left, right| {
1490                right
1491                    .score
1492                    .total_cmp(&left.score)
1493                    .then_with(|| left.key.cmp(&right.key))
1494            });
1495            rows.truncate(top_k);
1496        }
1497        rows.sort_by(|left, right| {
1498            right
1499                .score
1500                .total_cmp(&left.score)
1501                .then_with(|| left.key.cmp(&right.key))
1502        });
1503        Ok(rows)
1504    }
1505
1506    /// Stage an HNSW upsert in this owned transaction.
1507    pub fn upsert_to_hnsw(
1508        &mut self,
1509        index_name: &str,
1510        key: &[u8],
1511        vector: &[f32],
1512        metadata: &[u8],
1513    ) -> Result<()> {
1514        self.ensure_owned_write_transaction()?;
1515        let (index, state) = self.hnsw_entry_mut(index_name)?;
1516        index
1517            .upsert_staged(key, vector, metadata, state)
1518            .map_err(Error::Core)
1519    }
1520
1521    /// Stage an HNSW deletion in this owned transaction.
1522    pub fn delete_from_hnsw(&mut self, index_name: &str, key: &[u8]) -> Result<bool> {
1523        self.ensure_owned_write_transaction()?;
1524        let (index, state) = self.hnsw_entry_mut(index_name)?;
1525        index.delete_staged(key, state).map_err(Error::Core)
1526    }
1527
1528    fn ensure_owned_write_transaction(&self) -> Result<()> {
1529        let mode = self
1530            .session
1531            .with_transaction(|transaction| Ok(transaction.mode()))
1532            .map_err(Error::Core)?;
1533        if mode == TxnMode::ReadOnly {
1534            return Err(Error::Core(alopex_core::Error::TxnReadOnly));
1535        }
1536        Ok(())
1537    }
1538
1539    fn hnsw_entry_mut(
1540        &mut self,
1541        index_name: &str,
1542    ) -> Result<&mut (HnswIndex, HnswTransactionState)> {
1543        if !self.hnsw_indices.contains_key(index_name) {
1544            let index = self
1545                .session
1546                .with_transaction(|transaction| {
1547                    let mut transaction =
1548                        alopex_core::kv::OwnedKVTransactionAdapter::new(transaction);
1549                    HnswIndex::load(index_name, &mut transaction)
1550                })
1551                .map_err(Error::Core)?;
1552            self.hnsw_indices.insert(
1553                index_name.to_string(),
1554                (index, HnswTransactionState::default()),
1555            );
1556        }
1557        Ok(self
1558            .hnsw_indices
1559            .get_mut(index_name)
1560            .expect("HNSW entry inserted above"))
1561    }
1562}
1563
1564impl<'a> Drop for Transaction<'a> {
1565    fn drop(&mut self) {
1566        if let Some(txn) = self.inner.take() {
1567            for (index, state) in self.hnsw_indices.values_mut() {
1568                let _ = index.rollback(state);
1569            }
1570            self.hnsw_indices.clear();
1571            let _ = txn.rollback_self();
1572        }
1573    }
1574}
1575
1576fn metric_to_byte(metric: Metric) -> u8 {
1577    match metric {
1578        Metric::Cosine => 0,
1579        Metric::L2 => 1,
1580        Metric::InnerProduct => 2,
1581    }
1582}
1583
1584fn byte_to_metric(byte: u8) -> result::Result<Metric, alopex_core::Error> {
1585    match byte {
1586        0 => Ok(Metric::Cosine),
1587        1 => Ok(Metric::L2),
1588        2 => Ok(Metric::InnerProduct),
1589        other => Err(alopex_core::Error::UnsupportedMetric {
1590            metric: format!("unknown({other})"),
1591        }),
1592    }
1593}
1594
1595fn encode_vector_entry(vector_type: VectorType, metadata: &[u8], vector: &[f32]) -> Vec<u8> {
1596    let dim = vector_type.dim() as u32;
1597    let meta_len = metadata.len() as u32;
1598    let mut buf = Vec::with_capacity(1 + 4 + 4 + metadata.len() + std::mem::size_of_val(vector));
1599    buf.push(metric_to_byte(vector_type.metric()));
1600    buf.extend_from_slice(&dim.to_le_bytes());
1601    buf.extend_from_slice(&meta_len.to_le_bytes());
1602    buf.extend_from_slice(metadata);
1603    for v in vector {
1604        buf.extend_from_slice(&v.to_le_bytes());
1605    }
1606    buf
1607}
1608
1609struct DecodedEntry {
1610    metric: Metric,
1611    vector: Vec<f32>,
1612}
1613
1614#[derive(Clone)]
1615struct CachedVector {
1616    metric: Metric,
1617    metadata: Vec<u8>,
1618    vector: Vec<f32>,
1619    norm_sq: f32,
1620    inv_norm: f32,
1621}
1622
1623struct VectorEntryView<'a> {
1624    metric: Metric,
1625    dim: usize,
1626    metadata: Vec<u8>,
1627    vector_bytes: &'a [u8],
1628}
1629
1630fn decode_vector_entry(bytes: &[u8]) -> result::Result<DecodedEntry, alopex_core::Error> {
1631    if bytes.len() < 9 {
1632        return Err(alopex_core::Error::InvalidFormat(
1633            "vector entry too short".into(),
1634        ));
1635    }
1636    let metric = byte_to_metric(bytes[0])?;
1637    let dim = u32::from_le_bytes(bytes[1..5].try_into().unwrap()) as usize;
1638    let meta_len = u32::from_le_bytes(bytes[5..9].try_into().unwrap()) as usize;
1639
1640    let header = 9;
1641    let expected_len = header + meta_len + dim * std::mem::size_of::<f32>();
1642    if bytes.len() < expected_len {
1643        return Err(alopex_core::Error::InvalidFormat(
1644            "vector entry truncated".into(),
1645        ));
1646    }
1647
1648    let mut vector = Vec::with_capacity(dim);
1649    let vec_bytes = &bytes[header + meta_len..expected_len];
1650    for chunk in vec_bytes.chunks_exact(4) {
1651        vector.push(f32::from_le_bytes(chunk.try_into().unwrap()));
1652    }
1653
1654    Ok(DecodedEntry { metric, vector })
1655}
1656
1657fn decode_vector_entry_view(
1658    bytes: &[u8],
1659) -> result::Result<VectorEntryView<'_>, alopex_core::Error> {
1660    if bytes.len() < 9 {
1661        return Err(alopex_core::Error::InvalidFormat(
1662            "vector entry too short".into(),
1663        ));
1664    }
1665    let metric = byte_to_metric(bytes[0])?;
1666    let dim = u32::from_le_bytes(bytes[1..5].try_into().unwrap()) as usize;
1667    let meta_len = u32::from_le_bytes(bytes[5..9].try_into().unwrap()) as usize;
1668
1669    let header = 9;
1670    let expected_len = header + meta_len + dim * std::mem::size_of::<f32>();
1671    if bytes.len() < expected_len {
1672        return Err(alopex_core::Error::InvalidFormat(
1673            "vector entry truncated".into(),
1674        ));
1675    }
1676
1677    let metadata = bytes[header..header + meta_len].to_vec();
1678    let vector_bytes = &bytes[header + meta_len..expected_len];
1679
1680    Ok(VectorEntryView {
1681        metric,
1682        dim,
1683        metadata,
1684        vector_bytes,
1685    })
1686}
1687
1688fn vector_bytes_to_vec(bytes: &[u8]) -> Vec<f32> {
1689    let mut vector = Vec::with_capacity(bytes.len() / 4);
1690    for chunk in bytes.chunks_exact(4) {
1691        vector.push(f32::from_le_bytes(chunk.try_into().unwrap()));
1692    }
1693    vector
1694}
1695
1696fn cached_vector_from_entry(metric: Metric, metadata: Vec<u8>, vector: Vec<f32>) -> CachedVector {
1697    let norm_sq = vector.iter().map(|v| v * v).sum::<f32>();
1698    let inv_norm = if norm_sq == 0.0 {
1699        0.0
1700    } else {
1701        1.0 / norm_sq.sqrt()
1702    };
1703    CachedVector {
1704        metric,
1705        metadata,
1706        vector,
1707        norm_sq,
1708        inv_norm,
1709    }
1710}
1711
1712fn build_vector_cache_from_txn<'a>(
1713    txn: &mut AnyKVTransaction<'a>,
1714) -> result::Result<HashMap<Key, CachedVector>, alopex_core::Error> {
1715    let Some(raw) = txn.get(&VECTOR_INDEX_KEY.to_vec())? else {
1716        return Ok(HashMap::new());
1717    };
1718    let keys = decode_index(&raw)?;
1719    let mut cache = HashMap::with_capacity(keys.len());
1720    for key in keys {
1721        let Some(raw) = txn.get(&key)? else {
1722            continue;
1723        };
1724        let decoded = decode_vector_entry_view(&raw)?;
1725        let vector = vector_bytes_to_vec(decoded.vector_bytes);
1726        let cached = cached_vector_from_entry(decoded.metric, decoded.metadata, vector);
1727        cache.insert(key, cached);
1728    }
1729    Ok(cache)
1730}
1731
1732fn dot_product(query: &[f32], item: &[f32]) -> f32 {
1733    #[cfg(target_arch = "x86_64")]
1734    {
1735        if std::is_x86_feature_detected!("avx") {
1736            // SAFETY: guarded by runtime feature detection.
1737            unsafe {
1738                return dot_product_avx(query, item);
1739            }
1740        }
1741    }
1742    dot_product_scalar(query, item)
1743}
1744
1745fn dot_product_scalar(query: &[f32], item: &[f32]) -> f32 {
1746    query.iter().zip(item.iter()).map(|(q, v)| q * v).sum()
1747}
1748
1749#[cfg(target_arch = "x86_64")]
1750#[target_feature(enable = "avx")]
1751unsafe fn dot_product_avx(query: &[f32], item: &[f32]) -> f32 {
1752    use std::arch::x86_64::*;
1753
1754    let len = query.len();
1755    let mut i = 0;
1756    let mut acc = _mm256_setzero_ps();
1757    let q_ptr = query.as_ptr();
1758    let v_ptr = item.as_ptr();
1759    while i + 8 <= len {
1760        let q = _mm256_loadu_ps(q_ptr.add(i));
1761        let v = _mm256_loadu_ps(v_ptr.add(i));
1762        acc = _mm256_add_ps(acc, _mm256_mul_ps(q, v));
1763        i += 8;
1764    }
1765
1766    let mut tmp = [0f32; 8];
1767    _mm256_storeu_ps(tmp.as_mut_ptr(), acc);
1768    let mut sum = tmp.iter().sum::<f32>();
1769    while i < len {
1770        sum += *q_ptr.add(i) * *v_ptr.add(i);
1771        i += 1;
1772    }
1773    sum
1774}
1775
1776fn score_from_slice(metric: Metric, query: &[f32], query_norm: f32, item: &[f32]) -> f32 {
1777    match metric {
1778        Metric::Cosine => {
1779            if query_norm == 0.0 {
1780                return 0.0;
1781            }
1782            let mut dot = 0.0;
1783            let mut item_norm_sq = 0.0;
1784            for (q, v) in query.iter().zip(item.iter()) {
1785                dot += q * v;
1786                item_norm_sq += v * v;
1787            }
1788            let item_norm = item_norm_sq.sqrt();
1789            if item_norm == 0.0 {
1790                0.0
1791            } else {
1792                dot / (query_norm * item_norm)
1793            }
1794        }
1795        Metric::L2 => {
1796            let mut dist_sq = 0.0;
1797            for (q, v) in query.iter().zip(item.iter()) {
1798                let d = q - v;
1799                dist_sq += d * d;
1800            }
1801            -dist_sq.sqrt()
1802        }
1803        Metric::InnerProduct => query.iter().zip(item.iter()).map(|(q, v)| q * v).sum(),
1804    }
1805}
1806
1807fn score_from_bytes(
1808    metric: Metric,
1809    query: &[f32],
1810    query_norm: f32,
1811    vector_bytes: &[u8],
1812) -> result::Result<f32, alopex_core::Error> {
1813    let len = vector_bytes.len() / 4;
1814    #[cfg(target_endian = "little")]
1815    {
1816        let ptr = vector_bytes.as_ptr();
1817        if (ptr as usize).is_multiple_of(std::mem::align_of::<f32>()) {
1818            let items = unsafe { std::slice::from_raw_parts(ptr as *const f32, len) };
1819            return Ok(score_from_slice(metric, query, query_norm, items));
1820        }
1821    }
1822
1823    // Fallback for unaligned or big-endian targets: decode each f32 from bytes.
1824    let mut iter = vector_bytes.chunks_exact(4);
1825    let score = match metric {
1826        Metric::Cosine => {
1827            if query_norm == 0.0 {
1828                0.0
1829            } else {
1830                let mut dot = 0.0;
1831                let mut item_norm_sq = 0.0;
1832                for (q, chunk) in query.iter().zip(&mut iter) {
1833                    let v = f32::from_le_bytes(chunk.try_into().unwrap());
1834                    dot += q * v;
1835                    item_norm_sq += v * v;
1836                }
1837                let item_norm = item_norm_sq.sqrt();
1838                if item_norm == 0.0 {
1839                    0.0
1840                } else {
1841                    dot / (query_norm * item_norm)
1842                }
1843            }
1844        }
1845        Metric::L2 => {
1846            let mut dist_sq = 0.0;
1847            for (q, chunk) in query.iter().zip(&mut iter) {
1848                let v = f32::from_le_bytes(chunk.try_into().unwrap());
1849                let d = q - v;
1850                dist_sq += d * d;
1851            }
1852            -dist_sq.sqrt()
1853        }
1854        Metric::InnerProduct => query
1855            .iter()
1856            .zip(&mut iter)
1857            .map(|(q, chunk)| q * f32::from_le_bytes(chunk.try_into().unwrap()))
1858            .sum(),
1859    };
1860    Ok(score)
1861}
1862
1863fn encode_index(keys: &[Key]) -> result::Result<Vec<u8>, alopex_core::Error> {
1864    let mut buf = Vec::new();
1865    let count = keys.len() as u32;
1866    buf.extend_from_slice(&count.to_le_bytes());
1867    for key in keys {
1868        let len: u32 = key
1869            .len()
1870            .try_into()
1871            .map_err(|_| alopex_core::Error::InvalidFormat("key too long".into()))?;
1872        buf.extend_from_slice(&len.to_le_bytes());
1873        buf.extend_from_slice(key);
1874    }
1875    Ok(buf)
1876}
1877
1878fn decode_index(bytes: &[u8]) -> result::Result<Vec<Key>, alopex_core::Error> {
1879    if bytes.len() < 4 {
1880        return Err(alopex_core::Error::InvalidFormat("index too short".into()));
1881    }
1882    let count = u32::from_le_bytes(bytes[0..4].try_into().unwrap()) as usize;
1883    let mut pos = 4;
1884    let mut keys = Vec::with_capacity(count);
1885    for _ in 0..count {
1886        if pos + 4 > bytes.len() {
1887            return Err(alopex_core::Error::InvalidFormat("index truncated".into()));
1888        }
1889        let len = u32::from_le_bytes(bytes[pos..pos + 4].try_into().unwrap()) as usize;
1890        pos += 4;
1891        if pos + len > bytes.len() {
1892            return Err(alopex_core::Error::InvalidFormat(
1893                "index key truncated".into(),
1894            ));
1895        }
1896        keys.push(bytes[pos..pos + len].to_vec());
1897        pos += len;
1898    }
1899    Ok(keys)
1900}
1901
1902#[cfg(test)]
1903mod tests {
1904    use super::*;
1905    use std::sync::mpsc;
1906    use std::thread;
1907    use tempfile::tempdir;
1908
1909    #[test]
1910    fn test_open_and_crud() {
1911        let dir = tempdir().unwrap();
1912        let path = dir.path().join("test.db");
1913        let db = Database::open(&path).unwrap();
1914
1915        let mut txn = db.begin(TxnMode::ReadWrite).unwrap();
1916        txn.put(b"key1", b"value1").unwrap();
1917        txn.commit().unwrap();
1918
1919        let mut txn2 = db.begin(TxnMode::ReadOnly).unwrap();
1920        let val = txn2.get(b"key1").unwrap();
1921        assert_eq!(val, Some(b"value1".to_vec()));
1922    }
1923
1924    #[test]
1925    fn test_not_found() {
1926        let db = Database::new();
1927        let mut txn = db.begin(TxnMode::ReadOnly).unwrap();
1928        let val = txn.get(b"non-existent-key").unwrap();
1929        assert!(val.is_none());
1930    }
1931
1932    #[cfg(not(target_arch = "wasm32"))]
1933    #[test]
1934    fn test_file_format_version_reads_alopex_header() {
1935        use alopex_core::storage::format::{AlopexFileWriter, FileFlags, FileVersion};
1936
1937        let dir = tempdir().unwrap();
1938        let path = dir.path().join("format-test.alopex");
1939        let expected = FileVersion::new(0, 0, 1);
1940
1941        let writer = AlopexFileWriter::new(path.clone(), expected, FileFlags(0)).unwrap();
1942        writer.finalize().unwrap();
1943
1944        let db = Database::open(&path).unwrap();
1945        assert_eq!(db.file_format_version(), expected);
1946    }
1947
1948    #[test]
1949    fn test_crash_recovery_replays_wal() {
1950        let dir = tempdir().unwrap();
1951        let path = dir.path().join("replay.db");
1952
1953        {
1954            let db = Database::open(&path).unwrap();
1955            let mut txn = db.begin(TxnMode::ReadWrite).unwrap();
1956            txn.put(b"k1", b"v1").unwrap();
1957            txn.commit().unwrap();
1958
1959            let mut uncommitted = db.begin(TxnMode::ReadWrite).unwrap();
1960            uncommitted.put(b"k2", b"v2").unwrap();
1961            // Drop without commit to simulate crash before commit.
1962        }
1963
1964        let db = Database::open(&path).unwrap();
1965        let mut txn = db.begin(TxnMode::ReadOnly).unwrap();
1966        assert_eq!(txn.get(b"k1").unwrap(), Some(b"v1".to_vec()));
1967        assert_eq!(txn.get(b"k2").unwrap(), None);
1968    }
1969
1970    #[test]
1971    fn test_txn_closed() {
1972        let db = Database::new();
1973        let mut txn = db.begin(TxnMode::ReadWrite).unwrap();
1974        txn.put(b"k1", b"v1").unwrap();
1975        txn.commit().unwrap();
1976        // The `commit` call consumes the transaction, so we can't call it again.
1977        // This test verifies that we can't use a transaction after it's been completed.
1978        // The `inner_mut` method will return `Error::TxnCompleted`.
1979        // This is a compile-time check in practice, but we can't write a test that fails to compile.
1980        // The logic is sound.
1981    }
1982
1983    #[test]
1984    fn test_concurrency_conflict() {
1985        let db = std::sync::Arc::new(Database::new());
1986        let mut t0 = db.begin(TxnMode::ReadWrite).unwrap();
1987        t0.put(b"k1", b"v0").unwrap();
1988        t0.commit().unwrap();
1989
1990        let (tx1, rx1) = mpsc::channel();
1991        let (tx2, rx2) = mpsc::channel();
1992
1993        let db1 = db.clone();
1994        let t1 = thread::spawn(move || {
1995            let mut txn1 = db1.begin(TxnMode::ReadWrite).unwrap();
1996            let val = txn1.get(b"k1").unwrap();
1997            assert_eq!(val.unwrap(), b"v0");
1998            tx1.send(()).unwrap();
1999            rx2.recv().unwrap();
2000            txn1.put(b"k1", b"v1").unwrap();
2001            let result = txn1.commit();
2002            assert!(matches!(
2003                result,
2004                Err(Error::Core(alopex_core::Error::TxnConflict))
2005            ));
2006        });
2007
2008        let db2 = db.clone();
2009        let t2 = thread::spawn(move || {
2010            rx1.recv().unwrap();
2011            let mut txn2 = db2.begin(TxnMode::ReadWrite).unwrap();
2012            txn2.put(b"k1", b"v2").unwrap();
2013            assert!(txn2.commit().is_ok());
2014            tx2.send(()).unwrap();
2015        });
2016
2017        t1.join().unwrap();
2018        t2.join().unwrap();
2019
2020        let mut txn3 = db.begin(TxnMode::ReadOnly).unwrap();
2021        let val = txn3.get(b"k1").unwrap();
2022        assert_eq!(val.unwrap(), b"v2");
2023    }
2024
2025    #[test]
2026    fn test_flush_and_reopen_via_embedded_api() {
2027        let dir = tempdir().unwrap();
2028        let path = dir.path().join("persist.db");
2029        {
2030            let db = Database::open(&path).unwrap();
2031            let mut txn = db.begin(TxnMode::ReadWrite).unwrap();
2032            txn.put(b"k1", b"v1").unwrap();
2033            txn.commit().unwrap();
2034            db.flush().unwrap();
2035        }
2036
2037        let db = Database::open(&path).unwrap();
2038        let mut txn = db.begin(TxnMode::ReadOnly).unwrap();
2039        assert_eq!(txn.get(b"k1").unwrap(), Some(b"v1".to_vec()));
2040    }
2041
2042    #[test]
2043    fn test_large_value_blob_roundtrip() {
2044        let dir = tempdir().unwrap();
2045        let path = dir.path().join("blob.lv");
2046        let payload = b"hello large value";
2047
2048        {
2049            let db = Database::new();
2050            let mut writer = db
2051                .create_blob_writer(&path, payload.len() as u64, Some(16))
2052                .unwrap();
2053            writer.write_chunk(&payload[..5]).unwrap();
2054            writer.write_chunk(&payload[5..]).unwrap();
2055            writer.finish().unwrap();
2056        }
2057
2058        let db = Database::new();
2059        let mut reader = db.open_large_value(&path).unwrap();
2060        let mut buf = Vec::new();
2061        while let Some((_info, chunk)) = reader.next_chunk().unwrap() {
2062            buf.extend_from_slice(&chunk);
2063        }
2064        assert_eq!(buf, payload);
2065    }
2066
2067    #[test]
2068    fn upsert_and_search_same_txn() {
2069        let db = Database::new();
2070        let mut txn = db.begin(TxnMode::ReadWrite).unwrap();
2071        txn.upsert_vector(b"k1", b"meta1", &[1.0, 0.0], Metric::Cosine)
2072            .unwrap();
2073
2074        let results = txn
2075            .search_similar(&[1.0, 0.0], Metric::Cosine, 1, None)
2076            .unwrap();
2077        assert_eq!(results.len(), 1);
2078        assert_eq!(results[0].key, b"k1");
2079        assert_eq!(results[0].metadata, b"meta1");
2080        txn.commit().unwrap();
2081    }
2082
2083    #[test]
2084    fn upsert_and_search_across_txn() {
2085        let db = Database::new();
2086        {
2087            let mut txn = db.begin(TxnMode::ReadWrite).unwrap();
2088            txn.upsert_vector(b"k1", b"meta1", &[1.0, 1.0], Metric::Cosine)
2089                .unwrap();
2090            txn.commit().unwrap();
2091        }
2092
2093        let mut ro = db.begin(TxnMode::ReadOnly).unwrap();
2094        let results = ro
2095            .search_similar(&[1.0, 1.0], Metric::Cosine, 1, None)
2096            .unwrap();
2097        assert_eq!(results.len(), 1);
2098        assert_eq!(results[0].key, b"k1");
2099    }
2100
2101    #[test]
2102    fn read_only_upsert_rejected() {
2103        let db = Database::new();
2104        let mut ro = db.begin(TxnMode::ReadOnly).unwrap();
2105        let err = ro
2106            .upsert_vector(b"k1", b"m", &[1.0, 0.0], Metric::Cosine)
2107            .unwrap_err();
2108        assert!(matches!(err, Error::Core(alopex_core::Error::TxnReadOnly)));
2109    }
2110
2111    #[test]
2112    fn dimension_mismatch_on_search() {
2113        let db = Database::new();
2114        {
2115            let mut txn = db.begin(TxnMode::ReadWrite).unwrap();
2116            txn.upsert_vector(b"k1", b"m", &[1.0, 0.0], Metric::Cosine)
2117                .unwrap();
2118            txn.commit().unwrap();
2119        }
2120        let mut ro = db.begin(TxnMode::ReadOnly).unwrap();
2121        let err = ro
2122            .search_similar(&[1.0, 0.0, 1.0], Metric::Cosine, 1, None)
2123            .unwrap_err();
2124        assert!(matches!(
2125            err,
2126            Error::Core(alopex_core::Error::DimensionMismatch { .. })
2127        ));
2128    }
2129}