Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions crates/codegraph-graph/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -84,3 +84,8 @@ criterion = { version = "0.5", features = ["async_tokio"] }
[[bench]]
name = "search_bloom"
harness = false

[[bench]]
name = "lmdb_batch_read"
harness = false
required-features = ["lmdb"]
128 changes: 128 additions & 0 deletions crates/codegraph-graph/benches/lmdb_batch_read.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
//! Benchmark đọc LMDB: so sánh đọc TỪNG node (mỗi lần 1 read-txn) với đọc BATCH
//! (nhiều node trong 1 read-txn) — đo tác động của `get_nodes`/`get_childrens`.
//!
//! ```bash
//! cargo bench -p codegraph-graph --bench lmdb_batch_read --features lmdb
//! ```
//!
//! Nhóm đo:
//! - `single_read_{1,16,64}` — gọi `get_node` lặp lại cho từng id.
//! - `batch_read_{1,16,64}` — gọi `get_nodes` 1 lần cho cả id.
//! - `search_dfs` — DFS thật trên trie đã dựng (đo end-to-end).

use codegraph_graph::{CategoryStorage, LmdbStorage};
use criterion::{Criterion, black_box, criterion_group, criterion_main};

const N_NODES: usize = 4096;

/// Bộ id hợp lệ để đọc — dựng sẵn node thật trong LMDB.
fn setup(rt: &tokio::runtime::Runtime) -> (tempfile::TempDir, String, Vec<usize>) {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("bench.lmdb").to_string_lossy().into_owned();
let ids = rt.block_on(async {
let mut s = LmdbStorage::open(&path).await.unwrap();
let mut ids = Vec::with_capacity(N_NODES);
for i in 0..N_NODES {
// Prefix ngắn, giá trị record tăng dần — mô phỏng trie phẳng.
let id = s
.new_node(format!("n{i:05}").into_bytes(), i + 1)
.await
.unwrap();
ids.push(id);
}
ids
});
(dir, path, ids)
}

fn runtime() -> tokio::runtime::Runtime {
tokio::runtime::Builder::new_current_thread()
.build()
.unwrap()
}

fn bench_batch_vs_single(c: &mut Criterion) {
let rt = runtime();
let (_dir, path, ids) = setup(&rt);
let st = rt.block_on(async { LmdbStorage::open(&path).await.unwrap() });

let mut group = c.benchmark_group("lmdb_node_read");

for &k in &[1usize, 16, 64] {
let ids: Vec<usize> = ids.iter().copied().take(k).collect();

group.bench_function(format!("single_{k}"), |b| {
b.iter(|| {
rt.block_on(async {
for &id in &ids {
black_box(st.get_node(id).await.unwrap());
}
});
});
});

group.bench_function(format!("batch_{k}"), |b| {
b.iter(|| {
rt.block_on(async {
black_box(st.get_nodes(&ids).await.unwrap());
});
});
});
}

group.finish();
}

/// Đọc children: `get_children` lặp vs `get_childrens` batch.
fn bench_children_batch(c: &mut Criterion) {
let rt = runtime();
let (_dir, path, _ids) = setup(&rt);
let parents = rt.block_on(async {
let mut s = LmdbStorage::open(&path).await.unwrap();
let mut parents = Vec::new();
for i in 0..64usize {
let p = s
.new_node(format!("p{i:03}").into_bytes(), i)
.await
.unwrap();
// Mỗi parent có 32 child.
for j in 0..32usize {
let c = s
.new_node(format!("c{i:03}_{j:03}").into_bytes(), j)
.await
.unwrap();
let mut tx = s.new_tx();
tx.add_child(p, c).await.unwrap();
tx.commit().await.unwrap();
}
parents.push(p);
}
parents
});
let st = rt.block_on(async { LmdbStorage::open(&path).await.unwrap() });

let mut group = c.benchmark_group("lmdb_children_read");
for &k in &[1usize, 16, 64] {
let ps: Vec<usize> = parents.iter().copied().take(k).collect();
group.bench_function(format!("single_{k}"), |b| {
b.iter(|| {
rt.block_on(async {
for &p in &ps {
black_box(st.get_children(p).await.unwrap());
}
});
});
});
group.bench_function(format!("batch_{k}"), |b| {
b.iter(|| {
rt.block_on(async {
black_box(st.get_childrens(&ps).await.unwrap());
});
});
});
}
group.finish();
}

criterion_group!(benches, bench_batch_vs_single, bench_children_batch);
criterion_main!(benches);
50 changes: 37 additions & 13 deletions crates/codegraph-graph/src/radix.rs
Original file line number Diff line number Diff line change
Expand Up @@ -677,8 +677,16 @@ impl<T: Element> Radix<T> {
break;
};

let (prefix_bytes, _record) =
{ self.storage.read().await.get_node(frame.node_id).await? };
// Đọc node của frame + children của frame trong MỘT
// read-txn (`get_nodes`/`get_childrens`): xem
// `CategoryStorage::get_nodes` — giảm số read-txn mỗi
// vòng DFS, tránh chạm trần reader-slot của LMDB khi nhiều
// request search song song.
let (prefix_bytes, _record) = {
let st = self.storage.read().await;
let mut out = st.get_nodes(&[frame.node_id]).await?;
out.pop().expect("get_nodes trả về đúng 1 phần tử")
};
let prefix = Self::to_vec(&prefix_bytes);
let result = matcher(&prefix, pattern, frame.pattern_pos);

Expand All @@ -691,11 +699,9 @@ impl<T: Element> Radix<T> {
}

let children = {
self.storage
.read()
.await
.get_children(frame.node_id)
.await?
let st = self.storage.read().await;
let mut out = st.get_childrens(&[frame.node_id]).await?;
out.pop().expect("get_childrens trả về đúng 1 phần tử")
};
let mut descended = false;
while frame.cont_idx < result.continuations.len() {
Expand All @@ -707,12 +713,24 @@ impl<T: Element> Radix<T> {
}

let next_elem = pattern[pp];
// Pre-fetch node của MỌI child còn lại trong 1
// read-txn thay vì mỗi child 1 txn. Chỉ fetch phần chưa
// duyệt; `base` ánh xạ chỉ số tương đối về `children`
// (giống hệt hành vi cũ từng child một).
let base = frame.child_idx.min(children.len());
let remaining = &children[base..];
let child_nodes = if remaining.is_empty() {
Vec::new()
} else {
let st = self.storage.read().await;
st.get_nodes(remaining).await?
};
while frame.child_idx < children.len() {
let child = children[frame.child_idx];
let rel = frame.child_idx - base;
frame.child_idx += 1;
let (cp_bytes, _) =
{ self.storage.read().await.get_node(child).await? };
let cp = Self::to_vec(&cp_bytes);
let (cp_bytes, _) = &child_nodes[rel];
let cp = Self::to_vec(cp_bytes);
if cp.is_empty() || cp[0] != next_elem {
continue;
}
Expand Down Expand Up @@ -758,12 +776,18 @@ impl<T: Element> Radix<T> {
}
DfsState::Collect { root, mut stack } => {
if let Some((node_id, child_idx)) = stack.pop() {
let (_prefix_bytes, record) =
{ self.storage.read().await.get_node(node_id).await? };
// Gom node + children của cùng `node_id` — 2 read-txn
// điển hình xuống 2 lần gọi batch (mỗi lần 1 txn).
let (record, children) = {
let st = self.storage.read().await;
let mut node = st.get_nodes(&[node_id]).await?;
let node = node.pop().expect("get_nodes trả về đúng 1 phần tử");
let mut children = st.get_childrens(&[node_id]).await?;
(node.1, children.pop().expect("get_childrens trả 1 phần tử"))
};
if record != storage::EMPTY {
records.push(record);
}
let children = { self.storage.read().await.get_children(node_id).await? };
if child_idx < children.len() {
stack.push((node_id, child_idx + 1));
stack.push((children[child_idx], 0));
Expand Down
49 changes: 48 additions & 1 deletion crates/codegraph-graph/src/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,10 @@ enum TxOp {
/// chưa lộ ra cho reader cho tới khi `commit` hoàn tất.
///
/// `commit(self: Box<Self>)` tiêu thụ chính transaction — không thể commit 2 lần.
// `async_trait` tự sinh `#[must_use]` cho future; double_must_use (clippy 1.99)
// báo trùng vì `Pin<Box<dyn Future>>` vốn đã `must_use`. Lỗi nằm trong macro
// của async_trait, không phải code viết tay → allow tại chỗ sinh ra.
#[allow(clippy::double_must_use)]
#[async_trait]
pub trait Tx: Send {
async fn new_node(&mut self, prefix: Vec<u8>, record: usize) -> Result<usize>;
Expand All @@ -183,6 +187,7 @@ pub trait Tx: Send {
/// khi feature bật → method gọi được qua `dyn CategoryStorage` như cũ.
/// Backend không override → default no-op.
#[cfg(feature = "bloom-search")]
#[allow(clippy::double_must_use)] // async_trait sinh `must_use` trùng (clippy 1.99)
#[async_trait]
pub trait BloomStorage: Send + Sync {
async fn set_node_bloom(&mut self, _id: usize, _bloom: &[u8]) -> Result<()> {
Expand All @@ -199,6 +204,7 @@ pub trait BloomStorage: Send + Sync {
/// element id, cùng `clear`. Tách riêng để trait lõi gọn. `CategoryStorage`
/// super-bound trait này (luôn) → method gọi được qua `dyn CategoryStorage`.
/// Mặc định no-op.
#[allow(clippy::double_must_use)] // async_trait sinh `must_use` trùng (clippy 1.99)
#[async_trait]
pub trait NodeMetaStorage: Send + Sync {
/// Lưu metadata của node (opaque bytes, VD Node JSON) keyed theo element id.
Expand Down Expand Up @@ -236,6 +242,7 @@ pub trait NodeMetaStorage: Send + Sync {
/// lookup (KMP + DFS). Tách riêng để trait lõi gọn. `CategoryStorage`
/// super-bound trait này (luôn) → method gọi được qua `dyn CategoryStorage`.
/// Mặc định no-op.
#[allow(clippy::double_must_use)] // async_trait sinh `must_use` trùng (clippy 1.99)
#[async_trait]
pub trait ShortcutsStorage: Send + Sync {
/// Thêm `node_id` vào shortcut set của element `elem` (encoded bytes).
Expand All @@ -262,6 +269,7 @@ pub trait ShortcutsStorage: Send + Sync {
/// Edge-data storage: lưu/đọc metadata của mỗi edge id (opaque bytes) keyed
/// theo edge id. Tách riêng để trait lõi gọn. `CategoryStorage` super-bound
/// trait này (luôn) → method gọi được qua `dyn CategoryStorage`. Mặc định no-op.
#[allow(clippy::double_must_use)] // async_trait sinh `must_use` trùng (clippy 1.99)
#[async_trait]
pub trait EdgeDataStorage: Send + Sync {
/// Lưu dữ liệu edge (opaque bytes, VD CallEdgeMeta JSON) keyed theo edge id.
Expand All @@ -284,6 +292,7 @@ pub trait EdgeDataStorage: Send + Sync {
/// encode u64 LE 8-byte/element. Tách riêng để trait lõi gọn.
/// `CategoryStorage` super-bound trait này (luôn) → method gọi được qua
/// `dyn CategoryStorage`. Mặc định no-op.
#[allow(clippy::double_must_use)] // async_trait sinh `must_use` trùng (clippy 1.99)
#[async_trait]
pub trait ChainStorage: Send + Sync {
/// Lưu chain của owner (keyed theo record của owner; u64 LE 8-byte/element).
Expand All @@ -310,7 +319,8 @@ pub trait ChainStorage: Send + Sync {
macro_rules! declare_category_storage {
($($bounds:tt)*) => {
/// Radix-node storage: node management + transaction + 5 stream phụ.
#[async_trait]
#[allow(clippy::double_must_use)] // async_trait sinh `must_use` trùng (clippy 1.99)
#[async_trait]
pub trait CategoryStorage: $($bounds)* {
// ── Node management ──
async fn new_node(&mut self, prefix: Vec<u8>, record: usize) -> Result<usize>;
Expand All @@ -323,6 +333,40 @@ macro_rules! declare_category_storage {
async fn get_node(&self, id: usize) -> Result<(Vec<u8>, usize)>;
async fn get_children(&self, id: usize) -> Result<Vec<usize>>;

// ── Batch reads (MỐT read-txn cho nhiều id) ──
//
// Radix DFS (`radix.rs::search_dfs`) gọi `get_node`/`get_children`
// cho TỪNG node/child → mỗi lần đọc là một `begin_ro_txn` riêng nên
// một lần search tạo O(nodes) read-txn. Khi nhiều request search chạy
// song song (MCP/GraphQL dùng chung `Arc<GraphIndex>`), số read-txn
// đồng thời nhân lên và có thể chạm trần reader-slot của LMDB →
// `MDB_BAD_RSLOT`.
//
// 2 hàm này gom nhiều id vào MỘT read-txn. Backend nền (in-memory,
// sqlite, redis, rdbms) không override — default impl gọi lại hàm
// đơn lẻ nên hành vi giữ nguyên. LMDB override để dùng 1 txn.

/// Đọc nhiều node trong MỘT read-txn — thứ tự khớp `ids`.
///
/// Node không tồn tại → `Err(BranchOutOfRange)` (giống `get_node`).
async fn get_nodes(&self, ids: &[usize]) -> Result<Vec<(Vec<u8>, usize)>> {
let mut out = Vec::with_capacity(ids.len());
for &id in ids {
out.push(self.get_node(id).await?);
}
Ok(out)
}

/// Đọc children của nhiều node trong MỐT read-txn — `out[i]` là
/// children của `ids[i]`, đã sort. Node không tồn tại → `Err`.
async fn get_childrens(&self, ids: &[usize]) -> Result<Vec<Vec<usize>>> {
let mut out = Vec::with_capacity(ids.len());
for &id in ids {
out.push(self.get_children(id).await?);
}
Ok(out)
}

// ── Shard roots (endpoint) ──
async fn set_root(&mut self, shard: usize, root: usize) -> Result<()>;
async fn get_root(&self, shard: usize) -> Result<usize>;
Expand Down Expand Up @@ -357,6 +401,8 @@ declare_category_storage!(
/// - Chỉ `GraphIndex` / `SharedGraphIndex` dùng (`Radix` / `Search` không cần).
/// - Backend tối giản có thể bỏ qua (vd: chỉ cần `CategoryStorage` cho test).
/// - Cho phép phát triển/scale entity layer độc lập với radix.
// Xem `Tx` — cùng lý do `double_must_use` (clippy 1.99).
#[allow(clippy::double_must_use)]
#[async_trait]
pub trait EntityStorage: Send + Sync {
// ── Symbol registry ──
Expand Down Expand Up @@ -492,6 +538,7 @@ pub trait EntityStorage: Send + Sync {
///
/// Gộp `CategoryStorage` + 5 trait phụ + `EntityStorage`. Backend implement
/// 7 `impl` block riêng biệt — review từng phần độc lập được.
#[allow(clippy::double_must_use)] // async_trait sinh `must_use` trùng (clippy 1.99)
#[async_trait]
pub trait Storage:
CategoryStorage
Expand Down
Loading
Loading