1use datafusion::error::Result as DataFusionResult;
2use std::collections::HashMap;
3
4pub struct TemporalGraphProcessor;
7
8#[derive(Debug, Clone, PartialEq)]
10pub struct TemporalEdge {
11 pub source: String,
12 pub target: String,
13 pub weight: f64,
14 pub timestamp: i64, pub properties: HashMap<String, String>,
16}
17
18#[derive(Debug, Clone)]
20pub struct TimeWindow {
21 pub start_time: i64, pub end_time: i64, pub window_size: Option<i64>, }
25
26#[derive(Debug, Clone)]
28pub struct TemporalQueryConfig {
29 pub time_window: TimeWindow,
30 pub include_edge_properties: bool,
31 pub aggregate_by_time_bucket: Option<i64>, }
33
34impl TemporalGraphProcessor {
35 pub fn new() -> Self {
36 Self
37 }
38
39 pub fn filter_edges_by_time(
42 &self,
43 edges: &[TemporalEdge],
44 time_window: &TimeWindow,
45 ) -> DataFusionResult<Vec<TemporalEdge>> {
46 let filtered_edges = edges
47 .iter()
48 .filter(|edge| {
49 edge.timestamp >= time_window.start_time && edge.timestamp <= time_window.end_time
50 })
51 .cloned()
52 .collect();
53
54 Ok(filtered_edges)
55 }
56
57 pub fn get_graph_snapshot(
60 &self,
61 edges: &[TemporalEdge],
62 snapshot_time: i64,
63 ) -> DataFusionResult<Vec<TemporalEdge>> {
64 let snapshot_edges = edges
65 .iter()
66 .filter(|edge| edge.timestamp <= snapshot_time)
67 .cloned()
68 .collect();
69
70 Ok(snapshot_edges)
71 }
72
73 pub fn find_temporal_paths(
76 &self,
77 edges: &[TemporalEdge],
78 start_node: &str,
79 end_node: &str,
80 time_window: &TimeWindow,
81 ) -> DataFusionResult<Vec<Vec<TemporalEdge>>> {
82 let window_edges = self.filter_edges_by_time(edges, time_window)?;
84
85 let mut adjacency_map: HashMap<String, Vec<&TemporalEdge>> = HashMap::new();
87 for edge in &window_edges {
88 adjacency_map
89 .entry(edge.source.clone())
90 .or_insert_with(Vec::new)
91 .push(edge);
92 }
93
94 for edges_list in adjacency_map.values_mut() {
96 edges_list.sort_by_key(|edge| edge.timestamp);
97 }
98
99 let mut paths = Vec::new();
100 let mut current_path = Vec::new();
101
102 self.dfs_temporal_paths(
103 &adjacency_map,
104 start_node,
105 end_node,
106 &mut current_path,
107 &mut paths,
108 time_window.start_time,
109 );
110
111 Ok(paths)
112 }
113
114 fn dfs_temporal_paths(
116 &self,
117 adjacency_map: &HashMap<String, Vec<&TemporalEdge>>,
118 current_node: &str,
119 target_node: &str,
120 current_path: &mut Vec<TemporalEdge>,
121 all_paths: &mut Vec<Vec<TemporalEdge>>,
122 min_timestamp: i64,
123 ) {
124 if current_node == target_node {
125 all_paths.push(current_path.clone());
126 return;
127 }
128
129 if let Some(edges) = adjacency_map.get(current_node) {
130 for edge in edges {
131 if edge.timestamp >= min_timestamp {
133 current_path.push((*edge).clone());
134
135 self.dfs_temporal_paths(
136 adjacency_map,
137 &edge.target,
138 target_node,
139 current_path,
140 all_paths,
141 edge.timestamp,
142 );
143
144 current_path.pop();
145 }
146 }
147 }
148 }
149
150 pub fn aggregate_by_time_bucket(
153 &self,
154 edges: &[TemporalEdge],
155 bucket_size_seconds: i64,
156 ) -> DataFusionResult<HashMap<i64, usize>> {
157 let mut buckets: HashMap<i64, usize> = HashMap::new();
158
159 for edge in edges {
160 let bucket = (edge.timestamp / bucket_size_seconds) * bucket_size_seconds;
161 *buckets.entry(bucket).or_insert(0) += 1;
162 }
163
164 Ok(buckets)
165 }
166
167 pub fn get_active_nodes_at_time(
170 &self,
171 edges: &[TemporalEdge],
172 timestamp: i64,
173 lookback_window: i64, ) -> DataFusionResult<Vec<String>> {
175 let window_start = timestamp - lookback_window;
176 let time_window = TimeWindow {
177 start_time: window_start,
178 end_time: timestamp,
179 window_size: Some(lookback_window),
180 };
181
182 let active_edges = self.filter_edges_by_time(edges, &time_window)?;
183
184 let mut active_nodes: std::collections::HashSet<String> = std::collections::HashSet::new();
185 for edge in active_edges {
186 active_nodes.insert(edge.source);
187 active_nodes.insert(edge.target);
188 }
189
190 Ok(active_nodes.into_iter().collect())
191 }
192
193 pub fn calculate_temporal_centrality(
195 &self,
196 edges: &[TemporalEdge],
197 node_id: &str,
198 time_window: &TimeWindow,
199 ) -> DataFusionResult<f64> {
200 let window_edges = self.filter_edges_by_time(edges, time_window)?;
201
202 let mut degree = 0;
203 let total_edges = window_edges.len();
204
205 for edge in &window_edges {
206 if edge.source == node_id || edge.target == node_id {
207 degree += 1;
208 }
209 }
210
211 if total_edges == 0 {
212 Ok(0.0)
213 } else {
214 Ok(degree as f64 / total_edges as f64)
215 }
216 }
217
218 pub fn parse_and_execute_temporal_query(
221 &self,
222 _query: &str,
223 edges: &[TemporalEdge],
224 config: &TemporalQueryConfig,
225 ) -> DataFusionResult<Vec<TemporalEdge>> {
226 self.filter_edges_by_time(edges, &config.time_window)
230 }
231}
232
233impl Default for TemporalGraphProcessor {
234 fn default() -> Self {
235 Self::new()
236 }
237}
238
239pub struct TemporalGraphQueries;
241
242impl TemporalGraphQueries {
243 pub fn edges_in_last_n_days(
245 processor: &TemporalGraphProcessor,
246 edges: &[TemporalEdge],
247 days: i64,
248 current_time: i64,
249 ) -> DataFusionResult<Vec<TemporalEdge>> {
250 let seconds_per_day = 86400;
251 let window_start = current_time - (days * seconds_per_day);
252
253 let time_window = TimeWindow {
254 start_time: window_start,
255 end_time: current_time,
256 window_size: Some(days * seconds_per_day),
257 };
258
259 processor.filter_edges_by_time(edges, &time_window)
260 }
261
262 pub fn temporal_shortest_path(
264 processor: &TemporalGraphProcessor,
265 edges: &[TemporalEdge],
266 start: &str,
267 end: &str,
268 time_window: &TimeWindow,
269 ) -> DataFusionResult<Option<Vec<TemporalEdge>>> {
270 let paths = processor.find_temporal_paths(edges, start, end, time_window)?;
271
272 let shortest_path = paths
274 .into_iter()
275 .min_by_key(|path| {
276 if path.is_empty() {
277 0
278 } else {
279 path.last().unwrap().timestamp - path.first().unwrap().timestamp
280 }
281 });
282
283 Ok(shortest_path)
284 }
285
286 pub fn most_active_nodes_in_period(
288 processor: &TemporalGraphProcessor,
289 edges: &[TemporalEdge],
290 time_window: &TimeWindow,
291 top_n: usize,
292 ) -> DataFusionResult<Vec<(String, usize)>> {
293 let window_edges = processor.filter_edges_by_time(edges, time_window)?;
294
295 let mut node_activity: HashMap<String, usize> = HashMap::new();
296
297 for edge in window_edges {
298 *node_activity.entry(edge.source).or_insert(0) += 1;
299 *node_activity.entry(edge.target).or_insert(0) += 1;
300 }
301
302 let mut activity_vec: Vec<(String, usize)> = node_activity.into_iter().collect();
303 activity_vec.sort_by(|a, b| b.1.cmp(&a.1)); activity_vec.truncate(top_n);
305
306 Ok(activity_vec)
307 }
308}
309
310#[cfg(test)]
311mod tests {
312 use super::*;
313
314 fn create_test_temporal_edges() -> Vec<TemporalEdge> {
315 vec![
316 TemporalEdge {
317 source: "A".to_string(),
318 target: "B".to_string(),
319 weight: 1.0,
320 timestamp: 1704067200, properties: HashMap::new(),
322 },
323 TemporalEdge {
324 source: "B".to_string(),
325 target: "C".to_string(),
326 weight: 2.0,
327 timestamp: 1704153600, properties: HashMap::new(),
329 },
330 TemporalEdge {
331 source: "C".to_string(),
332 target: "D".to_string(),
333 weight: 1.0,
334 timestamp: 1704240000, properties: HashMap::new(),
336 },
337 TemporalEdge {
338 source: "A".to_string(),
339 target: "C".to_string(),
340 weight: 4.0,
341 timestamp: 1704326400, properties: HashMap::new(),
343 },
344 ]
345 }
346
347 #[test]
348 fn test_filter_edges_by_time() {
349 let processor = TemporalGraphProcessor::new();
350 let edges = create_test_temporal_edges();
351
352 let time_window = TimeWindow {
353 start_time: 1704067200, end_time: 1704240000, window_size: None,
356 };
357
358 let filtered = processor.filter_edges_by_time(&edges, &time_window).unwrap();
359 assert_eq!(filtered.len(), 3); assert!(!filtered.iter().any(|e| e.timestamp == 1704326400));
363 }
364
365 #[test]
366 fn test_get_graph_snapshot() {
367 let processor = TemporalGraphProcessor::new();
368 let edges = create_test_temporal_edges();
369
370 let snapshot = processor.get_graph_snapshot(&edges, 1704153600).unwrap();
371 assert_eq!(snapshot.len(), 2); let snapshot_later = processor.get_graph_snapshot(&edges, 1704326400).unwrap();
374 assert_eq!(snapshot_later.len(), 4); }
376
377 #[test]
378 fn test_find_temporal_paths() {
379 let processor = TemporalGraphProcessor::new();
380 let edges = create_test_temporal_edges();
381
382 let time_window = TimeWindow {
383 start_time: 1704067200,
384 end_time: 1704326400,
385 window_size: None,
386 };
387
388 let paths = processor.find_temporal_paths(&edges, "A", "D", &time_window).unwrap();
389 assert!(!paths.is_empty());
390
391 for path in &paths {
393 for i in 1..path.len() {
394 assert!(path[i].timestamp >= path[i-1].timestamp);
395 }
396 }
397 }
398
399 #[test]
400 fn test_aggregate_by_time_bucket() {
401 let processor = TemporalGraphProcessor::new();
402 let edges = create_test_temporal_edges();
403
404 let buckets = processor.aggregate_by_time_bucket(&edges, 86400).unwrap();
406
407 assert!(buckets.len() >= 3); for count in buckets.values() {
411 assert!(*count >= 1);
412 }
413 }
414
415 #[test]
416 fn test_get_active_nodes_at_time() {
417 let processor = TemporalGraphProcessor::new();
418 let edges = create_test_temporal_edges();
419
420 let active_nodes = processor.get_active_nodes_at_time(
421 &edges,
422 1704240000, 172800, ).unwrap();
425
426 assert!(active_nodes.contains(&"A".to_string()));
428 assert!(active_nodes.contains(&"B".to_string()));
429 assert!(active_nodes.contains(&"C".to_string()));
430 }
431
432 #[test]
433 fn test_calculate_temporal_centrality() {
434 let processor = TemporalGraphProcessor::new();
435 let edges = create_test_temporal_edges();
436
437 let time_window = TimeWindow {
438 start_time: 1704067200,
439 end_time: 1704326400,
440 window_size: None,
441 };
442
443 let centrality_a = processor.calculate_temporal_centrality(&edges, "A", &time_window).unwrap();
444 let centrality_d = processor.calculate_temporal_centrality(&edges, "D", &time_window).unwrap();
445
446 assert!(centrality_a > centrality_d);
448 }
449
450 #[test]
451 fn test_temporal_query_utilities() {
452 let processor = TemporalGraphProcessor::new();
453 let edges = create_test_temporal_edges();
454
455 let recent_edges = TemporalGraphQueries::edges_in_last_n_days(
457 &processor,
458 &edges,
459 5, 1704326400, ).unwrap();
462 assert_eq!(recent_edges.len(), 4); let time_window = TimeWindow {
466 start_time: 1704067200,
467 end_time: 1704326400,
468 window_size: None,
469 };
470
471 let shortest_path = TemporalGraphQueries::temporal_shortest_path(
472 &processor,
473 &edges,
474 "A",
475 "D",
476 &time_window,
477 ).unwrap();
478 assert!(shortest_path.is_some());
479
480 let active_nodes = TemporalGraphQueries::most_active_nodes_in_period(
482 &processor,
483 &edges,
484 &time_window,
485 3, ).unwrap();
487 assert!(active_nodes.len() <= 3);
488 assert!(!active_nodes.is_empty());
489 }
490
491 #[test]
492 fn test_empty_time_window() {
493 let processor = TemporalGraphProcessor::new();
494 let edges = create_test_temporal_edges();
495
496 let empty_window = TimeWindow {
498 start_time: 1500000000, end_time: 1600000000, window_size: None,
501 };
502
503 let filtered = processor.filter_edges_by_time(&edges, &empty_window).unwrap();
504 assert!(filtered.is_empty());
505
506 let paths = processor.find_temporal_paths(&edges, "A", "D", &empty_window).unwrap();
507 assert!(paths.is_empty());
508 }
509}