Skip to content
Open
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
67 changes: 62 additions & 5 deletions src/backend.rs
Original file line number Diff line number Diff line change
@@ -1,24 +1,62 @@
// 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};
use crate::xvalue_index::{new_dict_index, new_ext_index, new_index};

pub struct Backend<S: KVStore> {
store: S,
deadlines: RefCell<HashMap<String, Instant>>,
}

impl<S: KVStore> Backend<S> {
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<String> = normalize_children(nodes).into_iter().filter(|n| n != name).collect();
self.store.set(parent, &new_index(&filtered).encode());
}
XValue::Dict(nodes) => {
let filtered: Vec<String> = 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<String> = 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,
}
}

// ── 目录与路径工具 ──────────────────────────────────────────────
Expand Down Expand Up @@ -259,8 +297,10 @@ impl<S: KVStore> KVSpace for Backend<S> {
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])));
Expand Down Expand Up @@ -556,4 +596,21 @@ impl<S: KVStore> KVSpace for Backend<S> {
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(())
}
}
231 changes: 230 additions & 1 deletion src/ffi.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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_CHAR_UTF8, 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};
Expand Down Expand Up @@ -501,3 +502,231 @@ 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<u8> {
let v = get_one(kv, qk);
if is_none(&v) {
return Vec::new();
}
body_bytes(&v)
}

fn nq_frames(body: &[u8]) -> Vec<Vec<u8>> {
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<u8>]) -> 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 deadline = std::time::Instant::now() + Duration::from_nanos(timeout_ns);
loop {
if let Some(item) = nq_try_pop(kv, key) {
return alloc(item, out, out_len);
}
if std::time::Instant::now() >= deadline {
return 0;
}
std::thread::sleep(Duration::from_millis(1));
}
}

fn nq_try_pop(kv: &mut dyn KVSpace, key: &str) -> Option<Vec<u8>> {
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,
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::<i64>() {
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
}
}
}

#[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
}
}
}
Loading