wedb_embed 0.1.0

Embedded Kvrocks-compatible storage engine for WeDb
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
use rapidhash::v3::rapidhash_v3;
use serde::{Deserialize, Serialize};
use std::sync::atomic::{AtomicU64, Ordering};

/// Redis 数据类型枚举(对标 Apache Kvrocks RedisType)
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[repr(u8)]
pub enum RedisType {
    None = 0,
    String = 1,
    Hash = 2,
    List = 3,
    Set = 4,
    ZSet = 5,
    Bitmap = 6,
    SortedInt = 7,
    Stream = 8,
    Bloom = 9,
    Json = 10,
    HyperLogLog = 11,
    TDigest = 12,
    TimeSeries = 13,
    CuckooFilter = 14,
}

impl RedisType {
    #[inline]
    pub const fn name(&self) -> &'static str {
        match self {
            Self::None => "none",
            Self::String => "string",
            Self::Hash => "hash",
            Self::List => "list",
            Self::Set => "set",
            Self::ZSet => "zset",
            Self::Bitmap => "bitmap",
            Self::SortedInt => "sortedint",
            Self::Stream => "stream",
            Self::Bloom => "MBbloom--",
            Self::Json => "ReJSON-RL",
            Self::HyperLogLog => "hyperloglog",
            Self::TDigest => "TDIS-TYPE",
            Self::TimeSeries => "timeseries",
            Self::CuckooFilter => "MBbloomCF",
        }
    }

    #[inline]
    pub const fn from_u8(val: u8) -> Self {
        match val {
            1 => Self::String,
            2 => Self::Hash,
            3 => Self::List,
            4 => Self::Set,
            5 => Self::ZSet,
            6 => Self::Bitmap,
            7 => Self::SortedInt,
            8 => Self::Stream,
            9 => Self::Bloom,
            10 => Self::Json,
            11 => Self::HyperLogLog,
            12 => Self::TDigest,
            13 => Self::TimeSeries,
            14 => Self::CuckooFilter,
            _ => Self::None,
        }
    }

    #[inline]
    pub const fn is_single_kv_type(&self) -> bool {
        matches!(self, Self::String | Self::Json)
    }

    #[inline]
    pub const fn is_emptyable_type(&self) -> bool {
        matches!(
            self,
            Self::String
                | Self::Json
                | Self::Stream
                | Self::Bloom
                | Self::HyperLogLog
                | Self::TDigest
                | Self::TimeSeries
                | Self::CuckooFilter
        )
    }
}

// ================= Version 生成机制(对标 Apache Kvrocks 53-bit 时间戳 + 11-bit 计数器) =================

pub const VERSION_COUNTER_BITS: u32 = 11;
pub const VERSION_COUNTER_MASK: u64 = (1 << VERSION_COUNTER_BITS) - 1;

static VERSION_COUNTER: AtomicU64 = AtomicU64::new(0);

/// 初始化版本计数器(基于微秒与 rapidhash 生成随机初始偏移,避免主从切换时时钟回退冲突)
pub fn init_version_counter() {
    let now_nanos = coarsetime::Clock::now_since_epoch().as_nanos();
    let seed = rapidhash_v3(&now_nanos.to_be_bytes());
    VERSION_COUNTER.store(seed, Ordering::Relaxed);
}

/// 生成唯一递增版本号:高 53 位微秒时间戳 + 低 11 位原子计数器
#[inline]
pub fn generate_version() -> u64 {
    let ts_us = coarsetime::Clock::now_since_epoch().as_micros();
    let counter = VERSION_COUNTER.fetch_add(1, Ordering::Relaxed);
    (ts_us << VERSION_COUNTER_BITS) | (counter & VERSION_COUNTER_MASK)
}

/// 从版本号中解析出创建时间(秒与微秒,对标 Kvrocks Metadata::Time)
#[inline]
pub fn version_to_time(version: u64) -> (u64, u32) {
    let ts_us = version >> VERSION_COUNTER_BITS;
    let sec = ts_us / 1_000_000;
    let usec = (ts_us % 1_000_000) as u32;
    (sec, usec)
}

/// 基础通用 26 字节元数据结构(对标 Apache Kvrocks KeyMetadata)
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct KeyMeta {
    pub rtype: RedisType,
    pub flags: u8,
    pub expire_at: u64,
    pub version: u64,
    pub size: u64,
}

impl KeyMeta {
    pub const META_64BIT_ENCODING_MASK: u8 = 0x80;
    pub const META_TYPE_MASK: u8 = 0x0F;
    pub const ENCODED_SIZE: usize = 26; // 1B rtype + 1B flags + 8B expire + 8B version + 8B size
    pub const KVROCKS_COMPLEX_ENCODED_SIZE: usize = 25; // 1B flags(type+64bit) + 8B expire + 8B version + 8B size
    pub const KVROCKS_SINGLE_KV_ENCODED_SIZE: usize = 9; // 1B flags + 8B expire

    #[inline]
    pub const fn new(rtype: RedisType, expire_at: u64, version: u64, size: u64) -> Self {
        Self {
            rtype,
            flags: 0,
            expire_at,
            version,
            size,
        }
    }

    #[inline]
    pub fn new_with_version(rtype: RedisType, expire_at: u64, size: u64) -> Self {
        Self {
            rtype,
            flags: 0,
            expire_at,
            version: generate_version(),
            size,
        }
    }

    #[inline]
    pub const fn is_expired(&self, now_ms: u64) -> bool {
        if !self.is_emptyable_type() && self.size == 0 {
            return true;
        }
        self.expire_at > 0 && self.expire_at <= now_ms
    }

    #[inline]
    pub const fn is_single_kv_type(&self) -> bool {
        self.rtype.is_single_kv_type()
    }

    #[inline]
    pub const fn is_emptyable_type(&self) -> bool {
        self.rtype.is_emptyable_type()
    }

    #[inline]
    pub const fn ttl(&self, now_ms: u64) -> i64 {
        if self.expire_at == 0 {
            -1
        } else if self.expire_at < now_ms {
            -2
        } else {
            (self.expire_at - now_ms) as i64
        }
    }

    #[inline]
    pub fn expire_at_ms_to_sec(ms: u64) -> u64 {
        if ms == 0 {
            0
        } else if ms < 1000 {
            1
        } else {
            (ms + 499) / 1000
        }
    }

    #[inline]
    pub const fn is_64bit_encoded_flags(flags: u8) -> bool {
        flags & Self::META_64BIT_ENCODING_MASK != 0
    }

    #[inline]
    pub const fn is_64bit_encoded(&self) -> bool {
        Self::is_64bit_encoded_flags(self.flags)
    }

    #[inline]
    pub const fn common_encoded_size(&self) -> usize {
        if self.is_64bit_encoded() { 8 } else { 4 }
    }

    /// 获取元数据过期时间之后的字节偏移量(对标 Kvrocks Metadata::GetOffsetAfterExpire)
    #[inline]
    pub const fn get_offset_after_expire(flags: u8) -> usize {
        if Self::is_64bit_encoded_flags(flags) {
            1 + 8 // 1B flags + 8B expire
        } else {
            1 + 4 // 1B flags + 4B expire
        }
    }

    /// 获取复合元数据大小之后的字节偏移量(对标 Kvrocks Metadata::GetOffsetAfterSize)
    #[inline]
    pub const fn get_offset_after_size(flags: u8) -> usize {
        if Self::is_64bit_encoded_flags(flags) {
            1 + 8 + 8 + 8 // 1B flags + 8B expire + 8B version + 8B size
        } else {
            1 + 4 + 8 + 4 // 1B flags + 4B expire + 8B version + 4B size
        }
    }

    /// 编码为标准 26 字节元数据头
    #[inline]
    pub fn encode(&self) -> [u8; Self::ENCODED_SIZE] {
        let mut buf = [0u8; Self::ENCODED_SIZE];
        buf[0] = self.rtype as u8;
        buf[1] = self.flags;
        buf[2..10].copy_from_slice(&self.expire_at.to_be_bytes());
        buf[10..18].copy_from_slice(&self.version.to_be_bytes());
        buf[18..26].copy_from_slice(&self.size.to_be_bytes());
        buf
    }

    /// 编码为 Kvrocks 1:1 紧凑二进制格式(25字节复合类型 / 9字节SingleKV)
    #[inline]
    pub fn encode_kvrocks(&self) -> Vec<u8> {
        let flags = Self::META_64BIT_ENCODING_MASK | (self.rtype as u8 & Self::META_TYPE_MASK);
        if self.is_single_kv_type() {
            let mut out = Vec::with_capacity(Self::KVROCKS_SINGLE_KV_ENCODED_SIZE);
            out.push(flags);
            out.extend_from_slice(&self.expire_at.to_be_bytes());
            out
        } else {
            let mut out = Vec::with_capacity(Self::KVROCKS_COMPLEX_ENCODED_SIZE);
            out.push(flags);
            out.extend_from_slice(&self.expire_at.to_be_bytes());
            out.extend_from_slice(&self.version.to_be_bytes());
            out.extend_from_slice(&self.size.to_be_bytes());
            out
        }
    }

    /// 解码元数据头(自适应支持 26 字节标准头、Kvrocks 25 字节复合头及 9 字节 SingleKV 头)
    #[inline]
    pub fn decode(bytes: &[u8]) -> Option<Self> {
        if bytes.len() >= Self::ENCODED_SIZE
            && bytes[0] <= 14
            && (bytes[1] == 0 || bytes[1] == 0x80)
        {
            let rtype = RedisType::from_u8(bytes[0]);
            let flags = bytes[1];
            let mut exp_buf = [0u8; 8];
            exp_buf.copy_from_slice(&bytes[2..10]);
            let expire_at = u64::from_be_bytes(exp_buf);

            let mut ver_buf = [0u8; 8];
            ver_buf.copy_from_slice(&bytes[10..18]);
            let version = u64::from_be_bytes(ver_buf);

            let mut size_buf = [0u8; 8];
            size_buf.copy_from_slice(&bytes[18..26]);
            let size = u64::from_be_bytes(size_buf);

            return Some(Self {
                rtype,
                flags,
                expire_at,
                version,
                size,
            });
        }

        // Kvrocks 64位格式解析 (flags 最高位为 1)
        if !bytes.is_empty() && (bytes[0] & Self::META_64BIT_ENCODING_MASK != 0) {
            let flags = bytes[0];
            let rtype = RedisType::from_u8(flags & Self::META_TYPE_MASK);
            if rtype.is_single_kv_type() {
                if bytes.len() < Self::KVROCKS_SINGLE_KV_ENCODED_SIZE {
                    return None;
                }
                let mut exp_buf = [0u8; 8];
                exp_buf.copy_from_slice(&bytes[1..9]);
                let expire_at = u64::from_be_bytes(exp_buf);
                return Some(Self {
                    rtype,
                    flags,
                    expire_at,
                    version: 0,
                    size: 0,
                });
            } else if bytes.len() >= Self::KVROCKS_COMPLEX_ENCODED_SIZE {
                let mut exp_buf = [0u8; 8];
                exp_buf.copy_from_slice(&bytes[1..9]);
                let expire_at = u64::from_be_bytes(exp_buf);

                let mut ver_buf = [0u8; 8];
                ver_buf.copy_from_slice(&bytes[9..17]);
                let version = u64::from_be_bytes(ver_buf);

                let mut size_buf = [0u8; 8];
                size_buf.copy_from_slice(&bytes[17..25]);
                let size = u64::from_be_bytes(size_buf);

                return Some(Self {
                    rtype,
                    flags,
                    expire_at,
                    version,
                    size,
                });
            }
        }

        // 默认尝试 26 字节回退解码
        if bytes.len() >= Self::ENCODED_SIZE {
            let rtype = RedisType::from_u8(bytes[0]);
            let flags = bytes[1];
            let mut exp_buf = [0u8; 8];
            exp_buf.copy_from_slice(&bytes[2..10]);
            let expire_at = u64::from_be_bytes(exp_buf);

            let mut ver_buf = [0u8; 8];
            ver_buf.copy_from_slice(&bytes[10..18]);
            let version = u64::from_be_bytes(ver_buf);

            let mut size_buf = [0u8; 8];
            size_buf.copy_from_slice(&bytes[18..26]);
            let size = u64::from_be_bytes(size_buf);

            Some(Self {
                rtype,
                flags,
                expire_at,
                version,
                size,
            })
        } else {
            None
        }
    }
}

/// 归一化 Redis 索引范围(支持负数索引)
#[inline]
pub fn normalize_range(start: i64, stop: i64, len: i64) -> (i64, i64) {
    if len <= 0 {
        return (0, -1);
    }
    let mut s = if start < 0 { len + start } else { start };
    let mut e = if stop < 0 { len + stop } else { stop };
    if s < 0 {
        s = 0;
    }
    if e >= len {
        e = len - 1;
    }
    (s, e)
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn test_version_generation_and_time() {
        init_version_counter();
        let v1 = generate_version();
        let v2 = generate_version();
        assert!(v2 > v1);

        let (sec, usec) = version_to_time(v1);
        assert!(sec > 1_700_000_000);
        assert!(usec < 1_000_000);
    }

    #[test]
    fn test_key_meta_encode_decode_roundtrip() {
        let meta = KeyMeta::new(RedisType::Hash, 1_800_000_000_000, 123456789, 42);
        let encoded = meta.encode();
        assert_eq!(encoded.len(), KeyMeta::ENCODED_SIZE);

        let decoded = KeyMeta::decode(&encoded).expect("decode failed");
        assert_eq!(decoded.rtype, RedisType::Hash);
        assert_eq!(decoded.expire_at, 1_800_000_000_000);
        assert_eq!(decoded.version, 123456789);
        assert_eq!(decoded.size, 42);
    }

    #[test]
    fn test_key_meta_kvrocks_compatibility() {
        let meta = KeyMeta::new(RedisType::Set, 2_000_000_000_000, 9999, 10);
        let kvrocks_enc = meta.encode_kvrocks();
        assert_eq!(kvrocks_enc.len(), KeyMeta::KVROCKS_COMPLEX_ENCODED_SIZE);

        let decoded = KeyMeta::decode(&kvrocks_enc).expect("decode kvrocks failed");
        assert_eq!(decoded.rtype, RedisType::Set);
        assert_eq!(decoded.expire_at, 2_000_000_000_000);
        assert_eq!(decoded.version, 9999);
        assert_eq!(decoded.size, 10);
    }
}