Skip to main content

faucet_source_redshift/
convert.rs

1//! Row decoding and parameter binding for the Redshift source.
2//!
3//! Redshift is PostgreSQL wire-compatible, so this mirrors the native Postgres
4//! source: values decode through `sqlx`'s Postgres row API, and JSON bind values
5//! are classified before binding so large integers keep full precision.
6
7use serde_json::Value;
8use sqlx::{Column, Row};
9
10/// Convert a raw Redshift/Postgres column value to a `serde_json::Value`.
11///
12/// Tries progressively broader decodings and falls back to `Value::Null` for
13/// unsupported or SQL-NULL columns.
14pub(crate) fn pg_value_to_json(row: &sqlx::postgres::PgRow, col_name: &str) -> Value {
15    if let Ok(v) = row.try_get::<Value, _>(col_name) {
16        return v;
17    }
18    if let Ok(v) = row.try_get::<String, _>(col_name) {
19        return Value::String(v);
20    }
21    if let Ok(v) = row.try_get::<i64, _>(col_name) {
22        return Value::Number(v.into());
23    }
24    if let Ok(v) = row.try_get::<i32, _>(col_name) {
25        return Value::Number(v.into());
26    }
27    if let Ok(v) = row.try_get::<i16, _>(col_name) {
28        return Value::Number(v.into());
29    }
30    if let Ok(v) = row.try_get::<f64, _>(col_name) {
31        return serde_json::Number::from_f64(v)
32            .map(Value::Number)
33            .unwrap_or(Value::Null);
34    }
35    if let Ok(v) = row.try_get::<f32, _>(col_name) {
36        return serde_json::Number::from_f64(v as f64)
37            .map(Value::Number)
38            .unwrap_or(Value::Null);
39    }
40    if let Ok(v) = row.try_get::<bool, _>(col_name) {
41        return Value::Bool(v);
42    }
43    // Timestamps → RFC3339 / ISO-8601 strings.
44    if let Ok(v) =
45        row.try_get::<sqlx::types::chrono::DateTime<sqlx::types::chrono::Utc>, _>(col_name)
46    {
47        return Value::String(v.to_rfc3339());
48    }
49    if let Ok(v) = row.try_get::<sqlx::types::chrono::NaiveDateTime, _>(col_name) {
50        return Value::String(v.to_string());
51    }
52    if let Ok(v) = row.try_get::<sqlx::types::chrono::NaiveDate, _>(col_name) {
53        return Value::String(v.to_string());
54    }
55    if let Ok(v) = row.try_get::<sqlx::types::chrono::NaiveTime, _>(col_name) {
56        return Value::String(v.to_string());
57    }
58    if let Ok(v) = row.try_get::<sqlx::types::Uuid, _>(col_name) {
59        return Value::String(v.to_string());
60    }
61    // NUMERIC / DECIMAL → string, preserving exact precision.
62    if let Ok(v) = row.try_get::<sqlx::types::BigDecimal, _>(col_name) {
63        return Value::String(v.to_string());
64    }
65    // Binary (VARBYTE / bytea) → base64 so it survives the JSON round-trip.
66    if let Ok(v) = row.try_get::<Vec<u8>, _>(col_name) {
67        use base64::Engine as _;
68        return Value::String(base64::engine::general_purpose::STANDARD.encode(v));
69    }
70    Value::Null
71}
72
73/// Convert a single row into a JSON object keyed by column name.
74pub(crate) fn row_to_json(row: &sqlx::postgres::PgRow) -> Value {
75    let mut map = serde_json::Map::new();
76    for col in row.columns() {
77        let name = col.name().to_string();
78        let value = pg_value_to_json(row, &name);
79        map.insert(name, value);
80    }
81    Value::Object(map)
82}
83
84/// How a numeric bind value should be bound onto a `sqlx` query. Classifying
85/// before binding keeps any integer in `[i64::MIN, i64::MAX]` exact (binding
86/// large integers as `f64` silently rounds them).
87#[derive(Debug, Clone, Copy, PartialEq, Eq)]
88pub(crate) enum NumberBind {
89    /// Exact `i64`.
90    I64,
91    /// Above `i64::MAX`; bind the `u64` reinterpreted as `i64` (two's complement).
92    U64,
93    /// Genuine floating-point value.
94    F64,
95}
96
97/// Classify a JSON number into the bind category to use.
98pub(crate) fn classify_number(n: &serde_json::Number) -> NumberBind {
99    if n.is_i64() {
100        NumberBind::I64
101    } else if n.is_u64() {
102        NumberBind::U64
103    } else {
104        NumberBind::F64
105    }
106}
107
108/// Bind a slice of JSON values onto a `sqlx` query as native scalar types, in
109/// positional order (`$1, $2, …`). Binding a raw `serde_json::Value` would
110/// encode as `jsonb` and break comparisons against typed columns, so scalars
111/// are bound as their native types.
112pub(crate) fn bind_params<'q>(
113    mut query: sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments>,
114    binds: &'q [Value],
115) -> sqlx::query::Query<'q, sqlx::Postgres, sqlx::postgres::PgArguments> {
116    for value in binds {
117        query = match value {
118            Value::String(s) => query.bind(s.clone()),
119            Value::Number(n) => match classify_number(n) {
120                NumberBind::I64 => query.bind(n.as_i64().unwrap()),
121                NumberBind::U64 => query.bind(n.as_u64().unwrap() as i64),
122                NumberBind::F64 => query.bind(n.as_f64().unwrap_or(0.0)),
123            },
124            Value::Bool(b) => query.bind(*b),
125            Value::Null => query.bind(None::<String>),
126            _ => query.bind(value.to_string()),
127        };
128    }
129    query
130}
131
132#[cfg(test)]
133mod tests {
134    use super::*;
135    use serde_json::json;
136
137    fn num(v: serde_json::Value) -> serde_json::Number {
138        match v {
139            serde_json::Value::Number(n) => n,
140            _ => panic!("not a number"),
141        }
142    }
143
144    #[test]
145    fn classify_small_int_is_i64() {
146        assert_eq!(classify_number(&num(json!(42))), NumberBind::I64);
147        assert_eq!(classify_number(&num(json!(-7))), NumberBind::I64);
148    }
149
150    #[test]
151    fn classify_above_2_pow_53_stays_i64() {
152        let v = 9_007_199_254_740_993i64; // 2^53 + 1
153        assert_eq!(classify_number(&num(json!(v))), NumberBind::I64);
154    }
155
156    #[test]
157    fn classify_above_i64_max_is_u64() {
158        let v: u64 = i64::MAX as u64 + 1;
159        assert_eq!(classify_number(&num(json!(v))), NumberBind::U64);
160    }
161
162    #[test]
163    fn classify_float_is_f64() {
164        assert_eq!(classify_number(&num(json!(3.5))), NumberBind::F64);
165    }
166}