1mod 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
41pub enum BinaryProtocol {}
43
44pub enum CSVProtocol {}
46
47pub enum CursorProtocol {}
49
50pub 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
97fn 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
208fn 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)?; 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)?; 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![])?; 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 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 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 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 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..])? }
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 "" => 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())?; 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) }
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}