use super::readiness::{ReadinessProbe, ReadinessSignaller};
use super::{ExtensionBundle, ExtensionLifecycle, ExtensionWrapper};
use crate::capability::factory::{LocalInstanceFactory, SharedInstanceFactory};
use crate::config::ExtensionConfig;
use crate::error::Error;
use crate::local::extension as local_ext;
use crate::local::message::{LocalReceiver, LocalSender};
use crate::message::Sender;
use crate::shared::extension as shared_ext;
use crate::shared::message::{SharedReceiver, SharedSender};
use otel_arrow_dfe_channel::mpsc;
use otel_arrow_dfe_config::ExtensionId;
use otel_arrow_dfe_config::extension::ExtensionUserConfig;
use otel_arrow_dfe_telemetry::{otel_debug, otel_info};
use std::any::{Any, TypeId};
use std::sync::Arc;
use std::time::Duration;
fn emit_readiness_timeout_override(extension: &str, timeout: Duration) {
let default = super::readiness::DEFAULT_READINESS_TIMEOUT;
if timeout != default {
otel_info!(
"extension.readiness.timeout_override",
extension = extension,
timeout_ms = u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
default_ms = u64::try_from(default.as_millis()).unwrap_or(u64::MAX),
);
}
}
#[doc(hidden)]
pub struct SharedDecomposed {
pub(crate) extension: Option<Box<dyn shared_ext::Extension>>,
pub(crate) instance_factory: SharedInstanceFactory,
pub(crate) type_id: TypeId,
}
#[doc(hidden)]
pub struct LocalDecomposed {
pub(crate) extension: Option<std::rc::Rc<dyn local_ext::Extension>>,
pub(crate) instance_factory: LocalInstanceFactory,
pub(crate) type_id: TypeId,
}
#[doc(hidden)]
pub struct ActiveStage {
parent: ExtensionBundleBuilder,
readiness_timeout: Option<Duration>,
}
impl ActiveStage {
#[must_use]
pub fn with_readiness_probe(mut self) -> Self {
self.readiness_timeout = Some(super::readiness::DEFAULT_READINESS_TIMEOUT);
self
}
#[must_use]
pub fn with_readiness_probe_timeout_override(mut self, timeout: Duration) -> Self {
emit_readiness_timeout_override(self.parent.name.as_ref(), timeout);
self.readiness_timeout = Some(timeout);
self
}
#[must_use]
pub fn shared<E>(mut self, extension: E) -> ActiveSharedStage
where
E: shared_ext::Extension + Clone + Send + 'static,
{
self.parent.set_shared_active(extension);
if let Some(t) = self.readiness_timeout {
let (sig, probe) = ReadinessSignaller::pair(t);
self.parent.set_shared_readiness(probe, sig);
}
ActiveSharedStage {
parent: self.parent,
readiness_timeout: self.readiness_timeout,
}
}
#[must_use]
pub fn local<E>(mut self, extension: std::rc::Rc<E>) -> ActiveLocalStage
where
E: local_ext::Extension + Clone + 'static,
{
self.parent.set_local_active(extension);
if let Some(t) = self.readiness_timeout {
let (sig, probe) = ReadinessSignaller::pair(t);
self.parent.set_local_readiness(probe, sig);
}
ActiveLocalStage {
parent: self.parent,
readiness_timeout: self.readiness_timeout,
}
}
}
#[doc(hidden)]
pub struct ActiveSharedStage {
parent: ExtensionBundleBuilder,
readiness_timeout: Option<Duration>,
}
impl ActiveSharedStage {
#[must_use]
pub fn local<E>(mut self, extension: std::rc::Rc<E>) -> ActiveCompleteStage
where
E: local_ext::Extension + Clone + 'static,
{
self.parent.set_local_active(extension);
if let Some(t) = self.readiness_timeout {
let (sig, probe) = ReadinessSignaller::pair(t);
self.parent.set_local_readiness(probe, sig);
}
ActiveCompleteStage {
parent: self.parent,
}
}
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
#[doc(hidden)]
pub struct ActiveLocalStage {
parent: ExtensionBundleBuilder,
readiness_timeout: Option<Duration>,
}
impl ActiveLocalStage {
#[must_use]
pub fn shared<E>(mut self, extension: E) -> ActiveCompleteStage
where
E: shared_ext::Extension + Clone + Send + 'static,
{
self.parent.set_shared_active(extension);
if let Some(t) = self.readiness_timeout {
let (sig, probe) = ReadinessSignaller::pair(t);
self.parent.set_shared_readiness(probe, sig);
}
ActiveCompleteStage {
parent: self.parent,
}
}
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
#[doc(hidden)]
pub struct ActiveCompleteStage {
parent: ExtensionBundleBuilder,
}
impl ActiveCompleteStage {
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
#[doc(hidden)]
pub struct PassiveStage {
parent: ExtensionBundleBuilder,
}
impl PassiveStage {
#[must_use]
pub fn cloned(self) -> PassiveClonedStage {
PassiveClonedStage {
parent: self.parent,
}
}
#[must_use]
pub fn constructed(self) -> PassiveConstructedStage {
PassiveConstructedStage {
parent: self.parent,
}
}
}
#[doc(hidden)]
pub struct PassiveClonedStage {
parent: ExtensionBundleBuilder,
}
impl PassiveClonedStage {
#[must_use]
pub fn shared<E>(mut self, extension: E) -> PassiveClonedSharedStage
where
E: Clone + Send + 'static,
{
self.parent.set_shared_cloned(extension);
PassiveClonedSharedStage {
parent: self.parent,
}
}
#[must_use]
pub fn local<E>(mut self, extension: E) -> PassiveClonedLocalStage
where
E: Clone + 'static,
{
self.parent.set_local_cloned(extension);
PassiveClonedLocalStage {
parent: self.parent,
}
}
}
#[doc(hidden)]
pub struct PassiveClonedSharedStage {
parent: ExtensionBundleBuilder,
}
impl PassiveClonedSharedStage {
#[must_use]
pub fn local<E>(mut self, extension: E) -> PassiveClonedCompleteStage
where
E: Clone + 'static,
{
self.parent.set_local_cloned(extension);
PassiveClonedCompleteStage {
parent: self.parent,
}
}
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
#[doc(hidden)]
pub struct PassiveClonedLocalStage {
parent: ExtensionBundleBuilder,
}
impl PassiveClonedLocalStage {
#[must_use]
pub fn shared<E>(mut self, extension: E) -> PassiveClonedCompleteStage
where
E: Clone + Send + 'static,
{
self.parent.set_shared_cloned(extension);
PassiveClonedCompleteStage {
parent: self.parent,
}
}
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
#[doc(hidden)]
pub struct PassiveClonedCompleteStage {
parent: ExtensionBundleBuilder,
}
impl PassiveClonedCompleteStage {
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
#[doc(hidden)]
pub struct PassiveConstructedStage {
parent: ExtensionBundleBuilder,
}
impl PassiveConstructedStage {
#[must_use]
pub fn shared<E, F>(mut self, produce: F) -> PassiveConstructedSharedStage
where
E: Send + 'static,
F: Fn() -> E + Clone + Send + 'static,
{
self.parent.set_shared_constructed::<E, F>(produce);
PassiveConstructedSharedStage {
parent: self.parent,
}
}
#[must_use]
pub fn local<E, F>(mut self, produce: F) -> PassiveConstructedLocalStage
where
E: 'static,
F: Fn() -> E + Clone + 'static,
{
self.parent.set_local_constructed::<E, F>(produce);
PassiveConstructedLocalStage {
parent: self.parent,
}
}
}
#[doc(hidden)]
pub struct PassiveConstructedSharedStage {
parent: ExtensionBundleBuilder,
}
impl PassiveConstructedSharedStage {
#[must_use]
pub fn local<E, F>(mut self, produce: F) -> PassiveConstructedCompleteStage
where
E: 'static,
F: Fn() -> E + Clone + 'static,
{
self.parent.set_local_constructed::<E, F>(produce);
PassiveConstructedCompleteStage {
parent: self.parent,
}
}
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
#[doc(hidden)]
pub struct PassiveConstructedLocalStage {
parent: ExtensionBundleBuilder,
}
impl PassiveConstructedLocalStage {
#[must_use]
pub fn shared<E, F>(mut self, produce: F) -> PassiveConstructedCompleteStage
where
E: Send + 'static,
F: Fn() -> E + Clone + Send + 'static,
{
self.parent.set_shared_constructed::<E, F>(produce);
PassiveConstructedCompleteStage {
parent: self.parent,
}
}
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
#[doc(hidden)]
pub struct PassiveConstructedCompleteStage {
parent: ExtensionBundleBuilder,
}
impl PassiveConstructedCompleteStage {
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
#[doc(hidden)]
pub struct BackgroundEmptyStage {
parent: ExtensionBundleBuilder,
readiness_timeout: Option<Duration>,
}
impl BackgroundEmptyStage {
#[must_use]
pub fn with_readiness_probe(mut self) -> Self {
self.readiness_timeout = Some(super::readiness::DEFAULT_READINESS_TIMEOUT);
self
}
#[must_use]
pub fn with_readiness_probe_timeout_override(mut self, timeout: Duration) -> Self {
emit_readiness_timeout_override(self.parent.name.as_ref(), timeout);
self.readiness_timeout = Some(timeout);
self
}
#[must_use]
pub fn shared<E>(mut self, extension: E) -> BackgroundCompleteStage
where
E: shared_ext::Extension + Clone + Send + 'static,
{
self.parent.set_shared_active(extension);
if let Some(t) = self.readiness_timeout {
let (sig, probe) = ReadinessSignaller::pair(t);
self.parent.set_shared_readiness(probe, sig);
}
BackgroundCompleteStage {
parent: self.parent,
}
}
#[must_use]
pub fn local<E>(mut self, extension: std::rc::Rc<E>) -> BackgroundCompleteStage
where
E: local_ext::Extension + Clone + 'static,
{
self.parent.set_local_active(extension);
if let Some(t) = self.readiness_timeout {
let (sig, probe) = ReadinessSignaller::pair(t);
self.parent.set_local_readiness(probe, sig);
}
BackgroundCompleteStage {
parent: self.parent,
}
}
}
#[doc(hidden)]
pub struct BackgroundCompleteStage {
parent: ExtensionBundleBuilder,
}
impl BackgroundCompleteStage {
pub fn build(self) -> Result<ExtensionBundle, Error> {
self.parent.build()
}
}
pub struct ExtensionBundleBuilder {
pub(super) name: ExtensionId,
pub(super) user_config: Arc<ExtensionUserConfig>,
pub(super) runtime_config: ExtensionConfig,
shared: Option<SharedDecomposed>,
local: Option<LocalDecomposed>,
shared_probe: Option<ReadinessProbe>,
local_probe: Option<ReadinessProbe>,
shared_signaller: Option<ReadinessSignaller>,
local_signaller: Option<ReadinessSignaller>,
}
impl ExtensionBundleBuilder {
pub(super) fn new(
name: ExtensionId,
user_config: Arc<ExtensionUserConfig>,
runtime_config: ExtensionConfig,
) -> Self {
Self {
name,
user_config,
runtime_config,
shared: None,
local: None,
shared_probe: None,
local_probe: None,
shared_signaller: None,
local_signaller: None,
}
}
fn set_shared_active<E>(&mut self, extension: E)
where
E: shared_ext::Extension + Clone + Send + 'static,
{
let for_factory = extension.clone();
self.shared = Some(SharedDecomposed {
extension: Some(Box::new(extension)),
instance_factory: SharedInstanceFactory::new(move || {
Box::new(for_factory.clone()) as Box<dyn Any + Send>
}),
type_id: TypeId::of::<E>(),
});
}
fn set_local_active<E>(&mut self, extension: std::rc::Rc<E>)
where
E: local_ext::Extension + Clone + 'static,
{
let for_factory = std::rc::Rc::clone(&extension);
self.local = Some(LocalDecomposed {
extension: Some(extension),
instance_factory: LocalInstanceFactory::new(move || {
Box::new((*for_factory).clone()) as Box<dyn Any>
}),
type_id: TypeId::of::<E>(),
});
}
fn set_shared_cloned<E>(&mut self, extension: E)
where
E: Clone + Send + 'static,
{
self.shared = Some(SharedDecomposed {
extension: None,
instance_factory: SharedInstanceFactory::new(move || {
Box::new(extension.clone()) as Box<dyn Any + Send>
}),
type_id: TypeId::of::<E>(),
});
}
fn set_local_cloned<E>(&mut self, extension: E)
where
E: Clone + 'static,
{
self.local = Some(LocalDecomposed {
extension: None,
instance_factory: LocalInstanceFactory::new(move || {
Box::new(extension.clone()) as Box<dyn Any>
}),
type_id: TypeId::of::<E>(),
});
}
fn set_shared_constructed<E, F>(&mut self, produce: F)
where
E: Send + 'static,
F: Fn() -> E + Clone + Send + 'static,
{
self.shared = Some(SharedDecomposed {
extension: None,
instance_factory: SharedInstanceFactory::new(move || {
Box::new(produce()) as Box<dyn Any + Send>
}),
type_id: TypeId::of::<E>(),
});
}
fn set_local_constructed<E, F>(&mut self, produce: F)
where
E: 'static,
F: Fn() -> E + Clone + 'static,
{
self.local = Some(LocalDecomposed {
extension: None,
instance_factory: LocalInstanceFactory::new(move || {
Box::new(produce()) as Box<dyn Any>
}),
type_id: TypeId::of::<E>(),
});
}
fn set_shared_readiness(&mut self, probe: ReadinessProbe, signaller: ReadinessSignaller) {
debug_assert!(
self.shared_probe.is_none() && self.shared_signaller.is_none(),
"readiness probe already registered for the shared variant of extension '{}'",
self.name.as_ref()
);
self.shared_probe = Some(probe);
self.shared_signaller = Some(signaller);
}
fn set_local_readiness(&mut self, probe: ReadinessProbe, signaller: ReadinessSignaller) {
debug_assert!(
self.local_probe.is_none() && self.local_signaller.is_none(),
"readiness probe already registered for the local variant of extension '{}'",
self.name.as_ref()
);
self.local_probe = Some(probe);
self.local_signaller = Some(signaller);
}
#[must_use]
pub fn active(self) -> ActiveStage {
ActiveStage {
parent: self,
readiness_timeout: None,
}
}
#[must_use]
pub fn passive(self) -> PassiveStage {
PassiveStage { parent: self }
}
#[must_use]
pub fn background(self) -> BackgroundEmptyStage {
BackgroundEmptyStage {
parent: self,
readiness_timeout: None,
}
}
pub(super) fn build(self) -> Result<ExtensionBundle, Error> {
if let (Some(local), Some(shared)) = (&self.local, &self.shared)
&& local.type_id == shared.type_id
{
return Err(Error::InternalError {
message: "local and shared variants must use different concrete types; \
register only one execution model for single-variant extensions"
.into(),
});
}
let cap = self.runtime_config.control_channel.capacity;
let Self {
name,
user_config,
runtime_config,
shared: shared_decomposed,
local: local_decomposed,
shared_probe,
local_probe,
shared_signaller,
local_signaller,
} = self;
if let Some(probe) = local_probe.as_ref().or(shared_probe.as_ref())
&& probe.timeout().is_zero()
{
return Err(Error::ExtensionReadinessZeroTimeout {
extension: name.as_ref().to_owned(),
});
}
debug_assert!(
!(local_probe.is_some()
&& local_decomposed
.as_ref()
.is_none_or(|d| d.extension.is_none())),
"local readiness probe registered but no active local extension \u{2014} typestate broken"
);
debug_assert!(
!(shared_probe.is_some()
&& shared_decomposed
.as_ref()
.is_none_or(|d| d.extension.is_none())),
"shared readiness probe registered but no active shared extension \u{2014} typestate broken"
);
let local = local_decomposed.map(|l| {
let lifecycle = match l.extension {
Some(ext) => {
let (tx, rx) = mpsc::Channel::new(cap);
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
ExtensionLifecycle::Active {
extension: ext,
control_sender: Sender::Local(LocalSender::mpsc(tx)),
control_receiver: LocalReceiver::mpsc(rx),
shutdown_sender: Some(shutdown_tx),
shutdown_receiver: Some(shutdown_rx),
readiness_probe: local_probe,
readiness_signaller: local_signaller,
}
}
None => ExtensionLifecycle::Passive,
};
ExtensionWrapper::Local {
name: name.clone(),
user_config: user_config.clone(),
runtime_config: runtime_config.clone(),
telemetry: None,
lifecycle,
instance_factory: l.instance_factory,
}
});
let shared = shared_decomposed.map(|s| {
let lifecycle = match s.extension {
Some(ext) => {
let (tx, rx) = tokio::sync::mpsc::channel(cap);
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
ExtensionLifecycle::Active {
extension: ext,
control_sender: Sender::Shared(SharedSender::mpsc(tx)),
control_receiver: SharedReceiver::mpsc(rx),
shutdown_sender: Some(shutdown_tx),
shutdown_receiver: Some(shutdown_rx),
readiness_probe: shared_probe,
readiness_signaller: shared_signaller,
}
}
None => ExtensionLifecycle::Passive,
};
ExtensionWrapper::Shared {
name: name.clone(),
user_config: user_config.clone(),
runtime_config: runtime_config.clone(),
telemetry: None,
lifecycle,
instance_factory: s.instance_factory,
}
});
if local.is_none() && shared.is_none() {
return Err(Error::InternalError {
message: "ExtensionBundle must have at least one variant (local or shared)".into(),
});
}
for w in local.iter().chain(shared.iter()) {
let name = w.name();
otel_debug!(
"extension.builder.build",
name = name.as_ref(),
variant = match w {
ExtensionWrapper::Local { .. } => "local",
ExtensionWrapper::Shared { .. } => "shared",
},
lifecycle = if w.is_passive() { "passive" } else { "active" },
);
}
Ok(ExtensionBundle::from_parts(local, shared))
}
}
impl ExtensionWrapper {
#[must_use]
pub fn builder(
name: ExtensionId,
user_config: Arc<ExtensionUserConfig>,
config: &ExtensionConfig,
) -> ExtensionBundleBuilder {
ExtensionBundleBuilder::new(name, user_config, config.clone())
}
}