feat(media): observe sharded migration reads

Classify sharded payload keys in storage sweeps, preserve physical totals, and deduplicate logical legacy/sharded copies. Export read resolution, fallback, and duplicate-layout metrics.

Co-authored-by: npub128x7j3pwgm4vs8yra3c42fcgcwcvh94g3luwzkqa376du2q6l0esqcrwch <51cde9442e46eac81c83ec71552708c3b0cb96a88ff8e1581d8fb4de281afbf3@buzz.block.builderlab.xyz>
Signed-off-by: npub128x7j3pwgm4vs8yra3c42fcgcwcvh94g3luwzkqa376du2q6l0esqcrwch <51cde9442e46eac81c83ec71552708c3b0cb96a88ff8e1581d8fb4de281afbf3@buzz.block.builderlab.xyz>
This commit is contained in:
npub128x7j3pwgm4vs8yra3c42fcgcwcvh94g3luwzkqa376du2q6l0esqcrwch
2026-07-31 19:37:20 -04:00
parent a8219bfeb4
commit e5bd12bd46
5 changed files with 291 additions and 37 deletions
Generated
+1
View File
@@ -1043,6 +1043,7 @@ dependencies = [
"image",
"imagesize",
"infer",
"metrics",
"mp4",
"nostr",
"rust-s3",
+1
View File
@@ -32,6 +32,7 @@ tempfile = "3"
tokio-util = { version = "0.7", features = ["io"] }
futures-util = "0.3"
futures-core = "0.3"
metrics = { workspace = true }
[dev-dependencies]
tokio = { workspace = true, features = ["test-util"] }
+227 -29
View File
@@ -13,8 +13,8 @@
//!
//! | Class | Shape |
//! |---|---|
//! | thumb | `{sha256}.thumb.jpg` |
//! | blob | `{sha256}.{ext}` (ext: 1-8 mixed-case alphanumeric) |
//! | thumb | `{sha256}.thumb.jpg` or `m/{hh}/{hh}/{community-uuid}/{sha256}.thumb.jpg` |
//! | blob | `{sha256}.{ext}` or `m/{hh}/{hh}/{community-uuid}/{sha256}.{ext}` (ext: 1-8 mixed-case alphanumeric) |
//! | sidecar | `_meta/{community-uuid}/{sha256}.json` |
//! | auxiliary | `_uploads/{community-uuid}/{sha256}/{ulid}.json` |
//! | unknown | everything else |
@@ -32,10 +32,17 @@ use crate::error::MediaError;
/// `Auxiliary`, so visibility gauges stay loud instead of silently wrong.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum KeyClass {
/// `{sha256}.thumb.jpg` — attributed to the blob's sha.
Thumb { sha256: String },
/// `{sha256}.{ext}` — physical bytes, logical join key.
Blob { sha256: String, ext: String },
/// Legacy or sharded thumbnail, attributed to the blob's sha.
Thumb {
community: Option<Uuid>,
sha256: String,
},
/// Legacy or sharded blob; sharded keys carry direct community attribution.
Blob {
community: Option<Uuid>,
sha256: String,
ext: String,
},
/// `_meta/{community}/{sha256}.json` — the (community, sha) binding.
Sidecar { community: Uuid, sha256: String },
/// `_uploads/{community}/{sha256}/{event_id}.json` — fleet physical only.
@@ -52,11 +59,34 @@ pub enum KeyClass {
/// shape of the blob pattern's segment count), then blob, sidecar,
/// auxiliary, and finally unknown. See module docs for the exact shapes.
pub fn classify_key(key: &str) -> KeyClass {
if let Some((community, filename)) = parse_sharded_prefix(key) {
if let Some(parsed_sha) = parse_thumb_key(filename) {
return KeyClass::Thumb {
community: Some(community),
sha256: parsed_sha,
};
}
if let Some((parsed_sha, ext)) = parse_blob_key(filename) {
return KeyClass::Blob {
community: Some(community),
sha256: parsed_sha,
ext,
};
}
return KeyClass::Unknown;
}
if let Some(sha256) = parse_thumb_key(key) {
return KeyClass::Thumb { sha256 };
return KeyClass::Thumb {
community: None,
sha256,
};
}
if let Some((sha256, ext)) = parse_blob_key(key) {
return KeyClass::Blob { sha256, ext };
return KeyClass::Blob {
community: None,
sha256,
ext,
};
}
if let Some((community, sha256)) = parse_sidecar_key(key) {
return KeyClass::Sidecar { community, sha256 };
@@ -125,6 +155,27 @@ fn parse_canonical_uuid(s: &str) -> Option<Uuid> {
Uuid::parse_str(s).ok()
}
/// `m/{sha[0:2]}/{sha[2:4]}/{community}/{filename}`. The filename's digest
/// must agree with both shard segments; malformed migration keys stay unknown.
fn parse_sharded_prefix(key: &str) -> Option<(Uuid, &str)> {
let mut segments = key.split('/');
if segments.next()? != "m" {
return None;
}
let shard_1 = segments.next()?;
let shard_2 = segments.next()?;
let community = parse_canonical_uuid(segments.next()?)?;
let filename = segments.next()?;
if segments.next().is_some() || shard_1.len() != 2 || shard_2.len() != 2 {
return None;
}
let sha256 = filename.split('.').next()?;
if !is_sha256(sha256) || shard_1 != &sha256[..2] || shard_2 != &sha256[2..4] {
return None;
}
Some((community, filename))
}
/// `{sha256}.thumb.jpg`
fn parse_thumb_key(key: &str) -> Option<String> {
let mut parts = key.split('.');
@@ -221,6 +272,10 @@ pub struct BucketSnapshot {
pub multi_variant_shas: u64,
/// Total bytes of ALL blob variants belonging to anomalous shas.
pub multi_variant_bytes: u64,
/// Logical blob variants present in both legacy and sharded layouts.
pub duplicate_layout_variants: u64,
/// Physical bytes across both layouts for duplicate logical variants.
pub duplicate_layout_bytes: u64,
pub unknown_key_bytes: u64,
pub unknown_key_objects: u64,
}
@@ -228,14 +283,44 @@ pub struct BucketSnapshot {
/// Pure, incremental fold over classified bucket keys. Never retains a full
/// object listing — only per-sha/per-binding running totals, bounded by the
/// number of distinct shas and sidecar bindings actually present.
#[derive(Debug, Default)]
struct LayoutCopies {
legacy: Option<u64>,
sharded: HashMap<Uuid, u64>,
}
impl LayoutCopies {
fn insert(&mut self, community: Option<Uuid>, size: u64) {
match community {
Some(community) => {
self.sharded.insert(community, size);
}
None => {
self.legacy = Some(size);
}
}
}
fn physical_bytes(&self) -> u64 {
self.legacy.unwrap_or(0) + self.sharded.values().sum::<u64>()
}
fn logical_bytes(&self, community: Uuid) -> u64 {
self.sharded
.get(&community)
.copied()
.or(self.legacy)
.unwrap_or(0)
}
}
#[derive(Debug, Default)]
pub struct BucketAggregate {
/// sha -> bytes of every blob variant seen for that sha (D-EXT: multiple
/// entries is the multi-variant anomaly).
blob_variant_bytes: HashMap<String, Vec<u64>>,
/// sha -> thumb bytes. At most one thumb key per sha, so a plain insert
/// is correct (no accumulation needed).
thumb_bytes: HashMap<String, u64>,
/// (sha, ext) -> physical copies by layout. Layout copies are one logical
/// variant and must not double bill during migration.
blob_variants: HashMap<(String, String), LayoutCopies>,
/// sha -> physical thumbnail copies by layout.
thumb_copies: HashMap<String, LayoutCopies>,
/// (community, sha) -> sidecar object's own byte size (informational;
/// not part of logical bytes).
sidecar_bindings: HashMap<(Uuid, String), u64>,
@@ -251,14 +336,21 @@ impl BucketAggregate {
self.physical_objects += 1;
self.physical_bytes += size;
match classify_key(key) {
KeyClass::Thumb { sha256 } => {
self.thumb_bytes.insert(sha256, size);
}
KeyClass::Blob { sha256, .. } => {
self.blob_variant_bytes
KeyClass::Thumb { community, sha256 } => {
self.thumb_copies
.entry(sha256)
.or_default()
.push(size);
.insert(community, size);
}
KeyClass::Blob {
community,
sha256,
ext,
} => {
self.blob_variants
.entry((sha256, ext))
.or_default()
.insert(community, size);
}
KeyClass::Sidecar { community, sha256 } => {
self.sidecar_bindings.insert((community, sha256), size);
@@ -285,13 +377,28 @@ impl BucketAggregate {
let mut multi_variant_bytes = 0u64;
let mut orphan_blob_count = 0u64;
let mut orphan_blob_bytes = 0u64;
for (sha256, variants) in &self.blob_variant_bytes {
let variant_bytes: u64 = variants.iter().sum();
let mut duplicate_layout_variants = 0u64;
let mut duplicate_layout_bytes = 0u64;
for copies in self.blob_variants.values() {
if copies.legacy.is_some() && !copies.sharded.is_empty() {
duplicate_layout_variants += 1;
duplicate_layout_bytes += copies.physical_bytes();
}
}
let mut variants_by_sha: HashMap<&str, Vec<&LayoutCopies>> = HashMap::new();
for ((sha256, _), copies) in &self.blob_variants {
variants_by_sha
.entry(sha256.as_str())
.or_default()
.push(copies);
}
for (sha256, variants) in &variants_by_sha {
let variant_bytes: u64 = variants.iter().map(|copies| copies.physical_bytes()).sum();
if variants.len() > 1 {
multi_variant_shas += 1;
multi_variant_bytes += variant_bytes;
}
if !bound_shas.contains(sha256.as_str()) {
if !bound_shas.contains(*sha256) {
orphan_blob_count += 1;
orphan_blob_bytes += variant_bytes;
}
@@ -300,17 +407,25 @@ impl BucketAggregate {
let orphan_sidecar_count = self
.sidecar_bindings
.keys()
.filter(|(_, sha256)| !self.blob_variant_bytes.contains_key(sha256))
.filter(|(_, sha256)| !variants_by_sha.contains_key(sha256.as_str()))
.count() as u64;
let mut per_community: HashMap<Uuid, CommunityStorage> = HashMap::new();
for (community, sha256) in self.sidecar_bindings.keys() {
let blob_bytes: u64 = self
.blob_variant_bytes
let blob_bytes: u64 = variants_by_sha
.get(sha256.as_str())
.map(|variants| {
variants
.iter()
.map(|copies| copies.logical_bytes(*community))
.sum()
})
.unwrap_or(0);
let thumb_bytes = self
.thumb_copies
.get(sha256)
.map(|v| v.iter().sum())
.map(|copies| copies.logical_bytes(*community))
.unwrap_or(0);
let thumb_bytes = self.thumb_bytes.get(sha256).copied().unwrap_or(0);
let entry = per_community.entry(*community).or_default();
entry.bytes += blob_bytes + thumb_bytes;
entry.objects += 1;
@@ -329,6 +444,8 @@ impl BucketAggregate {
orphan_sidecar_count,
multi_variant_shas,
multi_variant_bytes,
duplicate_layout_variants,
duplicate_layout_bytes,
unknown_key_bytes: self.unknown_bytes,
unknown_key_objects: self.unknown_objects,
}
@@ -429,7 +546,10 @@ mod tests {
let s = sha(0xaa);
assert_eq!(
classify_key(&format!("{s}.thumb.jpg")),
KeyClass::Thumb { sha256: s }
KeyClass::Thumb {
community: None,
sha256: s,
}
);
}
@@ -439,6 +559,7 @@ mod tests {
assert_eq!(
classify_key(&format!("{s}.png")),
KeyClass::Blob {
community: None,
sha256: s,
ext: "png".to_string()
}
@@ -453,12 +574,51 @@ mod tests {
assert_eq!(
classify_key(&format!("{s}.Z")),
KeyClass::Blob {
community: None,
sha256: s,
ext: "Z".to_string()
}
);
}
#[test]
fn classifies_sharded_blob_and_thumb_keys_with_community() {
let s = sha(0xab);
let c = community(10);
assert_eq!(
classify_key(&format!("m/ab/ab/{c}/{s}.png")),
KeyClass::Blob {
community: Some(c),
sha256: s.clone(),
ext: "png".to_string(),
}
);
assert_eq!(
classify_key(&format!("m/ab/ab/{c}/{s}.thumb.jpg")),
KeyClass::Thumb {
community: Some(c),
sha256: s,
}
);
}
#[test]
fn malformed_sharded_keys_are_unknown() {
let s = sha(0xab);
let c = community(11);
for key in [
format!("m/ff/ab/{c}/{s}.png"),
format!("m/ab/ff/{c}/{s}.png"),
format!("m/a/ab/{c}/{s}.png"),
format!("m/ab/ab/not-a-uuid/{s}.png"),
format!("m/ab/ab/{c}/{s}.png/extra"),
format!("m/ab/ab/{c}/{}.png", s.to_uppercase()),
format!("m/ab/ab/{c}/{s}.tar.gz"),
] {
assert_eq!(classify_key(&key), KeyClass::Unknown, "key: {key}");
}
}
#[test]
fn classifies_sidecar_key() {
let s = sha(0xdd);
@@ -559,6 +719,44 @@ mod tests {
assert_eq!(snap.per_community[&c].objects, 1);
}
#[test]
fn dual_layout_copies_count_physically_but_dedupe_logical_usage() {
let s = sha(0xab);
let c = community(12);
let mut agg = BucketAggregate::default();
agg.fold(&format!("{s}.jpg"), 100);
agg.fold(&format!("m/ab/ab/{c}/{s}.jpg"), 100);
agg.fold(&format!("{s}.thumb.jpg"), 20);
agg.fold(&format!("m/ab/ab/{c}/{s}.thumb.jpg"), 20);
agg.fold(&format!("_meta/{c}/{s}.json"), 10);
let snap = agg.finish();
assert_eq!(snap.physical_objects, 5);
assert_eq!(snap.physical_bytes, 250);
assert_eq!(snap.logical_objects, 1);
assert_eq!(snap.logical_bytes, 120);
assert_eq!(snap.per_community[&c].bytes, 120);
assert_eq!(snap.multi_variant_shas, 0);
assert_eq!(snap.duplicate_layout_variants, 1);
assert_eq!(snap.duplicate_layout_bytes, 200);
assert_eq!(snap.unknown_key_objects, 0);
}
#[test]
fn sharded_copy_is_attributed_only_to_its_community() {
let s = sha(0xcd);
let sharded_community = community(13);
let other_community = community(14);
let mut agg = BucketAggregate::default();
agg.fold(&format!("m/cd/cd/{sharded_community}/{s}.jpg"), 200);
agg.fold(&format!("_meta/{sharded_community}/{s}.json"), 10);
agg.fold(&format!("_meta/{other_community}/{s}.json"), 10);
let snap = agg.finish();
assert_eq!(snap.per_community[&sharded_community].bytes, 200);
assert_eq!(snap.per_community[&other_community].bytes, 0);
}
#[test]
fn orphan_blob_has_no_sidecar_binding() {
let s = sha(0x44);
+58 -8
View File
@@ -187,14 +187,64 @@ impl MediaStorage {
ctx: &TenantContext,
payload_name: &str,
) -> Result<String, MediaError> {
let candidates =
crate::keys::read_candidates(ctx, payload_name).map_err(|_| MediaError::NotFound)?;
match self.head_with_metadata(&candidates.sharded).await? {
Some(_) => Ok(candidates.sharded),
None => match self.head_with_metadata(&candidates.legacy).await? {
Some(_) => Ok(candidates.legacy),
None => Err(MediaError::NotFound),
},
let candidates = match crate::keys::read_candidates(ctx, payload_name) {
Ok(candidates) => candidates,
Err(_) => {
metrics::counter!(
"buzz_media_s3_read_resolutions_total",
"result" => "missing"
)
.increment(1);
return Err(MediaError::NotFound);
}
};
match self.head_with_metadata(&candidates.sharded).await {
Ok(Some(_)) => {
metrics::counter!(
"buzz_media_s3_read_resolutions_total",
"result" => "sharded"
)
.increment(1);
Ok(candidates.sharded)
}
Ok(None) => {
metrics::counter!("buzz_media_s3_read_fallbacks_total").increment(1);
match self.head_with_metadata(&candidates.legacy).await {
Ok(Some(_)) => {
metrics::counter!(
"buzz_media_s3_read_resolutions_total",
"result" => "legacy"
)
.increment(1);
Ok(candidates.legacy)
}
Ok(None) => {
metrics::counter!(
"buzz_media_s3_read_resolutions_total",
"result" => "missing"
)
.increment(1);
Err(MediaError::NotFound)
}
Err(error) => {
metrics::counter!(
"buzz_media_s3_read_resolutions_total",
"result" => "storage_error"
)
.increment(1);
Err(error)
}
}
}
Err(error) => {
metrics::counter!(
"buzz_media_s3_read_resolutions_total",
"result" => "storage_error"
)
.increment(1);
Err(error)
}
}
}
+4
View File
@@ -320,6 +320,10 @@ pub async fn emit_storage_metrics(
metrics::gauge!("buzz_storage_orphan_sidecars").set(snapshot.orphan_sidecar_count as f64);
metrics::gauge!("buzz_storage_multi_variant_shas").set(snapshot.multi_variant_shas as f64);
metrics::gauge!("buzz_storage_multi_variant_bytes").set(snapshot.multi_variant_bytes as f64);
metrics::gauge!("buzz_storage_duplicate_layout_variants")
.set(snapshot.duplicate_layout_variants as f64);
metrics::gauge!("buzz_storage_duplicate_layout_bytes")
.set(snapshot.duplicate_layout_bytes as f64);
metrics::gauge!("buzz_storage_unknown_key_bytes").set(snapshot.unknown_key_bytes as f64);
metrics::gauge!("buzz_storage_unknown_key_objects").set(snapshot.unknown_key_objects as f64);