use base64::Engine as _;
use faucet_core::FaucetError;
use mysql_async::Value;
use mysql_async::binlog::row::BinlogRow;
use mysql_async::binlog::value::BinlogValue;
use serde_json::{Map, Value as Json};
pub fn value_to_json(v: &Value) -> Json {
match v {
Value::NULL => Json::Null,
Value::Int(i) => Json::from(*i),
Value::UInt(u) => Json::from(*u),
Value::Float(f) => serde_json::Number::from_f64(f64::from(*f))
.map(Json::Number)
.unwrap_or(Json::Null),
Value::Double(d) => serde_json::Number::from_f64(*d)
.map(Json::Number)
.unwrap_or(Json::Null),
Value::Bytes(b) => match std::str::from_utf8(b) {
Ok(s) => Json::String(s.to_owned()),
Err(_) => Json::String(base64::engine::general_purpose::STANDARD.encode(b)),
},
Value::Date(y, mo, d, h, mi, s, micro) => Json::String(format!(
"{y:04}-{mo:02}-{d:02} {h:02}:{mi:02}:{s:02}.{micro:06}"
)),
Value::Time(neg, days, h, mi, s, micro) => {
let sign = if *neg { "-" } else { "" };
let total_hours = u64::from(*days) * 24 + u64::from(*h);
Json::String(format!("{sign}{total_hours:02}:{mi:02}:{s:02}.{micro:06}"))
}
}
}
pub fn binlog_value_to_json(v: &BinlogValue<'_>) -> Json {
match v {
BinlogValue::Value(val) => value_to_json(val),
BinlogValue::Jsonb(jsonb_val) => {
match serde_json::Value::try_from(jsonb_val.clone().into_owned()) {
Ok(j) => j,
Err(e) => {
tracing::warn!(
error = %e,
"mysql-cdc: JSON column holds an opaque value with no JSON \
representation; emitting null for this column"
);
Json::Null
}
}
}
BinlogValue::JsonDiff(_) => Json::String("<JsonDiff>".into()),
}
}
pub fn binlog_row_to_json(row: &BinlogRow) -> Result<Json, FaucetError> {
let cols = row.columns_ref();
let mut obj = Map::with_capacity(cols.len());
for (i, col) in cols.iter().enumerate() {
let name = {
let n = col.name_str();
if n.is_empty() {
format!("col_{i}")
} else {
n.into_owned()
}
};
let val = match row.as_ref(i) {
Some(bv) => binlog_value_to_json(bv),
None => Json::Null,
};
obj.insert(name, val);
}
Ok(Json::Object(obj))
}
#[cfg(test)]
mod tests {
use super::*;
use mysql_async::Value;
#[test]
fn scalars() {
assert_eq!(value_to_json(&Value::NULL), Json::Null);
assert_eq!(value_to_json(&Value::Int(-5)), Json::from(-5));
assert_eq!(value_to_json(&Value::UInt(7)), Json::from(7u64));
assert_eq!(
value_to_json(&Value::Bytes(b"hello".to_vec())),
Json::String("hello".into())
);
}
#[test]
fn double() {
assert_eq!(value_to_json(&Value::Double(1.5)), Json::from(1.5));
}
#[test]
fn date_formats() {
let v = value_to_json(&Value::Date(2026, 6, 6, 12, 30, 0, 0));
assert_eq!(v, Json::String("2026-06-06 12:30:00.000000".into()));
}
#[test]
fn time_positive() {
let v = value_to_json(&Value::Time(false, 0, 1, 30, 0, 0));
assert_eq!(v, Json::String("01:30:00.000000".into()));
}
#[test]
fn time_negative() {
let v = value_to_json(&Value::Time(true, 0, 2, 0, 0, 500_000));
assert_eq!(v, Json::String("-02:00:00.500000".into()));
}
#[test]
fn non_utf8_bytes_base64() {
let v = value_to_json(&Value::Bytes(vec![0xff, 0xfe]));
assert!(matches!(v, Json::String(_)));
if let Json::String(s) = v {
base64::engine::general_purpose::STANDARD
.decode(&s)
.expect("should be valid base64");
}
}
#[test]
fn time_with_days_accumulates_hours() {
let v = value_to_json(&Value::Time(false, 1, 2, 3, 4, 0));
assert_eq!(v, Json::String("26:03:04.000000".into()));
}
#[test]
fn float_nan_is_null() {
let v = value_to_json(&Value::Float(f32::NAN));
assert_eq!(v, Json::Null);
}
#[test]
fn double_infinity_is_null() {
assert_eq!(value_to_json(&Value::Double(f64::INFINITY)), Json::Null);
assert_eq!(value_to_json(&Value::Double(f64::NEG_INFINITY)), Json::Null);
}
}