drasi-lib 0.9.2

Embedded Drasi for in-process data change processing using continuous queries
// Copyright 2025 The Drasi Authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
//     http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

pub mod base;
pub mod component_graph_source;
pub mod durability_config;
pub mod future_queue_source;
pub(crate) mod graph_elements;
pub mod manager;
mod traits;

#[cfg(test)]
pub(crate) mod tests;

use async_trait::async_trait;
use drasi_core::models::SourceChange;

/// Trait for publishing source changes to the Drasi server
#[async_trait]
pub trait Publisher: Send + Sync {
    /// Publish a source change event
    async fn publish(
        &self,
        change: SourceChange,
    ) -> Result<(), Box<dyn std::error::Error + Send + Sync>>;
}

/// Resource or attribute key that marks the origin of telemetry Drasi emitted.
pub const SOURCE_ORIGIN_ATTRIBUTE: &str = "drasi.source.origin";

/// Value of [`SOURCE_ORIGIN_ATTRIBUTE`] for Drasi-derived telemetry.
///
/// Sources that ingest telemetry SHOULD drop records with this marker so a
/// reaction that exports OTLP cannot feed back into the same graph.
pub const SOURCE_ORIGIN_DERIVED: &str = "derived";

// Re-export the Source trait and error types
pub use traits::Source;
pub use traits::SourceError;
pub use traits::{ByteLexPositionComparator, PositionComparator};

pub use base::{SourceBase, SourceBaseParams};
pub use component_graph_source::{ComponentGraphSource, COMPONENT_GRAPH_SOURCE_ID};
pub use durability_config::{CapacityPolicy, DurabilityConfig};
pub use future_queue_source::{FutureQueueSource, FUTURE_QUEUE_SOURCE_ID};
pub use manager::SourceManager;
pub use manager::{convert_json_to_element_properties, convert_json_to_element_value};