1use std::cmp::Ordering;
27use std::fmt;
28use std::hash::{Hash, Hasher};
29use std::sync::Arc;
30
31use async_trait::async_trait;
32use datafusion::common::{DFSchemaRef, Result};
33use datafusion::execution::context::{QueryPlanner, SessionState};
34use datafusion::logical_expr::{
35 Expr, Extension, InvariantLevel, LogicalPlan, UserDefinedLogicalNode,
36 UserDefinedLogicalNodeCore,
37};
38use datafusion::physical_plan::ExecutionPlan;
39use datafusion::physical_planner::{DefaultPhysicalPlanner, ExtensionPlanner, PhysicalPlanner};
40use uuid::Uuid;
41
42use crate::builder::{complete_event, fail_event, start_event};
43use crate::client::OpenLineageClient;
44use crate::config::OpenLineageConfig;
45use crate::context::LineageContextProvider;
46use crate::event::RunEvent;
47use crate::exec::OpenLineageExec;
48use crate::extract::extract;
49
50#[derive(Clone)]
60pub struct LineageMarker {
61 input: LogicalPlan,
62 complete: RunEvent,
65 client: OpenLineageClient,
66 producer: String,
67}
68
69impl fmt::Debug for LineageMarker {
70 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
71 f.debug_struct("LineageMarker").finish_non_exhaustive()
72 }
73}
74
75impl PartialEq for LineageMarker {
78 fn eq(&self, other: &Self) -> bool {
79 self.complete.run.run_id == other.complete.run.run_id && self.input == other.input
80 }
81}
82impl Eq for LineageMarker {}
83impl PartialOrd for LineageMarker {
84 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
85 self.complete
86 .run
87 .run_id
88 .partial_cmp(&other.complete.run.run_id)
89 }
90}
91impl Hash for LineageMarker {
92 fn hash<H: Hasher>(&self, state: &mut H) {
93 self.complete.run.run_id.hash(state);
94 }
95}
96
97impl UserDefinedLogicalNodeCore for LineageMarker {
98 fn name(&self) -> &str {
99 "LineageMarker"
100 }
101
102 fn inputs(&self) -> Vec<&LogicalPlan> {
103 vec![&self.input]
104 }
105
106 fn schema(&self) -> &DFSchemaRef {
107 self.input.schema()
108 }
109
110 fn check_invariants(&self, _check: InvariantLevel) -> Result<()> {
111 Ok(())
112 }
113
114 fn expressions(&self) -> Vec<Expr> {
115 vec![]
116 }
117
118 fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
119 write!(f, "LineageMarker")
120 }
121
122 fn with_exprs_and_inputs(
123 &self,
124 _exprs: Vec<Expr>,
125 mut inputs: Vec<LogicalPlan>,
126 ) -> Result<Self> {
127 Ok(Self {
128 input: inputs.pop().expect("LineageMarker has one input"),
129 complete: self.complete.clone(),
130 client: self.client.clone(),
131 producer: self.producer.clone(),
132 })
133 }
134}
135
136#[derive(Debug, Default)]
144pub struct LineageExtensionPlanner;
145
146#[async_trait]
147impl ExtensionPlanner for LineageExtensionPlanner {
148 async fn plan_extension(
149 &self,
150 _planner: &dyn PhysicalPlanner,
151 node: &dyn UserDefinedLogicalNode,
152 _logical_inputs: &[&LogicalPlan],
153 physical_inputs: &[Arc<dyn ExecutionPlan>],
154 _session_state: &SessionState,
155 ) -> Result<Option<Arc<dyn ExecutionPlan>>> {
156 let Some(marker) = node.as_any().downcast_ref::<LineageMarker>() else {
158 return Ok(None);
159 };
160 let inner = physical_inputs
161 .first()
162 .expect("LineageMarker has one physical input")
163 .clone();
164 Ok(Some(OpenLineageExec::new(
165 inner,
166 marker.client.clone(),
167 marker.complete.clone(),
168 marker.producer.clone(),
169 )))
170 }
171}
172
173pub struct OpenLineageQueryPlanner {
186 client: OpenLineageClient,
187 context: Arc<dyn LineageContextProvider>,
188 config: OpenLineageConfig,
189 physical: Arc<DefaultPhysicalPlanner>,
192}
193
194impl OpenLineageQueryPlanner {
195 pub fn new(
198 client: OpenLineageClient,
199 context: Arc<dyn LineageContextProvider>,
200 config: OpenLineageConfig,
201 extra_extension_planners: Vec<Arc<dyn ExtensionPlanner + Send + Sync>>,
202 ) -> Self {
203 let mut planners: Vec<Arc<dyn ExtensionPlanner + Send + Sync>> =
204 vec![Arc::new(LineageExtensionPlanner)];
205 planners.extend(extra_extension_planners);
206 Self {
207 client,
208 context,
209 config,
210 physical: Arc::new(DefaultPhysicalPlanner::with_extension_planners(planners)),
211 }
212 }
213}
214
215impl fmt::Debug for OpenLineageQueryPlanner {
216 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
218 f.debug_struct("OpenLineageQueryPlanner")
219 .finish_non_exhaustive()
220 }
221}
222
223#[async_trait]
224impl QueryPlanner for OpenLineageQueryPlanner {
225 async fn create_physical_plan(
226 &self,
227 logical_plan: &LogicalPlan,
228 session_state: &SessionState,
229 ) -> Result<Arc<dyn ExecutionPlan>> {
230 let mut lineage = extract(logical_plan, &self.config);
231 let cx = self.context.context(session_state).await;
232 lineage.sql = cx.sql.clone();
235
236 if lineage.inputs.is_empty() && lineage.outputs.is_empty() {
241 return self
242 .physical
243 .create_physical_plan(logical_plan, session_state)
244 .await;
245 }
246
247 let run_id = cx.run_id.unwrap_or_else(Uuid::now_v7);
248 self.client
249 .emit(start_event(run_id, &lineage, &cx, &self.config));
250
251 let marker = LineageMarker {
255 input: logical_plan.clone(),
256 complete: complete_event(run_id, &lineage, &cx, &self.config),
257 client: self.client.clone(),
258 producer: self.config.producer.clone(),
259 };
260 let wrapped = LogicalPlan::Extension(Extension {
261 node: Arc::new(marker),
262 });
263
264 match self
265 .physical
266 .create_physical_plan(&wrapped, session_state)
267 .await
268 {
269 Ok(plan) => Ok(plan),
270 Err(err) => {
271 self.client.emit(fail_event(
273 run_id,
274 &lineage,
275 &cx,
276 &self.config,
277 &err.to_string(),
278 ));
279 Err(err)
280 }
281 }
282 }
283}
284
285#[cfg(test)]
286mod tests {
287 use std::collections::hash_map::DefaultHasher;
288
289 use datafusion::logical_expr::LogicalPlanBuilder;
290
291 use super::*;
292 use crate::QueryLineage;
293 use crate::context::LineageContext;
294 use crate::transport::NoopTransport;
295
296 use datafusion::logical_expr::UserDefinedLogicalNodeCore as NodeCore;
300
301 fn marker(run_id: Uuid) -> LineageMarker {
304 let input = LogicalPlanBuilder::empty(false).build().unwrap();
305 let config = OpenLineageConfig::default();
306 let complete = complete_event(
307 run_id,
308 &QueryLineage::default(),
309 &LineageContext::default(),
310 &config,
311 );
312 LineageMarker {
313 input,
314 complete,
315 client: OpenLineageClient::new(Arc::new(NoopTransport)),
316 producer: config.producer,
317 }
318 }
319
320 fn hash_of(m: &LineageMarker) -> u64 {
321 let mut h = DefaultHasher::new();
322 m.hash(&mut h);
323 h.finish()
324 }
325
326 #[tokio::test]
328 async fn node_core_is_schema_transparent_and_expr_free() {
329 let m = marker(Uuid::now_v7());
330 assert_eq!(NodeCore::name(&m), "LineageMarker");
331 assert_eq!(NodeCore::inputs(&m).len(), 1);
332 assert_eq!(NodeCore::schema(&m), NodeCore::inputs(&m)[0].schema());
334 assert!(NodeCore::expressions(&m).is_empty());
335 assert!(NodeCore::check_invariants(&m, InvariantLevel::Always).is_ok());
336 assert_eq!(format!("{m:?}"), "LineageMarker { .. }");
337 }
338
339 #[tokio::test]
340 async fn with_exprs_and_inputs_rebuilds_preserving_payload() {
341 let run_id = Uuid::now_v7();
342 let m = marker(run_id);
343 let new_input = LogicalPlanBuilder::empty(true).build().unwrap();
344 let rebuilt = NodeCore::with_exprs_and_inputs(&m, vec![], vec![new_input.clone()]).unwrap();
345 assert_eq!(NodeCore::inputs(&rebuilt)[0], &new_input);
347 assert_eq!(rebuilt.complete.run.run_id, run_id);
348 }
349
350 #[tokio::test]
351 async fn identity_keys_on_run_id_and_input() {
352 let run_id = Uuid::now_v7();
353 let a = marker(run_id);
355 let b = marker(run_id);
356 assert_eq!(a, b);
357 assert_eq!(hash_of(&a), hash_of(&b));
358 assert_eq!(a.partial_cmp(&b), Some(Ordering::Equal));
359
360 let c = marker(Uuid::now_v7());
362 assert_ne!(a, c);
363 assert_eq!(
364 a.partial_cmp(&c),
365 a.complete.run.run_id.partial_cmp(&c.complete.run.run_id)
366 );
367 }
368}