termesh_test_support/
scripted_search.rs1use std::collections::VecDeque;
4use std::sync::atomic::{AtomicBool, Ordering};
5use std::sync::{Arc, Mutex};
6
7use termesh_core::{SearchEvent, SearchMode, SearchRequest};
8use termesh_search::{SearchEventSink, SearchResult, SearchService};
9
10struct Script {
11 mode: SearchMode,
12 query: String,
13 events: Vec<SearchEvent>,
14}
15
16#[derive(Default)]
17struct State {
18 scripts: VecDeque<Script>,
19 requests: Vec<SearchRequest>,
20}
21
22type Shared = Arc<Mutex<State>>;
23
24#[derive(Clone, Default)]
25pub struct ScriptedSearch {
26 shared: Shared,
27}
28
29#[derive(Clone)]
30pub struct ScriptedSearchControl {
31 shared: Shared,
32}
33
34impl ScriptedSearch {
35 pub fn new() -> Self {
36 Self::default()
37 }
38
39 pub fn control(&self) -> ScriptedSearchControl {
40 ScriptedSearchControl { shared: self.shared.clone() }
41 }
42
43 pub fn with_script(
44 self,
45 mode: SearchMode,
46 query: impl Into<String>,
47 events: Vec<SearchEvent>,
48 ) -> Self {
49 self.queue_script(mode, query, events);
50 self
51 }
52
53 pub fn queue_script(
54 &self,
55 mode: SearchMode,
56 query: impl Into<String>,
57 events: Vec<SearchEvent>,
58 ) {
59 self.shared.lock().expect("scripted search state poisoned").scripts.push_back(Script {
60 mode,
61 query: query.into(),
62 events,
63 });
64 }
65}
66
67impl ScriptedSearchControl {
68 pub fn requests(&self) -> Vec<SearchRequest> {
69 self.shared.lock().expect("scripted search state poisoned").requests.clone()
70 }
71}
72
73impl SearchService for ScriptedSearch {
74 fn search(
75 &mut self,
76 request: &SearchRequest,
77 cancelled: &AtomicBool,
78 sink: &SearchEventSink,
79 ) -> SearchResult<()> {
80 let events = {
81 let mut state = self.shared.lock().expect("scripted search state poisoned");
82 state.requests.push(request.clone());
83 let script = state
84 .scripts
85 .iter()
86 .position(|script| script.mode == request.mode && script.query == request.query)
87 .and_then(|index| state.scripts.remove(index));
88 script.map(|script| script.events).unwrap_or_default()
89 };
90
91 for event in events {
92 if cancelled.load(Ordering::Acquire) {
93 break;
94 }
95 sink(event);
96 }
97 Ok(())
98 }
99}