use alloc::string::String;
pub use maxwell_cdc::{
ColumnDefinition, ControlMessage, DatabaseChange, DatabaseDefinition, DatabaseDropChange,
DdlMetadata, Message, OpType, RowChange, TableAlterChange, TableCreateChange, TableDefinition,
TableDropChange, parse,
};
pub use crate::wire::ConversionError;
use crate::ChangesetFormat;
use crate::builders::{ChangeDelete, Insert, PatchDelete, Update};
use crate::schema::NamedColumns;
use crate::wire::{
Sealed, WireColumnItem, WireSource, WireType, build_changeset_delete, build_insert,
build_patch_delete, build_patchset_update, resolve_table,
};
use crate::builders::{DiffOps, DiffSetBuilder, PatchsetFormat};
use crate::wire::WireAdapter;
use core::fmt::Debug;
use core::hash::Hash;
#[derive(Debug, Clone, Copy, Default)]
pub struct Maxwell;
impl Sealed for Maxwell {}
impl WireSource for Maxwell {
type Payload<'a> = MaxwellColumn<'a>;
fn wire_type(payload: &Self::Payload<'_>) -> WireType {
payload.wire_type
}
fn column_name<'a>(payload: &'a Self::Payload<'_>) -> &'a str {
payload.column_name
}
}
#[derive(Debug, Clone, Copy)]
pub struct MaxwellColumn<'a> {
pub column_name: &'a str,
pub wire_type: WireType,
pub value: &'a serde_json::Value,
}
struct MaxwellItem<'a> {
name: &'a str,
value: &'a serde_json::Value,
}
impl WireColumnItem<Maxwell> for MaxwellItem<'_> {
fn name(&self) -> &str {
self.name
}
fn payload(&self, wire_type: WireType) -> MaxwellColumn<'_> {
MaxwellColumn {
column_name: self.name,
wire_type,
value: self.value,
}
}
}
impl MaxwellColumn<'_> {
pub fn decoded_by<D, S, B>(
self,
decoder: &D,
) -> Result<crate::encoding::Value<S, B>, crate::wire::DecodeError>
where
D: crate::wire::Decoder<Maxwell, S, B>,
{
decoder.decode(self)
}
}
use crate::wire::{Digestable, WireColumnTypes, WireSchema};
impl<T, S, B> Digestable<ChangesetFormat, T, S, B> for Message
where
T: NamedColumns + WireColumnTypes,
S: Clone + Debug + Hash + Eq + AsRef<str> + Default,
B: Clone + Debug + Hash + Eq + AsRef<[u8]> + Default,
{
type Src = Maxwell;
type Error = ConversionError;
fn digest_into<Sch, A>(
&self,
builder: DiffSetBuilder<ChangesetFormat, T, S, B>,
schema: &Sch,
adapter: &A,
) -> Result<DiffSetBuilder<ChangesetFormat, T, S, B>, ConversionError>
where
Sch: WireSchema<Table = T>,
A: WireAdapter<Maxwell, S, B>,
{
match self {
Message::Insert(row) | Message::BootstrapInsert(row) => {
let table = resolve_table(schema, row.table.as_str())?;
let insert = build_insert_from_maxwell(&row.data, table, adapter)?;
Ok(DiffOps::insert(builder, insert))
}
Message::Update(row) => {
let table = resolve_table(schema, row.table.as_str())?;
let update = build_changeset_update_from_maxwell(
&row.data,
row.old.as_ref(),
table,
adapter,
)?;
Ok(DiffOps::update(builder, update))
}
Message::Delete(row) => {
let table = resolve_table(schema, row.table.as_str())?;
let delete = build_changeset_delete_from_maxwell(&row.data, table, adapter)?;
Ok(DiffOps::delete(builder, delete))
}
_ => Ok(builder),
}
}
}
impl<T, S, B> Digestable<PatchsetFormat, T, S, B> for Message
where
T: NamedColumns + WireColumnTypes,
S: Clone + Debug + Hash + Eq + AsRef<str> + Default,
B: Clone + Debug + Hash + Eq + AsRef<[u8]> + Default,
{
type Src = Maxwell;
type Error = ConversionError;
fn digest_into<Sch, A>(
&self,
builder: DiffSetBuilder<PatchsetFormat, T, S, B>,
schema: &Sch,
adapter: &A,
) -> Result<DiffSetBuilder<PatchsetFormat, T, S, B>, ConversionError>
where
Sch: WireSchema<Table = T>,
A: WireAdapter<Maxwell, S, B>,
{
match self {
Message::Insert(row) | Message::BootstrapInsert(row) => {
let table = resolve_table(schema, row.table.as_str())?;
let insert = build_insert_from_maxwell(&row.data, table, adapter)?;
Ok(DiffOps::insert(builder, insert))
}
Message::Update(row) => {
let table = resolve_table(schema, row.table.as_str())?;
let update = build_patchset_update_from_maxwell(&row.data, table, adapter)?;
Ok(DiffOps::update(builder, update))
}
Message::Delete(row) => {
let table = resolve_table(schema, row.table.as_str())?;
let delete = build_patch_delete_from_maxwell(&row.data, table, adapter)?;
Ok(DiffOps::delete(builder, delete))
}
_ => Ok(builder),
}
}
}
fn build_insert_from_maxwell<T, S, B, A>(
data: &serde_json::Map<String, serde_json::Value>,
table: &T,
adapter: &A,
) -> Result<Insert<T, S, B>, ConversionError>
where
T: NamedColumns + WireColumnTypes,
S: Clone + AsRef<str>,
B: Clone + AsRef<[u8]>,
A: WireAdapter<Maxwell, S, B>,
{
build_insert(
data.iter().map(|(name, value)| MaxwellItem {
name: name.as_str(),
value,
}),
table,
adapter,
)
}
fn build_changeset_update_from_maxwell<T, S, B, A>(
data: &serde_json::Map<String, serde_json::Value>,
old: Option<&serde_json::Map<String, serde_json::Value>>,
table: &T,
adapter: &A,
) -> Result<Update<T, ChangesetFormat, S, B>, ConversionError>
where
T: NamedColumns + WireColumnTypes,
S: Clone + Debug + AsRef<str>,
B: Clone + Debug + AsRef<[u8]>,
A: WireAdapter<Maxwell, S, B>,
{
let mut update: Update<T, ChangesetFormat, S, B> = Update::from(table.clone());
for (name, new_value) in data {
let col_idx = table
.column_index(name)
.ok_or_else(|| ConversionError::ColumnNotFound(name.clone()))?;
let wire_type = table.column_type(col_idx);
let new_payload = MaxwellColumn {
column_name: name.as_str(),
wire_type,
value: new_value,
};
let new = adapter.decode(new_payload)?;
let old = if let Some(old_value) = old.and_then(|old_map| old_map.get(name)) {
let old_payload = MaxwellColumn {
column_name: name.as_str(),
wire_type,
value: old_value,
};
adapter.decode(old_payload)?
} else {
new.clone()
};
update = update
.set(col_idx, old, new)
.map_err(|_| ConversionError::ColumnNotFound(name.clone()))?;
}
Ok(update)
}
fn build_patchset_update_from_maxwell<T, S, B, A>(
data: &serde_json::Map<String, serde_json::Value>,
table: &T,
adapter: &A,
) -> Result<Update<T, PatchsetFormat, S, B>, ConversionError>
where
T: NamedColumns + WireColumnTypes,
S: Clone + AsRef<str>,
B: Clone + AsRef<[u8]>,
A: WireAdapter<Maxwell, S, B>,
{
build_patchset_update(
data.iter().map(|(name, value)| MaxwellItem {
name: name.as_str(),
value,
}),
table,
adapter,
)
}
fn build_changeset_delete_from_maxwell<T, S, B, A>(
data: &serde_json::Map<String, serde_json::Value>,
table: &T,
adapter: &A,
) -> Result<ChangeDelete<T, S, B>, ConversionError>
where
T: NamedColumns + WireColumnTypes,
S: Clone + Default + AsRef<str>,
B: Clone + Default + AsRef<[u8]>,
A: WireAdapter<Maxwell, S, B>,
{
build_changeset_delete(
data.iter().map(|(name, value)| MaxwellItem {
name: name.as_str(),
value,
}),
table,
adapter,
)
}
fn build_patch_delete_from_maxwell<T, S, B, A>(
data: &serde_json::Map<String, serde_json::Value>,
table: &T,
adapter: &A,
) -> Result<PatchDelete<T, S, B>, ConversionError>
where
T: NamedColumns + WireColumnTypes,
S: Clone + AsRef<str>,
B: Clone + AsRef<[u8]>,
A: WireAdapter<Maxwell, S, B>,
{
build_patch_delete(
data.iter().map(|(name, value)| MaxwellItem {
name: name.as_str(),
value,
}),
table,
adapter,
)
}