| 235 | } |
| 236 | |
| 237 | fn extract_key_value(event: &CdcEvent, key_field: &str) -> String { |
| 238 | let value = event.new_value.as_ref().or(event.old_value.as_ref()); |
| 239 | |
| 240 | if let Some(obj) = value.and_then(|v| v.as_object()) |
| 241 | && let Some(val) = obj.get(key_field) |
| 242 | { |
| 243 | return match val { |
| 244 | serde_json::Value::String(s) => s.clone(), |
| 245 | other => other.to_string(), |
| 246 | }; |
| 247 | } |
| 248 | |
| 249 | tracing::warn!( |
| 250 | collection = %event.collection, |
| 251 | row_id = %event.row_id, |
| 252 | key_field, |
| 253 | "compaction key field not found in event, falling back to row_id" |
| 254 | ); |
| 255 | event.row_id.clone() |
| 256 | } |
| 257 | |
| 258 | #[cfg(test)] |
| 259 | mod tests { |