From 9d3be1e5d77c724ec5095ef81cd7bb464ed94b00 Mon Sep 17 00:00:00 2001 From: Mia Date: Thu, 19 Jun 2025 19:40:54 +0000 Subject: [PATCH] feat(index): Switch Index to RocksDB --- Cargo.lock | 57 ++++++- consumer/src/db/record.rs | 11 +- consumer/src/db/sql/feedgen_upsert.sql | 3 +- consumer/src/db/sql/list_upsert.sql | 3 +- consumer/src/db/sql/starterpack_upsert.sql | 3 +- consumer/src/indexer/mod.rs | 25 ++-- parakeet-index/Cargo.toml | 4 +- parakeet-index/Dockerfile | 2 +- parakeet-index/src/server/db.rs | 165 ++++++++++++++------- parakeet-index/src/server/service.rs | 75 +++------- parakeet-index/src/server/utils.rs | 32 +--- 11 files changed, 226 insertions(+), 154 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index f2d77f48..f0c854ec 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -544,6 +544,16 @@ version = "1.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "325918d6fe32f23b19878fe4b34794ae41fc19ddbe53b10571a4874d44ffd39b" +[[package]] +name = "bzip2-sys" +version = "0.1.13+1.0.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "225bff33b2141874fe80d71e07d6eec4f85c5c216453dd96388240f96e1acc14" +dependencies = [ + "cc", + "pkg-config", +] + [[package]] name = "cbor4ii" version = "0.2.14" @@ -2224,6 +2234,31 @@ version = "0.2.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8355be11b20d696c8f18f6cc018c4e372165b1fa8126cef092399c9951984ffa" +[[package]] +name = "librocksdb-sys" +version = "0.17.1+9.9.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b7869a512ae9982f4d46ba482c2a304f1efd80c6412a3d4bf57bb79a619679f" +dependencies = [ + "bindgen", + "bzip2-sys", + "cc", + "libc", + "libz-sys", + "lz4-sys", +] + +[[package]] +name = "libz-sys" +version = "1.1.22" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b70e7a7df205e92a1a4cd9aaae7898dac0aa555503cc0a649494d0d60e7651d" +dependencies = [ + "cc", + "pkg-config", + "vcpkg", +] + [[package]] name = "linked-hash-map" version = "0.5.6" @@ -2270,6 +2305,16 @@ dependencies = [ "linked-hash-map", ] +[[package]] +name = "lz4-sys" +version = "1.11.1+lz4-1.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6bd8c0d6c6ed0cd30b3652886bb8711dc4bb01d637a68105a3d5158039b418e6" +dependencies = [ + "cc", + "libc", +] + [[package]] name = "match_cfg" version = "0.1.0" @@ -2651,8 +2696,8 @@ dependencies = [ "figment", "itertools 0.14.0", "prost", + "rocksdb", "serde", - "sled", "tokio", "tonic", "tonic-build", @@ -3306,6 +3351,16 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rocksdb" +version = "0.23.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "26ec73b20525cb235bad420f911473b69f9fe27cc856c5461bccd7e4af037f43" +dependencies = [ + "libc", + "librocksdb-sys", +] + [[package]] name = "rsa" version = "0.9.8" diff --git a/consumer/src/db/record.rs b/consumer/src/db/record.rs index 7bfcc580..0fde4652 100644 --- a/consumer/src/db/record.rs +++ b/consumer/src/db/record.rs @@ -1,4 +1,4 @@ -use super::{PgExecResult, PgOptResult}; +use super::{PgExecResult, PgOptResult, PgResult}; use crate::indexer::records::*; use crate::utils::{blob_ref, strongref_to_parts}; use chrono::prelude::*; @@ -61,7 +61,7 @@ pub async fn feedgen_upsert( repo: &str, cid: Cid, rec: AppBskyFeedGenerator, -) -> PgExecResult { +) -> PgResult { let cid = cid.to_string(); let description_facets = rec .description_facets @@ -84,6 +84,7 @@ pub async fn feedgen_upsert( ], ) .await + .map(|v| v == 0) } pub async fn feedgen_delete(conn: &mut C, at_uri: &str) -> PgExecResult { @@ -181,7 +182,7 @@ pub async fn list_upsert( repo: &str, cid: Cid, rec: AppBskyGraphList, -) -> PgExecResult { +) -> PgResult { let cid = cid.to_string(); let description_facets = rec .description_facets @@ -203,6 +204,7 @@ pub async fn list_upsert( ], ) .await + .map(|v| v == 0) } pub async fn list_delete(conn: &mut C, at_uri: &str) -> PgExecResult { @@ -559,7 +561,7 @@ pub async fn starter_pack_upsert( repo: &str, cid: Cid, rec: AppBskyGraphStarterPack, -) -> PgExecResult { +) -> PgResult { let cid = cid.to_string(); let record = serde_json::to_value(&rec).unwrap(); let description_facets = rec @@ -585,6 +587,7 @@ pub async fn starter_pack_upsert( ], ) .await + .map(|v| v == 0) } pub async fn starter_pack_delete(conn: &mut C, at_uri: &str) -> PgExecResult { diff --git a/consumer/src/db/sql/feedgen_upsert.sql b/consumer/src/db/sql/feedgen_upsert.sql index 6f731fef..376ec57d 100644 --- a/consumer/src/db/sql/feedgen_upsert.sql +++ b/consumer/src/db/sql/feedgen_upsert.sql @@ -8,4 +8,5 @@ ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, description=EXCLUDED.description, description_facets=EXCLUDED.description_facets, avatar_cid=EXCLUDED.avatar_cid, - indexed_at=NOW() \ No newline at end of file + indexed_at=NOW() +RETURNING XMAX \ No newline at end of file diff --git a/consumer/src/db/sql/list_upsert.sql b/consumer/src/db/sql/list_upsert.sql index 12ea3523..47070520 100644 --- a/consumer/src/db/sql/list_upsert.sql +++ b/consumer/src/db/sql/list_upsert.sql @@ -6,4 +6,5 @@ ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, description=EXCLUDED.description, description_facets=EXCLUDED.description_facets, avatar_cid=EXCLUDED.avatar_cid, - indexed_at=NOW() \ No newline at end of file + indexed_at=NOW() +RETURNING XMAX \ No newline at end of file diff --git a/consumer/src/db/sql/starterpack_upsert.sql b/consumer/src/db/sql/starterpack_upsert.sql index 12b49b0b..d7a9a93b 100644 --- a/consumer/src/db/sql/starterpack_upsert.sql +++ b/consumer/src/db/sql/starterpack_upsert.sql @@ -7,4 +7,5 @@ ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, description_facets=EXCLUDED.description_facets, list=EXCLUDED.list, feeds=EXCLUDED.feeds, - indexed_at=NOW() \ No newline at end of file + indexed_at=NOW() +RETURNING XMAX \ No newline at end of file diff --git a/consumer/src/indexer/mod.rs b/consumer/src/indexer/mod.rs index 85f650b1..a062d962 100644 --- a/consumer/src/indexer/mod.rs +++ b/consumer/src/indexer/mod.rs @@ -527,15 +527,15 @@ pub async fn index_op( } RecordTypes::AppBskyFeedGenerator(record) => { let labels = record.labels.clone(); - let count = db::feedgen_upsert(conn, at_uri, repo, cid, record).await?; + let did_insert = db::feedgen_upsert(conn, at_uri, repo, cid, record).await?; if let Some(labels) = labels { db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; } - deltas - .add_delta(repo, AggregateType::ProfileFeed, count as i32) - .await; + if did_insert { + deltas.incr(repo, AggregateType::ProfileFeed).await; + } } RecordTypes::AppBskyFeedLike(record) => { let subject = record.subject.uri.clone(); @@ -635,15 +635,15 @@ pub async fn index_op( } RecordTypes::AppBskyGraphList(record) => { let labels = record.labels.clone(); - let count = db::list_upsert(conn, at_uri, repo, cid, record).await?; + let did_insert = db::list_upsert(conn, at_uri, repo, cid, record).await?; if let Some(labels) = labels { db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; } - deltas - .add_delta(repo, AggregateType::ProfileList, count as i32) - .await; + if did_insert { + deltas.incr(repo, AggregateType::ProfileList).await; + } } RecordTypes::AppBskyGraphListBlock(record) => { db::list_block_insert(conn, at_uri, repo, record).await?; @@ -659,10 +659,11 @@ pub async fn index_op( db::list_item_insert(conn, at_uri, record).await?; } RecordTypes::AppBskyGraphStarterPack(record) => { - let count = db::starter_pack_upsert(conn, at_uri, repo, cid, record).await?; - deltas - .add_delta(repo, AggregateType::ProfileStarterpack, count as i32) - .await; + let did_insert = db::starter_pack_upsert(conn, at_uri, repo, cid, record).await?; + + if did_insert { + deltas.incr(repo, AggregateType::ProfileStarterpack).await; + } } RecordTypes::AppBskyGraphVerification(record) => { db::verification_insert(conn, at_uri, repo, cid, record).await?; diff --git a/parakeet-index/Cargo.toml b/parakeet-index/Cargo.toml index 9653d503..7f39b99c 100644 --- a/parakeet-index/Cargo.toml +++ b/parakeet-index/Cargo.toml @@ -14,8 +14,8 @@ prost = "0.13.5" eyre = { version = "0.6.12", optional = true } figment = { version = "0.10.19", features = ["env", "toml"], optional = true } itertools = { version = "0.14.0", optional = true } +rocksdb = { version = "0.23", default-features = false, features = ["lz4", "bindgen-runtime"], optional = true } serde = { version = "1.0.217", features = ["derive"], optional = true } -sled = { version = "0.34.7", optional = true } tokio = { version = "1.42.0", features = ["full"], optional = true } tonic-health = { version = "0.13.0", optional = true } tracing = { version = "0.1.40", optional = true } @@ -25,4 +25,4 @@ tracing-subscriber = { version = "0.3.18", optional = true } tonic-build = "0.13.0" [features] -server = ["dep:eyre", "dep:figment", "dep:itertools", "dep:serde", "dep:sled", "dep:tokio", "dep:tonic-health", "dep:tracing", "dep:tracing-subscriber"] \ No newline at end of file +server = ["dep:eyre", "dep:figment", "dep:itertools", "dep:rocksdb", "dep:serde", "dep:tokio", "dep:tonic-health", "dep:tracing", "dep:tracing-subscriber"] \ No newline at end of file diff --git a/parakeet-index/Dockerfile b/parakeet-index/Dockerfile index b568f999..bb22b83d 100644 --- a/parakeet-index/Dockerfile +++ b/parakeet-index/Dockerfile @@ -1,6 +1,6 @@ FROM rust:1.85-slim-bookworm AS builder WORKDIR /work -RUN apt-get update && apt-get install -y --no-install-recommends wget libssl-dev protobuf-compiler pkg-config && rm -rf /var/lib/apt/lists/* +RUN apt-get update && apt-get install -y --no-install-recommends wget libssl-dev protobuf-compiler pkg-config clang && rm -rf /var/lib/apt/lists/* RUN wget -qO /bin/grpc_health_probe https://github.com/grpc-ecosystem/grpc-health-probe/releases/download/v0.4.38/grpc_health_probe-linux-amd64 && \ chmod +x /bin/grpc_health_probe COPY . . diff --git a/parakeet-index/src/server/db.rs b/parakeet-index/src/server/db.rs index 2fa86ced..7aa8001a 100644 --- a/parakeet-index/src/server/db.rs +++ b/parakeet-index/src/server/db.rs @@ -1,51 +1,35 @@ -use crate::all_none; -use crate::server::utils::{ToIntExt, TreeExt, slice_as_i32}; -use sled::{Db, MergeOperator, Tree}; +use crate::server::utils::{ToIntExt, slice_as_i32, to_key_name}; +use crate::{AggregateType, all_none}; +use itertools::izip; +use rocksdb::{DB, MergeOperands}; +use std::collections::HashMap; use std::path::PathBuf; pub struct DbStore { - pub agg_db: Db, - pub label_db: Db, - - pub follows: Tree, - pub followers: Tree, - pub likes: Tree, - pub replies: Tree, - pub reposts: Tree, - pub embeds: Tree, - pub profile_posts: Tree, - pub profile_lists: Tree, - pub profile_feeds: Tree, - pub profile_starterpacks: Tree, + pub agg_db: DB, + pub label_db: DB, } impl DbStore { pub fn new(db_root: PathBuf) -> eyre::Result { - let agg_db = sled::open(db_root.join("aggdb"))?; - let label_db = sled::open(db_root.join("labeldb"))?; - - Ok(DbStore { - follows: open_tree(&agg_db, "follows", merge_delta)?, - followers: open_tree(&agg_db, "followers", merge_delta)?, - likes: open_tree(&agg_db, "likes", merge_delta)?, - replies: open_tree(&agg_db, "replies", merge_delta)?, - reposts: open_tree(&agg_db, "reposts", merge_delta)?, - embeds: open_tree(&agg_db, "embeds", merge_delta)?, - profile_posts: open_tree(&agg_db, "profile_posts", merge_delta)?, - profile_lists: open_tree(&agg_db, "profile_lists", merge_delta)?, - profile_feeds: open_tree(&agg_db, "profile_feeds", merge_delta)?, - profile_starterpacks: open_tree(&agg_db, "profile_starterpacks", merge_delta)?, - - agg_db, - label_db, - }) + let mut opts = rocksdb::Options::default(); + opts.create_if_missing(true); + opts.set_compression_type(rocksdb::DBCompressionType::Lz4); + + let mut agg_opts = opts.clone(); + agg_opts.set_merge_operator_associative("pk_i32_merge", pk_i32_merge); + let agg_db = DB::open(&agg_opts, db_root.join("aggdb"))?; + + let label_db = DB::open(&opts, db_root.join("labeldb"))?; + + Ok(DbStore { agg_db, label_db }) } pub fn get_post_stats(&self, post: &str) -> Option { - let replies = self.replies.get_i32(post); - let likes = self.likes.get_i32(post); - let reposts = self.reposts.get_i32(post); - let quotes = self.embeds.get_i32(post); + let replies = self.get_aggregate(post, AggregateType::Reply); + let likes = self.get_aggregate(post, AggregateType::Like); + let reposts = self.get_aggregate(post, AggregateType::Repost); + let quotes = self.get_aggregate(post, AggregateType::Embed); if all_none![replies, likes, reposts, quotes] { return None; @@ -59,13 +43,37 @@ impl DbStore { }) } + pub fn get_post_stats_many(&self, posts: Vec) -> HashMap { + let replies = self.get_aggregate_many(&posts, AggregateType::Reply); + let likes = self.get_aggregate_many(&posts, AggregateType::Like); + let reposts = self.get_aggregate_many(&posts, AggregateType::Repost); + let quotes = self.get_aggregate_many(&posts, AggregateType::Embed); + + izip!(posts, replies, likes, reposts, quotes) + .filter_map(|(key, replies, likes, reposts, quotes)| { + if all_none![replies, likes, reposts, quotes] { + return None; + } + + let stats = crate::PostStats { + replies: replies.unwrap_or_default(), + likes: likes.unwrap_or_default(), + reposts: reposts.unwrap_or_default(), + quotes: quotes.unwrap_or_default(), + }; + + Some((key, stats)) + }) + .collect() + } + pub fn get_profile_stats(&self, did: &str) -> Option { - let followers = self.followers.get_i32(did); - let following = self.follows.get_i32(did); - let posts = self.profile_posts.get_i32(did); - let lists = self.profile_lists.get_i32(did); - let feeds = self.profile_feeds.get_i32(did); - let starterpacks = self.profile_starterpacks.get_i32(did); + let followers = self.get_aggregate(did, AggregateType::Follower); + let following = self.get_aggregate(did, AggregateType::Follow); + let posts = self.get_aggregate(did, AggregateType::ProfilePost); + let lists = self.get_aggregate(did, AggregateType::ProfileList); + let feeds = self.get_aggregate(did, AggregateType::ProfileFeed); + let starterpacks = self.get_aggregate(did, AggregateType::ProfileStarterpack); if all_none![followers, following, posts, lists, feeds, starterpacks] { return None; @@ -80,19 +88,74 @@ impl DbStore { starterpacks: starterpacks.unwrap_or_default(), }) } -} -fn open_tree(db: &Db, name: &str, merge: impl MergeOperator + 'static) -> eyre::Result { - let tree = db.open_tree(name)?; + pub fn get_profile_stats_many( + &self, + dids: Vec, + ) -> HashMap { + let followers = self.get_aggregate_many(&dids, AggregateType::Follower); + let following = self.get_aggregate_many(&dids, AggregateType::Follow); + let posts = self.get_aggregate_many(&dids, AggregateType::ProfilePost); + let lists = self.get_aggregate_many(&dids, AggregateType::ProfileList); + let feeds = self.get_aggregate_many(&dids, AggregateType::ProfileFeed); + let starterpacks = self.get_aggregate_many(&dids, AggregateType::ProfileStarterpack); + + izip!( + dids, + followers, + following, + posts, + lists, + feeds, + starterpacks + ) + .filter_map( + |(key, followers, following, posts, lists, feeds, starterpacks)| { + if all_none![followers, following, posts, lists, feeds, starterpacks] { + return None; + } + + let stats = crate::ProfileStats { + followers: followers.unwrap_or_default(), + following: following.unwrap_or_default(), + posts: posts.unwrap_or_default(), + lists: lists.unwrap_or_default(), + feeds: feeds.unwrap_or_default(), + starterpacks: starterpacks.unwrap_or_default(), + }; + + Some((key, stats)) + }, + ) + .collect() + } + + pub fn get_aggregate(&self, uri: &str, typ: AggregateType) -> Option { + let key = to_key_name(uri, typ); + + self.agg_db + .get_pinned(&key) + .ok() + .flatten() + .and_then(|data| slice_as_i32(&data)) + } - tree.set_merge_operator(merge); + pub fn get_aggregate_many(&self, uris: &[String], typ: AggregateType) -> Vec> { + let keys = uris.iter().map(|uri| to_key_name(uri, typ)); - Ok(tree) + self.agg_db + .multi_get(keys) + .into_iter() + .filter_map(|v| v.ok()) + .map(|v| v.and_then(|data| slice_as_i32(&data))) + .collect() + } } -fn merge_delta(_key: &[u8], old: Option<&[u8]>, new: &[u8]) -> Option> { +fn pk_i32_merge(_key: &[u8], old: Option<&[u8]>, op: &MergeOperands) -> Option> { let old = old.and_then(slice_as_i32); - let new = slice_as_i32(new)?; + + let new = op.iter().map(slice_as_i32).sum::>()?; let res = match old { Some(old) => old + new, diff --git a/parakeet-index/src/server/service.rs b/parakeet-index/src/server/service.rs index cd91963f..5a41cc18 100644 --- a/parakeet-index/src/server/service.rs +++ b/parakeet-index/src/server/service.rs @@ -1,7 +1,6 @@ use crate::index::*; use crate::server::GlobalState; -use crate::server::utils::TreeExt; -use std::collections::HashMap; +use crate::server::utils::to_key_name; use std::ops::Deref; use std::sync::Arc; use tonic::codegen::tokio_stream::StreamExt; @@ -14,27 +13,11 @@ impl Service { Service(state) } - fn apply_delta( - &self, - uri: &str, - typ: AggregateType, - delta: i32, - ) -> sled::Result> { + fn apply_delta(&self, uri: &str, typ: AggregateType, delta: i32) -> Result<(), rocksdb::Error> { let val = delta.to_le_bytes(); + let key = to_key_name(uri, typ); - match typ { - AggregateType::Unknown => todo!(), - AggregateType::Follow => self.dbs.follows.merge(uri, val), - AggregateType::Follower => self.dbs.followers.merge(uri, val), - AggregateType::Like => self.dbs.likes.merge(uri, val), - AggregateType::Reply => self.dbs.replies.merge(uri, val), - AggregateType::Repost => self.dbs.reposts.merge(uri, val), - AggregateType::Embed => self.dbs.embeds.merge(uri, val), - AggregateType::ProfilePost => self.dbs.profile_posts.merge(uri, val), - AggregateType::ProfileList => self.dbs.profile_lists.merge(uri, val), - AggregateType::ProfileFeed => self.dbs.profile_feeds.merge(uri, val), - AggregateType::ProfileStarterpack => self.dbs.profile_starterpacks.merge(uri, val), - } + self.dbs.agg_db.merge(key, val) } } @@ -69,14 +52,18 @@ impl index_server::Index for Service { request: Request, ) -> Result, Status> { let inner = request.into_inner(); + let mut batch = rocksdb::WriteBatch::default(); for data in inner.deltas { - let res = self.apply_delta(&data.uri, data.typ(), data.delta); + let val = data.delta.to_le_bytes(); + let key = to_key_name(&data.uri, data.typ()); - if let Err(e) = res { - tracing::error!("failed to update stats DB: {e}"); - return Err(Status::unknown("failed to update stats DB")); - } + batch.merge(key, val); + } + + if let Err(e) = self.dbs.agg_db.write(batch) { + tracing::error!("failed to update stats DB: {e}"); + return Err(Status::unknown("failed to update stats DB")); } Ok(Response::new(AggregateDeltaRes {})) @@ -121,16 +108,7 @@ impl index_server::Index for Service { ) -> Result, Status> { let inner = request.into_inner(); - // idk if this is the best way of doing this???? - let entries = inner - .uris - .into_iter() - .filter_map(|uri| { - let stats = self.dbs.get_profile_stats(&uri)?; - - Some((uri, stats)) - }) - .collect::>(); + let entries = self.dbs.get_profile_stats_many(inner.uris); Ok(Response::new(GetProfileStatsManyRes { entries })) } @@ -152,15 +130,7 @@ impl index_server::Index for Service { ) -> Result, Status> { let inner = request.into_inner(); - let entries = inner - .uris - .into_iter() - .filter_map(|uri| { - let stats = self.dbs.get_post_stats(&uri)?; - - Some((uri, stats)) - }) - .collect::>(); + let entries = self.dbs.get_post_stats_many(inner.uris); Ok(Response::new(GetPostStatsManyRes { entries })) } @@ -173,8 +143,7 @@ impl index_server::Index for Service { let likes = self .dbs - .likes - .get_i32(inner.uri) + .get_aggregate(&inner.uri, AggregateType::Like) .map(|likes| LikeCount { likes }); Ok(Response::new(GetLikeCountRes { likes })) @@ -186,14 +155,12 @@ impl index_server::Index for Service { ) -> Result, Status> { let inner = request.into_inner(); - let entries = inner - .uris + let entries = self + .dbs + .get_aggregate_many(&inner.uris, AggregateType::Like) .into_iter() - .filter_map(|uri| { - let likes = self.dbs.likes.get_i32(&uri)?; - - Some((uri, LikeCount { likes })) - }) + .zip(inner.uris) + .filter_map(|(val, key)| val.map(|likes| (key, LikeCount { likes }))) .collect(); Ok(Response::new(GetLikeCountManyRes { entries })) diff --git a/parakeet-index/src/server/utils.rs b/parakeet-index/src/server/utils.rs index 94a5c276..85c30578 100644 --- a/parakeet-index/src/server/utils.rs +++ b/parakeet-index/src/server/utils.rs @@ -1,25 +1,10 @@ -use sled::{IVec, Tree}; +use crate::AggregateType; pub trait ToIntExt { fn as_i32(&self) -> Option; fn from_i32(i: i32) -> Self; } -impl ToIntExt for IVec { - fn as_i32(&self) -> Option { - if self.len() == 4 { - let bytes = self[0..4].try_into().ok()?; - Some(i32::from_le_bytes(bytes)) - } else { - None - } - } - - fn from_i32(i: i32) -> Self { - IVec::from(&i.to_le_bytes()) - } -} - impl ToIntExt for Vec { fn as_i32(&self) -> Option { if self.len() == 4 { @@ -41,19 +26,14 @@ pub fn slice_as_i32(data: &[u8]) -> Option { Some(i32::from_le_bytes(bytes)) } -pub trait TreeExt { - fn get_i32(&self, key: impl AsRef<[u8]>) -> Option; -} - -impl TreeExt for Tree { - fn get_i32(&self, key: impl AsRef<[u8]>) -> Option { - self.get(key).ok().flatten().and_then(|v| v.as_i32()) - } -} - #[macro_export] macro_rules! all_none { ($var0:ident, $($var:ident),*) => { $var0.is_none() $(&& $var.is_none())* }; } + +#[inline(always)] +pub fn to_key_name(uri: &str, typ: AggregateType) -> String { + format!("{uri}#{}", typ.as_str_name()) +} -- 2.51.2