Skip to main content

ahri_tre_pgmeta/
execution.rs

1use crate::{PgMetaError, PgMetadataConnection};
2use ahri_tre_libpq_oauth::{LibpqOAuthConnection, LibpqOAuthParameter, LibpqOAuthRow};
3use bytes::BytesMut;
4use postgres::types::{IsNull, ToSql, Type};
5use postgres::{Client, Row, Transaction};
6
7/// Parameters accepted by PostgreSQL metadata executor methods.
8pub type MetadataParameters<'a> = &'a [&'a (dyn ToSql + Sync)];
9
10/// Transaction-scoped metadata execution interface.
11pub trait MetadataTransaction {
12    /// Row representation returned by the concrete executor.
13    type Row;
14
15    /// Runs a query expected to return exactly one row.
16    ///
17    /// # Errors
18    ///
19    /// Returns an error when PostgreSQL/libpq rejects the query or the result
20    /// shape does not match the executor contract.
21    fn query_one(
22        &mut self,
23        operation: &'static str,
24        statement: &str,
25        params: MetadataParameters<'_>,
26    ) -> Result<Self::Row, PgMetaError>;
27
28    /// Runs a query expected to return zero or one row.
29    ///
30    /// # Errors
31    ///
32    /// Returns an error when PostgreSQL/libpq rejects the query or the result
33    /// shape contains more than one row.
34    fn query_optional(
35        &mut self,
36        operation: &'static str,
37        statement: &str,
38        params: MetadataParameters<'_>,
39    ) -> Result<Option<Self::Row>, PgMetaError>;
40
41    /// Runs a query that can return many rows.
42    ///
43    /// # Errors
44    ///
45    /// Returns an error when PostgreSQL/libpq rejects the query or parameters.
46    fn query_many(
47        &mut self,
48        operation: &'static str,
49        statement: &str,
50        params: MetadataParameters<'_>,
51    ) -> Result<Vec<Self::Row>, PgMetaError>;
52
53    /// Executes a metadata write command.
54    ///
55    /// # Errors
56    ///
57    /// Returns an error when PostgreSQL/libpq rejects the command or
58    /// parameters.
59    fn execute_command(
60        &mut self,
61        operation: &'static str,
62        statement: &str,
63        params: MetadataParameters<'_>,
64    ) -> Result<u64, PgMetaError>;
65}
66
67/// Session-scoped metadata execution interface.
68pub trait MetadataExecutor {
69    /// Row representation returned by the concrete executor.
70    type Row;
71
72    /// Runs a query expected to return exactly one row.
73    ///
74    /// # Errors
75    ///
76    /// Returns an error when PostgreSQL/libpq rejects the query or the result
77    /// shape does not match the executor contract.
78    fn query_one(
79        &mut self,
80        operation: &'static str,
81        statement: &str,
82        params: MetadataParameters<'_>,
83    ) -> Result<Self::Row, PgMetaError>;
84
85    /// Runs a query expected to return zero or one row.
86    ///
87    /// # Errors
88    ///
89    /// Returns an error when PostgreSQL/libpq rejects the query or the result
90    /// shape contains more than one row.
91    fn query_optional(
92        &mut self,
93        operation: &'static str,
94        statement: &str,
95        params: MetadataParameters<'_>,
96    ) -> Result<Option<Self::Row>, PgMetaError>;
97
98    /// Runs a query that can return many rows.
99    ///
100    /// # Errors
101    ///
102    /// Returns an error when PostgreSQL/libpq rejects the query or parameters.
103    fn query_many(
104        &mut self,
105        operation: &'static str,
106        statement: &str,
107        params: MetadataParameters<'_>,
108    ) -> Result<Vec<Self::Row>, PgMetaError>;
109
110    /// Executes a metadata write command.
111    ///
112    /// # Errors
113    ///
114    /// Returns an error when PostgreSQL/libpq rejects the command or
115    /// parameters.
116    fn execute_command(
117        &mut self,
118        operation: &'static str,
119        statement: &str,
120        params: MetadataParameters<'_>,
121    ) -> Result<u64, PgMetaError>;
122
123    /// Runs caller-supplied work inside a metadata transaction.
124    ///
125    /// # Errors
126    ///
127    /// Returns an error if the transaction cannot be opened or committed, or if
128    /// the callback returns an error.
129    fn with_transaction<T>(
130        &mut self,
131        operation: &'static str,
132        f: impl FnOnce(&mut dyn MetadataTransaction<Row = Self::Row>) -> Result<T, PgMetaError>,
133    ) -> Result<T, PgMetaError>;
134}
135
136pub(crate) struct PgClientExecutor<'a> {
137    client: &'a mut Client,
138}
139
140impl<'a> PgClientExecutor<'a> {
141    pub(crate) fn new(client: &'a mut Client) -> Self {
142        Self { client }
143    }
144}
145
146impl MetadataExecutor for PgClientExecutor<'_> {
147    type Row = Row;
148
149    fn query_one(
150        &mut self,
151        _operation: &'static str,
152        statement: &str,
153        params: MetadataParameters<'_>,
154    ) -> Result<Self::Row, PgMetaError> {
155        self.client
156            .query_one(statement, params)
157            .map_err(PgMetaError::Query)
158    }
159
160    fn query_optional(
161        &mut self,
162        _operation: &'static str,
163        statement: &str,
164        params: MetadataParameters<'_>,
165    ) -> Result<Option<Self::Row>, PgMetaError> {
166        self.client
167            .query_opt(statement, params)
168            .map_err(PgMetaError::Query)
169    }
170
171    fn query_many(
172        &mut self,
173        _operation: &'static str,
174        statement: &str,
175        params: MetadataParameters<'_>,
176    ) -> Result<Vec<Self::Row>, PgMetaError> {
177        self.client
178            .query(statement, params)
179            .map_err(PgMetaError::Query)
180    }
181
182    fn execute_command(
183        &mut self,
184        operation: &'static str,
185        statement: &str,
186        params: MetadataParameters<'_>,
187    ) -> Result<u64, PgMetaError> {
188        self.client
189            .execute(statement, params)
190            .map_err(|source| PgMetaError::Write { operation, source })
191    }
192
193    fn with_transaction<T>(
194        &mut self,
195        operation: &'static str,
196        f: impl FnOnce(&mut dyn MetadataTransaction<Row = Self::Row>) -> Result<T, PgMetaError>,
197    ) -> Result<T, PgMetaError> {
198        let mut transaction = self
199            .client
200            .transaction()
201            .map_err(PgMetaError::Transaction)?;
202        let result = {
203            let mut executor = PgTransactionExecutor {
204                transaction: &mut transaction,
205            };
206            f(&mut executor)?
207        };
208        transaction
209            .commit()
210            .map_err(|source| PgMetaError::Write { operation, source })?;
211        Ok(result)
212    }
213}
214
215impl MetadataExecutor for PgMetadataConnection {
216    type Row = Row;
217
218    fn query_one(
219        &mut self,
220        operation: &'static str,
221        statement: &str,
222        params: MetadataParameters<'_>,
223    ) -> Result<Self::Row, PgMetaError> {
224        let result =
225            PgClientExecutor::new(&mut self.client).query_one(operation, statement, params);
226        if let Err(error) = &result {
227            self.observe_error(error);
228        }
229        result
230    }
231
232    fn query_optional(
233        &mut self,
234        operation: &'static str,
235        statement: &str,
236        params: MetadataParameters<'_>,
237    ) -> Result<Option<Self::Row>, PgMetaError> {
238        let result =
239            PgClientExecutor::new(&mut self.client).query_optional(operation, statement, params);
240        if let Err(error) = &result {
241            self.observe_error(error);
242        }
243        result
244    }
245
246    fn query_many(
247        &mut self,
248        operation: &'static str,
249        statement: &str,
250        params: MetadataParameters<'_>,
251    ) -> Result<Vec<Self::Row>, PgMetaError> {
252        let result =
253            PgClientExecutor::new(&mut self.client).query_many(operation, statement, params);
254        if let Err(error) = &result {
255            self.observe_error(error);
256        }
257        result
258    }
259
260    fn execute_command(
261        &mut self,
262        operation: &'static str,
263        statement: &str,
264        params: MetadataParameters<'_>,
265    ) -> Result<u64, PgMetaError> {
266        let result =
267            PgClientExecutor::new(&mut self.client).execute_command(operation, statement, params);
268        if let Err(error) = &result {
269            self.observe_error(error);
270        }
271        result
272    }
273
274    fn with_transaction<T>(
275        &mut self,
276        operation: &'static str,
277        f: impl FnOnce(&mut dyn MetadataTransaction<Row = Self::Row>) -> Result<T, PgMetaError>,
278    ) -> Result<T, PgMetaError> {
279        if !self.admission_active {
280            return PgClientExecutor::new(&mut self.client).with_transaction(operation, f);
281        }
282        self.client
283            .batch_execute("SAVEPOINT dataset_metadata_helper")
284            .map_err(PgMetaError::Transaction)?;
285        let result = f(&mut PgAdmissionTransaction {
286            client: &mut self.client,
287        });
288        match result {
289            Ok(value) => {
290                self.client
291                    .batch_execute("RELEASE SAVEPOINT dataset_metadata_helper")
292                    .map_err(PgMetaError::Transaction)?;
293                Ok(value)
294            }
295            Err(error) => {
296                let _ = self.client.batch_execute("ROLLBACK TO SAVEPOINT dataset_metadata_helper; RELEASE SAVEPOINT dataset_metadata_helper");
297                Err(error)
298            }
299        }
300    }
301}
302
303impl MetadataExecutor for LibpqOAuthConnection {
304    type Row = LibpqOAuthRow;
305
306    fn query_one(
307        &mut self,
308        operation: &'static str,
309        statement: &str,
310        params: MetadataParameters<'_>,
311    ) -> Result<Self::Row, PgMetaError> {
312        let params = encode_libpq_parameters(operation, params)?;
313        self.query_one_raw(statement, &params)
314            .map_err(|source| PgMetaError::OAuth { operation, source })
315    }
316
317    fn query_optional(
318        &mut self,
319        operation: &'static str,
320        statement: &str,
321        params: MetadataParameters<'_>,
322    ) -> Result<Option<Self::Row>, PgMetaError> {
323        let params = encode_libpq_parameters(operation, params)?;
324        self.query_optional_raw(statement, &params)
325            .map_err(|source| PgMetaError::OAuth { operation, source })
326    }
327
328    fn query_many(
329        &mut self,
330        operation: &'static str,
331        statement: &str,
332        params: MetadataParameters<'_>,
333    ) -> Result<Vec<Self::Row>, PgMetaError> {
334        let params = encode_libpq_parameters(operation, params)?;
335        self.query_raw(statement, &params)
336            .map_err(|source| PgMetaError::OAuth { operation, source })
337    }
338
339    fn execute_command(
340        &mut self,
341        operation: &'static str,
342        statement: &str,
343        params: MetadataParameters<'_>,
344    ) -> Result<u64, PgMetaError> {
345        let params = encode_libpq_parameters(operation, params)?;
346        self.execute_raw(statement, &params)
347            .map_err(|source| PgMetaError::OAuth { operation, source })
348    }
349
350    fn with_transaction<T>(
351        &mut self,
352        operation: &'static str,
353        f: impl FnOnce(&mut dyn MetadataTransaction<Row = Self::Row>) -> Result<T, PgMetaError>,
354    ) -> Result<T, PgMetaError> {
355        let nested = self.is_in_transaction();
356        self.execute_raw(
357            if nested {
358                "SAVEPOINT dataset_metadata_helper"
359            } else {
360                "BEGIN"
361            },
362            &[],
363        )
364        .map_err(|source| PgMetaError::OAuth { operation, source })?;
365
366        let result = {
367            let mut transaction = LibpqOAuthTransaction { connection: self };
368            f(&mut transaction)
369        };
370
371        match result {
372            Ok(value) => {
373                if let Err(source) = self.execute_raw(
374                    if nested {
375                        "RELEASE SAVEPOINT dataset_metadata_helper"
376                    } else {
377                        "COMMIT"
378                    },
379                    &[],
380                ) {
381                    // A pre-dispatch deadline/cancel rejection leaves our own
382                    // transaction active. An uncertain lost COMMIT reply does
383                    // not report an active transaction and is not compensated here.
384                    if self.is_in_transaction() {
385                        let _ = self.execute_raw(if nested {
386                            "ROLLBACK TO SAVEPOINT dataset_metadata_helper; RELEASE SAVEPOINT dataset_metadata_helper"
387                        } else { "ROLLBACK" }, &[]);
388                    }
389                    return Err(PgMetaError::OAuth { operation, source });
390                }
391                Ok(value)
392            }
393            Err(error) => {
394                let _ = self.execute_raw(if nested {
395                    "ROLLBACK TO SAVEPOINT dataset_metadata_helper; RELEASE SAVEPOINT dataset_metadata_helper"
396                } else { "ROLLBACK" }, &[]);
397                Err(error)
398            }
399        }
400    }
401}
402
403pub(crate) struct PgTransactionExecutor<'tx, 'conn> {
404    transaction: &'tx mut Transaction<'conn>,
405}
406
407impl MetadataTransaction for PgTransactionExecutor<'_, '_> {
408    type Row = Row;
409
410    fn query_one(
411        &mut self,
412        _operation: &'static str,
413        statement: &str,
414        params: MetadataParameters<'_>,
415    ) -> Result<Self::Row, PgMetaError> {
416        self.transaction
417            .query_one(statement, params)
418            .map_err(PgMetaError::Query)
419    }
420
421    fn query_optional(
422        &mut self,
423        _operation: &'static str,
424        statement: &str,
425        params: MetadataParameters<'_>,
426    ) -> Result<Option<Self::Row>, PgMetaError> {
427        self.transaction
428            .query_opt(statement, params)
429            .map_err(PgMetaError::Query)
430    }
431
432    fn query_many(
433        &mut self,
434        _operation: &'static str,
435        statement: &str,
436        params: MetadataParameters<'_>,
437    ) -> Result<Vec<Self::Row>, PgMetaError> {
438        self.transaction
439            .query(statement, params)
440            .map_err(PgMetaError::Query)
441    }
442
443    fn execute_command(
444        &mut self,
445        operation: &'static str,
446        statement: &str,
447        params: MetadataParameters<'_>,
448    ) -> Result<u64, PgMetaError> {
449        self.transaction
450            .execute(statement, params)
451            .map_err(|source| PgMetaError::Write { operation, source })
452    }
453}
454
455struct LibpqOAuthTransaction<'a> {
456    connection: &'a mut LibpqOAuthConnection,
457}
458
459impl MetadataTransaction for LibpqOAuthTransaction<'_> {
460    type Row = LibpqOAuthRow;
461
462    fn query_one(
463        &mut self,
464        operation: &'static str,
465        statement: &str,
466        params: MetadataParameters<'_>,
467    ) -> Result<Self::Row, PgMetaError> {
468        let params = encode_libpq_parameters(operation, params)?;
469        self.connection
470            .query_one_raw(statement, &params)
471            .map_err(|source| PgMetaError::OAuth { operation, source })
472    }
473
474    fn query_optional(
475        &mut self,
476        operation: &'static str,
477        statement: &str,
478        params: MetadataParameters<'_>,
479    ) -> Result<Option<Self::Row>, PgMetaError> {
480        let params = encode_libpq_parameters(operation, params)?;
481        self.connection
482            .query_optional_raw(statement, &params)
483            .map_err(|source| PgMetaError::OAuth { operation, source })
484    }
485
486    fn query_many(
487        &mut self,
488        operation: &'static str,
489        statement: &str,
490        params: MetadataParameters<'_>,
491    ) -> Result<Vec<Self::Row>, PgMetaError> {
492        let params = encode_libpq_parameters(operation, params)?;
493        self.connection
494            .query_raw(statement, &params)
495            .map_err(|source| PgMetaError::OAuth { operation, source })
496    }
497
498    fn execute_command(
499        &mut self,
500        operation: &'static str,
501        statement: &str,
502        params: MetadataParameters<'_>,
503    ) -> Result<u64, PgMetaError> {
504        let params = encode_libpq_parameters(operation, params)?;
505        self.connection
506            .execute_raw(statement, &params)
507            .map_err(|source| PgMetaError::OAuth { operation, source })
508    }
509}
510
511fn encode_libpq_parameters(
512    operation: &'static str,
513    params: MetadataParameters<'_>,
514) -> Result<Vec<LibpqOAuthParameter>, PgMetaError> {
515    params
516        .iter()
517        .map(|param| encode_libpq_parameter(operation, *param))
518        .collect()
519}
520
521fn encode_libpq_parameter(
522    operation: &'static str,
523    param: &(dyn ToSql + Sync),
524) -> Result<LibpqOAuthParameter, PgMetaError> {
525    for ty in libpq_parameter_candidate_types() {
526        let mut value = BytesMut::new();
527        match param.to_sql_checked(ty, &mut value) {
528            Ok(IsNull::No) => {
529                return Ok(LibpqOAuthParameter::binary(ty.oid(), Some(value.to_vec())));
530            }
531            Ok(IsNull::Yes) => return Ok(LibpqOAuthParameter::binary(ty.oid(), None)),
532            Err(_) => {}
533        }
534    }
535
536    Err(PgMetaError::Decode {
537        field: "metadata parameter",
538        value: "<redacted>".to_string(),
539        message: format!("unsupported parameter type during {operation}"),
540    })
541}
542
543fn libpq_parameter_candidate_types() -> &'static [Type] {
544    &[
545        Type::BOOL,
546        Type::INT2,
547        Type::INT4,
548        Type::INT8,
549        Type::FLOAT4,
550        Type::FLOAT8,
551        Type::TEXT,
552        Type::VARCHAR,
553        Type::UUID,
554        Type::UUID_ARRAY,
555        Type::DATE,
556        Type::TIMESTAMP,
557        Type::TIMESTAMPTZ,
558        Type::BYTEA,
559        Type::JSON,
560        Type::JSONB,
561    ]
562}
563
564struct PgAdmissionTransaction<'a> {
565    client: &'a mut Client,
566}
567
568impl MetadataTransaction for PgAdmissionTransaction<'_> {
569    type Row = Row;
570    fn query_one(
571        &mut self,
572        operation: &'static str,
573        statement: &str,
574        params: MetadataParameters<'_>,
575    ) -> Result<Row, PgMetaError> {
576        PgClientExecutor::new(self.client).query_one(operation, statement, params)
577    }
578    fn query_optional(
579        &mut self,
580        operation: &'static str,
581        statement: &str,
582        params: MetadataParameters<'_>,
583    ) -> Result<Option<Row>, PgMetaError> {
584        PgClientExecutor::new(self.client).query_optional(operation, statement, params)
585    }
586    fn query_many(
587        &mut self,
588        operation: &'static str,
589        statement: &str,
590        params: MetadataParameters<'_>,
591    ) -> Result<Vec<Row>, PgMetaError> {
592        PgClientExecutor::new(self.client).query_many(operation, statement, params)
593    }
594    fn execute_command(
595        &mut self,
596        operation: &'static str,
597        statement: &str,
598        params: MetadataParameters<'_>,
599    ) -> Result<u64, PgMetaError> {
600        PgClientExecutor::new(self.client).execute_command(operation, statement, params)
601    }
602}
603
604#[cfg(test)]
605mod tests {
606    use super::{MetadataExecutor, MetadataParameters, MetadataTransaction};
607    use crate::PgMetaError;
608    use std::collections::VecDeque;
609
610    #[test]
611    fn oauth_parameters_preserve_governed_input_uuid_arrays() {
612        use postgres::types::{ToSql, Type};
613        let inputs = vec![uuid::Uuid::new_v4(), uuid::Uuid::new_v4()];
614        let encoded = super::encode_libpq_parameter("validate derivation", &inputs)
615            .expect("governed input references must cross the OAuth adapter");
616        let mut expected = bytes::BytesMut::new();
617        inputs
618            .to_sql_checked(&Type::UUID_ARRAY, &mut expected)
619            .unwrap();
620        assert_eq!(
621            encoded,
622            ahri_tre_libpq_oauth::LibpqOAuthParameter::binary(
623                Type::UUID_ARRAY.oid(),
624                Some(expected.to_vec())
625            )
626        );
627    }
628
629    #[derive(Debug, Clone, PartialEq, Eq)]
630    struct FakeRow(&'static str);
631
632    #[derive(Debug, Clone, PartialEq, Eq)]
633    enum FakeCall {
634        QueryOne(&'static str),
635        QueryOptional(&'static str),
636        QueryMany(&'static str),
637        Execute(&'static str),
638        Begin(&'static str),
639        Commit(&'static str),
640    }
641
642    #[derive(Debug, Default)]
643    struct FakeExecutor {
644        calls: Vec<FakeCall>,
645        one_rows: VecDeque<FakeRow>,
646        optional_rows: VecDeque<Option<FakeRow>>,
647        many_rows: VecDeque<Vec<FakeRow>>,
648        command_counts: VecDeque<u64>,
649        transaction: FakeTransaction,
650    }
651
652    impl FakeExecutor {
653        fn decode_error(message: impl Into<String>) -> PgMetaError {
654            PgMetaError::Decode {
655                field: "fake",
656                value: "metadata execution test".to_string(),
657                message: message.into(),
658            }
659        }
660    }
661
662    impl MetadataExecutor for FakeExecutor {
663        type Row = FakeRow;
664
665        fn query_one(
666            &mut self,
667            operation: &'static str,
668            _statement: &str,
669            _params: MetadataParameters<'_>,
670        ) -> Result<Self::Row, PgMetaError> {
671            self.calls.push(FakeCall::QueryOne(operation));
672            self.one_rows
673                .pop_front()
674                .ok_or_else(|| Self::decode_error("missing single row"))
675        }
676
677        fn query_optional(
678            &mut self,
679            operation: &'static str,
680            _statement: &str,
681            _params: MetadataParameters<'_>,
682        ) -> Result<Option<Self::Row>, PgMetaError> {
683            self.calls.push(FakeCall::QueryOptional(operation));
684            self.optional_rows
685                .pop_front()
686                .ok_or_else(|| Self::decode_error("missing optional row"))
687        }
688
689        fn query_many(
690            &mut self,
691            operation: &'static str,
692            _statement: &str,
693            _params: MetadataParameters<'_>,
694        ) -> Result<Vec<Self::Row>, PgMetaError> {
695            self.calls.push(FakeCall::QueryMany(operation));
696            self.many_rows
697                .pop_front()
698                .ok_or_else(|| Self::decode_error("missing row batch"))
699        }
700
701        fn execute_command(
702            &mut self,
703            operation: &'static str,
704            _statement: &str,
705            _params: MetadataParameters<'_>,
706        ) -> Result<u64, PgMetaError> {
707            self.calls.push(FakeCall::Execute(operation));
708            self.command_counts
709                .pop_front()
710                .ok_or_else(|| Self::decode_error("missing command count"))
711        }
712
713        fn with_transaction<T>(
714            &mut self,
715            operation: &'static str,
716            f: impl FnOnce(&mut dyn MetadataTransaction<Row = Self::Row>) -> Result<T, PgMetaError>,
717        ) -> Result<T, PgMetaError> {
718            self.calls.push(FakeCall::Begin(operation));
719            let result = f(&mut self.transaction)?;
720            self.calls.push(FakeCall::Commit(operation));
721            Ok(result)
722        }
723    }
724
725    #[derive(Debug, Default)]
726    struct FakeTransaction {
727        calls: Vec<FakeCall>,
728        one_rows: VecDeque<FakeRow>,
729        optional_rows: VecDeque<Option<FakeRow>>,
730        many_rows: VecDeque<Vec<FakeRow>>,
731        command_counts: VecDeque<u64>,
732    }
733
734    impl MetadataTransaction for FakeTransaction {
735        type Row = FakeRow;
736
737        fn query_one(
738            &mut self,
739            operation: &'static str,
740            _statement: &str,
741            _params: MetadataParameters<'_>,
742        ) -> Result<Self::Row, PgMetaError> {
743            self.calls.push(FakeCall::QueryOne(operation));
744            self.one_rows
745                .pop_front()
746                .ok_or_else(|| FakeExecutor::decode_error("missing transaction single row"))
747        }
748
749        fn query_optional(
750            &mut self,
751            operation: &'static str,
752            _statement: &str,
753            _params: MetadataParameters<'_>,
754        ) -> Result<Option<Self::Row>, PgMetaError> {
755            self.calls.push(FakeCall::QueryOptional(operation));
756            self.optional_rows
757                .pop_front()
758                .ok_or_else(|| FakeExecutor::decode_error("missing transaction optional row"))
759        }
760
761        fn query_many(
762            &mut self,
763            operation: &'static str,
764            _statement: &str,
765            _params: MetadataParameters<'_>,
766        ) -> Result<Vec<Self::Row>, PgMetaError> {
767            self.calls.push(FakeCall::QueryMany(operation));
768            self.many_rows
769                .pop_front()
770                .ok_or_else(|| FakeExecutor::decode_error("missing transaction row batch"))
771        }
772
773        fn execute_command(
774            &mut self,
775            operation: &'static str,
776            _statement: &str,
777            _params: MetadataParameters<'_>,
778        ) -> Result<u64, PgMetaError> {
779            self.calls.push(FakeCall::Execute(operation));
780            self.command_counts
781                .pop_front()
782                .ok_or_else(|| FakeExecutor::decode_error("missing transaction command count"))
783        }
784    }
785
786    #[test]
787    fn metadata_executor_covers_read_command_and_transaction_shapes() {
788        let mut executor = FakeExecutor::default();
789        executor.one_rows.push_back(FakeRow("single"));
790        executor.optional_rows.push_back(None);
791        executor
792            .many_rows
793            .push_back(vec![FakeRow("a"), FakeRow("b")]);
794        executor.command_counts.push_back(3);
795        executor
796            .transaction
797            .one_rows
798            .push_back(FakeRow("tx-single"));
799        executor
800            .transaction
801            .optional_rows
802            .push_back(Some(FakeRow("tx-optional")));
803        executor
804            .transaction
805            .many_rows
806            .push_back(vec![FakeRow("tx-a"), FakeRow("tx-b")]);
807        executor.transaction.command_counts.push_back(1);
808
809        assert_eq!(
810            executor
811                .query_one("single", "select 1", &[])
812                .expect("single row should be returned"),
813            FakeRow("single")
814        );
815        assert_eq!(
816            executor
817                .query_optional("optional", "select 1 where false", &[])
818                .expect("optional query should run"),
819            None
820        );
821        assert_eq!(
822            executor
823                .query_many("many", "select * from metadata", &[])
824                .expect("batch query should run"),
825            vec![FakeRow("a"), FakeRow("b")]
826        );
827        assert_eq!(
828            executor
829                .execute_command("command", "update metadata set touched = true", &[])
830                .expect("command should run"),
831            3
832        );
833
834        let transaction_result = executor
835            .with_transaction("transaction", |tx| {
836                let single = tx.query_one("tx-single", "select 1", &[])?;
837                let optional = tx.query_optional("tx-optional", "select 1", &[])?;
838                let many = tx.query_many("tx-many", "select 1", &[])?;
839                let updated = tx.execute_command("tx-command", "update metadata", &[])?;
840                Ok((single, optional, many, updated))
841            })
842            .expect("transaction should run");
843
844        assert_eq!(transaction_result.0, FakeRow("tx-single"));
845        assert_eq!(transaction_result.1, Some(FakeRow("tx-optional")));
846        assert_eq!(transaction_result.2, vec![FakeRow("tx-a"), FakeRow("tx-b")]);
847        assert_eq!(transaction_result.3, 1);
848        assert_eq!(
849            executor.calls,
850            vec![
851                FakeCall::QueryOne("single"),
852                FakeCall::QueryOptional("optional"),
853                FakeCall::QueryMany("many"),
854                FakeCall::Execute("command"),
855                FakeCall::Begin("transaction"),
856                FakeCall::Commit("transaction"),
857            ]
858        );
859        assert_eq!(
860            executor.transaction.calls,
861            vec![
862                FakeCall::QueryOne("tx-single"),
863                FakeCall::QueryOptional("tx-optional"),
864                FakeCall::QueryMany("tx-many"),
865                FakeCall::Execute("tx-command"),
866            ]
867        );
868    }
869
870    #[test]
871    fn metadata_executor_surfaces_common_error_shape() {
872        let mut executor = FakeExecutor::default();
873        let error = executor
874            .query_one("single", "select 1", &[])
875            .expect_err("missing fake row should fail");
876
877        assert!(matches!(error, PgMetaError::Decode { .. }));
878    }
879}