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
10 changes: 10 additions & 0 deletions crates/ember-core/src/memory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,16 @@
//! Provides byte-level accounting of memory used by entries. Updated
//! on every mutation so the engine can enforce memory limits and
//! report stats without scanning the entire keyspace.
//!
//! # Platform notes
//!
//! Overhead constants are empirical estimates for 64-bit platforms (x86-64,
//! aarch64). On 32-bit systems these would be smaller; the effect is that
//! we'd overestimate memory usage, which triggers eviction earlier than
//! necessary but doesn't cause correctness issues.
//!
//! The constants assume Rust's standard library allocator. Custom allocators
//! (jemalloc, mimalloc) may have different per-allocation overhead.

use crate::types::Value;

Expand Down
156 changes: 156 additions & 0 deletions crates/ember-core/src/shard.rs
Original file line number Diff line number Diff line change
Expand Up @@ -764,6 +764,39 @@ fn to_aof_record(req: &ShardRequest, resp: &ShardResponse) -> Option<AofRecord>
milliseconds: *milliseconds,
})
}
// Hash commands
(ShardRequest::HSet { key, fields }, ShardResponse::Len(_)) => Some(AofRecord::HSet {
key: key.clone(),
fields: fields.clone(),
}),
(ShardRequest::HDel { key, .. }, ShardResponse::HDelLen { removed, .. })
if !removed.is_empty() =>
{
Some(AofRecord::HDel {
key: key.clone(),
fields: removed.clone(),
})
}
(ShardRequest::HIncrBy { key, field, delta }, ShardResponse::Integer(_)) => {
Some(AofRecord::HIncrBy {
key: key.clone(),
field: field.clone(),
delta: *delta,
})
}
// Set commands
(ShardRequest::SAdd { key, members }, ShardResponse::Len(count)) if *count > 0 => {
Some(AofRecord::SAdd {
key: key.clone(),
members: members.clone(),
})
}
(ShardRequest::SRem { key, members }, ShardResponse::Len(count)) if *count > 0 => {
Some(AofRecord::SRem {
key: key.clone(),
members: members.clone(),
})
}
_ => None,
}
}
Expand Down Expand Up @@ -1501,4 +1534,127 @@ mod tests {
_ => panic!("expected Scan response"),
}
}

#[test]
fn to_aof_record_for_hset() {
let req = ShardRequest::HSet {
key: "h".into(),
fields: vec![("f1".into(), Bytes::from("v1"))],
};
let resp = ShardResponse::Len(1);
let record = to_aof_record(&req, &resp).unwrap();
match record {
AofRecord::HSet { key, fields } => {
assert_eq!(key, "h");
assert_eq!(fields.len(), 1);
}
_ => panic!("expected HSet record"),
}
}

#[test]
fn to_aof_record_for_hdel() {
let req = ShardRequest::HDel {
key: "h".into(),
fields: vec!["f1".into(), "f2".into()],
};
let resp = ShardResponse::HDelLen {
count: 2,
removed: vec!["f1".into(), "f2".into()],
};
let record = to_aof_record(&req, &resp).unwrap();
match record {
AofRecord::HDel { key, fields } => {
assert_eq!(key, "h");
assert_eq!(fields.len(), 2);
}
_ => panic!("expected HDel record"),
}
}

#[test]
fn to_aof_record_skips_hdel_when_none_removed() {
let req = ShardRequest::HDel {
key: "h".into(),
fields: vec!["f1".into()],
};
let resp = ShardResponse::HDelLen {
count: 0,
removed: vec![],
};
assert!(to_aof_record(&req, &resp).is_none());
}

#[test]
fn to_aof_record_for_hincrby() {
let req = ShardRequest::HIncrBy {
key: "h".into(),
field: "counter".into(),
delta: 5,
};
let resp = ShardResponse::Integer(10);
let record = to_aof_record(&req, &resp).unwrap();
match record {
AofRecord::HIncrBy { key, field, delta } => {
assert_eq!(key, "h");
assert_eq!(field, "counter");
assert_eq!(delta, 5);
}
_ => panic!("expected HIncrBy record"),
}
}

#[test]
fn to_aof_record_for_sadd() {
let req = ShardRequest::SAdd {
key: "s".into(),
members: vec!["m1".into(), "m2".into()],
};
let resp = ShardResponse::Len(2);
let record = to_aof_record(&req, &resp).unwrap();
match record {
AofRecord::SAdd { key, members } => {
assert_eq!(key, "s");
assert_eq!(members.len(), 2);
}
_ => panic!("expected SAdd record"),
}
}

#[test]
fn to_aof_record_skips_sadd_when_none_added() {
let req = ShardRequest::SAdd {
key: "s".into(),
members: vec!["m1".into()],
};
let resp = ShardResponse::Len(0);
assert!(to_aof_record(&req, &resp).is_none());
}

#[test]
fn to_aof_record_for_srem() {
let req = ShardRequest::SRem {
key: "s".into(),
members: vec!["m1".into()],
};
let resp = ShardResponse::Len(1);
let record = to_aof_record(&req, &resp).unwrap();
match record {
AofRecord::SRem { key, members } => {
assert_eq!(key, "s");
assert_eq!(members.len(), 1);
}
_ => panic!("expected SRem record"),
}
}

#[test]
fn to_aof_record_skips_srem_when_none_removed() {
let req = ShardRequest::SRem {
key: "s".into(),
members: vec!["m1".into()],
};
let resp = ShardResponse::Len(0);
assert!(to_aof_record(&req, &resp).is_none());
}
}
164 changes: 164 additions & 0 deletions crates/ember-persistence/src/aof.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,11 @@ const TAG_PERSIST: u8 = 10;
const TAG_PEXPIRE: u8 = 11;
const TAG_INCR: u8 = 12;
const TAG_DECR: u8 = 13;
const TAG_HSET: u8 = 14;
const TAG_HDEL: u8 = 15;
const TAG_HINCRBY: u8 = 16;
const TAG_SADD: u8 = 17;
const TAG_SREM: u8 = 18;

/// A single mutation record stored in the AOF.
#[derive(Debug, Clone, PartialEq)]
Expand Down Expand Up @@ -86,6 +91,23 @@ pub enum AofRecord {
Incr { key: String },
/// DECR key.
Decr { key: String },
/// HSET key field value [field value ...].
HSet {
key: String,
fields: Vec<(String, Bytes)>,
},
/// HDEL key field [field ...].
HDel { key: String, fields: Vec<String> },
/// HINCRBY key field delta.
HIncrBy {
key: String,
field: String,
delta: i64,
},
/// SADD key member [member ...].
SAdd { key: String, members: Vec<String> },
/// SREM key member [member ...].
SRem { key: String, members: Vec<String> },
}

impl AofRecord {
Expand Down Expand Up @@ -170,6 +192,45 @@ impl AofRecord {
format::write_u8(&mut buf, TAG_DECR).expect("vec write");
format::write_bytes(&mut buf, key.as_bytes()).expect("vec write");
}
AofRecord::HSet { key, fields } => {
format::write_u8(&mut buf, TAG_HSET).expect("vec write");
format::write_bytes(&mut buf, key.as_bytes()).expect("vec write");
format::write_u32(&mut buf, fields.len() as u32).expect("vec write");
for (field, value) in fields {
format::write_bytes(&mut buf, field.as_bytes()).expect("vec write");
format::write_bytes(&mut buf, value).expect("vec write");
}
}
AofRecord::HDel { key, fields } => {
format::write_u8(&mut buf, TAG_HDEL).expect("vec write");
format::write_bytes(&mut buf, key.as_bytes()).expect("vec write");
format::write_u32(&mut buf, fields.len() as u32).expect("vec write");
for field in fields {
format::write_bytes(&mut buf, field.as_bytes()).expect("vec write");
}
}
AofRecord::HIncrBy { key, field, delta } => {
format::write_u8(&mut buf, TAG_HINCRBY).expect("vec write");
format::write_bytes(&mut buf, key.as_bytes()).expect("vec write");
format::write_bytes(&mut buf, field.as_bytes()).expect("vec write");
format::write_i64(&mut buf, *delta).expect("vec write");
}
AofRecord::SAdd { key, members } => {
format::write_u8(&mut buf, TAG_SADD).expect("vec write");
format::write_bytes(&mut buf, key.as_bytes()).expect("vec write");
format::write_u32(&mut buf, members.len() as u32).expect("vec write");
for member in members {
format::write_bytes(&mut buf, member.as_bytes()).expect("vec write");
}
}
AofRecord::SRem { key, members } => {
format::write_u8(&mut buf, TAG_SREM).expect("vec write");
format::write_bytes(&mut buf, key.as_bytes()).expect("vec write");
format::write_u32(&mut buf, members.len() as u32).expect("vec write");
for member in members {
format::write_bytes(&mut buf, member.as_bytes()).expect("vec write");
}
}
}
buf
}
Expand Down Expand Up @@ -256,6 +317,50 @@ impl AofRecord {
let key = read_string(&mut cursor, "key")?;
Ok(AofRecord::Decr { key })
}
TAG_HSET => {
let key = read_string(&mut cursor, "key")?;
let count = format::read_u32(&mut cursor)?;
let mut fields = Vec::with_capacity(count as usize);
for _ in 0..count {
let field = read_string(&mut cursor, "field")?;
let value = Bytes::from(format::read_bytes(&mut cursor)?);
fields.push((field, value));
}
Ok(AofRecord::HSet { key, fields })
}
TAG_HDEL => {
let key = read_string(&mut cursor, "key")?;
let count = format::read_u32(&mut cursor)?;
let mut fields = Vec::with_capacity(count as usize);
for _ in 0..count {
fields.push(read_string(&mut cursor, "field")?);
}
Ok(AofRecord::HDel { key, fields })
}
TAG_HINCRBY => {
let key = read_string(&mut cursor, "key")?;
let field = read_string(&mut cursor, "field")?;
let delta = format::read_i64(&mut cursor)?;
Ok(AofRecord::HIncrBy { key, field, delta })
}
TAG_SADD => {
let key = read_string(&mut cursor, "key")?;
let count = format::read_u32(&mut cursor)?;
let mut members = Vec::with_capacity(count as usize);
for _ in 0..count {
members.push(read_string(&mut cursor, "member")?);
}
Ok(AofRecord::SAdd { key, members })
}
TAG_SREM => {
let key = read_string(&mut cursor, "key")?;
let count = format::read_u32(&mut cursor)?;
let mut members = Vec::with_capacity(count as usize);
for _ in 0..count {
members.push(read_string(&mut cursor, "member")?);
}
Ok(AofRecord::SRem { key, members })
}
_ => Err(FormatError::UnknownTag(tag)),
}
}
Expand Down Expand Up @@ -893,4 +998,63 @@ mod tests {
let p = aof_path(Path::new("/data"), 3);
assert_eq!(p, PathBuf::from("/data/shard-3.aof"));
}

#[test]
fn record_round_trip_hset() {
let rec = AofRecord::HSet {
key: "hash".into(),
fields: vec![
("f1".into(), Bytes::from("v1")),
("f2".into(), Bytes::from("v2")),
],
};
let bytes = rec.to_bytes();
let decoded = AofRecord::from_bytes(&bytes).unwrap();
assert_eq!(rec, decoded);
}

#[test]
fn record_round_trip_hdel() {
let rec = AofRecord::HDel {
key: "hash".into(),
fields: vec!["f1".into(), "f2".into()],
};
let bytes = rec.to_bytes();
let decoded = AofRecord::from_bytes(&bytes).unwrap();
assert_eq!(rec, decoded);
}

#[test]
fn record_round_trip_hincrby() {
let rec = AofRecord::HIncrBy {
key: "hash".into(),
field: "counter".into(),
delta: -42,
};
let bytes = rec.to_bytes();
let decoded = AofRecord::from_bytes(&bytes).unwrap();
assert_eq!(rec, decoded);
}

#[test]
fn record_round_trip_sadd() {
let rec = AofRecord::SAdd {
key: "set".into(),
members: vec!["m1".into(), "m2".into(), "m3".into()],
};
let bytes = rec.to_bytes();
let decoded = AofRecord::from_bytes(&bytes).unwrap();
assert_eq!(rec, decoded);
}

#[test]
fn record_round_trip_srem() {
let rec = AofRecord::SRem {
key: "set".into(),
members: vec!["m1".into()],
};
let bytes = rec.to_bytes();
let decoded = AofRecord::from_bytes(&bytes).unwrap();
assert_eq!(rec, decoded);
}
}
Loading