1use std::collections::{BTreeMap, BTreeSet};
11use std::net::IpAddr;
12use std::pin::Pin;
13use std::sync::Arc;
14use std::time::Duration;
15
16use chrono::{DateTime, Utc};
17use derivative::Derivative;
18use enum_kinds::EnumKind;
19use futures::Stream;
20use mz_adapter_types::connection::{ConnectionId, ConnectionIdType};
21use mz_auth::password::Password;
22use mz_cluster_client::ReplicaId;
23use mz_compute_types::ComputeInstanceId;
24use mz_compute_types::dataflows::DataflowDescription;
25use mz_controller_types::ClusterId;
26use mz_expr::RowSetFinishing;
27use mz_ore::collections::CollectionExt;
28use mz_ore::soft_assert_no_log;
29use mz_ore::tracing::OpenTelemetryContext;
30use mz_persist_client::PersistClient;
31use mz_pgcopy::CopyFormatParams;
32use mz_repr::global_id::TransientIdGen;
33use mz_repr::role_id::RoleId;
34use mz_repr::{CatalogItemId, ColumnIndex, Diff, GlobalId, Row, RowIterator, SqlRelationType};
35use mz_sql::ast::{FetchDirection, Raw, Statement};
36use mz_sql::catalog::ObjectType;
37use mz_sql::optimizer_metrics::OptimizerMetrics;
38use mz_sql::plan;
39use mz_sql::plan::{ExecuteTimeout, Plan, PlanKind, SideEffectingFunc};
40use mz_sql::session::user::User;
41use mz_sql::session::vars::{OwnedVarInput, SystemVars};
42use mz_sql_parser::ast::{AlterObjectRenameStatement, AlterOwnerStatement, DropObjectsStatement};
43use mz_storage_types::sources::Timeline;
44use mz_timestamp_oracle::TimestampOracle;
45use tokio::sync::{Semaphore, mpsc, oneshot, watch};
46use uuid::Uuid;
47
48use crate::active_compute_sink::ActiveSubscribeOwner;
49use crate::catalog::Catalog;
50use crate::config::{ScopedParameters, ScopedParametersScope, SystemParameterFrontend};
51use crate::coord::appends::{BuiltinTableAppendNotify, WriteResult};
52use crate::coord::consistency::CoordinatorInconsistencies;
53use crate::coord::peek::{PeekDataflowPlan, PeekResponseUnary};
54use crate::coord::timestamp_selection::TimestampDetermination;
55use crate::coord::{ExecuteContextExtra, ExecuteContextGuard};
56use crate::error::AdapterError;
57use crate::optimize::LirDataflowDescription;
58use crate::session::{EndTransactionAction, RowBatchStream, Session};
59use crate::statement_logging::{
60 FrontendStatementLoggingEvent, StatementEndedExecutionReason, StatementExecutionStrategy,
61 StatementLoggingFrontend,
62};
63use crate::statement_logging::{StatementLoggingId, WatchSetCreation};
64use crate::util::Transmittable;
65use crate::webhook::AppendWebhookResponse;
66use crate::{
67 AdapterNotice, AppendWebhookError, CollectionIdBundle, ReadHolds, TimestampExplanation,
68};
69
70#[derive(Debug)]
73pub struct CopyFromStdinWriter {
74 pub batch_txs: Vec<mpsc::Sender<Vec<u8>>>,
77 pub completion_rx: oneshot::Receiver<
80 Result<(Vec<mz_persist_client::batch::ProtoBatch>, u64), crate::AdapterError>,
81 >,
82}
83
84#[derive(Debug)]
85pub struct CatalogSnapshot {
86 pub catalog: Arc<Catalog>,
87}
88
89#[derive(Debug)]
90pub enum Command {
91 CatalogSnapshot {
92 tx: oneshot::Sender<CatalogSnapshot>,
93 },
94
95 Startup {
96 tx: oneshot::Sender<Result<StartupResponse, AdapterError>>,
97 user: User,
98 conn_id: ConnectionId,
99 client_ip: Option<IpAddr>,
100 secret_key: u32,
101 uuid: Uuid,
102 application_name: String,
103 notice_tx: mpsc::UnboundedSender<AdapterNotice>,
104 },
105
106 AuthenticatePassword {
107 tx: oneshot::Sender<Result<(), AdapterError>>,
108 role_name: String,
109 password: Option<Password>,
110 },
111
112 AuthenticateGetSASLChallenge {
113 tx: oneshot::Sender<Result<SASLChallengeResponse, AdapterError>>,
114 role_name: String,
115 nonce: String,
116 },
117
118 AuthenticateVerifySASLProof {
119 tx: oneshot::Sender<Result<SASLVerifyProofResponse, AdapterError>>,
120 role_name: String,
121 proof: String,
122 auth_message: String,
123 mock_hash: String,
124 },
125
126 CheckRoleCanLogin {
127 tx: oneshot::Sender<Result<(), AdapterError>>,
128 role_name: String,
129 },
130
131 Execute {
132 portal_name: String,
133 session: Session,
134 tx: oneshot::Sender<Response<ExecuteResponse>>,
135 outer_ctx_extra: Option<ExecuteContextExtra>,
140 },
141
142 Commit {
147 action: EndTransactionAction,
148 session: Session,
149 tx: oneshot::Sender<Response<ExecuteResponse>>,
150 },
151
152 CancelRequest {
153 conn_id: ConnectionIdType,
154 secret_key: u32,
155 },
156
157 PrivilegedCancelRequest {
158 conn_id: ConnectionId,
159 },
160
161 GetWebhook {
162 database: String,
163 schema: String,
164 name: String,
165 tx: oneshot::Sender<Result<AppendWebhookResponse, AppendWebhookError>>,
166 },
167
168 GetSystemVars {
169 tx: oneshot::Sender<SystemVars>,
170 },
171
172 SetSystemVars {
173 vars: BTreeMap<String, String>,
174 conn_id: ConnectionId,
175 tx: oneshot::Sender<Result<(), AdapterError>>,
176 },
177
178 UpdateScopedSystemParameters {
185 overrides: ScopedParameters,
186 prune_scope: ScopedParametersScope,
189 tx: oneshot::Sender<()>,
190 },
191
192 InstallScopedSystemParameterFrontend {
199 frontend: Arc<SystemParameterFrontend>,
200 },
201
202 InjectAuditEvents {
203 events: Vec<crate::catalog::InjectedAuditEvent>,
204 conn_id: ConnectionId,
205 tx: oneshot::Sender<Result<(), AdapterError>>,
206 },
207
208 Terminate {
209 conn_id: ConnectionId,
210 tx: Option<oneshot::Sender<Result<(), AdapterError>>>,
211 },
212
213 StartCopyFromStdin {
217 target_id: CatalogItemId,
218 target_name: String,
219 columns: Vec<ColumnIndex>,
220 row_desc: mz_repr::RelationDesc,
222 params: mz_pgcopy::CopyFormatParams<'static>,
224 session: Session,
225 tx: oneshot::Sender<Response<CopyFromStdinWriter>>,
226 },
227
228 RetireExecute {
235 data: ExecuteContextExtra,
236 reason: StatementEndedExecutionReason,
237 },
238
239 CheckConsistency {
240 tx: oneshot::Sender<Result<(), CoordinatorInconsistencies>>,
241 },
242
243 Dump {
244 tx: oneshot::Sender<Result<serde_json::Value, anyhow::Error>>,
245 },
246
247 GetComputeInstanceClient {
248 instance_id: ComputeInstanceId,
249 tx: oneshot::Sender<
250 Result<
251 mz_compute_client::controller::instance_client::InstanceClient,
252 mz_compute_client::controller::error::InstanceMissing,
253 >,
254 >,
255 },
256
257 GetOracle {
258 timeline: Timeline,
259 tx: oneshot::Sender<
260 Result<Arc<dyn TimestampOracle<mz_repr::Timestamp> + Send + Sync>, AdapterError>,
261 >,
262 },
263
264 DetermineRealTimeRecentTimestamp {
265 source_ids: BTreeSet<GlobalId>,
266 real_time_recency_timeout: Duration,
267 tx: oneshot::Sender<Result<Option<mz_repr::Timestamp>, AdapterError>>,
268 },
269
270 GetTransactionReadHoldsBundle {
271 conn_id: ConnectionId,
272 tx: oneshot::Sender<Option<ReadHolds>>,
273 },
274
275 StoreTransactionReadHolds {
277 conn_id: ConnectionId,
278 read_holds: ReadHolds,
279 tx: oneshot::Sender<()>,
280 },
281
282 ExecuteSlowPathPeek {
283 dataflow_plan: Box<PeekDataflowPlan>,
284 determination: TimestampDetermination,
285 finishing: RowSetFinishing,
286 compute_instance: ComputeInstanceId,
287 target_replica: Option<ReplicaId>,
288 intermediate_result_type: SqlRelationType,
289 source_ids: BTreeSet<GlobalId>,
290 conn_id: ConnectionId,
291 max_result_size: u64,
292 max_query_result_size: Option<u64>,
293 watch_set: Option<WatchSetCreation>,
296 tx: oneshot::Sender<Result<ExecuteResponse, AdapterError>>,
297 },
298
299 ExecuteSubscribe {
300 df_desc: DataflowDescription<mz_compute_types::plan::LirRelationExpr>,
301 dependency_ids: BTreeSet<GlobalId>,
302 cluster_id: ComputeInstanceId,
303 replica_id: Option<ReplicaId>,
304 conn_id: ConnectionId,
305 session_uuid: Uuid,
306 read_holds: ReadHolds,
307 plan: plan::SubscribePlan,
308 statement_logging_id: Option<StatementLoggingId>,
309 tx: oneshot::Sender<Result<ExecuteResponse, AdapterError>>,
310 },
311
312 CopyToPreflight {
316 s3_sink_connection: mz_compute_types::sinks::CopyToS3OneshotSinkConnection,
318 sink_id: GlobalId,
320 tx: oneshot::Sender<Result<(), AdapterError>>,
322 },
323
324 ExecuteCopyTo {
325 df_desc: Box<DataflowDescription<mz_compute_types::plan::LirRelationExpr>>,
326 compute_instance: ComputeInstanceId,
327 target_replica: Option<ReplicaId>,
328 source_ids: BTreeSet<GlobalId>,
329 conn_id: ConnectionId,
330 watch_set: Option<WatchSetCreation>,
333 tx: oneshot::Sender<Result<ExecuteResponse, AdapterError>>,
334 },
335
336 ExecuteSideEffectingFunc {
338 plan: SideEffectingFunc,
339 conn_id: ConnectionId,
340 tx: oneshot::Sender<Result<ExecuteResponse, AdapterError>>,
341 },
342
343 LookupConnection {
353 connection_id: u32,
354 tx: oneshot::Sender<Option<(ConnectionId, RoleId)>>,
355 },
356
357 RegisterFrontendPeek {
361 uuid: Uuid,
362 conn_id: ConnectionId,
363 cluster_id: mz_controller_types::ClusterId,
364 depends_on: BTreeSet<GlobalId>,
365 is_fast_path: bool,
366 watch_set: Option<WatchSetCreation>,
369 tx: oneshot::Sender<Result<(), AdapterError>>,
370 },
371
372 UnregisterFrontendPeek {
381 uuid: Uuid,
382 reason: StatementEndedExecutionReason,
383 tx: oneshot::Sender<()>,
384 },
385
386 ExplainTimestamp {
389 conn_id: ConnectionId,
390 session_wall_time: DateTime<Utc>,
391 cluster_id: ClusterId,
392 id_bundle: CollectionIdBundle,
393 determination: TimestampDetermination,
394 tx: oneshot::Sender<TimestampExplanation>,
395 },
396
397 FrontendStatementLogging(FrontendStatementLoggingEvent),
400
401 RegisterConnectionCancelWatch {
407 conn_id: ConnectionId,
408 tx: oneshot::Sender<watch::Receiver<bool>>,
409 },
410
411 CreateInternalSubscribe {
416 df_desc: Box<LirDataflowDescription>,
417 cluster_id: ComputeInstanceId,
418 replica_id: Option<ReplicaId>,
419 depends_on: BTreeSet<GlobalId>,
420 as_of: mz_repr::Timestamp,
421 arity: usize,
422 sink_id: GlobalId,
423 owner: ActiveSubscribeOwner,
424 start_time: mz_ore::now::EpochMillis,
425 read_holds: ReadHolds,
426 tx: oneshot::Sender<Result<mpsc::UnboundedReceiver<PeekResponseUnary>, AdapterError>>,
427 },
428
429 AttemptWrite {
440 attempt: WriteAttemptKind,
441 target_id: CatalogItemId,
442 target_global_id: GlobalId,
443 diffs: Vec<(Row, Diff)>,
444 tx: oneshot::Sender<WriteResult>,
445 },
446
447 DropInternalSubscribe {
450 sink_id: GlobalId,
451 },
452}
453
454#[derive(Debug)]
459pub enum WriteAttemptKind {
460 Session {
464 conn_id: ConnectionId,
465 write_ts: Option<mz_repr::Timestamp>,
466 },
467 Background { write_ts: mz_repr::Timestamp },
470}
471
472impl Command {
473 pub fn session(&self) -> Option<&Session> {
474 match self {
475 Command::Execute { session, .. }
476 | Command::Commit { session, .. }
477 | Command::StartCopyFromStdin { session, .. } => Some(session),
478 Command::CancelRequest { .. }
479 | Command::Startup { .. }
480 | Command::AuthenticatePassword { .. }
481 | Command::AuthenticateGetSASLChallenge { .. }
482 | Command::AuthenticateVerifySASLProof { .. }
483 | Command::CheckRoleCanLogin { .. }
484 | Command::CatalogSnapshot { .. }
485 | Command::PrivilegedCancelRequest { .. }
486 | Command::GetWebhook { .. }
487 | Command::Terminate { .. }
488 | Command::GetSystemVars { .. }
489 | Command::SetSystemVars { .. }
490 | Command::UpdateScopedSystemParameters { .. }
491 | Command::InstallScopedSystemParameterFrontend { .. }
492 | Command::RetireExecute { .. }
493 | Command::CheckConsistency { .. }
494 | Command::Dump { .. }
495 | Command::GetComputeInstanceClient { .. }
496 | Command::GetOracle { .. }
497 | Command::DetermineRealTimeRecentTimestamp { .. }
498 | Command::GetTransactionReadHoldsBundle { .. }
499 | Command::StoreTransactionReadHolds { .. }
500 | Command::ExecuteSlowPathPeek { .. }
501 | Command::ExecuteSubscribe { .. }
502 | Command::CopyToPreflight { .. }
503 | Command::ExecuteCopyTo { .. }
504 | Command::ExecuteSideEffectingFunc { .. }
505 | Command::LookupConnection { .. }
506 | Command::RegisterFrontendPeek { .. }
507 | Command::UnregisterFrontendPeek { .. }
508 | Command::ExplainTimestamp { .. }
509 | Command::FrontendStatementLogging(..)
510 | Command::InjectAuditEvents { .. }
511 | Command::RegisterConnectionCancelWatch { .. }
512 | Command::CreateInternalSubscribe { .. }
513 | Command::AttemptWrite { .. }
514 | Command::DropInternalSubscribe { .. } => None,
515 }
516 }
517
518 pub fn session_mut(&mut self) -> Option<&mut Session> {
519 match self {
520 Command::Execute { session, .. }
521 | Command::Commit { session, .. }
522 | Command::StartCopyFromStdin { session, .. } => Some(session),
523 Command::CancelRequest { .. }
524 | Command::Startup { .. }
525 | Command::AuthenticatePassword { .. }
526 | Command::AuthenticateGetSASLChallenge { .. }
527 | Command::AuthenticateVerifySASLProof { .. }
528 | Command::CheckRoleCanLogin { .. }
529 | Command::CatalogSnapshot { .. }
530 | Command::PrivilegedCancelRequest { .. }
531 | Command::GetWebhook { .. }
532 | Command::Terminate { .. }
533 | Command::GetSystemVars { .. }
534 | Command::SetSystemVars { .. }
535 | Command::UpdateScopedSystemParameters { .. }
536 | Command::InstallScopedSystemParameterFrontend { .. }
537 | Command::RetireExecute { .. }
538 | Command::CheckConsistency { .. }
539 | Command::Dump { .. }
540 | Command::GetComputeInstanceClient { .. }
541 | Command::GetOracle { .. }
542 | Command::DetermineRealTimeRecentTimestamp { .. }
543 | Command::GetTransactionReadHoldsBundle { .. }
544 | Command::StoreTransactionReadHolds { .. }
545 | Command::ExecuteSlowPathPeek { .. }
546 | Command::ExecuteSubscribe { .. }
547 | Command::CopyToPreflight { .. }
548 | Command::ExecuteCopyTo { .. }
549 | Command::ExecuteSideEffectingFunc { .. }
550 | Command::LookupConnection { .. }
551 | Command::RegisterFrontendPeek { .. }
552 | Command::UnregisterFrontendPeek { .. }
553 | Command::ExplainTimestamp { .. }
554 | Command::FrontendStatementLogging(..)
555 | Command::InjectAuditEvents { .. }
556 | Command::RegisterConnectionCancelWatch { .. }
557 | Command::CreateInternalSubscribe { .. }
558 | Command::AttemptWrite { .. }
559 | Command::DropInternalSubscribe { .. } => None,
560 }
561 }
562}
563
564#[derive(Debug)]
565pub struct Response<T> {
566 pub result: Result<T, AdapterError>,
567 pub session: Session,
568 pub otel_ctx: OpenTelemetryContext,
569}
570
571#[derive(Debug, Clone, Copy)]
572pub struct SuperuserAttribute(pub Option<bool>);
573
574#[derive(Derivative)]
576#[derivative(Debug)]
577pub struct StartupResponse {
578 pub role_id: RoleId,
580 pub superuser_attribute: SuperuserAttribute,
585 #[derivative(Debug = "ignore")]
587 pub write_notify: BuiltinTableAppendNotify,
588 pub session_defaults: BTreeMap<String, OwnedVarInput>,
590 pub catalog: Arc<Catalog>,
591 pub storage_collections:
592 Arc<dyn mz_storage_client::storage_collections::StorageCollections + Send + Sync>,
593 pub transient_id_gen: Arc<TransientIdGen>,
594 pub optimizer_metrics: OptimizerMetrics,
595 pub persist_client: PersistClient,
596 pub statement_logging_frontend: StatementLoggingFrontend,
597 pub occ_write_semaphore: Arc<Semaphore>,
600 pub frontend_read_then_write_enabled: bool,
603 pub group_commit_notifier: crate::coord::appends::GroupCommitNotifier,
606 pub read_only: bool,
609}
610
611#[derive(Derivative)]
612#[derivative(Debug)]
613pub struct SASLChallengeResponse {
614 pub iteration_count: usize,
615 pub salt: String,
617 pub nonce: String,
618}
619
620#[derive(Derivative)]
621#[derivative(Debug)]
622pub struct SASLVerifyProofResponse {
623 pub verifier: String,
624}
625
626impl Transmittable for StartupResponse {
629 type Allowed = bool;
630 fn to_allowed(&self) -> Self::Allowed {
631 true
632 }
633}
634
635#[derive(Debug, Clone)]
637pub struct CatalogDump(String);
638
639impl CatalogDump {
640 pub fn new(raw: String) -> Self {
641 CatalogDump(raw)
642 }
643
644 pub fn into_string(self) -> String {
645 self.0
646 }
647}
648
649impl Transmittable for CatalogDump {
650 type Allowed = bool;
651 fn to_allowed(&self) -> Self::Allowed {
652 true
653 }
654}
655
656impl Transmittable for SystemVars {
657 type Allowed = bool;
658 fn to_allowed(&self) -> Self::Allowed {
659 true
660 }
661}
662
663#[derive(EnumKind, Derivative)]
665#[derivative(Debug)]
666#[enum_kind(ExecuteResponseKind, derive(PartialOrd, Ord))]
667pub enum ExecuteResponse {
668 AlteredDefaultPrivileges,
670 AlteredObject(ObjectType),
672 AlteredRole,
674 AlteredSystemConfiguration,
676 ClosedCursor,
678 Comment,
680 Copied(usize),
682 CopyTo {
684 format: mz_sql::plan::CopyFormat,
685 resp: Box<ExecuteResponse>,
686 },
687 CopyFrom {
688 target_id: CatalogItemId,
690 target_name: String,
692 columns: Vec<ColumnIndex>,
693 params: CopyFormatParams<'static>,
694 ctx_extra: ExecuteContextGuard,
695 },
696 CreatedConnection,
698 CreatedDatabase,
700 CreatedSchema,
702 CreatedRole,
704 CreatedCluster,
706 CreatedClusterReplica,
708 CreatedIndex,
710 CreatedMetricSink,
712 CreatedIntrospectionSubscribe,
714 CreatedSecret,
716 CreatedSink,
718 CreatedSource,
720 CreatedTable,
722 CreatedView,
724 CreatedViews,
726 CreatedMaterializedView,
728 CreatedType,
730 CreatedNetworkPolicy,
732 Deallocate { all: bool },
734 DeclaredCursor,
736 Deleted(usize),
738 DiscardedTemp,
740 DiscardedAll {
742 params: BTreeMap<&'static str, String>,
744 },
745 DroppedObject(ObjectType),
747 DroppedOwned,
749 EmptyQuery,
751 Fetch {
753 name: String,
755 count: Option<FetchDirection>,
757 timeout: ExecuteTimeout,
759 ctx_extra: ExecuteContextGuard,
760 },
761 GrantedPrivilege,
763 GrantedRole,
765 Inserted(usize),
767 Prepare,
769 Raised,
771 ReassignOwned,
773 RevokedPrivilege,
775 RevokedRole,
777 SendingRowsStreaming {
779 #[derivative(Debug = "ignore")]
780 rows: Pin<Box<dyn Stream<Item = PeekResponseUnary> + Send + Sync>>,
781 instance_id: ComputeInstanceId,
782 strategy: StatementExecutionStrategy,
783 },
784 SendingRowsImmediate {
787 #[derivative(Debug = "ignore")]
788 rows: Box<dyn RowIterator + Send + Sync>,
789 },
790 SetVariable {
792 name: String,
793 reset: bool,
795 },
796 StartedTransaction,
798 Subscribing {
801 #[derivative(Debug = "ignore")]
802 rx: RowBatchStream,
803 ctx_extra: ExecuteContextGuard,
804 instance_id: ComputeInstanceId,
805 },
806 TransactionCommitted {
808 params: BTreeMap<&'static str, String>,
810 },
811 TransactionRolledBack {
813 params: BTreeMap<&'static str, String>,
815 },
816 Updated(usize),
818 ValidatedConnection,
820}
821
822impl TryFrom<&Statement<Raw>> for ExecuteResponse {
823 type Error = ();
824
825 fn try_from(stmt: &Statement<Raw>) -> Result<Self, Self::Error> {
827 let resp_kinds = Plan::generated_from(&stmt.into())
828 .iter()
829 .map(ExecuteResponse::generated_from)
830 .flatten()
831 .cloned()
832 .collect::<BTreeSet<ExecuteResponseKind>>();
833 let resps = resp_kinds
834 .iter()
835 .map(|r| (*r).try_into())
836 .collect::<Result<Vec<ExecuteResponse>, _>>();
837 if let Ok(resps) = resps {
839 if resps.len() == 1 {
840 return Ok(resps.into_element());
841 }
842 }
843 let resp = match stmt {
844 Statement::DropObjects(DropObjectsStatement { object_type, .. }) => {
845 ExecuteResponse::DroppedObject((*object_type).into())
846 }
847 Statement::AlterObjectRename(AlterObjectRenameStatement { object_type, .. })
848 | Statement::AlterOwner(AlterOwnerStatement { object_type, .. }) => {
849 ExecuteResponse::AlteredObject((*object_type).into())
850 }
851 _ => return Err(()),
852 };
853 soft_assert_no_log!(
855 resp_kinds.len() == 1
856 && resp_kinds.first().expect("must exist") == &ExecuteResponseKind::from(&resp),
857 "ExecuteResponses out of sync with planner"
858 );
859 Ok(resp)
860 }
861}
862
863impl TryInto<ExecuteResponse> for ExecuteResponseKind {
864 type Error = ();
865
866 fn try_into(self) -> Result<ExecuteResponse, Self::Error> {
869 match self {
870 ExecuteResponseKind::AlteredDefaultPrivileges => {
871 Ok(ExecuteResponse::AlteredDefaultPrivileges)
872 }
873 ExecuteResponseKind::AlteredObject => Err(()),
874 ExecuteResponseKind::AlteredRole => Ok(ExecuteResponse::AlteredRole),
875 ExecuteResponseKind::AlteredSystemConfiguration => {
876 Ok(ExecuteResponse::AlteredSystemConfiguration)
877 }
878 ExecuteResponseKind::ClosedCursor => Ok(ExecuteResponse::ClosedCursor),
879 ExecuteResponseKind::Comment => Ok(ExecuteResponse::Comment),
880 ExecuteResponseKind::Copied => Err(()),
881 ExecuteResponseKind::CopyTo => Err(()),
882 ExecuteResponseKind::CopyFrom => Err(()),
883 ExecuteResponseKind::CreatedConnection => Ok(ExecuteResponse::CreatedConnection),
884 ExecuteResponseKind::CreatedDatabase => Ok(ExecuteResponse::CreatedDatabase),
885 ExecuteResponseKind::CreatedSchema => Ok(ExecuteResponse::CreatedSchema),
886 ExecuteResponseKind::CreatedRole => Ok(ExecuteResponse::CreatedRole),
887 ExecuteResponseKind::CreatedCluster => Ok(ExecuteResponse::CreatedCluster),
888 ExecuteResponseKind::CreatedClusterReplica => {
889 Ok(ExecuteResponse::CreatedClusterReplica)
890 }
891 ExecuteResponseKind::CreatedIndex => Ok(ExecuteResponse::CreatedIndex),
892 ExecuteResponseKind::CreatedMetricSink => Ok(ExecuteResponse::CreatedMetricSink),
893 ExecuteResponseKind::CreatedSecret => Ok(ExecuteResponse::CreatedSecret),
894 ExecuteResponseKind::CreatedSink => Ok(ExecuteResponse::CreatedSink),
895 ExecuteResponseKind::CreatedSource => Ok(ExecuteResponse::CreatedSource),
896 ExecuteResponseKind::CreatedTable => Ok(ExecuteResponse::CreatedTable),
897 ExecuteResponseKind::CreatedView => Ok(ExecuteResponse::CreatedView),
898 ExecuteResponseKind::CreatedViews => Ok(ExecuteResponse::CreatedViews),
899 ExecuteResponseKind::CreatedMaterializedView => {
900 Ok(ExecuteResponse::CreatedMaterializedView)
901 }
902 ExecuteResponseKind::CreatedNetworkPolicy => Ok(ExecuteResponse::CreatedNetworkPolicy),
903 ExecuteResponseKind::CreatedType => Ok(ExecuteResponse::CreatedType),
904 ExecuteResponseKind::Deallocate => Err(()),
905 ExecuteResponseKind::DeclaredCursor => Ok(ExecuteResponse::DeclaredCursor),
906 ExecuteResponseKind::Deleted => Err(()),
907 ExecuteResponseKind::DiscardedTemp => Ok(ExecuteResponse::DiscardedTemp),
908 ExecuteResponseKind::DiscardedAll => Err(()),
909 ExecuteResponseKind::DroppedObject => Err(()),
910 ExecuteResponseKind::DroppedOwned => Ok(ExecuteResponse::DroppedOwned),
911 ExecuteResponseKind::EmptyQuery => Ok(ExecuteResponse::EmptyQuery),
912 ExecuteResponseKind::Fetch => Err(()),
913 ExecuteResponseKind::GrantedPrivilege => Ok(ExecuteResponse::GrantedPrivilege),
914 ExecuteResponseKind::GrantedRole => Ok(ExecuteResponse::GrantedRole),
915 ExecuteResponseKind::Inserted => Err(()),
916 ExecuteResponseKind::Prepare => Ok(ExecuteResponse::Prepare),
917 ExecuteResponseKind::Raised => Ok(ExecuteResponse::Raised),
918 ExecuteResponseKind::ReassignOwned => Ok(ExecuteResponse::ReassignOwned),
919 ExecuteResponseKind::RevokedPrivilege => Ok(ExecuteResponse::RevokedPrivilege),
920 ExecuteResponseKind::RevokedRole => Ok(ExecuteResponse::RevokedRole),
921 ExecuteResponseKind::SetVariable => Err(()),
922 ExecuteResponseKind::StartedTransaction => Ok(ExecuteResponse::StartedTransaction),
923 ExecuteResponseKind::Subscribing => Err(()),
924 ExecuteResponseKind::TransactionCommitted => Err(()),
925 ExecuteResponseKind::TransactionRolledBack => Err(()),
926 ExecuteResponseKind::Updated => Err(()),
927 ExecuteResponseKind::ValidatedConnection => Ok(ExecuteResponse::ValidatedConnection),
928 ExecuteResponseKind::SendingRowsStreaming => Err(()),
929 ExecuteResponseKind::SendingRowsImmediate => Err(()),
930 ExecuteResponseKind::CreatedIntrospectionSubscribe => {
931 Ok(ExecuteResponse::CreatedIntrospectionSubscribe)
932 }
933 }
934 }
935}
936
937impl ExecuteResponse {
938 pub fn tag(&self) -> Option<String> {
939 use ExecuteResponse::*;
940 match self {
941 AlteredDefaultPrivileges => Some("ALTER DEFAULT PRIVILEGES".into()),
942 AlteredObject(o) => Some(format!("ALTER {}", o)),
943 AlteredRole => Some("ALTER ROLE".into()),
944 AlteredSystemConfiguration => Some("ALTER SYSTEM".into()),
945 ClosedCursor => Some("CLOSE CURSOR".into()),
946 Comment => Some("COMMENT".into()),
947 Copied(n) => Some(format!("COPY {}", n)),
948 CopyTo { .. } => None,
949 CopyFrom { .. } => None,
950 CreatedConnection { .. } => Some("CREATE CONNECTION".into()),
951 CreatedDatabase { .. } => Some("CREATE DATABASE".into()),
952 CreatedSchema { .. } => Some("CREATE SCHEMA".into()),
953 CreatedRole => Some("CREATE ROLE".into()),
954 CreatedCluster { .. } => Some("CREATE CLUSTER".into()),
955 CreatedClusterReplica { .. } => Some("CREATE CLUSTER REPLICA".into()),
956 CreatedIndex { .. } => Some("CREATE INDEX".into()),
957 CreatedMetricSink { .. } => Some("CREATE METRIC SINK".into()),
958 CreatedSecret { .. } => Some("CREATE SECRET".into()),
959 CreatedSink { .. } => Some("CREATE SINK".into()),
960 CreatedSource { .. } => Some("CREATE SOURCE".into()),
961 CreatedTable { .. } => Some("CREATE TABLE".into()),
962 CreatedView { .. } => Some("CREATE VIEW".into()),
963 CreatedViews { .. } => Some("CREATE VIEWS".into()),
964 CreatedMaterializedView { .. } => Some("CREATE MATERIALIZED VIEW".into()),
965 CreatedType => Some("CREATE TYPE".into()),
966 CreatedNetworkPolicy => Some("CREATE NETWORKPOLICY".into()),
967 Deallocate { all } => Some(format!("DEALLOCATE{}", if *all { " ALL" } else { "" })),
968 DeclaredCursor => Some("DECLARE CURSOR".into()),
969 Deleted(n) => Some(format!("DELETE {}", n)),
970 DiscardedTemp => Some("DISCARD TEMP".into()),
971 DiscardedAll { .. } => Some("DISCARD ALL".into()),
972 DroppedObject(o) => Some(format!("DROP {o}")),
973 DroppedOwned => Some("DROP OWNED".into()),
974 EmptyQuery => None,
975 Fetch { .. } => None,
976 GrantedPrivilege => Some("GRANT".into()),
977 GrantedRole => Some("GRANT ROLE".into()),
978 Inserted(n) => {
979 Some(format!("INSERT 0 {}", n))
987 }
988 Prepare => Some("PREPARE".into()),
989 Raised => Some("RAISE".into()),
990 ReassignOwned => Some("REASSIGN OWNED".into()),
991 RevokedPrivilege => Some("REVOKE".into()),
992 RevokedRole => Some("REVOKE ROLE".into()),
993 SendingRowsStreaming { .. } | SendingRowsImmediate { .. } => None,
994 SetVariable { reset: true, .. } => Some("RESET".into()),
995 SetVariable { reset: false, .. } => Some("SET".into()),
996 StartedTransaction { .. } => Some("BEGIN".into()),
997 Subscribing { .. } => None,
998 TransactionCommitted { .. } => Some("COMMIT".into()),
999 TransactionRolledBack { .. } => Some("ROLLBACK".into()),
1000 Updated(n) => Some(format!("UPDATE {}", n)),
1001 ValidatedConnection => Some("VALIDATE CONNECTION".into()),
1002 CreatedIntrospectionSubscribe => Some("CREATE INTROSPECTION SUBSCRIBE".into()),
1003 }
1004 }
1005
1006 pub fn generated_from(plan: &PlanKind) -> &'static [ExecuteResponseKind] {
1010 use ExecuteResponseKind::*;
1011 use PlanKind::*;
1012
1013 match plan {
1014 AbortTransaction => &[TransactionRolledBack],
1015 AlterClusterRename
1016 | AlterClusterSwap
1017 | AlterCluster
1018 | AlterClusterReplicaRename
1019 | AlterOwner
1020 | AlterItemRename
1021 | AlterRetainHistory
1022 | AlterSourceTimestampInterval
1023 | AlterNoop
1024 | AlterSchemaRename
1025 | AlterSchemaSwap
1026 | AlterSecret
1027 | AlterConnection
1028 | AlterSource
1029 | AlterSink
1030 | AlterTableAddColumn
1031 | AlterMaterializedViewApplyReplacement
1032 | AlterNetworkPolicy => &[AlteredObject],
1033 AlterDefaultPrivileges => &[AlteredDefaultPrivileges],
1034 AlterSetCluster => &[AlteredObject],
1035 AlterRole => &[AlteredRole],
1036 AlterSystemSet | AlterSystemReset | AlterSystemResetAll => {
1037 &[AlteredSystemConfiguration]
1038 }
1039 Close => &[ClosedCursor],
1040 PlanKind::CopyFrom => &[ExecuteResponseKind::CopyFrom, ExecuteResponseKind::Copied],
1041 PlanKind::CopyTo => &[ExecuteResponseKind::Copied],
1042 PlanKind::Comment => &[ExecuteResponseKind::Comment],
1043 CommitTransaction => &[TransactionCommitted, TransactionRolledBack],
1044 CreateConnection => &[CreatedConnection],
1045 CreateDatabase => &[CreatedDatabase],
1046 CreateSchema => &[CreatedSchema],
1047 CreateRole => &[CreatedRole],
1048 CreateCluster => &[CreatedCluster],
1049 CreateClusterReplica => &[CreatedClusterReplica],
1050 CreateSource | CreateSources => &[CreatedSource],
1051 CreateSecret => &[CreatedSecret],
1052 CreateSink => &[CreatedSink],
1053 CreateTable => &[CreatedTable],
1054 CreateView => &[CreatedView],
1055 CreateMaterializedView => &[CreatedMaterializedView],
1056 CreateIndex => &[CreatedIndex],
1057 CreateMetricSink => &[CreatedMetricSink],
1058 CreateType => &[CreatedType],
1059 PlanKind::Deallocate => &[ExecuteResponseKind::Deallocate],
1060 CreateNetworkPolicy => &[CreatedNetworkPolicy],
1061 Declare => &[DeclaredCursor],
1062 DiscardTemp => &[DiscardedTemp],
1063 DiscardAll => &[DiscardedAll],
1064 DropObjects => &[DroppedObject],
1065 DropOwned => &[DroppedOwned],
1066 PlanKind::EmptyQuery => &[ExecuteResponseKind::EmptyQuery],
1067 ExplainPlan | ExplainPushdown | ExplainTimestamp | Select | ShowAllVariables
1068 | ShowCreate | ShowColumns | ShowVariable | InspectShard | ExplainSinkSchema => &[
1069 ExecuteResponseKind::CopyTo,
1070 SendingRowsStreaming,
1071 SendingRowsImmediate,
1072 ],
1073 Execute | ReadThenWrite => &[
1074 Deleted,
1075 Inserted,
1076 SendingRowsStreaming,
1077 SendingRowsImmediate,
1078 Updated,
1079 ],
1080 PlanKind::Fetch => &[ExecuteResponseKind::Fetch],
1081 GrantPrivileges => &[GrantedPrivilege],
1082 GrantRole => &[GrantedRole],
1083 Insert => &[Inserted, SendingRowsImmediate],
1084 PlanKind::Prepare => &[ExecuteResponseKind::Prepare],
1085 PlanKind::Raise => &[ExecuteResponseKind::Raised],
1086 PlanKind::ReassignOwned => &[ExecuteResponseKind::ReassignOwned],
1087 RevokePrivileges => &[RevokedPrivilege],
1088 RevokeRole => &[RevokedRole],
1089 PlanKind::SetVariable | ResetVariable | PlanKind::SetTransaction => {
1090 &[ExecuteResponseKind::SetVariable]
1091 }
1092 PlanKind::Subscribe => &[Subscribing, ExecuteResponseKind::CopyTo],
1093 StartTransaction => &[StartedTransaction],
1094 SideEffectingFunc => &[SendingRowsStreaming, SendingRowsImmediate],
1095 ValidateConnection => &[ExecuteResponseKind::ValidatedConnection],
1096 }
1097 }
1098}
1099
1100impl Transmittable for ExecuteResponse {
1104 type Allowed = ExecuteResponseKind;
1105 fn to_allowed(&self) -> Self::Allowed {
1106 ExecuteResponseKind::from(self)
1107 }
1108}