1use std::sync::Arc;
9
10use smol_str::SmolStr;
11
12use crate::capability::{Capability, CapabilitySet};
13use crate::errors::PluginError;
14use crate::plugin::PluginId;
15use crate::qname::QName;
16use crate::registry::PluginRegistry;
17use crate::surfaces::{
18 AggregateSurface, AlgorithmSurface, AppendReg, AuthSurface, AuthzSurface, BackgroundJobSurface,
19 CatalogSurface, CdcSurface, CollationSurface, CrdtSurface, DynPendingRegistration, HookSurface,
20 IndexKindSurface, KeyedUniqueReg, LabelStorageSurface, LocyAggregateSurface,
21 LocyGeneratorSurface, LocyPredicateSurface, LogicalTypeSurface, NamedUniqueReg,
22 OptimizerRuleSurface, ProcedureSurface, ReplacementScanSurface, ScalarSurface, TriggerSurface,
23 VersionedReg, WindowSurface,
24};
25use crate::traits::aggregate::{AggSignature, AggregatePluginFn};
26use crate::traits::algorithm::AlgorithmProvider;
27use crate::traits::background::BackgroundJobProvider;
28use crate::traits::catalog::{CatalogProvider, ReplacementScanProvider};
29use crate::traits::cdc::CdcOutputProvider;
30use crate::traits::collation::CollationProvider;
31use crate::traits::connector::{AuthProvider, AuthzPolicy};
32use crate::traits::crdt::{CrdtKind, CrdtKindProvider};
33use crate::traits::hook::SessionHook;
34use crate::traits::index::{IndexKind, IndexKindProvider};
35use crate::traits::locy::{
36 GenSignature, LocyAggregate, LocyGenerator, LocyPredicate, PredSignature,
37};
38use crate::traits::operator::OptimizerRuleProvider;
39use crate::traits::procedure::{ProcedurePlugin, ProcedureSignature};
40use crate::traits::scalar::{FnSignature, ScalarPluginFn};
41use crate::traits::trigger::TriggerPlugin;
42use crate::traits::types::LogicalTypeProvider;
43use crate::traits::window::{WindowPluginFn, WindowSignature};
44
45pub struct PluginRegistrar<'a> {
56 plugin_id: PluginId,
57 effective_caps: &'a CapabilitySet,
58 registry: &'a PluginRegistry,
59 pending: Vec<Box<dyn DynPendingRegistration>>,
60 aggregate_qnames: Vec<QName>,
67}
68
69impl<'a> std::fmt::Debug for PluginRegistrar<'a> {
70 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
71 f.debug_struct("PluginRegistrar")
72 .field("plugin_id", &self.plugin_id)
73 .field("pending", &self.pending.len())
74 .finish_non_exhaustive()
75 }
76}
77
78impl<'a> PluginRegistrar<'a> {
79 #[must_use]
84 pub fn new(
85 plugin_id: PluginId,
86 effective_caps: &'a CapabilitySet,
87 registry: &'a PluginRegistry,
88 ) -> Self {
89 Self {
90 plugin_id,
91 effective_caps,
92 registry,
93 pending: Vec::new(),
94 aggregate_qnames: Vec::new(),
95 }
96 }
97
98 #[must_use]
104 pub fn staged_aggregate_qnames(&self) -> &[QName] {
105 &self.aggregate_qnames
106 }
107
108 #[must_use]
110 pub fn plugin_id(&self) -> &PluginId {
111 &self.plugin_id
112 }
113
114 pub fn set_plugin_id(&mut self, plugin_id: PluginId) {
122 self.plugin_id = plugin_id;
123 }
124
125 fn require(&self, cap: &Capability) -> Result<(), PluginError> {
126 if self.effective_caps.contains_variant(cap) {
127 Ok(())
128 } else {
129 Err(PluginError::CapabilityRequired(cap.clone()))
130 }
131 }
132
133 fn validate_qname(&self, qname: &QName) -> Result<(), PluginError> {
134 if !qname.is_builtin() && qname.namespace() != self.plugin_id.as_str() {
135 return Err(PluginError::internal(format!(
136 "plugin `{}` cannot register qname `{}` outside its namespace",
137 self.plugin_id, qname
138 )));
139 }
140 Ok(())
141 }
142
143 pub fn scalar_fn(
151 &mut self,
152 qname: QName,
153 sig: FnSignature,
154 f: Arc<dyn ScalarPluginFn>,
155 ) -> Result<&mut Self, PluginError> {
156 self.require(&Capability::ScalarFn)?;
157 self.validate_qname(&qname)?;
158 self.pending.push(Box::new(NamedUniqueReg::<ScalarSurface> {
159 q: qname,
160 sig,
161 provider: f,
162 }));
163 Ok(self)
164 }
165
166 pub fn aggregate_fn(
172 &mut self,
173 qname: QName,
174 sig: AggSignature,
175 f: Arc<dyn AggregatePluginFn>,
176 ) -> Result<&mut Self, PluginError> {
177 self.require(&Capability::AggregateFn)?;
178 self.validate_qname(&qname)?;
179 self.aggregate_qnames.push(qname.clone());
180 self.pending
181 .push(Box::new(NamedUniqueReg::<AggregateSurface> {
182 q: qname,
183 sig,
184 provider: f,
185 }));
186 Ok(self)
187 }
188
189 pub fn window_fn(
195 &mut self,
196 qname: QName,
197 sig: WindowSignature,
198 f: Arc<dyn WindowPluginFn>,
199 ) -> Result<&mut Self, PluginError> {
200 self.require(&Capability::WindowFn)?;
201 self.validate_qname(&qname)?;
202 self.pending.push(Box::new(NamedUniqueReg::<WindowSurface> {
203 q: qname,
204 sig,
205 provider: f,
206 }));
207 Ok(self)
208 }
209
210 pub fn procedure(
217 &mut self,
218 qname: QName,
219 sig: ProcedureSignature,
220 p: Arc<dyn ProcedurePlugin>,
221 ) -> Result<&mut Self, PluginError> {
222 use crate::traits::procedure::ProcedureMode;
223 self.require(&Capability::Procedure)?;
224 match sig.mode {
225 ProcedureMode::Write => self.require(&Capability::ProcedureWrites)?,
226 ProcedureMode::Schema => self.require(&Capability::ProcedureSchema)?,
227 ProcedureMode::Dbms => self.require(&Capability::ProcedureDbms)?,
228 ProcedureMode::Read => {}
229 }
230 self.validate_qname(&qname)?;
231 self.pending
232 .push(Box::new(VersionedReg::<ProcedureSurface> {
233 q: qname,
234 sig,
235 provider: p,
236 }));
237 Ok(self)
238 }
239
240 pub fn locy_aggregate(
246 &mut self,
247 qname: QName,
248 a: Arc<dyn LocyAggregate>,
249 ) -> Result<&mut Self, PluginError> {
250 self.require(&Capability::LocyAggregate)?;
251 self.validate_qname(&qname)?;
252 self.pending
253 .push(Box::new(NamedUniqueReg::<LocyAggregateSurface> {
254 q: qname,
255 sig: (),
256 provider: a,
257 }));
258 Ok(self)
259 }
260
261 pub fn locy_predicate(
267 &mut self,
268 qname: QName,
269 sig: PredSignature,
270 p: Arc<dyn LocyPredicate>,
271 ) -> Result<&mut Self, PluginError> {
272 self.require(&Capability::LocyPredicate)?;
273 self.validate_qname(&qname)?;
274 self.pending
275 .push(Box::new(NamedUniqueReg::<LocyPredicateSurface> {
276 q: qname,
277 sig,
278 provider: p,
279 }));
280 Ok(self)
281 }
282
283 pub fn locy_generator(
289 &mut self,
290 qname: QName,
291 sig: GenSignature,
292 p: Arc<dyn LocyGenerator>,
293 ) -> Result<&mut Self, PluginError> {
294 self.require(&Capability::LocyGenerator)?;
295 self.validate_qname(&qname)?;
296 self.pending
297 .push(Box::new(NamedUniqueReg::<LocyGeneratorSurface> {
298 q: qname,
299 sig,
300 provider: p,
301 }));
302 Ok(self)
303 }
304
305 pub fn optimizer_rule(
311 &mut self,
312 r: Arc<dyn OptimizerRuleProvider>,
313 ) -> Result<&mut Self, PluginError> {
314 self.require(&Capability::Operator)?;
315 self.pending
316 .push(Box::new(AppendReg::<OptimizerRuleSurface> { provider: r }));
317 Ok(self)
318 }
319
320 pub fn index_kind(
326 &mut self,
327 kind: IndexKind,
328 p: Arc<dyn IndexKindProvider>,
329 ) -> Result<&mut Self, PluginError> {
330 self.require(&Capability::Index)?;
331 self.pending
332 .push(Box::new(KeyedUniqueReg::<IndexKindSurface> {
333 key_override: Some(kind),
334 provider: p,
335 }));
336 Ok(self)
337 }
338
339 pub fn label_storage(
349 &mut self,
350 label: impl Into<SmolStr>,
351 storage: Arc<dyn crate::traits::storage::Storage>,
352 ) -> Result<&mut Self, PluginError> {
353 self.require(&Capability::Storage)?;
354 self.pending
355 .push(Box::new(KeyedUniqueReg::<LabelStorageSurface> {
356 key_override: Some(label.into()),
357 provider: storage,
358 }));
359 Ok(self)
360 }
361
362 pub fn algorithm(
368 &mut self,
369 qname: QName,
370 p: Arc<dyn AlgorithmProvider>,
371 ) -> Result<&mut Self, PluginError> {
372 self.require(&Capability::Algorithm)?;
373 self.validate_qname(&qname)?;
374 p.signature()
378 .check_slices(crate::traits::algorithm::HOST_CAPABILITY_SLICES)
379 .map_err(|e| PluginError::SliceUnavailable(e.message))?;
380 let effective_caps = self.effective_caps.clone();
383 self.pending
384 .push(Box::new(NamedUniqueReg::<AlgorithmSurface> {
385 q: qname,
386 sig: effective_caps,
387 provider: p,
388 }));
389 Ok(self)
390 }
391
392 pub fn crdt_kind(
398 &mut self,
399 kind: CrdtKind,
400 p: Arc<dyn CrdtKindProvider>,
401 ) -> Result<&mut Self, PluginError> {
402 self.require(&Capability::Crdt)?;
403 self.pending.push(Box::new(KeyedUniqueReg::<CrdtSurface> {
404 key_override: Some(kind),
405 provider: p,
406 }));
407 Ok(self)
408 }
409
410 pub fn hook(&mut self, h: Arc<dyn SessionHook>) -> Result<&mut Self, PluginError> {
416 self.require(&Capability::Hook)?;
417 self.pending
418 .push(Box::new(AppendReg::<HookSurface> { provider: h }));
419 Ok(self)
420 }
421
422 pub fn logical_type(
428 &mut self,
429 t: Arc<dyn LogicalTypeProvider>,
430 ) -> Result<&mut Self, PluginError> {
431 self.require(&Capability::Type)?;
432 self.pending
433 .push(Box::new(KeyedUniqueReg::<LogicalTypeSurface> {
434 key_override: None,
435 provider: t,
436 }));
437 Ok(self)
438 }
439
440 pub fn auth_provider(&mut self, p: Arc<dyn AuthProvider>) -> Result<&mut Self, PluginError> {
446 self.require(&Capability::Auth)?;
447 self.pending
448 .push(Box::new(AppendReg::<AuthSurface> { provider: p }));
449 Ok(self)
450 }
451
452 pub fn authz_policy(&mut self, p: Arc<dyn AuthzPolicy>) -> Result<&mut Self, PluginError> {
458 self.require(&Capability::Authz)?;
459 self.pending
460 .push(Box::new(AppendReg::<AuthzSurface> { provider: p }));
461 Ok(self)
462 }
463
464 pub fn trigger(&mut self, t: Arc<dyn TriggerPlugin>) -> Result<&mut Self, PluginError> {
470 self.require(&Capability::Trigger)?;
471 self.pending
472 .push(Box::new(AppendReg::<TriggerSurface> { provider: t }));
473 Ok(self)
474 }
475
476 pub fn collation(&mut self, c: Arc<dyn CollationProvider>) -> Result<&mut Self, PluginError> {
482 self.require(&Capability::Collation)?;
483 self.pending
484 .push(Box::new(KeyedUniqueReg::<CollationSurface> {
485 key_override: None,
486 provider: c,
487 }));
488 Ok(self)
489 }
490
491 pub fn cdc_output(&mut self, c: Arc<dyn CdcOutputProvider>) -> Result<&mut Self, PluginError> {
497 self.require(&Capability::Cdc)?;
498 self.pending.push(Box::new(KeyedUniqueReg::<CdcSurface> {
499 key_override: None,
500 provider: c,
501 }));
502 Ok(self)
503 }
504
505 pub fn catalog(&mut self, c: Arc<dyn CatalogProvider>) -> Result<&mut Self, PluginError> {
511 self.require(&Capability::Catalog)?;
512 self.pending
513 .push(Box::new(KeyedUniqueReg::<CatalogSurface> {
514 key_override: None,
515 provider: c,
516 }));
517 Ok(self)
518 }
519
520 pub fn replacement_scan(
526 &mut self,
527 r: Arc<dyn ReplacementScanProvider>,
528 ) -> Result<&mut Self, PluginError> {
529 self.require(&Capability::Catalog)?;
530 self.pending
531 .push(Box::new(AppendReg::<ReplacementScanSurface> {
532 provider: r,
533 }));
534 Ok(self)
535 }
536
537 pub fn background_job(
544 &mut self,
545 j: Arc<dyn BackgroundJobProvider>,
546 ) -> Result<&mut Self, PluginError> {
547 self.require(&Capability::BackgroundJob { max_concurrent: 0 })?;
548 self.pending
549 .push(Box::new(AppendReg::<BackgroundJobSurface> { provider: j }));
550 Ok(self)
551 }
552
553 pub fn commit_to_registry(self) -> Result<(), PluginError> {
564 self.registry.apply_pending(&self.plugin_id, self.pending)
565 }
566
567 #[must_use]
573 pub fn pending_len(&self) -> usize {
574 self.pending.len()
575 }
576}