otel-arrow-dfe-engine 0.61.0

Async pipeline engine
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

//! Capability system for extensions.
//!
//! This module defines the type-safe capability resolution infrastructure.
//! Extensions register capabilities via [`ExtensionCapabilities`], and
//! node factories consume them via [`registry::Capabilities`].
//!
//! Capability traits are defined per-capability in submodules grouped by domain
//! (e.g., `auth::bearer_token_provider`), with local (!Send) and shared
//! (Send + Sync) variants re-exported from
//! [`local::capability`](crate::local::capability) and
//! [`shared::capability`](crate::shared::capability).

// Capability code is organized so every public item has exactly **one** export
// surface -- no item is reachable via two paths:
//
// - A capability's **data types and registration handle** are exposed only on
//   the **scoped** surface under its domain. Types shared across a domain's
//   capabilities are re-exported at the domain root, `capability::<domain>`
//   (e.g. `capability::auth::BearerToken`); a capability's own handle and
//   capability-specific types live in its submodule, `capability::<domain>::<name>`
//   (e.g. `capability::auth::bearer_token_provider::{BearerTokenProvider,
//   TokenStream}`). Their defining modules stay private, so each item has a
//   single public path.
//
// - A capability's **trait variants** (the `local` `!Send` / `shared`
//   `Send + Sync` traits the `#[capability]` macro generates) are exposed only on the
//   **execution-model** surface, `{local,shared}::capability::<domain>::<name>` --
//   the surface extensions implement against, alongside `{local,shared}::extension`,
//   etc.
//
// The two surfaces don't overlap because the `#[capability]` macro emits its
// `local`/`shared` trait modules as `pub(crate)` (an implementation detail), so
// they are never a public path under `capability::<domain>::<name>`; the traits
// become public only through the hand-written `{local,shared}::capability`
// re-exports.
//
// - Genuinely **shared** framework infrastructure -- used across capabilities
//   rather than owned by one -- is re-exported flat at `capability::` (below): the
//   error types (`error`), the instance-factory types (`factory`), and the
//   `ExtensionCapability` trait plus `KNOWN_CAPABILITIES` (defined in this
//   module). These are unique, collision-free vocabulary where a single short
//   canonical path is worth more than namespacing. `registry` stays `pub`
//   because `registry::Capabilities` is not re-exported at the root.
pub mod auth;
pub(crate) mod error;
pub(crate) mod factory;
pub mod registry;
pub mod vendor_bundle;

pub use error::{CapabilityError, CapabilityErrorSource};
pub use factory::{LocalInstanceFactory, SharedInstanceFactory};

// -- Sealed ExtensionCapability trait -----------------------------------------

/// Sealing module -- prevents external crates from implementing
/// [`ExtensionCapability`].
mod private {
    /// Sealed marker trait. Only the `#[capability]` proc macro (or
    /// hand-written impls inside the engine crate) can implement this.
    pub trait Sealed {}
}

/// Compile-time-sealed trait binding a capability registration struct
/// to its local and shared trait object types.
///
/// Each capability (e.g., `BearerTokenProvider`) has a zero-sized
/// registration struct that implements this trait. The associated types
/// tell [`Capabilities::require_local`](registry::Capabilities::require_local)
/// and [`Capabilities::require_shared`](registry::Capabilities::require_shared)
/// which concrete `dyn Trait` to downcast to.
///
/// This trait is sealed -- only the engine crate (via the
/// `#[capability]` proc macro or manual impls) can add new
/// capabilities.
pub trait ExtensionCapability: private::Sealed + 'static {
    /// Capability name as a const, usable in static contexts.
    /// Must match the `name` argument in the `#[capability]` attribute.
    const NAME: &'static str;

    /// The local (!Send) trait object type for this capability
    /// (e.g., `dyn local::capability::auth::bearer_token_provider::BearerTokenProvider`).
    type Local: ?Sized + 'static;

    /// The shared (Send + Sync) trait object type for this capability
    /// (e.g., `dyn shared::capability::auth::bearer_token_provider::BearerTokenProvider`).
    type Shared: ?Sized + Send + Sync + 'static;

    /// Human-readable name used in error messages and config validation.
    #[must_use]
    fn name() -> &'static str {
        Self::NAME
    }

    /// Wraps a freshly-produced shared trait object as a local trait object.
    ///
    /// The `#[capability]` proc macro generates an impl that constructs
    /// the capability's `SharedAsLocal` adapter. Because `local::Trait`
    /// and `shared::Trait` are generated from the same source trait,
    /// the adapter is always constructible -- there is no opt-out.
    /// A capability whose local and shared semantics diverge should be
    /// split into two distinct capabilities rather than expressed as an
    /// adapter refusal.
    ///
    /// Called from each capability entry's
    /// `adapt_as_local` fn pointer when a shared-only extension is
    /// consumed via `require_local` / the `SharedAsLocal` fallback
    /// path in `resolve_bindings`.
    ///
    /// # Per-node freshness
    ///
    /// Invoked **once per node** that binds the capability via the
    /// fallback path, with a new `Box<Self::Shared>` minted by the
    /// factory for that node. This matches the per-node instance
    /// semantics of [`Capabilities::require_shared`] -- every binding
    /// gets its own instance.
    ///
    /// If your shared impl relies on state that must be reset **per call**
    /// rather than per node, declare an explicit `local:` variant via
    /// the `extension_capabilities!` macro instead of relying on this
    /// fallback.
    ///
    /// Authors normally don't implement this by hand; the macro handles
    /// it. If you are hand-rolling an `ExtensionCapability` impl (only
    /// needed inside the engine crate for testing), return
    /// `Box::new(YourAdapter(shared))`.
    ///
    /// [`Capabilities::require_shared`]: registry::Capabilities::require_shared
    fn wrap_shared_as_local(shared: Box<Self::Shared>) -> Box<Self::Local>;
}

/// Re-export for use by the `#[capability]` proc macro's generated code.
/// `pub(crate)` (not `pub`) preserves the seal: the macro only expands
/// inside this crate, so external crates still can't reach `Sealed` to
/// forge an `ExtensionCapability` impl.
#[doc(hidden)]
pub(crate) use private::Sealed as CapabilitySealed;

// -- KNOWN_CAPABILITIES (link-time registration) ------------------------------

/// A link-time-registered capability descriptor.
///
/// Each `#[capability]` invocation produces a static entry in the
/// [`KNOWN_CAPABILITIES`] distributed slice. The engine uses this at
/// config validation time to map string names to `TypeId`s.
#[doc(hidden)]
pub struct KnownCapability {
    /// Human-readable name (e.g., `"bearer_token_provider"`).
    pub name: &'static str,
    /// Short description of what the capability does. Authored at the
    /// `#[capability(description = "...")]` site.
    ///
    /// TODO(extension-system): not yet read by any consumer.
    pub description: &'static str,
    /// `TypeId` of the zero-sized registration struct.
    pub type_id: fn() -> std::any::TypeId,
}

/// Link-time registry of all capabilities defined in the binary.
///
/// Populated by `#[capability]` proc macro entries. Used by
/// `resolve_bindings()` to validate capability names and retrieve
/// `TypeId`s.
//
// `linkme::distributed_slice` requires a `pub static`; `#[doc(hidden)]`
// excludes it from generated rustdoc so external crates don't see it in
// the public API surface.
#[doc(hidden)]
#[allow(unsafe_code)]
#[linkme::distributed_slice]
pub static KNOWN_CAPABILITIES: [KnownCapability] = [..];

// -- ExtensionCapabilities (factory metadata) ---------------------------------

/// Static metadata describing which capabilities an extension factory provides.
///
/// Carried on [`ExtensionFactory`](crate::ExtensionFactory) and used by:
/// - Config validation: checking that capability bindings reference
///   capabilities the extension actually provides.
/// - `resolve_bindings()`: knowing which registry slots to populate.
///
/// Constructed via the `extension_capabilities!` macro.
///
/// The `register_shared` / `register_local` fn pointers are the bridge
/// between the extension's type-erased instance factories and the
/// capability registry. The engine invokes them at bundle-registration
/// time, passing the extension's `ExtensionId` plus a clone of the
/// appropriate `*InstanceFactory`. The fn pointer internally builds one
/// [`SharedCapabilityEntry`](registry::SharedCapabilityEntry) per
/// listed capability and inserts it into the registry.
///
/// Returns [`registry::Error::InternalError`] on a duplicate
/// `(capability, extension)` insert.
#[derive(Clone)]
pub struct ExtensionCapabilities {
    /// Capability names provided by the **shared** variant.
    pub shared: &'static [&'static str],
    /// Capability names provided by the **local** variant.
    pub local: &'static [&'static str],
    /// Register all shared-variant capabilities into the registry.
    /// No-op when `shared` is empty.
    pub register_shared: fn(
        ext_id: otel_arrow_dfe_config::ExtensionId,
        factory: SharedInstanceFactory,
        registry: &mut registry::CapabilityRegistry,
    ) -> Result<(), registry::Error>,
    /// Register all local-variant capabilities into the registry.
    /// No-op when `local` is empty.
    pub register_local: fn(
        ext_id: otel_arrow_dfe_config::ExtensionId,
        factory: LocalInstanceFactory,
        registry: &mut registry::CapabilityRegistry,
    ) -> Result<(), registry::Error>,
}

/// Declares which capabilities an extension provides.
///
/// The left-hand side names the extension type(s) -- one or two,
/// depending on form -- and the right-hand side is a single capability
/// list shared by both execution models.
/// Three forms:
///
/// ```rust,ignore
/// // Shared-only (local consumers served via SharedAsLocal fallback).
/// extension_capabilities!(shared: MyExt => [BearerTokenProvider, KeyValueStore]);
///
/// // Local-only.
/// extension_capabilities!(local: MyLocalExt => [KeyValueStore]);
///
/// // Dual-type -- distinct shared/local types, same capability list.
/// extension_capabilities!(
///     (shared: MySharedKv, local: MyLocalKv) => [KeyValueStore]
/// );
/// ```
///
/// Each capability `$cap` in the list must have a `#[capability]`-generated
/// `shared_entry::<E>` and/or `local_entry::<E>` associated fn. The macro
/// invokes them per listed capability, passing a clone of the extension's
/// instance factory, and inserts the result into the registry.
///
/// In the dual form, `S` must implement `shared::$cap` and `L` must
/// implement `local::$cap` for every capability in the list -- mismatches
/// surface as standard trait-bound errors at the macro call site.
#[macro_export]
macro_rules! extension_capabilities {
    // Shared-only extension (automatic local fallback via SharedAsLocal).
    (shared: $ext:ty => [$($cap:ty),+ $(,)?]) => {
        $crate::capability::ExtensionCapabilities {
            shared: &[$(<$cap as $crate::capability::ExtensionCapability>::NAME),+],
            local: &[],
            register_shared: |ext_id, factory, registry| {
                $(
                    registry.register_shared(
                        ::std::any::TypeId::of::<$cap>(),
                        <$cap>::shared_entry::<$ext>(ext_id.clone(), factory.clone()),
                    )?;
                )+
                Ok(())
            },
            register_local: |_, _, _| Ok(()),
        }
    };
    // Local-only extension.
    (local: $ext:ty => [$($cap:ty),+ $(,)?]) => {
        $crate::capability::ExtensionCapabilities {
            shared: &[],
            local: &[$(<$cap as $crate::capability::ExtensionCapability>::NAME),+],
            register_shared: |_, _, _| Ok(()),
            register_local: |ext_id, factory, registry| {
                $(
                    registry.register_local(
                        ::std::any::TypeId::of::<$cap>(),
                        <$cap>::local_entry::<$ext>(ext_id.clone(), factory.clone()),
                    )?;
                )+
                Ok(())
            },
        }
    };
    // Dual-type extension -- distinct shared/local types, same capability list.
    ((shared: $sext:ty, local: $lext:ty) => [$($cap:ty),+ $(,)?]) => {
        $crate::capability::ExtensionCapabilities {
            shared: &[$(<$cap as $crate::capability::ExtensionCapability>::NAME),+],
            local: &[$(<$cap as $crate::capability::ExtensionCapability>::NAME),+],
            register_shared: |ext_id, factory, registry| {
                $(
                    registry.register_shared(
                        ::std::any::TypeId::of::<$cap>(),
                        <$cap>::shared_entry::<$sext>(ext_id.clone(), factory.clone()),
                    )?;
                )+
                Ok(())
            },
            register_local: |ext_id, factory, registry| {
                $(
                    registry.register_local(
                        ::std::any::TypeId::of::<$cap>(),
                        <$cap>::local_entry::<$lext>(ext_id.clone(), factory.clone()),
                    )?;
                )+
                Ok(())
            },
        }
    };
}

#[cfg(test)]
mod tests;