arrow_graph/sql/
temporal_queries.rs

1use datafusion::error::Result as DataFusionResult;
2use std::collections::HashMap;
3
4/// Temporal graph query processor for time-based edge filtering
5/// Supports queries like: SELECT * FROM edges WHERE edge_timestamp BETWEEN '2024-01-01' AND '2024-12-31'
6pub struct TemporalGraphProcessor;
7
8/// Represents a temporal edge with timestamp information
9#[derive(Debug, Clone, PartialEq)]
10pub struct TemporalEdge {
11    pub source: String,
12    pub target: String,
13    pub weight: f64,
14    pub timestamp: i64, // Unix timestamp
15    pub properties: HashMap<String, String>,
16}
17
18/// Time window configuration for temporal queries
19#[derive(Debug, Clone)]
20pub struct TimeWindow {
21    pub start_time: i64,  // Unix timestamp
22    pub end_time: i64,    // Unix timestamp
23    pub window_size: Option<i64>, // Optional sliding window size in seconds
24}
25
26/// Configuration for temporal graph queries
27#[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>, // Bucket size in seconds for temporal aggregation
32}
33
34impl TemporalGraphProcessor {
35    pub fn new() -> Self {
36        Self
37    }
38
39    /// Filter edges by time window
40    /// Example: Find all edges that occurred between two timestamps
41    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    /// Get graph snapshot at a specific timestamp
58    /// Returns all edges that were active at the given time
59    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    /// Find temporal paths - paths that respect time ordering
74    /// Each edge in the path must have a timestamp >= previous edge
75    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        // Filter edges within time window first
83        let window_edges = self.filter_edges_by_time(edges, time_window)?;
84        
85        // Build adjacency map with temporal ordering
86        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        // Sort edges by timestamp for each node
95        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    /// Depth-first search for temporal paths
115    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                // Ensure temporal ordering: each edge must happen after the previous one
132                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    /// Aggregate edges by time buckets
151    /// Example: Count edges per hour, day, week, etc.
152    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    /// Find active nodes at a specific time
168    /// Returns nodes that had at least one edge within the time window before the timestamp
169    pub fn get_active_nodes_at_time(
170        &self,
171        edges: &[TemporalEdge],
172        timestamp: i64,
173        lookback_window: i64, // How far back to look for activity
174    ) -> 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    /// Calculate temporal centrality - how central a node is within a time window
194    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    /// Process a temporal SQL query (simplified parser)
219    /// Example: "SELECT source, target FROM edges WHERE edge_timestamp BETWEEN 1704067200 AND 1735689600"
220    pub fn parse_and_execute_temporal_query(
221        &self,
222        _query: &str,
223        edges: &[TemporalEdge],
224        config: &TemporalQueryConfig,
225    ) -> DataFusionResult<Vec<TemporalEdge>> {
226        // This is a simplified implementation
227        // A full implementation would parse the actual SQL syntax
228        
229        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
239/// Utility functions for common temporal graph operations
240pub struct TemporalGraphQueries;
241
242impl TemporalGraphQueries {
243    /// Standard "edges in last N days" query
244    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    /// Standard "temporal shortest path" query
263    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        // Return the path with minimum total time duration
273        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    /// Standard "nodes active in period" query
287    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)); // Sort by activity descending
304        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, // 2024-01-01 00:00:00
321                properties: HashMap::new(),
322            },
323            TemporalEdge {
324                source: "B".to_string(),
325                target: "C".to_string(),
326                weight: 2.0,
327                timestamp: 1704153600, // 2024-01-02 00:00:00
328                properties: HashMap::new(),
329            },
330            TemporalEdge {
331                source: "C".to_string(),
332                target: "D".to_string(),
333                weight: 1.0,
334                timestamp: 1704240000, // 2024-01-03 00:00:00
335                properties: HashMap::new(),
336            },
337            TemporalEdge {
338                source: "A".to_string(),
339                target: "C".to_string(),
340                weight: 4.0,
341                timestamp: 1704326400, // 2024-01-04 00:00:00
342                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, // 2024-01-01
354            end_time: 1704240000,   // 2024-01-03
355            window_size: None,
356        };
357
358        let filtered = processor.filter_edges_by_time(&edges, &time_window).unwrap();
359        assert_eq!(filtered.len(), 3); // Should include first 3 edges
360        
361        // Last edge should be filtered out
362        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); // Should include first 2 edges
372        
373        let snapshot_later = processor.get_graph_snapshot(&edges, 1704326400).unwrap();
374        assert_eq!(snapshot_later.len(), 4); // Should include all edges
375    }
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        // Verify temporal ordering in paths
392        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        // Group by day (86400 seconds)
405        let buckets = processor.aggregate_by_time_bucket(&edges, 86400).unwrap();
406        
407        assert!(buckets.len() >= 3); // Should have at least 3 different days
408        
409        // Each bucket should have at least 1 edge
410        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, // 2024-01-03
423            172800,     // 2 days lookback
424        ).unwrap();
425        
426        // Should include nodes A, B, C, D
427        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        // A should have higher centrality than D (appears in more edges)
447        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        // Test edges in last N days
456        let recent_edges = TemporalGraphQueries::edges_in_last_n_days(
457            &processor,
458            &edges,
459            5, // last 5 days
460            1704326400, // from 2024-01-04
461        ).unwrap();
462        assert_eq!(recent_edges.len(), 4); // All edges should be included
463        
464        // Test temporal shortest path
465        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        // Test most active nodes
481        let active_nodes = TemporalGraphQueries::most_active_nodes_in_period(
482            &processor,
483            &edges,
484            &time_window,
485            3, // top 3
486        ).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        // Time window that excludes all edges
497        let empty_window = TimeWindow {
498            start_time: 1500000000, // Much earlier
499            end_time: 1600000000,   // Still before any edges
500            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}