use super::types::{Schema, ZSetBatch};
use gnitz_wire::control::{append_frame, ControlHeader};
use gnitz_wire::txn_frame::{encode_items, FrameItem};
use gnitz_wire::WireConflictMode;
use gnitz_wire::{ClientVerb, WireFlags};
use std::sync::Arc;
pub fn encode_frame(hdr: ControlHeader, blob: &[u8], schema: Option<&[u8]>, data: Option<&ZSetBatch>) -> Vec<u8> {
let mut out = Vec::new();
append_frame(
&mut out,
&hdr,
blob,
schema,
data.filter(|b| !b.is_empty()).map(|b| b.wire_regions()).as_deref(),
);
out
}
pub struct PushFamily {
pub target: crate::Target,
pub schema: Arc<Schema>,
pub batch: ZSetBatch,
pub mode: WireConflictMode,
pub basis: u64,
}
pub fn encode_push_txn(families: &[PushFamily]) -> Vec<u8> {
let schemas: Vec<Vec<u8>> = families.iter().map(|f| f.schema.to_block()).collect();
let items: Vec<FrameItem> = families
.iter()
.zip(&schemas)
.map(|(f, schema)| FrameItem {
hdr: ControlHeader {
flags: WireFlags {
conflict_mode: f.mode,
..Default::default()
},
..ControlHeader::naming(ClientVerb::PushTxn, f.target, f.basis)
},
schema: Some(schema),
data: Some(f.batch.wire_regions()),
})
.collect();
encode_items(ClientVerb::PushTxn, &items)
}
pub fn encode_ddl_txn(families: &[(u64, ZSetBatch)]) -> Vec<u8> {
let items: Vec<FrameItem> = families
.iter()
.map(|(tid, b)| FrameItem {
hdr: ControlHeader { target_id: *tid, ..Default::default() },
schema: None,
data: Some(b.wire_regions()),
})
.collect();
encode_items(ClientVerb::DdlTxn, &items)
}
#[cfg(test)]
#[path = "tests/message.rs"]
mod tests;