use std::path::PathBuf;
use std::time::Duration;
use storage_engines::engine::{open_engine, KvEngine};
use storage_engines::{EngineError, SearchQuery};
fn main() {
if let Err(e) = run() {
eprintln!("出错了: {e}");
std::process::exit(1);
}
}
fn run() -> Result<(), EngineError> {
let mut args = std::env::args().skip(1);
let name = args.next().unwrap_or_else(|| "memory".to_string());
let dir: Option<PathBuf> = match name.to_ascii_lowercase().as_str() {
"memory" | "mem" | "in_memory" => None,
_ => Some(
args.next()
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from("./data").join(&name)),
),
};
let db = open_engine(&name, dir.as_deref())?;
println!("== 引擎: {} ==", db.name());
match db.path() {
Some(p) => println!(" 库路径: {}", p.display()),
None => println!(" 纯内存,无库路径"),
}
println!();
basic_ops(db.as_ref());
batch_and_scan(db.as_ref());
ttl_ops(db.as_ref());
write_conflict(db.as_ref());
bulk_load(db.as_ref());
maintenance(db.as_ref())?;
println!("\n全部示例跑完。");
Ok(())
}
fn basic_ops(db: &dyn KvEngine) {
println!("-- 1. 基本读写 --");
let tx = db.begin();
println!(" 事务版本号 = {}", tx.version());
tx.set(b"name", b"Alice".to_vec());
tx.set(b"city", b"Beijing".to_vec());
println!(" get(name) = {:?}", as_str(tx.get(b"name")));
println!(" exists(城市) = {}", tx.exists(b"city"));
tx.set(b"visits", b"10".to_vec());
let n = tx.incr(b"visits", 5).expect("incr");
println!(" incr(visits, +5) = {n}");
tx.delete(b"city");
println!(" 删除后 get(city) = {:?}", tx.get(b"city"));
if let Some(meta) = tx.get_meta(b"name") {
println!(
" meta(name): 版本={} 长度={:?} 已删除={}",
meta.version, meta.value_len, meta.deleted
);
}
tx.commit();
println!();
}
fn batch_and_scan(db: &dyn KvEngine) {
println!("-- 2. 批量与扫描 --");
let tx = db.begin();
let users: Vec<(Vec<u8>, Vec<u8>)> = (1..=5)
.map(|i| {
(
format!("user:{i:02}").into_bytes(),
format!("用户{i}").into_bytes(),
)
})
.collect();
if let Err(bad_key) = tx.batch_set(&users) {
println!(" 批量写在 {:?} 上冲突", as_str(Some(bad_key)));
tx.rollback();
return;
}
let got = tx.multi_get(&[b"user:01".to_vec(), b"user:99".to_vec()]);
println!(" multi_get = {:?}", got.iter().map(as_ref_str).collect::<Vec<_>>());
let range = tx.scan(Some(b"user:02"), Some(b"user:04"));
println!(" scan[user:02, user:04) = {:?}", keys_of(&range));
println!(" prefix_scan(user:) = {:?}", keys_of(&tx.prefix_scan(b"user:")));
println!(" key_count(user:) = {}", tx.key_count(b"user:"));
let rev = tx.reverse_scan(None, None);
println!(" reverse_scan 前 3 个 = {:?}", keys_of(&rev[..rev.len().min(3)]));
println!(" seek(user:03) = {:?}", rec_key(tx.seek(b"user:03")));
println!(" seek_prev(user:03) = {:?}", rec_key(tx.seek_prev(b"user:03")));
let page = tx.search_keys(&SearchQuery::contains("user", 0, 3));
println!(
" 搜索 \"user\" 第0页: {:?} (共{}条 {}页, 有下一页={})",
keys_of(&page.items),
page.total,
page.total_pages,
page.has_next()
);
tx.commit();
println!();
}
fn ttl_ops(db: &dyn KvEngine) {
println!("-- 3. TTL --");
let tx = db.begin();
tx.set_with_ttl(b"session", b"token-abc".to_vec(), Duration::from_secs(60));
match tx.get_ttl(b"session") {
None => println!(" session 不存在或已过期"),
Some(None) => println!(" session 是永久的"),
Some(Some(d)) => println!(" session 还剩 {} 秒", d.as_secs()),
}
tx.refresh_ttl(b"session", Duration::from_secs(120));
println!(" 续期后还剩 {:?} 秒", tx.get_ttl(b"session").flatten().map(|d| d.as_secs()));
tx.persist_key(b"session");
println!(" persist 后 = {:?} (Some(None) 就是永久)", tx.get_ttl(b"session"));
println!(" purge_expired 清理了 {} 条", tx.purge_expired());
tx.commit();
println!();
}
fn write_conflict(db: &dyn KvEngine) {
println!("-- 4. 写冲突 --");
let t1 = db.begin();
let t2 = db.begin();
println!(" t1.set(counter) -> {}", t1.set(b"counter", b"1".to_vec()));
let ok = t2.set(b"counter", b"2".to_vec());
println!(" t2.set(counter) -> {ok} <- false 说明撞车了");
t1.commit();
t2.rollback();
let t3 = db.begin();
println!(" 最终值 = {:?} (t1 赢)", as_str(t3.get(b"counter")));
t3.commit();
println!();
}
fn bulk_load(db: &dyn KvEngine) {
println!("-- 5. 批量导入 --");
println!(" 注意:单写者假定,导入期间别开别的写事务");
{
let mut bulk = db.begin_bulk();
bulk.put(b"bulk:a", b"1".to_vec());
bulk.put_batch(&[(b"bulk:b".to_vec(), b"2".to_vec())]);
bulk.put_batch_owned(vec![(b"bulk:c".to_vec(), b"3".to_vec())]);
bulk.finish();
}
let tx = db.begin();
println!(" 导入结果 = {:?}", keys_of(&tx.prefix_scan(b"bulk:")));
tx.commit();
println!();
}
fn maintenance(db: &dyn KvEngine) -> Result<(), EngineError> {
println!("-- 6. 维护操作 --");
db.flush()?;
db.checkpoint()?;
println!(" flush + checkpoint 完成");
let stats = db.vacuum()?;
println!(
" vacuum: 水位={} 清理版本={} blob重写={:?}",
stats.xmin, stats.versions_removed, stats.blob_rewritten
);
match db.merge() {
Ok(()) => println!(" merge 完成"),
Err(EngineError::Unsupported { engine, op }) => {
println!(" {engine} 不支持 {op}(正常,只有 bitcask 有)")
}
Err(e) => return Err(e),
}
let all = db.export_latest_visible(false);
println!(" 当前共 {} 个可见 key", all.len());
println!();
Ok(())
}
fn as_str(v: Option<Vec<u8>>) -> Option<String> {
v.map(|b| String::from_utf8_lossy(&b).into_owned())
}
fn as_ref_str(v: &Option<Vec<u8>>) -> Option<String> {
v.as_ref().map(|b| String::from_utf8_lossy(b).into_owned())
}
fn keys_of(records: &[storage_engines::ExportRecord]) -> Vec<String> {
records
.iter()
.map(|r| String::from_utf8_lossy(&r.key).into_owned())
.collect()
}
fn rec_key(r: Option<storage_engines::ExportRecord>) -> Option<String> {
r.map(|r| String::from_utf8_lossy(&r.key).into_owned())
}