use std::sync::Arc;
use nodedb_types::sync::wire::SyncProvenance;
use tracing::warn;
use crate::event::types::{EventSource, RowId, WriteEvent, WriteOp};
use crate::types::{Lsn, TenantId, VShardId};
type LabelEventFields = (WriteOp, Option<Arc<[u8]>>, Option<Arc<[u8]>>);
pub(super) fn parse_put_record(
payload: &[u8],
tenant_id: TenantId,
vshard_id: VShardId,
lsn: Lsn,
sequence: &mut u64,
) -> Option<WriteEvent> {
if let Ok((disc, collection, key, value, _ttl_ms)) =
zerompk::from_msgpack::<(&str, String, Vec<u8>, Vec<u8>, u64)>(payload)
&& disc == "kv_put"
{
*sequence += 1;
let key_str = String::from_utf8_lossy(&key);
let (system_time_ms, valid_time_ms) =
crate::event::bitemporal_extract::extract_stamps(Some(&value));
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::Insert,
row_id: RowId::new(key_str.as_ref()),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: Some(Arc::from(value.as_slice())),
old_value: None,
system_time_ms,
valid_time_ms,
user_id: None,
statement_digest: None,
});
}
if let Ok((disc, collection, entries, _ttl_ms)) =
zerompk::from_msgpack::<(&str, String, Vec<(Vec<u8>, Vec<u8>)>, u64)>(payload)
&& disc == "kv_batch_put"
{
*sequence += 1;
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::BulkInsert {
count: entries.len() as u32,
},
row_id: RowId::new("_batch"),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: None,
old_value: None,
system_time_ms: None,
valid_time_ms: None,
user_id: None,
statement_digest: None,
});
}
if let Ok((collection, document_id, value, _prov, _surrogate)) =
zerompk::from_msgpack::<(String, String, Vec<u8>, Option<SyncProvenance>, u32)>(payload)
{
*sequence += 1;
let (system_time_ms, valid_time_ms) =
crate::event::bitemporal_extract::extract_stamps(Some(&value));
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::Insert,
row_id: RowId::new(document_id.as_str()),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: Some(Arc::from(value.as_slice())),
old_value: None,
system_time_ms,
valid_time_ms,
user_id: None,
statement_digest: None,
});
}
if let Ok((collection, document_id, value, _prov)) =
zerompk::from_msgpack::<(String, String, Vec<u8>, Option<SyncProvenance>)>(payload)
{
*sequence += 1;
let (system_time_ms, valid_time_ms) =
crate::event::bitemporal_extract::extract_stamps(Some(&value));
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::Insert,
row_id: RowId::new(document_id.as_str()),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: Some(Arc::from(value.as_slice())),
old_value: None,
system_time_ms,
valid_time_ms,
user_id: None,
statement_digest: None,
});
}
if let Ok((collection, document_id, value)) =
zerompk::from_msgpack::<(String, String, Vec<u8>)>(payload)
{
*sequence += 1;
let (system_time_ms, valid_time_ms) =
crate::event::bitemporal_extract::extract_stamps(Some(&value));
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::Insert,
row_id: RowId::new(document_id.as_str()),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: Some(Arc::from(value.as_slice())),
old_value: None,
system_time_ms,
valid_time_ms,
user_id: None,
statement_digest: None,
});
}
if let Ok((collection, src_id, label, dst_id, properties)) =
zerompk::from_msgpack::<(String, String, String, String, Vec<u8>)>(payload)
{
*sequence += 1;
let (system_time_ms, valid_time_ms) =
crate::event::bitemporal_extract::extract_stamps(Some(&properties));
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::Insert,
row_id: RowId::new(
crate::event::graph_cdc::edge_row_id(&src_id, &label, &dst_id).as_str(),
),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: Some(Arc::from(properties.as_slice())),
old_value: None,
system_time_ms,
valid_time_ms,
user_id: None,
statement_digest: None,
});
}
warn!(
lsn = lsn.as_u64(),
payload_len = payload.len(),
"WAL replay: unrecognized Put payload format, skipping"
);
None
}
pub(super) fn parse_graph_node_label_record(
payload: &[u8],
is_set: bool,
tenant_id: TenantId,
vshard_id: VShardId,
lsn: Lsn,
sequence: &mut u64,
) -> Option<WriteEvent> {
let (node_id, labels) = match zerompk::from_msgpack::<(String, Vec<String>)>(payload) {
Ok(decoded) => decoded,
Err(_) => {
warn!(
lsn = lsn.as_u64(),
payload_len = payload.len(),
"WAL replay: malformed graph node-label payload, skipping"
);
return None;
}
};
*sequence += 1;
let value = crate::event::graph_cdc::graph_label_delta_value(&labels);
let (op, new_value, old_value): LabelEventFields = if is_set {
(WriteOp::Insert, Some(Arc::from(value.as_slice())), None)
} else {
(WriteOp::Delete, None, Some(Arc::from(value.as_slice())))
};
Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(crate::event::graph_cdc::GRAPH_LABEL_STREAM),
op,
row_id: RowId::new(node_id.as_str()),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value,
old_value,
system_time_ms: None,
valid_time_ms: None,
user_id: None,
statement_digest: None,
})
}
pub(super) fn parse_delete_record(
payload: &[u8],
tenant_id: TenantId,
vshard_id: VShardId,
lsn: Lsn,
sequence: &mut u64,
) -> Option<WriteEvent> {
if let Ok((disc, collection, keys)) =
zerompk::from_msgpack::<(&str, String, Vec<Vec<u8>>)>(payload)
&& disc == "kv_delete"
{
*sequence += 1;
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::BulkDelete {
count: keys.len() as u32,
},
row_id: RowId::new("_batch"),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: None,
old_value: None,
system_time_ms: None,
valid_time_ms: None,
user_id: None,
statement_digest: None,
});
}
if let Ok((collection, document_id, _prov, _surrogate)) =
zerompk::from_msgpack::<(String, String, Option<SyncProvenance>, u32)>(payload)
{
*sequence += 1;
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::Delete,
row_id: RowId::new(document_id.as_str()),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: None,
old_value: None,
system_time_ms: None,
valid_time_ms: None,
user_id: None,
statement_digest: None,
});
}
if let Ok((collection, document_id, _prov)) =
zerompk::from_msgpack::<(String, String, Option<SyncProvenance>)>(payload)
{
*sequence += 1;
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::Delete,
row_id: RowId::new(document_id.as_str()),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: None,
old_value: None,
system_time_ms: None,
valid_time_ms: None,
user_id: None,
statement_digest: None,
});
}
if let Ok((collection, document_id)) = zerompk::from_msgpack::<(String, String)>(payload) {
*sequence += 1;
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::Delete,
row_id: RowId::new(document_id.as_str()),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: None,
old_value: None,
system_time_ms: None,
valid_time_ms: None,
user_id: None,
statement_digest: None,
});
}
if let Ok((collection, src_id, label, dst_id)) =
zerompk::from_msgpack::<(String, String, String, String)>(payload)
{
*sequence += 1;
return Some(WriteEvent {
sequence: *sequence,
collection: Arc::from(collection.as_str()),
op: WriteOp::Delete,
row_id: RowId::new(
crate::event::graph_cdc::edge_row_id(&src_id, &label, &dst_id).as_str(),
),
lsn,
tenant_id,
vshard_id,
source: EventSource::User,
new_value: None,
old_value: None,
system_time_ms: None,
valid_time_ms: None,
user_id: None,
statement_digest: None,
});
}
warn!(
lsn = lsn.as_u64(),
payload_len = payload.len(),
"WAL replay: unrecognized Delete payload format, skipping"
);
None
}