Skip to main content

connectorx/sources/postgres/
mod.rs

1//! Source implementation for Postgres database, including the TLS support (client only).
2
3mod connection;
4mod errors;
5mod typesystem;
6
7pub use self::errors::PostgresSourceError;
8pub use cidr_02::IpInet;
9pub use connection::rewrite_tls_args;
10pub use pgvector::{Bit, HalfVector, SparseVector, Vector};
11pub use typesystem::{PostgresTypePairs, PostgresTypeSystem};
12
13use crate::constants::DB_BUFFER_SIZE;
14use crate::{
15    data_order::DataOrder,
16    errors::ConnectorXError,
17    sources::{PartitionParser, Produce, Source, SourcePartition},
18    sql::{count_query, CXQuery},
19};
20use anyhow::anyhow;
21use chrono::{DateTime, FixedOffset, NaiveDate, NaiveDateTime, NaiveTime, Utc};
22use csv::{ReaderBuilder, StringRecord, StringRecordsIntoIter};
23use fehler::{throw, throws};
24use hex::decode;
25use postgres::{
26    binary_copy::{BinaryCopyOutIter, BinaryCopyOutRow},
27    fallible_iterator::FallibleIterator,
28    tls::{MakeTlsConnect, TlsConnect},
29    Config, CopyOutReader, Row, RowIter, SimpleQueryMessage, Socket,
30};
31use r2d2::{Pool, PooledConnection};
32use r2d2_postgres::PostgresConnectionManager;
33use rust_decimal::Decimal;
34use serde_json::{from_str, Value};
35use sqlparser::dialect::PostgreSqlDialect;
36use std::collections::HashMap;
37use std::convert::TryFrom;
38use std::marker::PhantomData;
39use uuid::Uuid;
40
41/// Protocol - Binary based bulk load
42pub enum BinaryProtocol {}
43
44/// Protocol - CSV based bulk load
45pub enum CSVProtocol {}
46
47/// Protocol - use Cursor
48pub enum CursorProtocol {}
49
50/// Protocol - use Simple Query
51pub enum SimpleProtocol {}
52
53type PgManager<C> = PostgresConnectionManager<C>;
54type PgConn<C> = PooledConnection<PgManager<C>>;
55
56macro_rules! impl_produce_unimplemented {
57    ($(($protocol: ty, $t: ty, $msg: expr),)+) => {
58        $(
59            impl<'r> Produce<'r, $t> for $protocol {
60                type Error = PostgresSourceError;
61
62                #[throws(PostgresSourceError)]
63                fn produce(&'r mut self) -> $t {
64                   unimplemented!($msg);
65                }
66            }
67
68            impl<'r> Produce<'r, Option<$t>> for $protocol {
69                type Error = PostgresSourceError;
70
71                #[throws(PostgresSourceError)]
72                fn produce(&'r mut self) -> Option<$t> {
73                   unimplemented!($msg);
74                }
75            }
76        )+
77    };
78}
79
80impl_produce_unimplemented!(
81    (PostgresCSVSourceParser<'_>, HashMap<String, Option<String>>, "Please use `cursor` protocol for hstore type"),
82    (PostgresCSVSourceParser<'_>, Vector, "Please use `binary` protocol for vector type"),
83    (PostgresCSVSourceParser<'_>, HalfVector, "Please use `binary` protocol for halfvector type"),
84    (PostgresCSVSourceParser<'_>, Bit, "Please use `binary` protocol for bit type"),
85    (PostgresCSVSourceParser<'_>, SparseVector, "Please use `binary` protocol for sparsevector type"),
86
87
88    (PostgresSimpleSourceParser,HashMap<String, Option<String>>, "unimplemented"),
89    (PostgresSimpleSourceParser,Value, "unimplemented"),
90    (PostgresSimpleSourceParser, Vector, "Please use `binary` protocol for vector type"),
91    (PostgresSimpleSourceParser, HalfVector, "Please use `binary` protocol for halfvector type"),
92    (PostgresSimpleSourceParser, Bit, "Please use `binary` protocol for bit type"),
93    (PostgresSimpleSourceParser, SparseVector, "Please use `binary` protocol for sparsevector type"),
94
95);
96
97/// Wrap queries containing range or tsvector columns to cast those columns
98/// to `::text`. For example:
99///   SELECT * FROM orders WHERE id < 2
100/// becomes:
101///   SELECT "id", "period"::text, "name" FROM (SELECT * FROM orders WHERE id < 2) AS _cx_sub
102///
103/// This makes binary and cursor protocols work without type-specific FromSql
104/// implementations, and preserves the text representation for CSV and simple.
105fn maybe_rewrite_text_query(
106    query: &CXQuery<String>,
107    names: &[String],
108    schema: &[PostgresTypeSystem],
109) -> CXQuery<String> {
110    let text_mask: Vec<bool> = schema
111        .iter()
112        .map(|ts| {
113            matches!(
114                ts,
115                PostgresTypeSystem::Range(_) | PostgresTypeSystem::TsVector(_)
116            )
117        })
118        .collect();
119
120    if !text_mask.iter().any(|&cast| cast) {
121        return query.clone();
122    }
123
124    let cols: String = names
125        .iter()
126        .zip(text_mask.iter())
127        .map(|(name, cast_to_text)| {
128            let quoted = quote_ident(name);
129            if *cast_to_text {
130                format!("{}::text", quoted)
131            } else {
132                quoted
133            }
134        })
135        .collect::<Vec<_>>()
136        .join(", ");
137
138    let rewritten = format!("SELECT {} FROM ({}) AS _cx_sub", cols, query.as_str());
139    CXQuery::Wrapped(rewritten)
140}
141
142fn quote_ident(ident: &str) -> String {
143    format!("\"{}\"", ident.replace('\"', "\"\""))
144}
145
146#[cfg(test)]
147mod text_rewrite_tests {
148    use super::{maybe_rewrite_text_query, quote_ident, PostgresTypeSystem};
149    use crate::sql::CXQuery;
150
151    #[test]
152    fn quote_ident_escapes_embedded_quotes() {
153        assert_eq!(quote_ident("a\"b"), "\"a\"\"b\"");
154    }
155
156    #[test]
157    fn rewrite_escapes_column_names_and_casts_only_ranges() {
158        let q = CXQuery::Naked("SELECT 1".to_string());
159        let names = vec!["plain".to_string(), "a\"b".to_string()];
160        let schema = vec![
161            PostgresTypeSystem::Int4(true),
162            PostgresTypeSystem::Range(true),
163        ];
164        let rewritten = maybe_rewrite_text_query(&q, &names, &schema);
165        assert_eq!(
166            rewritten.as_str(),
167            "SELECT \"plain\", \"a\"\"b\"::text FROM (SELECT 1) AS _cx_sub"
168        );
169    }
170
171    #[test]
172    fn rewrite_casts_tsvector_and_ranges() {
173        let q = CXQuery::Wrapped("SELECT * FROM documents WHERE id < 2".to_string());
174        let names = vec!["id".to_string(), "a\"b".to_string(), "period".to_string()];
175        for nullable in [false, true] {
176            let schema = vec![
177                PostgresTypeSystem::Int4(false),
178                PostgresTypeSystem::TsVector(nullable),
179                PostgresTypeSystem::Range(nullable),
180            ];
181            assert_eq!(
182                maybe_rewrite_text_query(&q, &names, &schema).as_str(),
183                "SELECT \"id\", \"a\"\"b\"::text, \"period\"::text FROM (SELECT * FROM documents WHERE id < 2) AS _cx_sub"
184            );
185        }
186    }
187
188    #[test]
189    fn rewrite_leaves_other_types_unchanged() {
190        for q in [
191            CXQuery::Naked("SELECT 'hello' AS text".to_string()),
192            CXQuery::Wrapped("SELECT 'hello' AS text".to_string()),
193        ] {
194            let rewritten = maybe_rewrite_text_query(
195                &q,
196                &["text".to_string()],
197                &[PostgresTypeSystem::Text(true)],
198            );
199            assert_eq!(rewritten.as_str(), q.as_str());
200            assert_eq!(
201                std::mem::discriminant(&rewritten),
202                std::mem::discriminant(&q)
203            );
204        }
205    }
206}
207
208// take a row and unwrap the interior field from column 0
209fn convert_row<'b, R: TryFrom<usize> + postgres::types::FromSql<'b> + Clone>(row: &'b Row) -> R {
210    let nrows: Option<R> = row.get(0);
211    nrows.expect("Could not parse int result from count_query")
212}
213
214#[throws(PostgresSourceError)]
215fn get_total_rows<C>(conn: &mut PgConn<C>, query: &CXQuery<String>) -> usize
216where
217    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
218    C::TlsConnect: Send,
219    C::Stream: Send,
220    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
221{
222    let dialect = PostgreSqlDialect {};
223
224    let row = conn.query_one(count_query(query, &dialect)?.as_str(), &[])?;
225    let col_type = PostgresTypeSystem::from(row.columns()[0].type_());
226    match col_type {
227        PostgresTypeSystem::Int2(_) => convert_row::<i16>(&row) as usize,
228        PostgresTypeSystem::Int4(_) => convert_row::<i32>(&row) as usize,
229        PostgresTypeSystem::Int8(_) => convert_row::<i64>(&row) as usize,
230        _ => throw!(anyhow!(
231            "The result of the count query was not an int, aborting."
232        )),
233    }
234}
235
236pub struct PostgresSource<P, C>
237where
238    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
239    C::TlsConnect: Send,
240    C::Stream: Send,
241    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
242{
243    pool: Pool<PgManager<C>>,
244    origin_query: Option<String>,
245    queries: Vec<CXQuery<String>>,
246    names: Vec<String>,
247    schema: Vec<PostgresTypeSystem>,
248    pg_schema: Vec<postgres::types::Type>,
249    pre_execution_queries: Option<Vec<String>>,
250    _protocol: PhantomData<P>,
251}
252
253impl<P, C> PostgresSource<P, C>
254where
255    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
256    C::TlsConnect: Send,
257    C::Stream: Send,
258    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
259{
260    #[throws(PostgresSourceError)]
261    pub fn new(config: Config, tls: C, nconn: usize) -> Self {
262        let manager = PostgresConnectionManager::new(config, tls);
263        let pool = Pool::builder().max_size(nconn as u32).build(manager)?;
264
265        Self {
266            pool,
267            origin_query: None,
268            queries: vec![],
269            names: vec![],
270            schema: vec![],
271            pg_schema: vec![],
272            pre_execution_queries: None,
273            _protocol: PhantomData,
274        }
275    }
276}
277
278impl<P, C> Source for PostgresSource<P, C>
279where
280    PostgresSourcePartition<P, C>:
281        SourcePartition<TypeSystem = PostgresTypeSystem, Error = PostgresSourceError>,
282    P: Send,
283    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
284    C::TlsConnect: Send,
285    C::Stream: Send,
286    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
287{
288    const DATA_ORDERS: &'static [DataOrder] = &[DataOrder::RowMajor];
289    type Partition = PostgresSourcePartition<P, C>;
290    type TypeSystem = PostgresTypeSystem;
291    type Error = PostgresSourceError;
292
293    #[throws(PostgresSourceError)]
294    fn set_data_order(&mut self, data_order: DataOrder) {
295        if !matches!(data_order, DataOrder::RowMajor) {
296            throw!(ConnectorXError::UnsupportedDataOrder(data_order));
297        }
298    }
299
300    fn set_queries<Q: ToString>(&mut self, queries: &[CXQuery<Q>]) {
301        self.queries = queries.iter().map(|q| q.map(Q::to_string)).collect();
302    }
303
304    fn set_origin_query(&mut self, query: Option<String>) {
305        self.origin_query = query;
306    }
307
308    fn set_pre_execution_queries(&mut self, pre_execution_queries: Option<&[String]>) {
309        self.pre_execution_queries = pre_execution_queries.map(|s| s.to_vec());
310    }
311
312    #[throws(PostgresSourceError)]
313    fn fetch_metadata(&mut self) {
314        assert!(!self.queries.is_empty());
315
316        let mut conn = self.pool.get()?;
317        let first_query = &self.queries[0];
318
319        let stmt = conn.prepare(first_query.as_str())?;
320
321        let (names, pg_types): (Vec<String>, Vec<postgres::types::Type>) = stmt
322            .columns()
323            .iter()
324            .map(|col| (col.name().to_string(), col.type_().clone()))
325            .unzip();
326
327        self.names = names;
328        self.schema = pg_types.iter().map(PostgresTypeSystem::from).collect();
329        self.pg_schema = self
330            .schema
331            .iter()
332            .zip(pg_types.iter())
333            .map(|(t1, t2)| PostgresTypePairs(t2, t1).into())
334            .collect();
335    }
336
337    #[throws(PostgresSourceError)]
338    fn result_rows(&mut self) -> Option<usize> {
339        match &self.origin_query {
340            Some(q) => {
341                let cxq = CXQuery::Naked(q.clone());
342                let mut conn = self.pool.get()?;
343                let nrows = get_total_rows(&mut conn, &cxq)?;
344                Some(nrows)
345            }
346            None => None,
347        }
348    }
349
350    fn names(&self) -> Vec<String> {
351        self.names.clone()
352    }
353
354    fn schema(&self) -> Vec<Self::TypeSystem> {
355        self.schema.clone()
356    }
357
358    #[throws(PostgresSourceError)]
359    fn partition(self) -> Vec<Self::Partition> {
360        let mut ret = vec![];
361        for query in self.queries {
362            let mut conn = self.pool.get()?;
363            let rewritten = maybe_rewrite_text_query(&query, &self.names, &self.schema);
364
365            if let Some(pre_queries) = &self.pre_execution_queries {
366                for pre_query in pre_queries {
367                    conn.query(pre_query, &[])?;
368                }
369            }
370
371            ret.push(PostgresSourcePartition::<P, C>::new(
372                conn,
373                &rewritten,
374                &self.schema,
375                &self.pg_schema,
376            ));
377        }
378        ret
379    }
380}
381
382pub struct PostgresSourcePartition<P, C>
383where
384    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
385    C::TlsConnect: Send,
386    C::Stream: Send,
387    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
388{
389    conn: PgConn<C>,
390    query: CXQuery<String>,
391    schema: Vec<PostgresTypeSystem>,
392    pg_schema: Vec<postgres::types::Type>,
393    nrows: usize,
394    ncols: usize,
395    _protocol: PhantomData<P>,
396}
397
398impl<P, C> PostgresSourcePartition<P, C>
399where
400    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
401    C::TlsConnect: Send,
402    C::Stream: Send,
403    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
404{
405    pub fn new(
406        conn: PgConn<C>,
407        query: &CXQuery<String>,
408        schema: &[PostgresTypeSystem],
409        pg_schema: &[postgres::types::Type],
410    ) -> Self {
411        Self {
412            conn,
413            query: query.clone(),
414            schema: schema.to_vec(),
415            pg_schema: pg_schema.to_vec(),
416            nrows: 0,
417            ncols: schema.len(),
418            _protocol: PhantomData,
419        }
420    }
421}
422
423impl<C> SourcePartition for PostgresSourcePartition<BinaryProtocol, C>
424where
425    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
426    C::TlsConnect: Send,
427    C::Stream: Send,
428    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
429{
430    type TypeSystem = PostgresTypeSystem;
431    type Parser<'a> = PostgresBinarySourcePartitionParser<'a>;
432    type Error = PostgresSourceError;
433
434    #[throws(PostgresSourceError)]
435    fn result_rows(&mut self) -> () {
436        self.nrows = get_total_rows(&mut self.conn, &self.query)?;
437    }
438
439    #[throws(PostgresSourceError)]
440    fn parser(&mut self) -> Self::Parser<'_> {
441        let query = format!("COPY ({}) TO STDOUT WITH BINARY", self.query);
442        let reader = self.conn.copy_out(&*query)?; // unless reading the data, it seems like issue the query is fast
443        let iter = BinaryCopyOutIter::new(reader, &self.pg_schema);
444
445        PostgresBinarySourcePartitionParser::new(iter, &self.schema)
446    }
447
448    fn nrows(&self) -> usize {
449        self.nrows
450    }
451
452    fn ncols(&self) -> usize {
453        self.ncols
454    }
455}
456
457impl<C> SourcePartition for PostgresSourcePartition<CSVProtocol, C>
458where
459    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
460    C::TlsConnect: Send,
461    C::Stream: Send,
462    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
463{
464    type TypeSystem = PostgresTypeSystem;
465    type Parser<'a> = PostgresCSVSourceParser<'a>;
466    type Error = PostgresSourceError;
467
468    #[throws(PostgresSourceError)]
469    fn result_rows(&mut self) {
470        self.nrows = get_total_rows(&mut self.conn, &self.query)?;
471    }
472
473    #[throws(PostgresSourceError)]
474    fn parser(&mut self) -> Self::Parser<'_> {
475        let query = format!("COPY ({}) TO STDOUT WITH CSV", self.query);
476        let reader = self.conn.copy_out(&*query)?; // unless reading the data, it seems like issue the query is fast
477        let iter = ReaderBuilder::new()
478            .has_headers(false)
479            .from_reader(reader)
480            .into_records();
481
482        PostgresCSVSourceParser::new(iter, &self.schema)
483    }
484
485    fn nrows(&self) -> usize {
486        self.nrows
487    }
488
489    fn ncols(&self) -> usize {
490        self.ncols
491    }
492}
493
494impl<C> SourcePartition for PostgresSourcePartition<CursorProtocol, C>
495where
496    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
497    C::TlsConnect: Send,
498    C::Stream: Send,
499    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
500{
501    type TypeSystem = PostgresTypeSystem;
502    type Parser<'a> = PostgresRawSourceParser<'a>;
503    type Error = PostgresSourceError;
504
505    #[throws(PostgresSourceError)]
506    fn result_rows(&mut self) {
507        self.nrows = get_total_rows(&mut self.conn, &self.query)?;
508    }
509
510    #[throws(PostgresSourceError)]
511    fn parser(&mut self) -> Self::Parser<'_> {
512        let iter = self
513            .conn
514            .query_raw::<_, bool, _>(self.query.as_str(), vec![])?; // unless reading the data, it seems like issue the query is fast
515        PostgresRawSourceParser::new(iter, &self.schema)
516    }
517
518    fn nrows(&self) -> usize {
519        self.nrows
520    }
521
522    fn ncols(&self) -> usize {
523        self.ncols
524    }
525}
526pub struct PostgresBinarySourcePartitionParser<'a> {
527    iter: BinaryCopyOutIter<'a>,
528    rowbuf: Vec<BinaryCopyOutRow>,
529    ncols: usize,
530    current_col: usize,
531    current_row: usize,
532    is_finished: bool,
533}
534
535impl<'a> PostgresBinarySourcePartitionParser<'a> {
536    pub fn new(iter: BinaryCopyOutIter<'a>, schema: &[PostgresTypeSystem]) -> Self {
537        Self {
538            iter,
539            rowbuf: Vec::with_capacity(DB_BUFFER_SIZE),
540            ncols: schema.len(),
541            current_row: 0,
542            current_col: 0,
543            is_finished: false,
544        }
545    }
546
547    #[throws(PostgresSourceError)]
548    fn next_loc(&mut self) -> (usize, usize) {
549        let ret = (self.current_row, self.current_col);
550        self.current_row += (self.current_col + 1) / self.ncols;
551        self.current_col = (self.current_col + 1) % self.ncols;
552        ret
553    }
554}
555
556impl<'a> PartitionParser<'a> for PostgresBinarySourcePartitionParser<'a> {
557    type TypeSystem = PostgresTypeSystem;
558    type Error = PostgresSourceError;
559
560    #[throws(PostgresSourceError)]
561    fn fetch_next(&mut self) -> (usize, bool) {
562        assert!(self.current_col == 0);
563        let remaining_rows = self.rowbuf.len() - self.current_row;
564        if remaining_rows > 0 {
565            return (remaining_rows, self.is_finished);
566        } else if self.is_finished {
567            return (0, self.is_finished);
568        }
569
570        // clear the buffer
571        if !self.rowbuf.is_empty() {
572            self.rowbuf.drain(..);
573        }
574        for _ in 0..DB_BUFFER_SIZE {
575            match self.iter.next()? {
576                Some(row) => {
577                    self.rowbuf.push(row);
578                }
579                None => {
580                    self.is_finished = true;
581                    break;
582                }
583            }
584        }
585
586        // reset current cursor positions
587        self.current_row = 0;
588        self.current_col = 0;
589
590        (self.rowbuf.len(), self.is_finished)
591    }
592}
593
594macro_rules! impl_produce {
595    ($($t: ty,)+) => {
596        $(
597            impl<'r, 'a> Produce<'r, $t> for PostgresBinarySourcePartitionParser<'a> {
598                type Error = PostgresSourceError;
599
600                #[throws(PostgresSourceError)]
601                fn produce(&'r mut self) -> $t {
602                    let (ridx, cidx) = self.next_loc()?;
603                    let row = &self.rowbuf[ridx];
604                    let val = row.try_get(cidx)?;
605                    val
606                }
607            }
608
609            impl<'r, 'a> Produce<'r, Option<$t>> for PostgresBinarySourcePartitionParser<'a> {
610                type Error = PostgresSourceError;
611
612                #[throws(PostgresSourceError)]
613                fn produce(&'r mut self) -> Option<$t> {
614                    let (ridx, cidx) = self.next_loc()?;
615                    let row = &self.rowbuf[ridx];
616                    let val = row.try_get(cidx)?;
617                    val
618                }
619            }
620        )+
621    };
622}
623
624impl_produce!(
625    i8,
626    i16,
627    i32,
628    i64,
629    u32,
630    f32,
631    f64,
632    Decimal,
633    bool,
634    &'r str,
635    Vec<u8>,
636    NaiveTime,
637    Uuid,
638    Value,
639    IpInet,
640    Vector,
641    HalfVector,
642    Bit,
643    SparseVector,
644    Vec<Option<bool>>,
645    Vec<Option<i16>>,
646    Vec<Option<i32>>,
647    Vec<Option<i64>>,
648    Vec<Option<Decimal>>,
649    Vec<Option<f32>>,
650    Vec<Option<f64>>,
651    Vec<Option<String>>,
652);
653
654impl<'r> Produce<'r, NaiveDateTime> for PostgresBinarySourcePartitionParser<'_> {
655    type Error = PostgresSourceError;
656
657    #[throws(PostgresSourceError)]
658    fn produce(&'r mut self) -> NaiveDateTime {
659        let (ridx, cidx) = self.next_loc()?;
660        let row = &self.rowbuf[ridx];
661        let val = row.try_get(cidx)?;
662        match val {
663            postgres::types::Timestamp::PosInfinity => NaiveDateTime::MAX,
664            postgres::types::Timestamp::NegInfinity => NaiveDateTime::MIN,
665            postgres::types::Timestamp::Value(t) => t,
666        }
667    }
668}
669
670impl<'r> Produce<'r, Option<NaiveDateTime>> for PostgresBinarySourcePartitionParser<'_> {
671    type Error = PostgresSourceError;
672
673    #[throws(PostgresSourceError)]
674    fn produce(&'r mut self) -> Option<NaiveDateTime> {
675        let (ridx, cidx) = self.next_loc()?;
676        let row = &self.rowbuf[ridx];
677        let val = row.try_get(cidx)?;
678        match val {
679            Some(postgres::types::Timestamp::PosInfinity) => Some(NaiveDateTime::MAX),
680            Some(postgres::types::Timestamp::NegInfinity) => Some(NaiveDateTime::MIN),
681            Some(postgres::types::Timestamp::Value(t)) => t,
682            None => None,
683        }
684    }
685}
686
687impl<'r> Produce<'r, DateTime<Utc>> for PostgresBinarySourcePartitionParser<'_> {
688    type Error = PostgresSourceError;
689
690    #[throws(PostgresSourceError)]
691    fn produce(&'r mut self) -> DateTime<Utc> {
692        let (ridx, cidx) = self.next_loc()?;
693        let row = &self.rowbuf[ridx];
694        let val = row.try_get(cidx)?;
695        match val {
696            postgres::types::Timestamp::PosInfinity => DateTime::<Utc>::MAX_UTC,
697            postgres::types::Timestamp::NegInfinity => DateTime::<Utc>::MIN_UTC,
698            postgres::types::Timestamp::Value(t) => t,
699        }
700    }
701}
702
703impl<'r> Produce<'r, Option<DateTime<Utc>>> for PostgresBinarySourcePartitionParser<'_> {
704    type Error = PostgresSourceError;
705
706    #[throws(PostgresSourceError)]
707    fn produce(&'r mut self) -> Option<DateTime<Utc>> {
708        let (ridx, cidx) = self.next_loc()?;
709        let row = &self.rowbuf[ridx];
710        let val = row.try_get(cidx)?;
711        match val {
712            Some(postgres::types::Timestamp::PosInfinity) => Some(DateTime::<Utc>::MAX_UTC),
713            Some(postgres::types::Timestamp::NegInfinity) => Some(DateTime::<Utc>::MIN_UTC),
714            Some(postgres::types::Timestamp::Value(t)) => t,
715            None => None,
716        }
717    }
718}
719
720impl<'r> Produce<'r, NaiveDate> for PostgresBinarySourcePartitionParser<'_> {
721    type Error = PostgresSourceError;
722
723    #[throws(PostgresSourceError)]
724    fn produce(&'r mut self) -> NaiveDate {
725        let (ridx, cidx) = self.next_loc()?;
726        let row = &self.rowbuf[ridx];
727        let val = row.try_get(cidx)?;
728        match val {
729            postgres::types::Date::PosInfinity => NaiveDate::MAX,
730            postgres::types::Date::NegInfinity => NaiveDate::MIN,
731            postgres::types::Date::Value(t) => t,
732        }
733    }
734}
735
736impl<'r> Produce<'r, Option<NaiveDate>> for PostgresBinarySourcePartitionParser<'_> {
737    type Error = PostgresSourceError;
738
739    #[throws(PostgresSourceError)]
740    fn produce(&'r mut self) -> Option<NaiveDate> {
741        let (ridx, cidx) = self.next_loc()?;
742        let row = &self.rowbuf[ridx];
743        let val = row.try_get(cidx)?;
744        match val {
745            Some(postgres::types::Date::PosInfinity) => Some(NaiveDate::MAX),
746            Some(postgres::types::Date::NegInfinity) => Some(NaiveDate::MIN),
747            Some(postgres::types::Date::Value(t)) => t,
748            None => None,
749        }
750    }
751}
752
753impl Produce<'_, HashMap<String, Option<String>>> for PostgresBinarySourcePartitionParser<'_> {
754    type Error = PostgresSourceError;
755    #[throws(PostgresSourceError)]
756    fn produce(&mut self) -> HashMap<String, Option<String>> {
757        unimplemented!("Please use `cursor` protocol for hstore type");
758    }
759}
760
761impl Produce<'_, Option<HashMap<String, Option<String>>>>
762    for PostgresBinarySourcePartitionParser<'_>
763{
764    type Error = PostgresSourceError;
765    #[throws(PostgresSourceError)]
766    fn produce(&mut self) -> Option<HashMap<String, Option<String>>> {
767        unimplemented!("Please use `cursor` protocol for hstore type");
768    }
769}
770
771pub struct PostgresCSVSourceParser<'a> {
772    iter: StringRecordsIntoIter<CopyOutReader<'a>>,
773    rowbuf: Vec<StringRecord>,
774    ncols: usize,
775    current_col: usize,
776    current_row: usize,
777    is_finished: bool,
778}
779
780impl<'a> PostgresCSVSourceParser<'a> {
781    pub fn new(
782        iter: StringRecordsIntoIter<CopyOutReader<'a>>,
783        schema: &[PostgresTypeSystem],
784    ) -> Self {
785        Self {
786            iter,
787            rowbuf: Vec::with_capacity(DB_BUFFER_SIZE),
788            ncols: schema.len(),
789            current_row: 0,
790            current_col: 0,
791            is_finished: false,
792        }
793    }
794
795    #[throws(PostgresSourceError)]
796    fn next_loc(&mut self) -> (usize, usize) {
797        let ret = (self.current_row, self.current_col);
798        self.current_row += (self.current_col + 1) / self.ncols;
799        self.current_col = (self.current_col + 1) % self.ncols;
800        ret
801    }
802}
803
804impl<'a> PartitionParser<'a> for PostgresCSVSourceParser<'a> {
805    type Error = PostgresSourceError;
806    type TypeSystem = PostgresTypeSystem;
807
808    #[throws(PostgresSourceError)]
809    fn fetch_next(&mut self) -> (usize, bool) {
810        assert!(self.current_col == 0);
811        let remaining_rows = self.rowbuf.len() - self.current_row;
812        if remaining_rows > 0 {
813            return (remaining_rows, self.is_finished);
814        } else if self.is_finished {
815            return (0, self.is_finished);
816        }
817
818        if !self.rowbuf.is_empty() {
819            self.rowbuf.drain(..);
820        }
821        for _ in 0..DB_BUFFER_SIZE {
822            if let Some(row) = self.iter.next() {
823                self.rowbuf.push(row?);
824            } else {
825                self.is_finished = true;
826                break;
827            }
828        }
829        self.current_row = 0;
830        self.current_col = 0;
831        (self.rowbuf.len(), self.is_finished)
832    }
833}
834
835macro_rules! impl_csv_produce {
836    ($($t: ty,)+) => {
837        $(
838            impl<'r, 'a> Produce<'r, $t> for PostgresCSVSourceParser<'a> {
839                type Error = PostgresSourceError;
840
841                #[throws(PostgresSourceError)]
842                fn produce(&'r mut self) -> $t {
843                    let (ridx, cidx) = self.next_loc()?;
844                    self.rowbuf[ridx][cidx].parse().map_err(|_| {
845                        ConnectorXError::cannot_produce::<$t>(Some(self.rowbuf[ridx][cidx].into()))
846                    })?
847                }
848            }
849
850            impl<'r, 'a> Produce<'r, Option<$t>> for PostgresCSVSourceParser<'a> {
851                type Error = PostgresSourceError;
852
853                #[throws(PostgresSourceError)]
854                fn produce(&'r mut self) -> Option<$t> {
855                    let (ridx, cidx) = self.next_loc()?;
856                    match &self.rowbuf[ridx][cidx][..] {
857                        "" => None,
858                        v => Some(v.parse().map_err(|_| {
859                            ConnectorXError::cannot_produce::<$t>(Some(self.rowbuf[ridx][cidx].into()))
860                        })?),
861                    }
862                }
863            }
864        )+
865    };
866}
867
868impl_csv_produce!(i8, i16, i32, i64, u32, f32, f64, Uuid, IpInet,);
869
870macro_rules! impl_csv_vec_produce {
871    ($($t: ty,)+) => {
872        $(
873            impl<'r, 'a> Produce<'r, Vec<Option<$t>>> for PostgresCSVSourceParser<'a> {
874                type Error = PostgresSourceError;
875
876                #[throws(PostgresSourceError)]
877                fn produce(&mut self) -> Vec<Option<$t>> {
878                    let (ridx, cidx) = self.next_loc()?;
879                    let s = &self.rowbuf[ridx][cidx][..];
880                    match s {
881                        "{}" => vec![],
882                        _ if s.len() < 3 => throw!(ConnectorXError::cannot_produce::<$t>(Some(s.into()))),
883                        s => s[1..s.len() - 1]
884                            .split(",")
885                            .map(|v| {
886                                if v == "NULL" {
887                                    Ok(None)
888                                } else {
889                                    match v.parse() {
890                                        Ok(v) => Ok(Some(v)),
891                                        Err(e) => Err(e).map_err(|_| ConnectorXError::cannot_produce::<$t>(Some(s.into())))
892                                    }
893                                }
894                            })
895                            .collect::<Result<Vec<Option<$t>>, ConnectorXError>>()?,
896                    }
897                }
898            }
899
900            impl<'r, 'a> Produce<'r, Option<Vec<Option<$t>>>> for PostgresCSVSourceParser<'a> {
901                type Error = PostgresSourceError;
902
903                #[throws(PostgresSourceError)]
904                fn produce(&mut self) -> Option<Vec<Option<$t>>> {
905                    let (ridx, cidx) = self.next_loc()?;
906                    let s = &self.rowbuf[ridx][cidx][..];
907                    match s {
908                        "" => None,
909                        "{}" => Some(vec![]),
910                        _ if s.len() < 3 => throw!(ConnectorXError::cannot_produce::<$t>(Some(s.into()))),
911                        s => Some(
912                            s[1..s.len() - 1]
913                                .split(",")
914                                .map(|v| {
915                                    if v == "NULL" {
916                                        Ok(None)
917                                    } else {
918                                        match v.parse() {
919                                            Ok(v) => Ok(Some(v)),
920                                            Err(e) => Err(e).map_err(|_| ConnectorXError::cannot_produce::<$t>(Some(s.into())))
921                                        }
922                                    }
923                                })
924                                .collect::<Result<Vec<Option<$t>>, ConnectorXError>>()?,
925                        ),
926                    }
927                }
928            }
929        )+
930    };
931}
932
933impl_csv_vec_produce!(i8, i16, i32, i64, f32, f64, Decimal, String,);
934
935impl Produce<'_, bool> for PostgresCSVSourceParser<'_> {
936    type Error = PostgresSourceError;
937
938    #[throws(PostgresSourceError)]
939    fn produce(&mut self) -> bool {
940        let (ridx, cidx) = self.next_loc()?;
941        let ret = match &self.rowbuf[ridx][cidx][..] {
942            "t" => true,
943            "f" => false,
944            _ => throw!(ConnectorXError::cannot_produce::<bool>(Some(
945                self.rowbuf[ridx][cidx].into()
946            ))),
947        };
948        ret
949    }
950}
951
952impl Produce<'_, Option<bool>> for PostgresCSVSourceParser<'_> {
953    type Error = PostgresSourceError;
954
955    #[throws(PostgresSourceError)]
956    fn produce(&mut self) -> Option<bool> {
957        let (ridx, cidx) = self.next_loc()?;
958        let ret = match &self.rowbuf[ridx][cidx][..] {
959            "" => None,
960            "t" => Some(true),
961            "f" => Some(false),
962            _ => throw!(ConnectorXError::cannot_produce::<bool>(Some(
963                self.rowbuf[ridx][cidx].into()
964            ))),
965        };
966        ret
967    }
968}
969
970impl Produce<'_, Vec<Option<bool>>> for PostgresCSVSourceParser<'_> {
971    type Error = PostgresSourceError;
972
973    #[throws(PostgresSourceError)]
974    fn produce(&mut self) -> Vec<Option<bool>> {
975        let (ridx, cidx) = self.next_loc()?;
976        let s = &self.rowbuf[ridx][cidx][..];
977        match s {
978            "{}" => vec![],
979            _ if s.len() < 3 => throw!(ConnectorXError::cannot_produce::<bool>(Some(s.into()))),
980            s => s[1..s.len() - 1]
981                .split(',')
982                .map(|v| match v {
983                    "NULL" => Ok(None),
984                    "t" => Ok(Some(true)),
985                    "f" => Ok(Some(false)),
986                    _ => throw!(ConnectorXError::cannot_produce::<bool>(Some(s.into()))),
987                })
988                .collect::<Result<Vec<Option<bool>>, ConnectorXError>>()?,
989        }
990    }
991}
992
993impl Produce<'_, Option<Vec<Option<bool>>>> for PostgresCSVSourceParser<'_> {
994    type Error = PostgresSourceError;
995
996    #[throws(PostgresSourceError)]
997    fn produce(&mut self) -> Option<Vec<Option<bool>>> {
998        let (ridx, cidx) = self.next_loc()?;
999        let s = &self.rowbuf[ridx][cidx][..];
1000        match s {
1001            "" => None,
1002            "{}" => Some(vec![]),
1003            _ if s.len() < 3 => throw!(ConnectorXError::cannot_produce::<bool>(Some(s.into()))),
1004            s => Some(
1005                s[1..s.len() - 1]
1006                    .split(',')
1007                    .map(|v| match v {
1008                        "NULL" => Ok(None),
1009                        "t" => Ok(Some(true)),
1010                        "f" => Ok(Some(false)),
1011                        _ => throw!(ConnectorXError::cannot_produce::<bool>(Some(s.into()))),
1012                    })
1013                    .collect::<Result<Vec<Option<bool>>, ConnectorXError>>()?,
1014            ),
1015        }
1016    }
1017}
1018
1019impl<'r> Produce<'r, Decimal> for PostgresCSVSourceParser<'_> {
1020    type Error = PostgresSourceError;
1021
1022    #[throws(PostgresSourceError)]
1023    fn produce(&'r mut self) -> Decimal {
1024        let (ridx, cidx) = self.next_loc()?;
1025        match &self.rowbuf[ridx][cidx][..] {
1026            "Infinity" => Decimal::MAX,
1027            "-Infinity" => Decimal::MIN,
1028            v => v
1029                .parse()
1030                .map_err(|_| ConnectorXError::cannot_produce::<Decimal>(Some(v.into())))?,
1031        }
1032    }
1033}
1034
1035impl<'r> Produce<'r, Option<Decimal>> for PostgresCSVSourceParser<'_> {
1036    type Error = PostgresSourceError;
1037
1038    #[throws(PostgresSourceError)]
1039    fn produce(&'r mut self) -> Option<Decimal> {
1040        let (ridx, cidx) = self.next_loc()?;
1041        match &self.rowbuf[ridx][cidx][..] {
1042            "" => None,
1043            "Infinity" => Some(Decimal::MAX),
1044            "-Infinity" => Some(Decimal::MIN),
1045            v => Some(
1046                v.parse()
1047                    .map_err(|_| ConnectorXError::cannot_produce::<Decimal>(Some(v.into())))?,
1048            ),
1049        }
1050    }
1051}
1052
1053impl Produce<'_, DateTime<Utc>> for PostgresCSVSourceParser<'_> {
1054    type Error = PostgresSourceError;
1055
1056    #[throws(PostgresSourceError)]
1057    fn produce(&mut self) -> DateTime<Utc> {
1058        let (ridx, cidx) = self.next_loc()?;
1059        match &self.rowbuf[ridx][cidx][..] {
1060            "infinity" => DateTime::<Utc>::MAX_UTC,
1061            "-infinity" => DateTime::<Utc>::MIN_UTC,
1062            // postgres csv return example: 1970-01-01 00:00:01+00
1063            v => format!("{}:00", v)
1064                .parse()
1065                .map_err(|_| ConnectorXError::cannot_produce::<DateTime<Utc>>(Some(v.into())))?,
1066        }
1067    }
1068}
1069
1070impl Produce<'_, Option<DateTime<Utc>>> for PostgresCSVSourceParser<'_> {
1071    type Error = PostgresSourceError;
1072
1073    #[throws(PostgresSourceError)]
1074    fn produce(&mut self) -> Option<DateTime<Utc>> {
1075        let (ridx, cidx) = self.next_loc()?;
1076        match &self.rowbuf[ridx][cidx][..] {
1077            "" => None,
1078            "infinity" => Some(DateTime::<Utc>::MAX_UTC),
1079            "-infinity" => Some(DateTime::<Utc>::MIN_UTC),
1080            v => {
1081                // postgres csv return example: 1970-01-01 00:00:01+00
1082                Some(format!("{}:00", v).parse().map_err(|_| {
1083                    ConnectorXError::cannot_produce::<DateTime<Utc>>(Some(v.into()))
1084                })?)
1085            }
1086        }
1087    }
1088}
1089
1090impl Produce<'_, NaiveDate> for PostgresCSVSourceParser<'_> {
1091    type Error = PostgresSourceError;
1092
1093    #[throws(PostgresSourceError)]
1094    fn produce(&mut self) -> NaiveDate {
1095        let (ridx, cidx) = self.next_loc()?;
1096        match &self.rowbuf[ridx][cidx][..] {
1097            "infinity" => NaiveDate::MAX,
1098            "-infinity" => NaiveDate::MIN,
1099            v => NaiveDate::parse_from_str(v, "%Y-%m-%d")
1100                .map_err(|_| ConnectorXError::cannot_produce::<NaiveDate>(Some(v.into())))?,
1101        }
1102    }
1103}
1104
1105impl Produce<'_, Option<NaiveDate>> for PostgresCSVSourceParser<'_> {
1106    type Error = PostgresSourceError;
1107
1108    #[throws(PostgresSourceError)]
1109    fn produce(&mut self) -> Option<NaiveDate> {
1110        let (ridx, cidx) = self.next_loc()?;
1111        match &self.rowbuf[ridx][cidx][..] {
1112            "" => None,
1113            "infinity" => Some(NaiveDate::MAX),
1114            "-infinity" => Some(NaiveDate::MIN),
1115            v => Some(
1116                NaiveDate::parse_from_str(v, "%Y-%m-%d")
1117                    .map_err(|_| ConnectorXError::cannot_produce::<NaiveDate>(Some(v.into())))?,
1118            ),
1119        }
1120    }
1121}
1122
1123impl Produce<'_, NaiveDateTime> for PostgresCSVSourceParser<'_> {
1124    type Error = PostgresSourceError;
1125
1126    #[throws(PostgresSourceError)]
1127    fn produce(&mut self) -> NaiveDateTime {
1128        let (ridx, cidx) = self.next_loc()?;
1129        match &self.rowbuf[ridx][cidx] {
1130            "infinity" => NaiveDateTime::MAX,
1131            "-infinity" => NaiveDateTime::MIN,
1132            v => NaiveDateTime::parse_from_str(v, "%Y-%m-%d %H:%M:%S%.f")
1133                .map_err(|_| ConnectorXError::cannot_produce::<NaiveDateTime>(Some(v.into())))?,
1134        }
1135    }
1136}
1137
1138impl Produce<'_, Option<NaiveDateTime>> for PostgresCSVSourceParser<'_> {
1139    type Error = PostgresSourceError;
1140
1141    #[throws(PostgresSourceError)]
1142    fn produce(&mut self) -> Option<NaiveDateTime> {
1143        let (ridx, cidx) = self.next_loc()?;
1144        match &self.rowbuf[ridx][cidx][..] {
1145            "" => None,
1146            "infinity" => Some(NaiveDateTime::MAX),
1147            "-infinity" => Some(NaiveDateTime::MIN),
1148            v => Some(
1149                NaiveDateTime::parse_from_str(v, "%Y-%m-%d %H:%M:%S%.f").map_err(|_| {
1150                    ConnectorXError::cannot_produce::<NaiveDateTime>(Some(v.into()))
1151                })?,
1152            ),
1153        }
1154    }
1155}
1156
1157impl Produce<'_, NaiveTime> for PostgresCSVSourceParser<'_> {
1158    type Error = PostgresSourceError;
1159
1160    #[throws(PostgresSourceError)]
1161    fn produce(&mut self) -> NaiveTime {
1162        let (ridx, cidx) = self.next_loc()?;
1163        NaiveTime::parse_from_str(&self.rowbuf[ridx][cidx], "%H:%M:%S%.f").map_err(|_| {
1164            ConnectorXError::cannot_produce::<NaiveTime>(Some(self.rowbuf[ridx][cidx].into()))
1165        })?
1166    }
1167}
1168
1169impl Produce<'_, Option<NaiveTime>> for PostgresCSVSourceParser<'_> {
1170    type Error = PostgresSourceError;
1171
1172    #[throws(PostgresSourceError)]
1173    fn produce(&mut self) -> Option<NaiveTime> {
1174        let (ridx, cidx) = self.next_loc()?;
1175        match &self.rowbuf[ridx][cidx][..] {
1176            "" => None,
1177            v => Some(
1178                NaiveTime::parse_from_str(v, "%H:%M:%S%.f")
1179                    .map_err(|_| ConnectorXError::cannot_produce::<NaiveTime>(Some(v.into())))?,
1180            ),
1181        }
1182    }
1183}
1184
1185impl<'r> Produce<'r, &'r str> for PostgresCSVSourceParser<'_> {
1186    type Error = PostgresSourceError;
1187
1188    #[throws(PostgresSourceError)]
1189    fn produce(&'r mut self) -> &'r str {
1190        let (ridx, cidx) = self.next_loc()?;
1191        &self.rowbuf[ridx][cidx]
1192    }
1193}
1194
1195impl<'r> Produce<'r, Option<&'r str>> for PostgresCSVSourceParser<'_> {
1196    type Error = PostgresSourceError;
1197
1198    #[throws(PostgresSourceError)]
1199    fn produce(&'r mut self) -> Option<&'r str> {
1200        let (ridx, cidx) = self.next_loc()?;
1201        match &self.rowbuf[ridx][cidx][..] {
1202            "" => None,
1203            v => Some(v),
1204        }
1205    }
1206}
1207
1208impl<'r> Produce<'r, Vec<u8>> for PostgresCSVSourceParser<'_> {
1209    type Error = PostgresSourceError;
1210
1211    #[throws(PostgresSourceError)]
1212    fn produce(&'r mut self) -> Vec<u8> {
1213        let (ridx, cidx) = self.next_loc()?;
1214        decode(&self.rowbuf[ridx][cidx][2..])? // escape \x in the beginning
1215    }
1216}
1217
1218impl<'r> Produce<'r, Option<Vec<u8>>> for PostgresCSVSourceParser<'_> {
1219    type Error = PostgresSourceError;
1220
1221    #[throws(PostgresSourceError)]
1222    fn produce(&'r mut self) -> Option<Vec<u8>> {
1223        let (ridx, cidx) = self.next_loc()?;
1224        match &self.rowbuf[ridx][cidx] {
1225            // escape \x in the beginning, empty if None
1226            "" => None,
1227            v => Some(decode(&v[2..])?),
1228        }
1229    }
1230}
1231
1232impl<'r> Produce<'r, Value> for PostgresCSVSourceParser<'_> {
1233    type Error = PostgresSourceError;
1234
1235    #[throws(PostgresSourceError)]
1236    fn produce(&'r mut self) -> Value {
1237        let (ridx, cidx) = self.next_loc()?;
1238        let v = &self.rowbuf[ridx][cidx];
1239        from_str(v).map_err(|_| ConnectorXError::cannot_produce::<Value>(Some(v.into())))?
1240    }
1241}
1242
1243impl<'r> Produce<'r, Option<Value>> for PostgresCSVSourceParser<'_> {
1244    type Error = PostgresSourceError;
1245
1246    #[throws(PostgresSourceError)]
1247    fn produce(&'r mut self) -> Option<Value> {
1248        let (ridx, cidx) = self.next_loc()?;
1249
1250        match &self.rowbuf[ridx][cidx][..] {
1251            "" => None,
1252            v => {
1253                from_str(v).map_err(|_| ConnectorXError::cannot_produce::<Value>(Some(v.into())))?
1254            }
1255        }
1256    }
1257}
1258
1259pub struct PostgresRawSourceParser<'a> {
1260    iter: RowIter<'a>,
1261    rowbuf: Vec<Row>,
1262    ncols: usize,
1263    current_col: usize,
1264    current_row: usize,
1265    is_finished: bool,
1266}
1267
1268impl<'a> PostgresRawSourceParser<'a> {
1269    pub fn new(iter: RowIter<'a>, schema: &[PostgresTypeSystem]) -> Self {
1270        Self {
1271            iter,
1272            rowbuf: Vec::with_capacity(DB_BUFFER_SIZE),
1273            ncols: schema.len(),
1274            current_row: 0,
1275            current_col: 0,
1276            is_finished: false,
1277        }
1278    }
1279
1280    #[throws(PostgresSourceError)]
1281    fn next_loc(&mut self) -> (usize, usize) {
1282        let ret = (self.current_row, self.current_col);
1283        self.current_row += (self.current_col + 1) / self.ncols;
1284        self.current_col = (self.current_col + 1) % self.ncols;
1285        ret
1286    }
1287}
1288
1289impl<'a> PartitionParser<'a> for PostgresRawSourceParser<'a> {
1290    type TypeSystem = PostgresTypeSystem;
1291    type Error = PostgresSourceError;
1292
1293    #[throws(PostgresSourceError)]
1294    fn fetch_next(&mut self) -> (usize, bool) {
1295        assert!(self.current_col == 0);
1296        let remaining_rows = self.rowbuf.len() - self.current_row;
1297        if remaining_rows > 0 {
1298            return (remaining_rows, self.is_finished);
1299        } else if self.is_finished {
1300            return (0, self.is_finished);
1301        }
1302
1303        if !self.rowbuf.is_empty() {
1304            self.rowbuf.drain(..);
1305        }
1306        for _ in 0..DB_BUFFER_SIZE {
1307            if let Some(row) = self.iter.next()? {
1308                self.rowbuf.push(row);
1309            } else {
1310                self.is_finished = true;
1311                break;
1312            }
1313        }
1314        self.current_row = 0;
1315        self.current_col = 0;
1316        (self.rowbuf.len(), self.is_finished)
1317    }
1318}
1319
1320macro_rules! impl_produce {
1321    ($($t: ty,)+) => {
1322        $(
1323            impl<'r, 'a> Produce<'r, $t> for PostgresRawSourceParser<'a> {
1324                type Error = PostgresSourceError;
1325
1326                #[throws(PostgresSourceError)]
1327                fn produce(&'r mut self) -> $t {
1328                    let (ridx, cidx) = self.next_loc()?;
1329                    let row = &self.rowbuf[ridx];
1330                    let val = row.try_get(cidx)?;
1331                    val
1332                }
1333            }
1334
1335            impl<'r, 'a> Produce<'r, Option<$t>> for PostgresRawSourceParser<'a> {
1336                type Error = PostgresSourceError;
1337
1338                #[throws(PostgresSourceError)]
1339                fn produce(&'r mut self) -> Option<$t> {
1340                    let (ridx, cidx) = self.next_loc()?;
1341                    let row = &self.rowbuf[ridx];
1342                    let val = row.try_get(cidx)?;
1343                    val
1344                }
1345            }
1346        )+
1347    };
1348}
1349
1350impl_produce!(
1351    i8,
1352    i16,
1353    i32,
1354    i64,
1355    u32,
1356    f32,
1357    f64,
1358    Decimal,
1359    bool,
1360    &'r str,
1361    Vec<u8>,
1362    NaiveTime,
1363    Uuid,
1364    Value,
1365    IpInet,
1366    Vector,
1367    HalfVector,
1368    Bit,
1369    SparseVector,
1370    HashMap<String, Option<String>>,
1371    Vec<Option<bool>>,
1372    Vec<Option<String>>,
1373    Vec<Option<i16>>,
1374    Vec<Option<i32>>,
1375    Vec<Option<i64>>,
1376    Vec<Option<f32>>,
1377    Vec<Option<f64>>,
1378    Vec<Option<Decimal>>,
1379);
1380
1381impl<'r> Produce<'r, DateTime<Utc>> for PostgresRawSourceParser<'_> {
1382    type Error = PostgresSourceError;
1383
1384    #[throws(PostgresSourceError)]
1385    fn produce(&'r mut self) -> DateTime<Utc> {
1386        let (ridx, cidx) = self.next_loc()?;
1387        let row = &self.rowbuf[ridx];
1388        let val: postgres::types::Timestamp<DateTime<Utc>> = row.try_get(cidx)?;
1389        match val {
1390            postgres::types::Timestamp::PosInfinity => DateTime::<Utc>::MAX_UTC,
1391            postgres::types::Timestamp::NegInfinity => DateTime::<Utc>::MIN_UTC,
1392            postgres::types::Timestamp::Value(t) => t,
1393        }
1394    }
1395}
1396
1397impl<'r> Produce<'r, Option<DateTime<Utc>>> for PostgresRawSourceParser<'_> {
1398    type Error = PostgresSourceError;
1399
1400    #[throws(PostgresSourceError)]
1401    fn produce(&'r mut self) -> Option<DateTime<Utc>> {
1402        let (ridx, cidx) = self.next_loc()?;
1403        let row = &self.rowbuf[ridx];
1404        let val = row.try_get(cidx)?;
1405        match val {
1406            Some(postgres::types::Timestamp::PosInfinity) => Some(DateTime::<Utc>::MAX_UTC),
1407            Some(postgres::types::Timestamp::NegInfinity) => Some(DateTime::<Utc>::MIN_UTC),
1408            Some(postgres::types::Timestamp::Value(t)) => t,
1409            None => None,
1410        }
1411    }
1412}
1413
1414impl<'r> Produce<'r, NaiveDateTime> for PostgresRawSourceParser<'_> {
1415    type Error = PostgresSourceError;
1416
1417    #[throws(PostgresSourceError)]
1418    fn produce(&'r mut self) -> NaiveDateTime {
1419        let (ridx, cidx) = self.next_loc()?;
1420        let row = &self.rowbuf[ridx];
1421        let val: postgres::types::Timestamp<NaiveDateTime> = row.try_get(cidx)?;
1422        match val {
1423            postgres::types::Timestamp::PosInfinity => NaiveDateTime::MAX,
1424            postgres::types::Timestamp::NegInfinity => NaiveDateTime::MIN,
1425            postgres::types::Timestamp::Value(t) => t,
1426        }
1427    }
1428}
1429
1430impl<'r> Produce<'r, Option<NaiveDateTime>> for PostgresRawSourceParser<'_> {
1431    type Error = PostgresSourceError;
1432
1433    #[throws(PostgresSourceError)]
1434    fn produce(&'r mut self) -> Option<NaiveDateTime> {
1435        let (ridx, cidx) = self.next_loc()?;
1436        let row = &self.rowbuf[ridx];
1437        let val = row.try_get(cidx)?;
1438        match val {
1439            Some(postgres::types::Timestamp::PosInfinity) => Some(NaiveDateTime::MAX),
1440            Some(postgres::types::Timestamp::NegInfinity) => Some(NaiveDateTime::MIN),
1441            Some(postgres::types::Timestamp::Value(t)) => t,
1442            None => None,
1443        }
1444    }
1445}
1446
1447impl<'r> Produce<'r, NaiveDate> for PostgresRawSourceParser<'_> {
1448    type Error = PostgresSourceError;
1449
1450    #[throws(PostgresSourceError)]
1451    fn produce(&'r mut self) -> NaiveDate {
1452        let (ridx, cidx) = self.next_loc()?;
1453        let row = &self.rowbuf[ridx];
1454        let val: postgres::types::Date<NaiveDate> = row.try_get(cidx)?;
1455        match val {
1456            postgres::types::Date::PosInfinity => NaiveDate::MAX,
1457            postgres::types::Date::NegInfinity => NaiveDate::MIN,
1458            postgres::types::Date::Value(t) => t,
1459        }
1460    }
1461}
1462
1463impl<'r> Produce<'r, Option<NaiveDate>> for PostgresRawSourceParser<'_> {
1464    type Error = PostgresSourceError;
1465
1466    #[throws(PostgresSourceError)]
1467    fn produce(&'r mut self) -> Option<NaiveDate> {
1468        let (ridx, cidx) = self.next_loc()?;
1469        let row = &self.rowbuf[ridx];
1470        let val = row.try_get(cidx)?;
1471        match val {
1472            Some(postgres::types::Date::PosInfinity) => Some(NaiveDate::MAX),
1473            Some(postgres::types::Date::NegInfinity) => Some(NaiveDate::MIN),
1474            Some(postgres::types::Date::Value(t)) => t,
1475            None => None,
1476        }
1477    }
1478}
1479
1480impl<C> SourcePartition for PostgresSourcePartition<SimpleProtocol, C>
1481where
1482    C: MakeTlsConnect<Socket> + Clone + 'static + Sync + Send,
1483    C::TlsConnect: Send,
1484    C::Stream: Send,
1485    <C::TlsConnect as TlsConnect<Socket>>::Future: Send,
1486{
1487    type TypeSystem = PostgresTypeSystem;
1488    type Parser<'a> = PostgresSimpleSourceParser;
1489    type Error = PostgresSourceError;
1490
1491    #[throws(PostgresSourceError)]
1492    fn result_rows(&mut self) {
1493        self.nrows = get_total_rows(&mut self.conn, &self.query)?;
1494    }
1495
1496    #[throws(PostgresSourceError)]
1497    fn parser(&mut self) -> Self::Parser<'_> {
1498        let rows = self.conn.simple_query(self.query.as_str())?; // unless reading the data, it seems like issue the query is fast
1499        PostgresSimpleSourceParser::new(rows, &self.schema)
1500    }
1501
1502    fn nrows(&self) -> usize {
1503        self.nrows
1504    }
1505
1506    fn ncols(&self) -> usize {
1507        self.ncols
1508    }
1509}
1510
1511pub struct PostgresSimpleSourceParser {
1512    rows: Vec<SimpleQueryMessage>,
1513    ncols: usize,
1514    current_col: usize,
1515    current_row: usize,
1516}
1517impl PostgresSimpleSourceParser {
1518    pub fn new(rows: Vec<SimpleQueryMessage>, schema: &[PostgresTypeSystem]) -> Self {
1519        Self {
1520            rows,
1521            ncols: schema.len(),
1522            current_row: 0,
1523            current_col: 0,
1524        }
1525    }
1526
1527    #[throws(PostgresSourceError)]
1528    fn next_loc(&mut self) -> (usize, usize) {
1529        let ret = (self.current_row, self.current_col);
1530        self.current_row += (self.current_col + 1) / self.ncols;
1531        self.current_col = (self.current_col + 1) % self.ncols;
1532        ret
1533    }
1534}
1535
1536impl PartitionParser<'_> for PostgresSimpleSourceParser {
1537    type TypeSystem = PostgresTypeSystem;
1538    type Error = PostgresSourceError;
1539
1540    #[throws(PostgresSourceError)]
1541    fn fetch_next(&mut self) -> (usize, bool) {
1542        self.current_row = 0;
1543        self.current_col = 0;
1544        if !self.rows.is_empty() {
1545            if let SimpleQueryMessage::RowDescription(_) = &self.rows[0] {
1546                self.current_row = 1;
1547            }
1548        }
1549
1550        (self.rows.len() - 1 - self.current_row, true) // last message is command complete
1551    }
1552}
1553
1554macro_rules! impl_simple_produce {
1555    ($($t: ty,)+) => {
1556        $(
1557            impl<'r> Produce<'r, $t> for PostgresSimpleSourceParser {
1558                type Error = PostgresSourceError;
1559
1560                #[throws(PostgresSourceError)]
1561                fn produce(&'r mut self) -> $t {
1562                    let (ridx, cidx) = self.next_loc()?;
1563                    let val = match &self.rows[ridx] {
1564                        SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1565                            Some(s) => s
1566                                .parse()
1567                                .map_err(|_| ConnectorXError::cannot_produce::<$t>(Some(s.into())))?,
1568                            None => throw!(anyhow!(
1569                                "Cannot parse NULL in NOT NULL column."
1570                            )),
1571                        },
1572                        SimpleQueryMessage::CommandComplete(c) => {
1573                            panic!("get command: {}", c);
1574                        }
1575                        _ => {
1576                            panic!("what?");
1577                        }
1578                    };
1579                    val
1580                }
1581            }
1582
1583            impl<'r, 'a> Produce<'r, Option<$t>> for PostgresSimpleSourceParser {
1584                type Error = PostgresSourceError;
1585
1586                #[throws(PostgresSourceError)]
1587                fn produce(&'r mut self) -> Option<$t> {
1588                    let (ridx, cidx) = self.next_loc()?;
1589                    let val = match &self.rows[ridx] {
1590                        SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1591                            Some(s) => Some(
1592                                s.parse()
1593                                    .map_err(|_| ConnectorXError::cannot_produce::<$t>(Some(s.into())))?,
1594                            ),
1595                            None => None,
1596                        },
1597                        SimpleQueryMessage::CommandComplete(c) => {
1598                            panic!("get command: {}", c);
1599                        }
1600                        _ => {
1601                            panic!("what?");
1602                        }
1603                    };
1604                    val
1605                }
1606            }
1607        )+
1608    };
1609}
1610
1611impl_simple_produce!(i8, i16, i32, i64, u32, f32, f64, Uuid, IpInet,);
1612
1613impl<'r> Produce<'r, bool> for PostgresSimpleSourceParser {
1614    type Error = PostgresSourceError;
1615
1616    #[throws(PostgresSourceError)]
1617    fn produce(&'r mut self) -> bool {
1618        let (ridx, cidx) = self.next_loc()?;
1619        let val = match &self.rows[ridx] {
1620            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1621                Some(s) => match s {
1622                    "t" => true,
1623                    "f" => false,
1624                    _ => throw!(ConnectorXError::cannot_produce::<bool>(Some(s.into()))),
1625                },
1626                None => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
1627            },
1628            SimpleQueryMessage::CommandComplete(c) => {
1629                panic!("get command: {}", c);
1630            }
1631            _ => {
1632                panic!("what?");
1633            }
1634        };
1635        val
1636    }
1637}
1638
1639impl<'r> Produce<'r, Option<bool>> for PostgresSimpleSourceParser {
1640    type Error = PostgresSourceError;
1641
1642    #[throws(PostgresSourceError)]
1643    fn produce(&'r mut self) -> Option<bool> {
1644        let (ridx, cidx) = self.next_loc()?;
1645        let val = match &self.rows[ridx] {
1646            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1647                Some(s) => match s {
1648                    "t" => Some(true),
1649                    "f" => Some(false),
1650                    _ => throw!(ConnectorXError::cannot_produce::<bool>(Some(s.into()))),
1651                },
1652                None => None,
1653            },
1654            SimpleQueryMessage::CommandComplete(c) => {
1655                panic!("get command: {}", c);
1656            }
1657            _ => {
1658                panic!("what?");
1659            }
1660        };
1661        val
1662    }
1663}
1664
1665impl<'r> Produce<'r, Decimal> for PostgresSimpleSourceParser {
1666    type Error = PostgresSourceError;
1667
1668    #[throws(PostgresSourceError)]
1669    fn produce(&'r mut self) -> Decimal {
1670        let (ridx, cidx) = self.next_loc()?;
1671        let val = match &self.rows[ridx] {
1672            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1673                Some("Infinity") => Decimal::MAX,
1674                Some("-Infinity") => Decimal::MIN,
1675                Some(s) => s
1676                    .parse()
1677                    .map_err(|_| ConnectorXError::cannot_produce::<Decimal>(Some(s.into())))?,
1678                None => throw!(anyhow!("Cannot parse NULL in NOT NULL column.")),
1679            },
1680            SimpleQueryMessage::CommandComplete(c) => {
1681                panic!("get command: {}", c);
1682            }
1683            _ => {
1684                panic!("what?");
1685            }
1686        };
1687        val
1688    }
1689}
1690
1691impl<'r> Produce<'r, Option<Decimal>> for PostgresSimpleSourceParser {
1692    type Error = PostgresSourceError;
1693
1694    #[throws(PostgresSourceError)]
1695    fn produce(&'r mut self) -> Option<Decimal> {
1696        let (ridx, cidx) = self.next_loc()?;
1697        let val = match &self.rows[ridx] {
1698            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1699                Some("Infinity") => Some(Decimal::MAX),
1700                Some("-Infinity") => Some(Decimal::MIN),
1701                Some(s) => Some(
1702                    s.parse()
1703                        .map_err(|_| ConnectorXError::cannot_produce::<Decimal>(Some(s.into())))?,
1704                ),
1705                None => None,
1706            },
1707            SimpleQueryMessage::CommandComplete(c) => {
1708                panic!("get command: {}", c);
1709            }
1710            _ => {
1711                panic!("what?");
1712            }
1713        };
1714        val
1715    }
1716}
1717
1718impl<'r> Produce<'r, &'r str> for PostgresSimpleSourceParser {
1719    type Error = PostgresSourceError;
1720
1721    #[throws(PostgresSourceError)]
1722    fn produce(&'r mut self) -> &'r str {
1723        let (ridx, cidx) = self.next_loc()?;
1724        let val = match &self.rows[ridx] {
1725            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1726                Some(s) => s,
1727                None => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
1728            },
1729            SimpleQueryMessage::CommandComplete(c) => {
1730                panic!("get command: {}", c);
1731            }
1732            _ => {
1733                panic!("what?");
1734            }
1735        };
1736        val
1737    }
1738}
1739
1740impl<'r> Produce<'r, Option<&'r str>> for PostgresSimpleSourceParser {
1741    type Error = PostgresSourceError;
1742
1743    #[throws(PostgresSourceError)]
1744    fn produce(&'r mut self) -> Option<&'r str> {
1745        let (ridx, cidx) = self.next_loc()?;
1746        let val = match &self.rows[ridx] {
1747            SimpleQueryMessage::Row(row) => row.try_get(cidx)?,
1748            SimpleQueryMessage::CommandComplete(c) => {
1749                panic!("get command: {}", c);
1750            }
1751            _ => {
1752                panic!("what?");
1753            }
1754        };
1755        val
1756    }
1757}
1758
1759impl<'r> Produce<'r, Vec<u8>> for PostgresSimpleSourceParser {
1760    type Error = PostgresSourceError;
1761
1762    #[throws(PostgresSourceError)]
1763    fn produce(&'r mut self) -> Vec<u8> {
1764        let (ridx, cidx) = self.next_loc()?;
1765        let val = match &self.rows[ridx] {
1766            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1767                Some(s) => {
1768                    let mut res = s.chars();
1769                    res.next();
1770                    res.next();
1771                    decode(
1772                        res.enumerate()
1773                            .fold(String::new(), |acc, (_i, c)| format!("{}{}", acc, c))
1774                            .chars()
1775                            .map(|c| c as u8)
1776                            .collect::<Vec<u8>>(),
1777                    )?
1778                }
1779                None => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
1780            },
1781            SimpleQueryMessage::CommandComplete(c) => {
1782                panic!("get command: {}", c);
1783            }
1784            _ => {
1785                panic!("what?");
1786            }
1787        };
1788        val
1789    }
1790}
1791
1792impl<'r> Produce<'r, Option<Vec<u8>>> for PostgresSimpleSourceParser {
1793    type Error = PostgresSourceError;
1794
1795    #[throws(PostgresSourceError)]
1796    fn produce(&'r mut self) -> Option<Vec<u8>> {
1797        let (ridx, cidx) = self.next_loc()?;
1798        let val = match &self.rows[ridx] {
1799            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1800                Some(s) => {
1801                    let mut res = s.chars();
1802                    res.next();
1803                    res.next();
1804                    Some(decode(
1805                        res.enumerate()
1806                            .fold(String::new(), |acc, (_i, c)| format!("{}{}", acc, c))
1807                            .chars()
1808                            .map(|c| c as u8)
1809                            .collect::<Vec<u8>>(),
1810                    )?)
1811                }
1812                None => None,
1813            },
1814            SimpleQueryMessage::CommandComplete(c) => {
1815                panic!("get command: {}", c);
1816            }
1817            _ => {
1818                panic!("what?");
1819            }
1820        };
1821        val
1822    }
1823}
1824
1825fn rem_first_and_last(value: &str) -> &str {
1826    let mut chars = value.chars();
1827    chars.next();
1828    chars.next_back();
1829    chars.as_str()
1830}
1831
1832macro_rules! impl_simple_vec_produce {
1833    ($($t: ty,)+) => {
1834        $(
1835            impl<'r> Produce<'r, Vec<Option<$t>>> for PostgresSimpleSourceParser {
1836                type Error = PostgresSourceError;
1837
1838                #[throws(PostgresSourceError)]
1839                fn produce(&'r mut self) -> Vec<Option<$t>> {
1840                    let (ridx, cidx) = self.next_loc()?;
1841                    let val = match &self.rows[ridx] {
1842                        SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1843                            Some(s) => match s{
1844                                "" => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
1845                                "{}" => vec![],
1846                                _ => rem_first_and_last(s).split(",").map(|v| {
1847                                    if v == "NULL" {
1848                                        Ok(None)
1849                                    } else {
1850                                        match v.parse() {
1851                                            Ok(v) => Ok(Some(v)),
1852                                            Err(e) => Err(e).map_err(|_| ConnectorXError::cannot_produce::<Vec<$t>>(Some(s.into())))
1853                                        }
1854                                    }
1855                                }).collect::<Result<Vec<Option<$t>>, ConnectorXError>>()?
1856                            },
1857                            None => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
1858                        },
1859                        SimpleQueryMessage::CommandComplete(c) => {
1860                            panic!("get command: {}", c);
1861                        }
1862                        _ => {
1863                            panic!("what?");
1864                        }
1865                    };
1866                    val
1867                }
1868            }
1869
1870            impl<'r, 'a> Produce<'r, Option<Vec<Option<$t>>>> for PostgresSimpleSourceParser {
1871                type Error = PostgresSourceError;
1872
1873                #[throws(PostgresSourceError)]
1874                fn produce(&'r mut self) -> Option<Vec<Option<$t>>> {
1875                    let (ridx, cidx) = self.next_loc()?;
1876                    let val = match &self.rows[ridx] {
1877
1878                        SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1879                            Some(s) => match s{
1880                                "" => None,
1881                                "{}" => Some(vec![]),
1882                                _ => Some(rem_first_and_last(s).split(",").map(|v| {
1883                                    if v == "NULL" {
1884                                        Ok(None)
1885                                    } else {
1886                                        match v.parse() {
1887                                            Ok(v) => Ok(Some(v)),
1888                                            Err(e) => Err(e).map_err(|_| ConnectorXError::cannot_produce::<Vec<$t>>(Some(s.into())))
1889                                        }
1890                                    }
1891                                }).collect::<Result<Vec<Option<$t>>, ConnectorXError>>()?)
1892                            },
1893                            None => None,
1894                        },
1895
1896                        SimpleQueryMessage::CommandComplete(c) => {
1897                            panic!("get command: {}", c);
1898                        }
1899                        _ => {
1900                            panic!("what?");
1901                        }
1902                    };
1903                    val
1904                }
1905            }
1906        )+
1907    };
1908}
1909impl_simple_vec_produce!(i16, i32, i64, f32, f64, Decimal, String,);
1910
1911impl<'r> Produce<'r, Vec<Option<bool>>> for PostgresSimpleSourceParser {
1912    type Error = PostgresSourceError;
1913
1914    #[throws(PostgresSourceError)]
1915    fn produce(&'r mut self) -> Vec<Option<bool>> {
1916        let (ridx, cidx) = self.next_loc()?;
1917        let val = match &self.rows[ridx] {
1918            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1919                Some(s) => match s {
1920                    "" => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
1921                    "{}" => vec![],
1922                    _ => rem_first_and_last(s)
1923                        .split(',')
1924                        .map(|token| match token {
1925                            "NULL" => Ok(None),
1926                            "t" => Ok(Some(true)),
1927                            "f" => Ok(Some(false)),
1928                            _ => {
1929                                throw!(ConnectorXError::cannot_produce::<Vec<bool>>(Some(s.into())))
1930                            }
1931                        })
1932                        .collect::<Result<Vec<Option<bool>>, ConnectorXError>>()?,
1933                },
1934                None => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
1935            },
1936            SimpleQueryMessage::CommandComplete(c) => {
1937                panic!("get command: {}", c);
1938            }
1939            _ => {
1940                panic!("what?");
1941            }
1942        };
1943        val
1944    }
1945}
1946
1947impl<'r> Produce<'r, Option<Vec<Option<bool>>>> for PostgresSimpleSourceParser {
1948    type Error = PostgresSourceError;
1949
1950    #[throws(PostgresSourceError)]
1951    fn produce(&'r mut self) -> Option<Vec<Option<bool>>> {
1952        let (ridx, cidx) = self.next_loc()?;
1953        let val = match &self.rows[ridx] {
1954            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1955                Some(s) => match s {
1956                    "" => None,
1957                    "{}" => Some(vec![]),
1958                    _ => Some(
1959                        rem_first_and_last(s)
1960                            .split(',')
1961                            .map(|token| match token {
1962                                "NULL" => Ok(None),
1963                                "t" => Ok(Some(true)),
1964                                "f" => Ok(Some(false)),
1965                                _ => {
1966                                    throw!(ConnectorXError::cannot_produce::<Vec<bool>>(Some(
1967                                        s.into()
1968                                    )))
1969                                }
1970                            })
1971                            .collect::<Result<Vec<Option<bool>>, ConnectorXError>>()?,
1972                    ),
1973                },
1974                None => None,
1975            },
1976            SimpleQueryMessage::CommandComplete(c) => {
1977                panic!("get command: {}", c);
1978            }
1979            _ => {
1980                panic!("what?");
1981            }
1982        };
1983        val
1984    }
1985}
1986
1987impl<'r> Produce<'r, NaiveDate> for PostgresSimpleSourceParser {
1988    type Error = PostgresSourceError;
1989
1990    #[throws(PostgresSourceError)]
1991    fn produce(&'r mut self) -> NaiveDate {
1992        let (ridx, cidx) = self.next_loc()?;
1993        let val = match &self.rows[ridx] {
1994            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
1995                Some(s) => match s {
1996                    "infinity" => NaiveDate::MAX,
1997                    "-infinity" => NaiveDate::MIN,
1998                    s => NaiveDate::parse_from_str(s, "%Y-%m-%d").map_err(|_| {
1999                        ConnectorXError::cannot_produce::<NaiveDate>(Some(s.into()))
2000                    })?,
2001                },
2002                None => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
2003            },
2004            SimpleQueryMessage::CommandComplete(c) => {
2005                panic!("get command: {}", c);
2006            }
2007            _ => {
2008                panic!("what?");
2009            }
2010        };
2011        val
2012    }
2013}
2014
2015impl<'r> Produce<'r, Option<NaiveDate>> for PostgresSimpleSourceParser {
2016    type Error = PostgresSourceError;
2017
2018    #[throws(PostgresSourceError)]
2019    fn produce(&'r mut self) -> Option<NaiveDate> {
2020        let (ridx, cidx) = self.next_loc()?;
2021        let val = match &self.rows[ridx] {
2022            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
2023                Some(s) => match s {
2024                    "infinity" => Some(NaiveDate::MAX),
2025                    "-infinity" => Some(NaiveDate::MIN),
2026                    s => Some(NaiveDate::parse_from_str(s, "%Y-%m-%d").map_err(|_| {
2027                        ConnectorXError::cannot_produce::<Option<NaiveDate>>(Some(s.into()))
2028                    })?),
2029                },
2030                None => None,
2031            },
2032            SimpleQueryMessage::CommandComplete(c) => {
2033                panic!("get command: {}", c);
2034            }
2035            _ => {
2036                panic!("what?");
2037            }
2038        };
2039        val
2040    }
2041}
2042
2043impl<'r> Produce<'r, NaiveTime> for PostgresSimpleSourceParser {
2044    type Error = PostgresSourceError;
2045
2046    #[throws(PostgresSourceError)]
2047    fn produce(&'r mut self) -> NaiveTime {
2048        let (ridx, cidx) = self.next_loc()?;
2049        let val = match &self.rows[ridx] {
2050            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
2051                Some(s) => NaiveTime::parse_from_str(s, "%H:%M:%S%.f")
2052                    .map_err(|_| ConnectorXError::cannot_produce::<NaiveTime>(Some(s.into())))?,
2053                None => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
2054            },
2055            SimpleQueryMessage::CommandComplete(c) => {
2056                panic!("get command: {}", c);
2057            }
2058            _ => {
2059                panic!("what?");
2060            }
2061        };
2062        val
2063    }
2064}
2065
2066impl<'r> Produce<'r, Option<NaiveTime>> for PostgresSimpleSourceParser {
2067    type Error = PostgresSourceError;
2068
2069    #[throws(PostgresSourceError)]
2070    fn produce(&'r mut self) -> Option<NaiveTime> {
2071        let (ridx, cidx) = self.next_loc()?;
2072        let val = match &self.rows[ridx] {
2073            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
2074                Some(s) => Some(NaiveTime::parse_from_str(s, "%H:%M:%S%.f").map_err(|_| {
2075                    ConnectorXError::cannot_produce::<Option<NaiveTime>>(Some(s.into()))
2076                })?),
2077                None => None,
2078            },
2079            SimpleQueryMessage::CommandComplete(c) => {
2080                panic!("get command: {}", c);
2081            }
2082            _ => {
2083                panic!("what?");
2084            }
2085        };
2086        val
2087    }
2088}
2089
2090impl<'r> Produce<'r, NaiveDateTime> for PostgresSimpleSourceParser {
2091    type Error = PostgresSourceError;
2092
2093    #[throws(PostgresSourceError)]
2094    fn produce(&'r mut self) -> NaiveDateTime {
2095        let (ridx, cidx) = self.next_loc()?;
2096        let val =
2097            match &self.rows[ridx] {
2098                SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
2099                    Some(s) => match s {
2100                        "infinity" => NaiveDateTime::MAX,
2101                        "-infinity" => NaiveDateTime::MIN,
2102                        s => NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S%.f").map_err(
2103                            |_| ConnectorXError::cannot_produce::<NaiveDateTime>(Some(s.into())),
2104                        )?,
2105                    },
2106                    None => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
2107                },
2108                SimpleQueryMessage::CommandComplete(c) => {
2109                    panic!("get command: {}", c);
2110                }
2111                _ => {
2112                    panic!("what?");
2113                }
2114            };
2115        val
2116    }
2117}
2118
2119impl<'r> Produce<'r, Option<NaiveDateTime>> for PostgresSimpleSourceParser {
2120    type Error = PostgresSourceError;
2121
2122    #[throws(PostgresSourceError)]
2123    fn produce(&'r mut self) -> Option<NaiveDateTime> {
2124        let (ridx, cidx) = self.next_loc()?;
2125        let val = match &self.rows[ridx] {
2126            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
2127                Some(s) => match s {
2128                    "infinity" => Some(NaiveDateTime::MAX),
2129                    "-infinity" => Some(NaiveDateTime::MIN),
2130                    s => Some(
2131                        NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S%.f").map_err(|_| {
2132                            ConnectorXError::cannot_produce::<Option<NaiveDateTime>>(Some(s.into()))
2133                        })?,
2134                    ),
2135                },
2136                None => None,
2137            },
2138            SimpleQueryMessage::CommandComplete(c) => {
2139                panic!("get command: {}", c);
2140            }
2141            _ => {
2142                panic!("what?");
2143            }
2144        };
2145        val
2146    }
2147}
2148
2149impl<'r> Produce<'r, DateTime<Utc>> for PostgresSimpleSourceParser {
2150    type Error = PostgresSourceError;
2151
2152    #[throws(PostgresSourceError)]
2153    fn produce(&'r mut self) -> DateTime<Utc> {
2154        let (ridx, cidx) = self.next_loc()?;
2155        let val = match &self.rows[ridx] {
2156            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
2157                Some("infinity") => DateTime::<Utc>::MAX_UTC,
2158                Some("-infinity") => DateTime::<Utc>::MIN_UTC,
2159                Some(s) => {
2160                    let time_string = format!("{}:00", s).to_owned();
2161                    let slice: &str = &time_string[..];
2162                    let time: DateTime<FixedOffset> =
2163                        DateTime::parse_from_str(slice, "%Y-%m-%d %H:%M:%S%.f%:z").unwrap();
2164
2165                    time.with_timezone(&Utc)
2166                }
2167                None => throw!(anyhow!("Cannot parse NULL in non-NULL column.")),
2168            },
2169            SimpleQueryMessage::CommandComplete(c) => {
2170                panic!("get command: {}", c);
2171            }
2172            _ => {
2173                panic!("what?");
2174            }
2175        };
2176        val
2177    }
2178}
2179
2180impl<'r> Produce<'r, Option<DateTime<Utc>>> for PostgresSimpleSourceParser {
2181    type Error = PostgresSourceError;
2182
2183    #[throws(PostgresSourceError)]
2184    fn produce(&'r mut self) -> Option<DateTime<Utc>> {
2185        let (ridx, cidx) = self.next_loc()?;
2186        let val = match &self.rows[ridx] {
2187            SimpleQueryMessage::Row(row) => match row.try_get(cidx)? {
2188                Some("infinity") => Some(DateTime::<Utc>::MAX_UTC),
2189                Some("-infinity") => Some(DateTime::<Utc>::MIN_UTC),
2190                Some(s) => {
2191                    let time_string = format!("{}:00", s).to_owned();
2192                    let slice: &str = &time_string[..];
2193                    let time: DateTime<FixedOffset> =
2194                        DateTime::parse_from_str(slice, "%Y-%m-%d %H:%M:%S%.f%:z").unwrap();
2195
2196                    Some(time.with_timezone(&Utc))
2197                }
2198                None => None,
2199            },
2200            SimpleQueryMessage::CommandComplete(c) => {
2201                panic!("get command: {}", c);
2202            }
2203            _ => {
2204                panic!("what?");
2205            }
2206        };
2207        val
2208    }
2209}