#![allow(dead_code)]
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use datafusion_openlineage::DataFusionConfig;
use datafusion_openlineage::config::OpenLineageConfig;
use datafusion_openlineage::event::RunEvent;
use datafusion_openlineage::transport::{Transport, TransportError};
#[derive(Debug, Default, Clone)]
pub struct RecordingTransport {
pub events: Arc<Mutex<Vec<RunEvent>>>,
}
impl RecordingTransport {
pub fn events(&self) -> Vec<RunEvent> {
self.events.lock().unwrap().clone()
}
}
#[async_trait]
impl Transport for RecordingTransport {
async fn emit(&self, event: &RunEvent) -> Result<(), TransportError> {
self.events.lock().unwrap().push(event.clone());
Ok(())
}
}
#[derive(Debug)]
pub struct FailingTransport;
#[async_trait]
impl Transport for FailingTransport {
async fn emit(&self, _event: &RunEvent) -> Result<(), TransportError> {
Err(TransportError::Other("boom".to_string()))
}
}
pub fn config(namespace: &str, datafusion: bool) -> OpenLineageConfig {
let base = if datafusion {
OpenLineageConfig::for_datafusion()
} else {
OpenLineageConfig::default()
};
OpenLineageConfig {
job_namespace: namespace.to_string(),
..base
}
}