import _bootstrap # noqa: F401 import unittest from fault_log_analyzer.filters import CaptureFilter from fault_log_analyzer.pipeline import CapturePipeline from fault_log_analyzer.storage import InMemoryDedupCache, InMemoryLogSink class TestCapturePipeline(unittest.TestCase): def test_process_raw_fault(self): produced = [] sink = InMemoryLogSink() p = CapturePipeline( capture_filter=CaptureFilter(), dedup=InMemoryDedupCache(), log_sink=sink, fault_producer=produced.append, ) raw = { "timestamp": "2025-01-01T10:00:00Z", "level": "ERROR", "message": "disk full", "host_id": "h-1", } fl = p.process_raw(raw) self.assertIsNotNone(fl) self.assertEqual(fl.host_id, "h-1") self.assertEqual(fl.level, "ERROR") self.assertEqual(len(sink.search()), 1) self.assertEqual(len(produced), 1) def test_process_raw_non_fault(self): p = CapturePipeline(capture_filter=CaptureFilter()) self.assertIsNone(p.process_raw({"level": "INFO", "message": "hello"})) def test_dedup_same_fingerprint(self): p = CapturePipeline(capture_filter=CaptureFilter(), dedup=InMemoryDedupCache()) raw = {"level": "ERROR", "message": "disk full"} self.assertIsNotNone(p.process_raw(raw)) self.assertIsNone(p.process_raw(raw)) def test_process_batch(self): p = CapturePipeline(capture_filter=CaptureFilter()) out = p.process_batch( [ {"level": "ERROR", "message": "disk full"}, {"level": "INFO", "message": "ok"}, ] ) self.assertEqual(len(out), 1) self.assertEqual(out[0].message, "disk full") if __name__ == "__main__": unittest.main()