From 7ec470168208c4cd082ef61fe55748f5eb8d361e Mon Sep 17 00:00:00 2001 From: junyao <1071307515@qq.com> Date: Wed, 19 Aug 2026 16:04:56 +0000 Subject: [PATCH 1/4] feat: durable Notify/Take queue --- src/ffi.rs | 109 ++++++++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 108 insertions(+), 1 deletion(-) diff --git a/src/ffi.rs b/src/ffi.rs index 28053cf..ba6cc0f 100644 --- a/src/ffi.rs +++ b/src/ffi.rs @@ -18,8 +18,9 @@ use std::time::Duration; use crate::conn::conn; use crate::kvspace::{KVSpace, KVPair}; use crate::kvspace_common::get_one; +use crate::r#const::KIND_UINT8; use crate::xvalue::{ - decode_xvalue, decode_xvalue_head, encode_head, new_ptr, + body_bytes, decode_xvalue, decode_xvalue_head, encode_head, is_none, new_ptr, }; use crate::xvalue_bool::new_bool; use crate::xvalue_byte::{new_char, new_char_byte}; @@ -501,3 +502,109 @@ pub extern "C" fn kvspaceNewFloat64( let x = new_float64(&[v]); alloc(x.encode(), out, out_len) } + +fn nq_key(key: &str) -> String { + format!("/\u{2025}notify{key}") +} + +fn nq_load(kv: &mut dyn KVSpace, qk: &str) -> Vec { + let v = get_one(kv, qk); + if is_none(&v) { + return Vec::new(); + } + body_bytes(&v) +} + +fn nq_frames(body: &[u8]) -> Vec> { + let mut out = Vec::new(); + let mut i = 0; + while i + 4 <= body.len() { + let n = u32::from_le_bytes([body[i], body[i + 1], body[i + 2], body[i + 3]]) as usize; + i += 4; + if i + n > body.len() { + break; + } + out.push(body[i..i + n].to_vec()); + i += n; + } + out +} + +fn nq_save(kv: &mut dyn KVSpace, qk: &str, frames: &[Vec]) -> Result<(), String> { + if frames.is_empty() { + return kv.del(&[qk.to_string()]); + } + let mut raw = Vec::new(); + for f in frames { + raw.extend_from_slice(&(f.len() as u32).to_le_bytes()); + raw.extend_from_slice(f); + } + let tlv = encode_head(KIND_UINT8, 0, &[raw.len() as i32], &raw); + kv.set(&[KVPair { + key: qk.to_string(), + val: decode_xvalue(&tlv), + }]) +} + +#[no_mangle] +pub extern "C" fn kvspaceNotify( + h: *mut Handle, + key: *const c_char, + val: *const u8, + len: u32, + err: *mut c_char, + err_cap: u32, +) -> c_int { + if h.is_null() || val.is_null() || len == 0 { + return 1; + } + let kv: &mut dyn KVSpace = unsafe { &mut **h }; + let key = unsafe { cstr(key) }; + let frame = unsafe { std::slice::from_raw_parts(val, len as usize) }.to_vec(); + let qk = nq_key(key); + let mut frames = nq_frames(&nq_load(kv, &qk)); + frames.push(frame); + match nq_save(kv, &qk, &frames) { + Ok(()) => 0, + Err(e) => { + write_err(err, err_cap, &e); + 1 + } + } +} + +#[no_mangle] +pub extern "C" fn kvspaceTake( + h: *mut Handle, + key: *const c_char, + timeout_ns: u64, + out: *mut *mut u8, + out_len: *mut u32, +) -> c_int { + if out.is_null() || out_len.is_null() { + return 1; + } + unsafe { + *out = std::ptr::null_mut(); + *out_len = 0; + } + if h.is_null() { + return 1; + } + let kv: &mut dyn KVSpace = unsafe { &mut **h }; + let key = unsafe { cstr(key) }; + let qk = nq_key(key); + let deadline = std::time::Instant::now() + Duration::from_nanos(timeout_ns); + loop { + let mut frames = nq_frames(&nq_load(kv, &qk)); + if !frames.is_empty() { + let item = frames.remove(0); + let _ = nq_save(kv, &qk, &frames); + return alloc(item, out, out_len); + } + if std::time::Instant::now() >= deadline { + return 0; + } + std::thread::sleep(Duration::from_millis(1)); + } +} From 5c956245bdef63b838c30158517d5278d2fd9f94 Mon Sep 17 00:00:00 2001 From: junyao <1071307515@qq.com> Date: Thu, 20 Aug 2026 04:11:12 +0000 Subject: [PATCH 2/4] feat: atomic Incr --- src/ffi.rs | 51 ++++++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 50 insertions(+), 1 deletion(-) diff --git a/src/ffi.rs b/src/ffi.rs index ba6cc0f..4dcb0c8 100644 --- a/src/ffi.rs +++ b/src/ffi.rs @@ -18,7 +18,7 @@ use std::time::Duration; use crate::conn::conn; use crate::kvspace::{KVSpace, KVPair}; use crate::kvspace_common::get_one; -use crate::r#const::KIND_UINT8; +use crate::r#const::{KIND_CHAR_UTF8, KIND_UINT8}; use crate::xvalue::{ body_bytes, decode_xvalue, decode_xvalue_head, encode_head, is_none, new_ptr, }; @@ -608,3 +608,52 @@ pub extern "C" fn kvspaceTake( std::thread::sleep(Duration::from_millis(1)); } } + +#[no_mangle] +pub extern "C" fn kvspaceIncr( + h: *mut Handle, + key: *const c_char, + out: *mut i64, + err: *mut c_char, + err_cap: u32, +) -> c_int { + if h.is_null() || out.is_null() { + write_err(err, err_cap, "Incr: bad args"); + return 1; + } + unsafe { *out = 0; } + let kv: &mut dyn KVSpace = unsafe { &mut **h }; + let key = unsafe { cstr(key) }; + let cur = get_one(kv, key); + let mut n: i64 = 0; + if !is_none(&cur) { + if !cur.kind().starts_with("char/") { + write_err(err, err_cap, "Incr: counter is not a Char"); + return 1; + } + let s = cur.value_string(); + match s.parse::() { + Ok(v) => n = v, + Err(_) => { + write_err(err, err_cap, "Incr: unparsable counter"); + return 1; + } + } + } + if n == i64::MAX { + write_err(err, err_cap, "Incr: overflow"); + return 1; + } + n += 1; + let val = new_char(KIND_CHAR_UTF8, &n.to_string()); + match kv.set(&[KVPair { key: key.to_string(), val }]) { + Ok(()) => { + unsafe { *out = n; } + 0 + } + Err(e) => { + write_err(err, err_cap, &e); + 1 + } + } +} From a463ca0a09055ffd87a75def725089dccaafb46e Mon Sep 17 00:00:00 2001 From: junyao <1071307515@qq.com> Date: Thu, 20 Aug 2026 05:30:52 +0000 Subject: [PATCH 3/4] feat: Expire TTL --- src/backend.rs | 67 +++++++++++++++++++++++++++++++++++++++---- src/ffi.rs | 22 ++++++++++++++ src/fs/kvspace.rs | 39 +++++++++++++++++++++++-- src/kvspace.rs | 3 ++ src/kvspace_common.rs | 15 ++++++++++ src/redis/store.rs | 8 ++++++ src/store.rs | 4 +++ 7 files changed, 150 insertions(+), 8 deletions(-) diff --git a/src/backend.rs b/src/backend.rs index f5cccb1..b0bd4c9 100644 --- a/src/backend.rs +++ b/src/backend.rs @@ -1,12 +1,14 @@ // backend.rs — 对齐 redis/kvspace.go 的 KVSpace 实现逻辑,参数化于 KVStore 原语。 // redis 与 fs 后端共用这份逻辑,只替换底层 store。 -use std::time::Duration; +use std::cell::RefCell; +use std::collections::HashMap; +use std::time::{Duration, Instant}; use crate::r#const::*; use crate::kvspace::{KVSpace, KVPair}; use crate::kvspace_common::{ - dir_exists, get_one, join_path, mk_index_recursive, sep_path, split_index, validate_ptr, watch_value, + dir_exists, expire_key_ok, get_one, join_path, mk_index_recursive, sep_path, split_index, validate_ptr, watch_value, }; use crate::store::KVStore; use crate::xvalue::{decode_xvalue, decode_xvalue_head, is_none, is_ptr, ptr_target, XValue}; @@ -14,11 +16,47 @@ use crate::xvalue_index::{new_dict_index, new_ext_index, new_index}; pub struct Backend { store: S, + deadlines: RefCell>, } impl Backend { pub fn new(store: S) -> Self { - Backend { store } + Backend { store, deadlines: RefCell::new(HashMap::new()) } + } + + fn remove_child_exact(&self, parent: &str, name: &str) { + match self.store.get(parent) { + None => {} + Some(data) => { + let v = decode_xvalue(&data); + match v { + XValue::Index(nodes) => { + let filtered: Vec = normalize_children(nodes).into_iter().filter(|n| n != name).collect(); + self.store.set(parent, &new_index(&filtered).encode()); + } + XValue::Dict(nodes) => { + let filtered: Vec = normalize_children(nodes).into_iter().filter(|n| n != name).collect(); + self.store.set(parent, &new_dict_index(&filtered).encode()); + } + XValue::ExtIndex(e) => { + let filtered: Vec = e.childs.into_iter().filter(|n| n != name).collect(); + self.store.set(parent, &new_ext_index(&filtered, &e.ext_path).encode()); + } + other => panic!("remove_child_exact: unexpected kind {}", other.kind()), + } + } + } + } + + fn expired_now(&self, key: &str) -> bool { + let mut d = self.deadlines.borrow_mut(); + match d.get(key).copied() { + Some(t) if Instant::now() >= t => { + d.remove(key); + true + } + _ => false, + } } // ── 目录与路径工具 ────────────────────────────────────────────── @@ -259,8 +297,10 @@ impl KVSpace for Backend { let full_refs: Vec<&str> = full_keys.iter().map(|(_, f)| f.as_str()).collect(); let full_vals = self.store.get_many(&full_refs); let mut ext_keys: Vec<(usize, String)> = Vec::new(); - for (idx, (i, _)) in full_keys.iter().enumerate() { - if let Some(data) = &full_vals[idx] { + for (idx, (i, full)) in full_keys.iter().enumerate() { + if self.expired_now(full) { + results[*i] = Some(XValue::None); + } else if let Some(data) = &full_vals[idx] { results[*i] = Some(decode_xvalue(data)); } else if !ext_t.is_empty() { ext_keys.push((*i, join_path(&ext_t, &keys[*i]))); @@ -556,4 +596,21 @@ impl KVSpace for Backend { fn dis_conn(&mut self) -> Result<(), String> { Ok(()) } + + fn expire(&mut self, key: &str, ttl: Duration) -> Result<(), String> { + expire_key_ok(key)?; + if ttl.is_zero() { + return Err("Expire: ttl must be > 0".into()); + } + let resolved = self.resolve_path(key); + if self.store.get(&resolved).is_none() || self.expired_now(&resolved) { + return Err("Expire: missing key".into()); + } + let (parent, name, _) = split_index(&resolved); + self.remove_child_exact(&parent, &name); + if !self.store.pexpire(&resolved, ttl) { + self.deadlines.borrow_mut().insert(resolved, Instant::now() + ttl); + } + Ok(()) + } } diff --git a/src/ffi.rs b/src/ffi.rs index 4dcb0c8..9e16261 100644 --- a/src/ffi.rs +++ b/src/ffi.rs @@ -657,3 +657,25 @@ pub extern "C" fn kvspaceIncr( } } } + +#[no_mangle] +pub extern "C" fn kvspaceExpire( + h: *mut Handle, + key: *const c_char, + ttl_ns: u64, + err: *mut c_char, + err_cap: u32, +) -> c_int { + if h.is_null() { + write_err(err, err_cap, "Expire: bad args"); + return 1; + } + let kv: &mut dyn KVSpace = unsafe { &mut **h }; + match kv.expire(unsafe { cstr(key) }, Duration::from_nanos(ttl_ns)) { + Ok(()) => 0, + Err(e) => { + write_err(err, err_cap, &e); + 1 + } + } +} diff --git a/src/fs/kvspace.rs b/src/fs/kvspace.rs index a88e72f..2d8f45e 100644 --- a/src/fs/kvspace.rs +++ b/src/fs/kvspace.rs @@ -2,12 +2,13 @@ // 编码:kvspace 的 '.' 一律替换为 './'(父目录名带尾点 + "/" 分隔成员),反向 './' → '.'。 // "/" 与 "." 的 index 都从 readdir 派生;ExtIndex 用目录内 __extindex__ 文件存 ext_target_path(第一行)。 +use std::collections::HashMap; use std::fs; use std::path::{Path, PathBuf}; -use std::time::Duration; +use std::time::{Duration, Instant}; use crate::kvspace::{KVSpace, KVPair}; -use crate::kvspace_common::{join_path, sep_path, split_index, validate_ptr, watch_value, SepKind}; +use crate::kvspace_common::{expire_key_ok, join_path, sep_path, split_index, validate_ptr, watch_value, SepKind}; use crate::r#const::*; use crate::xvalue::*; use crate::xvalue_index::{new_dict_index, new_ext_index, new_index}; @@ -18,6 +19,7 @@ const ORDER_MARKER: &str = "__order__"; pub struct FsKVSpace { root: PathBuf, + deadlines: HashMap, } pub fn connect(root: &str) -> FsKVSpace { @@ -28,7 +30,7 @@ pub fn connect(root: &str) -> FsKVSpace { impl FsKVSpace { pub fn new(root: &str) -> Self { fs::create_dir_all(root).unwrap_or_else(|e| panic!("kvspace-fs: create root {}: {}", root, e)); - FsKVSpace { root: PathBuf::from(root) } + FsKVSpace { root: PathBuf::from(root), deadlines: HashMap::new() } } /// kvspace key → fs 路径:'.'(成员分隔)→ './';段首 '.'(如 .todo)是字面量不替换。 @@ -61,6 +63,17 @@ impl FsKVSpace { key.ends_with(DIR_INDEX_SUF) || key.ends_with(DICT_SEP) } + fn expired_now(&mut self, key: &str) -> bool { + match self.deadlines.get(key).copied() { + Some(t) if Instant::now() >= t => { + self.deadlines.remove(key); + let _ = fs::remove_file(self.fs_path(key)); + true + } + _ => false, + } + } + fn read_leaf(&self, key: &str) -> Option> { if key.contains("//") { return None; @@ -314,6 +327,9 @@ impl KVSpace for FsKVSpace { if Self::is_dir_key(&full) { return self.dir_value(&full); } + if self.expired_now(&full) { + return XValue::None; + } if let Some(data) = self.read_leaf(&full) { return decode_xvalue(&data); } @@ -393,6 +409,10 @@ impl KVSpace for FsKVSpace { return Vec::new(); } let mut members = self.dir_children(&resolved); + members.retain(|m| { + let full = join_path(&resolved, m); + !self.deadlines.contains_key(&full) + }); if expand_ext { let ext_t = self.prefix_ext(&resolved); @@ -493,6 +513,19 @@ impl KVSpace for FsKVSpace { fn dis_conn(&mut self) -> Result<(), String> { Ok(()) } + + fn expire(&mut self, key: &str, ttl: Duration) -> Result<(), String> { + expire_key_ok(key)?; + if ttl.is_zero() { + return Err("Expire: ttl must be > 0".into()); + } + let resolved = self.resolve_path(key); + if self.read_leaf(&resolved).is_none() || self.expired_now(&resolved) { + return Err("Expire: missing key".into()); + } + self.deadlines.insert(resolved, Instant::now() + ttl); + Ok(()) + } } // 供测试用:返回 root 下的顶层条目数。 diff --git a/src/kvspace.rs b/src/kvspace.rs index 31bc5c4..c5454e9 100644 --- a/src/kvspace.rs +++ b/src/kvspace.rs @@ -44,4 +44,7 @@ pub trait KVSpace { fn clear(&mut self) -> Result<(), String>; fn dis_conn(&mut self) -> Result<(), String>; + + /// Hide key from List immediately. Get still returns it until ttl. + fn expire(&mut self, key: &str, ttl: Duration) -> Result<(), String>; } diff --git a/src/kvspace_common.rs b/src/kvspace_common.rs index d132c6c..1849e21 100644 --- a/src/kvspace_common.rs +++ b/src/kvspace_common.rs @@ -7,6 +7,21 @@ use crate::kvspace::{KVSpace, KVPair}; use crate::xvalue::{body_bytes, is_none, plain, XValue}; use crate::xvalue_index::new_index; +pub fn expire_key_ok(key: &str) -> Result<(), String> { + if key.is_empty() || !key.starts_with('/') { + return Err("Expire: key is not an absolute path".into()); + } + if key.ends_with('/') { + return Err("Expire: directory".into()); + } + for seg in key[1..].split('/') { + if seg.is_empty() || seg == "." || seg == ".." { + return Err("Expire: key is not an absolute path".into()); + } + } + Ok(()) +} + /// JoinPath 拼接父子路径。 pub fn join_path(parent: &str, child: &str) -> String { if parent == PATH_SEP { diff --git a/src/redis/store.rs b/src/redis/store.rs index 4019709..3b39eb7 100644 --- a/src/redis/store.rs +++ b/src/redis/store.rs @@ -176,6 +176,14 @@ impl KVStore for RedisStore { keys } + fn pexpire(&self, key: &str, ttl: std::time::Duration) -> bool { + let ms = ttl.as_millis().max(1).to_string(); + match self.cmd(&[b"PEXPIRE", key.as_bytes(), ms.as_bytes()]) { + Resp::Integer(n) => n == 1, + _ => false, + } + } + fn flush(&self) { let _ = self.cmd(&[b"FLUSHDB"]); } diff --git a/src/store.rs b/src/store.rs index d189d21..d79de0f 100644 --- a/src/store.rs +++ b/src/store.rs @@ -13,4 +13,8 @@ pub trait KVStore { /// 返回所有以 prefix 开头的 key(含 prefix 自身,若存在)。 fn scan_keys(&self, prefix: &str) -> Vec; fn flush(&self); + /// Native TTL. true if the store will drop Get after ttl. false → caller hides Get itself. + fn pexpire(&self, _key: &str, _ttl: std::time::Duration) -> bool { + false + } } From 1ab3c07c7526ab9145812ca31a1409a9402134b2 Mon Sep 17 00:00:00 2001 From: junyao <1071307515@qq.com> Date: Thu, 20 Aug 2026 05:39:09 +0000 Subject: [PATCH 4/4] feat: WatchAny multi-key wait --- src/ffi.rs | 61 +++++++++++++++++++++++++++++++++++++++++++++++++----- 1 file changed, 56 insertions(+), 5 deletions(-) diff --git a/src/ffi.rs b/src/ffi.rs index 9e16261..572dedb 100644 --- a/src/ffi.rs +++ b/src/ffi.rs @@ -593,13 +593,9 @@ pub extern "C" fn kvspaceTake( } let kv: &mut dyn KVSpace = unsafe { &mut **h }; let key = unsafe { cstr(key) }; - let qk = nq_key(key); let deadline = std::time::Instant::now() + Duration::from_nanos(timeout_ns); loop { - let mut frames = nq_frames(&nq_load(kv, &qk)); - if !frames.is_empty() { - let item = frames.remove(0); - let _ = nq_save(kv, &qk, &frames); + if let Some(item) = nq_try_pop(kv, key) { return alloc(item, out, out_len); } if std::time::Instant::now() >= deadline { @@ -609,6 +605,61 @@ pub extern "C" fn kvspaceTake( } } +fn nq_try_pop(kv: &mut dyn KVSpace, key: &str) -> Option> { + let qk = nq_key(key); + let mut frames = nq_frames(&nq_load(kv, &qk)); + if frames.is_empty() { + return None; + } + let item = frames.remove(0); + let _ = nq_save(kv, &qk, &frames); + Some(item) +} + +#[no_mangle] +pub extern "C" fn kvspaceWatchAny( + h: *mut Handle, + keys: *const *const c_char, + nkeys: u32, + timeout_ns: u64, + out_key: *mut *mut u8, + out_key_len: *mut u32, + out: *mut *mut u8, + out_len: *mut u32, +) -> c_int { + if out_key.is_null() || out_key_len.is_null() || out.is_null() || out_len.is_null() { + return 1; + } + unsafe { + *out_key = std::ptr::null_mut(); + *out_key_len = 0; + *out = std::ptr::null_mut(); + *out_len = 0; + } + if h.is_null() || keys.is_null() || nkeys == 0 { + return 1; + } + let kv: &mut dyn KVSpace = unsafe { &mut **h }; + let deadline = std::time::Instant::now() + Duration::from_nanos(timeout_ns); + let ks: &[*const c_char] = unsafe { std::slice::from_raw_parts(keys, nkeys as usize) }; + loop { + for p in ks { + let key = unsafe { cstr(*p) }; + if key.is_empty() { + continue; + } + if let Some(item) = nq_try_pop(kv, key) { + let _ = alloc(key.as_bytes().to_vec(), out_key, out_key_len); + return alloc(item, out, out_len); + } + } + if std::time::Instant::now() >= deadline { + return 0; + } + std::thread::sleep(Duration::from_millis(1)); + } +} + #[no_mangle] pub extern "C" fn kvspaceIncr( h: *mut Handle,