use serde_json::{Map, Value, json};
use crate::config::StartPosition;
use crate::lsn::Lsn;
pub const LSN_ALIAS: &str = "__faucet_lsn";
pub const SEQVAL_ALIAS: &str = "__faucet_seqval";
pub const OP_COLUMN: &str = "__$operation";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OpAction {
Emit(&'static str),
Skip,
}
pub fn op_action(operation: i64) -> Result<OpAction, faucet_core::FaucetError> {
match operation {
1 => Ok(OpAction::Emit("d")),
2 => Ok(OpAction::Emit("i")),
3 => Ok(OpAction::Skip),
4 => Ok(OpAction::Emit("u")),
other => Err(faucet_core::FaucetError::Source(format!(
"mssql-cdc: unrecognized __$operation code {other} (expected 1, 2, 3, or 4)"
))),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PollPlan {
NoChanges { set_bookmark: Option<Lsn> },
Query { from: Lsn, to: Lsn, gap: bool },
}
pub fn plan_poll(
bookmark: Option<Lsn>,
min_lsn: Option<Lsn>,
max_lsn: Option<Lsn>,
start: StartPosition,
) -> PollPlan {
let Some(to) = max_lsn else {
return PollPlan::NoChanges { set_bookmark: None };
};
let Some(min) = min_lsn else {
return PollPlan::NoChanges { set_bookmark: None };
};
let candidate = match bookmark {
Some(bm) => match bm.increment() {
Some(next) => next,
None => return PollPlan::NoChanges { set_bookmark: None },
},
None => match start {
StartPosition::Current => {
return PollPlan::NoChanges {
set_bookmark: Some(to),
};
}
StartPosition::Earliest => min,
},
};
if candidate > to {
return PollPlan::NoChanges { set_bookmark: None };
}
if candidate < min {
return PollPlan::Query {
from: min,
to,
gap: true,
};
}
PollPlan::Query {
from: candidate,
to,
gap: false,
}
}
pub fn business_columns(decoded: &Value) -> Map<String, Value> {
let mut out = Map::new();
if let Value::Object(map) = decoded {
for (k, v) in map {
if k == LSN_ALIAS || k == SEQVAL_ALIAS || k.starts_with("__$") {
continue;
}
out.insert(k.clone(), v.clone());
}
}
out
}
pub fn build_change_envelope(
op: &str,
schema: &str,
table: &str,
lsn_hex: &str,
seqval_hex: Option<&str>,
columns: Map<String, Value>,
) -> Value {
let is_delete = op == "d";
let (before, after) = if is_delete {
(Value::Object(columns), Value::Null)
} else {
(Value::Null, Value::Object(columns))
};
let mut obj = Map::new();
obj.insert("op".into(), json!(op));
obj.insert("schema".into(), json!(schema));
obj.insert("table".into(), json!(table));
obj.insert("before".into(), before);
obj.insert("after".into(), after);
obj.insert("lsn".into(), json!(lsn_hex));
if let Some(seq) = seqval_hex {
obj.insert("seqval".into(), json!(seq));
}
Value::Object(obj)
}
#[cfg(test)]
mod tests {
use super::*;
fn lsn(hex: &str) -> Lsn {
Lsn::from_hex(hex).unwrap()
}
#[test]
fn op_action_maps_all_codes() {
assert_eq!(op_action(1).unwrap(), OpAction::Emit("d"));
assert_eq!(op_action(2).unwrap(), OpAction::Emit("i"));
assert_eq!(op_action(3).unwrap(), OpAction::Skip);
assert_eq!(op_action(4).unwrap(), OpAction::Emit("u"));
}
#[test]
fn op_action_rejects_unknown_code() {
let err = op_action(7).unwrap_err();
assert!(err.to_string().contains("__$operation code 7"), "{err}");
assert!(op_action(0).is_err());
}
#[test]
fn plan_no_max_lsn_is_no_changes() {
let p = plan_poll(
None,
Some(lsn("00000000000000000001")),
None,
StartPosition::Earliest,
);
assert_eq!(p, PollPlan::NoChanges { set_bookmark: None });
}
#[test]
fn plan_fresh_current_anchors_at_max() {
let max = lsn("00000000000000000100");
let p = plan_poll(
None,
Some(lsn("00000000000000000001")),
Some(max),
StartPosition::Current,
);
assert_eq!(
p,
PollPlan::NoChanges {
set_bookmark: Some(max)
}
);
}
#[test]
fn plan_fresh_earliest_queries_from_min() {
let min = lsn("00000000000000000005");
let max = lsn("00000000000000000100");
let p = plan_poll(None, Some(min), Some(max), StartPosition::Earliest);
assert_eq!(
p,
PollPlan::Query {
from: min,
to: max,
gap: false
}
);
}
#[test]
fn plan_resume_queries_from_increment() {
let bm = lsn("00000000000000000010");
let min = lsn("00000000000000000001");
let max = lsn("00000000000000000100");
let p = plan_poll(Some(bm), Some(min), Some(max), StartPosition::Current);
assert_eq!(
p,
PollPlan::Query {
from: bm.increment().unwrap(),
to: max,
gap: false
}
);
}
#[test]
fn plan_resume_no_new_changes() {
let bm = lsn("00000000000000000100");
let min = lsn("00000000000000000001");
let max = lsn("00000000000000000100");
let p = plan_poll(Some(bm), Some(min), Some(max), StartPosition::Current);
assert_eq!(p, PollPlan::NoChanges { set_bookmark: None });
}
#[test]
fn plan_resume_before_min_flags_gap_and_clamps() {
let bm = lsn("00000000000000000001");
let min = lsn("00000000000000000050");
let max = lsn("00000000000000000100");
let p = plan_poll(Some(bm), Some(min), Some(max), StartPosition::Current);
assert_eq!(
p,
PollPlan::Query {
from: min,
to: max,
gap: true
}
);
}
#[test]
fn business_columns_strips_metadata() {
let decoded = json!({
"__faucet_lsn": "00000000000000000001",
"__faucet_seqval": "00000000000000000001",
"__$start_lsn": "AAAA",
"__$operation": 2,
"__$seqval": "BBBB",
"__$update_mask": "CC",
"id": 1,
"name": "alice"
});
let cols = business_columns(&decoded);
assert_eq!(cols.len(), 2);
assert_eq!(cols["id"], json!(1));
assert_eq!(cols["name"], json!("alice"));
assert!(!cols.contains_key("__$operation"));
assert!(!cols.contains_key("__faucet_lsn"));
assert!(!cols.contains_key("__faucet_seqval"));
}
#[test]
fn insert_envelope_puts_columns_in_after() {
let mut cols = Map::new();
cols.insert("id".into(), json!(1));
let env = build_change_envelope(
"i",
"dbo",
"Orders",
"00000000000000000003",
Some("0000000000000000000300"),
cols,
);
assert_eq!(env["op"], "i");
assert_eq!(env["schema"], "dbo");
assert_eq!(env["table"], "Orders");
assert_eq!(env["before"], Value::Null);
assert_eq!(env["after"]["id"], json!(1));
assert_eq!(env["lsn"], "00000000000000000003");
assert_eq!(env["seqval"], "0000000000000000000300");
}
#[test]
fn delete_envelope_puts_columns_in_before() {
let mut cols = Map::new();
cols.insert("id".into(), json!(42));
let env = build_change_envelope("d", "dbo", "Orders", "00000000000000000004", None, cols);
assert_eq!(env["op"], "d");
assert_eq!(env["before"]["id"], json!(42));
assert_eq!(env["after"], Value::Null);
assert!(env.as_object().unwrap().get("seqval").is_none());
}
#[test]
fn update_envelope_is_upsert_after() {
let mut cols = Map::new();
cols.insert("id".into(), json!(1));
cols.insert("name".into(), json!("bob"));
let env = build_change_envelope("u", "dbo", "Orders", "00000000000000000005", None, cols);
assert_eq!(env["op"], "u");
assert_eq!(env["before"], Value::Null);
assert_eq!(env["after"]["name"], "bob");
}
}