fault-log-analyzer/tests/test_pipeline.py

69 lines
2.6 KiB
Python

"""故障日志捕获管道测试。"""
import _bootstrap # noqa: F401
import unittest
from fault_log_analyzer.filters import FaultFilter
from fault_log_analyzer.models import FaultFilterRule
from fault_log_analyzer.pipeline import CapturePipeline
from fault_log_analyzer.storage import MemoryStorage
class TestCapturePipeline(unittest.TestCase):
def setUp(self):
self.storage = MemoryStorage()
self.pipeline = CapturePipeline(storage=self.storage)
def test_error_captured_and_classified(self):
log = self.pipeline.process(
{"level": "ERROR", "message": "disk full", "host_id": "h-1"}
)
self.assertIsNotNone(log)
self.assertEqual(log.fault_type, "disk_full")
self.assertEqual(log.host_id, "h-1")
self.assertEqual(self.storage.query_fault_logs()["total"], 1)
def test_info_skipped(self):
self.assertIsNone(self.pipeline.process({"level": "INFO", "message": "hello"}))
self.assertEqual(self.storage.query_fault_logs()["total"], 0)
def test_duplicate_skipped(self):
self.pipeline.process({"level": "ERROR", "message": "disk full"})
second = self.pipeline.process({"level": "ERROR", "message": "disk full"})
self.assertIsNone(second)
self.assertEqual(self.pipeline.duplicates, 1)
self.assertEqual(self.pipeline.captured, 1)
def test_producer_callback(self):
produced = []
pipeline = CapturePipeline(storage=MemoryStorage(), producer=produced.append)
log = pipeline.process({"level": "ERROR", "message": "boom"})
self.assertEqual(produced, [log])
def test_process_batch(self):
logs = self.pipeline.process_batch(
[
{"level": "ERROR", "message": "boom"},
{"level": "INFO", "message": "skip"},
{"level": "FATAL", "message": "fatal thing"},
]
)
self.assertEqual(len(logs), 2)
self.assertEqual(self.pipeline.captured, 2)
def test_custom_filter(self):
f = FaultFilter([FaultFilterRule(name="warn", level="WARN")])
pipeline = CapturePipeline(storage=MemoryStorage(), capture_filter=f)
self.assertIsNotNone(pipeline.process({"level": "WARN", "message": "x"}))
self.assertIsNone(pipeline.process({"level": "ERROR", "message": "x"}))
def test_raw_line_pipeline(self):
line = "2025-01-01 10:18:00 ERROR [svc] no space left on device"
log = self.pipeline.process(line)
self.assertIsNotNone(log)
self.assertEqual(log.service, "svc")
self.assertEqual(log.fault_type, "disk_full")
if __name__ == "__main__":
unittest.main()