(&mut self, messages: impl Iterator<Item = SharedMessage> + Send)
| 39 | #[async_trait] |
| 40 | impl Processor for Validator { |
| 41 | async fn batch(&mut self, messages: impl Iterator<Item = SharedMessage> + Send) -> Result<()> { |
| 42 | for message in messages { |
| 43 | match message.header().stream_key().name() { |
| 44 | INFO_STREAM => self.data.debugger_infos.push({ |
| 45 | let mut debugger_info = deser_info(&message); |
| 46 | debugger_info.redacted(); |
| 47 | debugger_info |
| 48 | }), |
| 49 | FILE_STREAM => self.data.files.push({ |
| 50 | let mut file: SourceFile = deser(&message); |
| 51 | file.redacted(); |
| 52 | file |
| 53 | }), |
| 54 | BREAKPOINT_STREAM => self.data.breakpoints.push(deser(&message)), |
| 55 | EVENT_STREAM => { |
| 56 | let mut event = EventStream::read_from(message.message().into_bytes().into()); |
| 57 | event.redacted(); |
| 58 | self.data.events.push(event); |
| 59 | } |
| 60 | ALLOCATION_STREAM => {} |
| 61 | _ => anyhow::bail!("Unexpected stream key {}", message.stream_key()), |
| 62 | } |
| 63 | } |
| 64 | Ok(()) |
| 65 | } |
| 66 | |
| 67 | async fn end(&mut self) -> Result<()> { |
| 68 | Ok(()) |
nothing calls this directly
no test coverage detected