Skip to main content

wpl/ast/
ann_func.rs

1use crate::WparseError;
2use crate::WplEvaluator;
3use crate::ast::AnnFun;
4use smol_str::SmolStr;
5use std::collections::BTreeMap;
6use wp_model_core::model::{DataField, DataRecord};
7use wp_model_core::raw::RawData;
8use wp_source_types::SourceEvent;
9
10pub trait AnnotationFunc {
11    fn proc(&self, src: &SourceEvent, data: &mut DataRecord) -> Result<(), WparseError>;
12}
13
14#[derive(Clone, Debug)]
15pub struct TagAnnotation {
16    args: BTreeMap<SmolStr, SmolStr>,
17}
18
19impl AnnotationFunc for TagAnnotation {
20    fn proc(&self, _src: &SourceEvent, data: &mut DataRecord) -> Result<(), WparseError> {
21        for (key, val) in &self.args {
22            data.append(DataField::from_chars(key.clone(), val.clone()));
23        }
24        Ok(())
25    }
26}
27
28#[derive(Clone, Debug)]
29pub struct NoopAnnotation;
30
31impl AnnotationFunc for NoopAnnotation {
32    fn proc(&self, _src: &SourceEvent, _data: &mut DataRecord) -> Result<(), WparseError> {
33        Ok(())
34    }
35}
36
37#[derive(Clone, Debug)]
38pub struct RawCopy {
39    raw_key: SmolStr,
40}
41
42impl AnnotationFunc for RawCopy {
43    fn proc(&self, src: &SourceEvent, data: &mut DataRecord) -> Result<(), WparseError> {
44        match &src.payload {
45            RawData::String(raw) => {
46                data.append(DataField::from_chars(self.raw_key.clone(), raw.clone()));
47            }
48            RawData::Bytes(raw) => {
49                data.append(DataField::from_chars(
50                    self.raw_key.clone(),
51                    String::from_utf8_lossy(raw).into_owned(),
52                ));
53            }
54            RawData::ArcBytes(raw) => {
55                data.append(DataField::from_chars(
56                    self.raw_key.clone(),
57                    String::from_utf8_lossy(raw).into_owned(),
58                ));
59            }
60        }
61        Ok(())
62    }
63}
64
65/// 将原始 payload 复制给指定 rule 的 parser 解析,把产出的字段并入当前 record。
66///
67/// `target` 在构建期由 motor 层注入(按 `rule_name` 解析同包 rule 的 parser);
68/// 未注入时(如 editor/station 解析路径)`proc` 直接 no-op。
69/// 目标 rule 只解析,不执行其自身注解。
70#[derive(Clone)]
71pub struct CopyEventParseAnnotation {
72    pub rule_name: SmolStr,
73    pub target: Option<WplEvaluator>,
74}
75
76impl std::fmt::Debug for CopyEventParseAnnotation {
77    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
78        f.debug_struct("CopyEventParseAnnotation")
79            .field("rule_name", &self.rule_name)
80            .field("target_set", &self.target.is_some())
81            .finish()
82    }
83}
84
85impl AnnotationFunc for CopyEventParseAnnotation {
86    fn proc(&self, src: &SourceEvent, data: &mut DataRecord) -> Result<(), WparseError> {
87        let Some(target) = &self.target else {
88            return Ok(());
89        };
90        // 复制 src.payload 喂给目标 rule 的 parser(proc_ref 避免每条事件 clone payload)
91        let (target_rec, _left) = target.proc_ref(src.event_id, &src.payload, 0)?;
92        // 把目标 rule 解析产出的字段并入当前 record
93        data.merge(target_rec);
94        Ok(())
95    }
96}
97
98#[derive(Clone, Debug)]
99pub enum AnnotationType {
100    Tag(TagAnnotation),
101    Copy(RawCopy),
102    Null(NoopAnnotation),
103    CopyEventParse(CopyEventParseAnnotation),
104}
105
106impl AnnotationFunc for AnnotationType {
107    fn proc(&self, src: &SourceEvent, data: &mut DataRecord) -> Result<(), WparseError> {
108        match self {
109            AnnotationType::Tag(func) => func.proc(src, data),
110            AnnotationType::Null(func) => func.proc(src, data),
111            AnnotationType::Copy(func) => func.proc(src, data),
112            AnnotationType::CopyEventParse(func) => func.proc(src, data),
113        }
114    }
115}
116
117impl AnnotationType {
118    pub fn convert(ann: &Option<AnnFun>) -> Vec<Self> {
119        let mut vec = vec![];
120        if let Some(ann) = ann {
121            if !ann.tags.is_empty() {
122                vec.push(AnnotationType::Tag(TagAnnotation {
123                    args: ann.tags.clone(),
124                }));
125            }
126
127            if let Some((k, v)) = &ann.copy_raw {
128                if k == "name" {
129                    vec.push(AnnotationType::Copy(RawCopy { raw_key: v.clone() }));
130                } else {
131                    vec.push(AnnotationType::Null(NoopAnnotation {}))
132                }
133            }
134
135            if let Some(rule) = &ann.copy_event_parse {
136                // target 留空,构建期由 motor 层按 rule_name 注入目标 rule 的 parser
137                vec.push(AnnotationType::CopyEventParse(CopyEventParseAnnotation {
138                    rule_name: rule.clone(),
139                    target: None,
140                }));
141            }
142        } else {
143            vec.push(AnnotationType::Null(NoopAnnotation {}))
144        }
145        vec
146    }
147}
148
149#[cfg(test)]
150mod tests {
151    use super::*;
152    use crate::pkg::DEFAULT_KEY;
153    use bytes::Bytes;
154    use orion_error::dev::testing::TestAssert;
155    use std::collections::BTreeMap;
156    use std::sync::Arc;
157    use wp_model_core::model::DataRecord;
158    use wp_model_core::raw::RawData;
159    use wp_source_types::{SourceEvent, Tags};
160
161    #[test]
162    fn test_tag_fun() {
163        let ann = AnnFun {
164            tags: BTreeMap::from([("tag_1".into(), "x".into())]),
165            copy_raw: None,
166            copy_event_parse: None,
167        };
168        let tag = AnnotationType::convert(&Some(ann));
169        let mut data = DataRecord::test_value();
170        let src = SourceEvent::new(
171            1,
172            DEFAULT_KEY.to_string(),
173            RawData::String("test".to_string()),
174            Tags::new().into(),
175        );
176        tag.first().unwrap().proc(&src, &mut data).assert();
177        let expected = DataField::from_chars("tag_1", "x");
178        assert_eq!(data.field("tag_1").map(|s| s.as_field()), Some(&expected));
179    }
180
181    #[test]
182    fn test_copy_fun() {
183        let ann = AnnFun {
184            tags: Default::default(),
185            copy_raw: Some(("name".into(), "raw".into())),
186            copy_event_parse: None,
187        };
188        let tag = AnnotationType::convert(&Some(ann));
189        let mut data = DataRecord::test_value();
190        let src = SourceEvent::new(
191            1,
192            DEFAULT_KEY.to_string(),
193            RawData::String("test".to_string()),
194            Tags::new().into(),
195        );
196        tag.first().unwrap().proc(&src, &mut data).unwrap();
197        let expected = DataField::from_chars("raw", "test");
198        assert_eq!(data.field("raw").map(|s| s.as_field()), Some(&expected));
199    }
200
201    #[test]
202    fn test_copy_fun_handles_invalid_utf8_bytes() {
203        let tag = copy_raw_tag("raw");
204        let raw = Bytes::from_static(b"hello \xff\xfe\xc0\xaf");
205        let expected_raw = String::from_utf8_lossy(&raw).into_owned();
206        let mut data = DataRecord::test_value();
207        let src = SourceEvent::new(
208            1,
209            DEFAULT_KEY.to_string(),
210            RawData::Bytes(raw),
211            Tags::new().into(),
212        );
213
214        tag.first().unwrap().proc(&src, &mut data).unwrap();
215
216        let expected = DataField::from_chars("raw", expected_raw);
217        assert_eq!(data.field("raw").map(|s| s.as_field()), Some(&expected));
218    }
219
220    #[test]
221    fn test_copy_fun_handles_invalid_utf8_arc_bytes() {
222        let tag = copy_raw_tag("raw");
223        let raw = Arc::new(b"hello \xff\xfe\xc0\xaf".to_vec());
224        let expected_raw = String::from_utf8_lossy(raw.as_slice()).into_owned();
225        let mut data = DataRecord::test_value();
226        let src = SourceEvent::new(
227            1,
228            DEFAULT_KEY.to_string(),
229            RawData::ArcBytes(raw),
230            Tags::new().into(),
231        );
232
233        tag.first().unwrap().proc(&src, &mut data).unwrap();
234
235        let expected = DataField::from_chars("raw", expected_raw);
236        assert_eq!(data.field("raw").map(|s| s.as_field()), Some(&expected));
237    }
238
239    fn copy_raw_tag(raw_key: &str) -> Vec<AnnotationType> {
240        AnnotationType::convert(&Some(AnnFun {
241            tags: Default::default(),
242            copy_raw: Some(("name".into(), raw_key.into())),
243            copy_event_parse: None,
244        }))
245    }
246
247    /// 构造一个注入了 target parser 的 copy_event_parse 注解,
248    /// target rule 把 payload 整体捕获为 raw 字段。
249    fn copy_event_parse_with_target(rule_code: &str, rule_name: &str) -> Vec<AnnotationType> {
250        let target = WplEvaluator::from_code(rule_code).expect("build target evaluator");
251        let mut funcs = AnnotationType::convert(&Some(AnnFun {
252            tags: Default::default(),
253            copy_raw: None,
254            copy_event_parse: Some(rule_name.into()),
255        }));
256        for ann in &mut funcs {
257            if let AnnotationType::CopyEventParse(c) = ann {
258                c.target = Some(target.clone());
259            }
260        }
261        funcs
262    }
263
264    #[test]
265    fn test_copy_event_parse_merges_target_fields() {
266        // 目标 rule:解析 JSON {"raw":"..."} 产出名为 raw 的 chars 字段
267        let funcs =
268            copy_event_parse_with_target(r#"rule raw_event { (json(chars@raw)) }"#, "raw_event");
269        let mut data = DataRecord::test_value();
270        let src = SourceEvent::new(
271            1,
272            DEFAULT_KEY.to_string(),
273            RawData::String(r#"{ "raw": "payload-content" }"#.to_string()),
274            Tags::new().into(),
275        );
276        funcs.first().unwrap().proc(&src, &mut data).unwrap();
277        let expected = DataField::from_chars("raw", "payload-content");
278        assert_eq!(data.field("raw").map(|s| s.as_field()), Some(&expected));
279    }
280
281    #[test]
282    fn test_copy_event_parse_noop_without_target() {
283        // 未注入 target(editor/station 路径)时应 no-op,不报错也不写字段
284        let funcs = AnnotationType::convert(&Some(AnnFun {
285            tags: Default::default(),
286            copy_raw: None,
287            copy_event_parse: Some("raw_event".into()),
288        }));
289        let mut data = DataRecord::test_value();
290        let src = SourceEvent::new(
291            1,
292            DEFAULT_KEY.to_string(),
293            RawData::String("payload-content".to_string()),
294            Tags::new().into(),
295        );
296        funcs.first().unwrap().proc(&src, &mut data).unwrap();
297        // 未注入 target 时 raw 字段不应存在
298        assert!(data.field("raw").is_none());
299    }
300}