Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,9 @@ testpaths = ["tests"]
python_files = ["test_*.py"]
python_classes = ["Test*"]
python_functions = ["test_*"]
markers = [
"integration: tests that load real models (slow, requires network)",
]
addopts = [
"--cov=cordon",
"--cov-report=term-missing",
Expand Down
15 changes: 7 additions & 8 deletions src/cordon/analysis/thresholder.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,14 +38,13 @@ def select_significant(

scores = np.array([sw.score for sw in scored_windows])

# Calculate percentile thresholds
# e.g., min=0.05 (exclude top 5%) -> 95th percentile
# e.g., max=0.15 (include up to 15%) -> 85th percentile
upper_percentile = (1 - config.anomaly_range_min) * 100
lower_percentile = (1 - config.anomaly_range_max) * 100

upper_threshold = np.percentile(scores, upper_percentile)
lower_threshold = np.percentile(scores, lower_percentile)
# anomaly_range_min = fraction to exclude from top (e.g., 0.05 = exclude top 5%)
# anomaly_range_max = cumulative fraction to include (e.g., 0.15 = include up to top 15%)
exclusion_cutoff = (1 - config.anomaly_range_min) * 100
inclusion_floor = (1 - config.anomaly_range_max) * 100

upper_threshold = np.percentile(scores, exclusion_cutoff)
lower_threshold = np.percentile(scores, inclusion_floor)

# Select windows in the range: lower <= score < upper
selected = [
Expand Down
7 changes: 7 additions & 0 deletions src/cordon/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,12 @@ def parse_args() -> argparse.Namespace:
default=None,
help="Device for embedding and scoring (default: auto-detect)",
)
config_group.add_argument(
"--max-line-length",
type=int,
default=None,
help="Maximum characters per line before splitting into virtual lines (default: no limit)",
)
config_group.add_argument(
"--scoring-batch-size",
type=int,
Expand Down Expand Up @@ -300,6 +306,7 @@ def _main_impl() -> None:
try:
config = AnalysisConfig(
window_size=args.window_size,
max_line_length=args.max_line_length,
k_neighbors=args.k_neighbors,
anomaly_percentile=anomaly_percentile,
anomaly_range_min=anomaly_range_min,
Expand Down
10 changes: 10 additions & 0 deletions src/cordon/core/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,15 @@
from cordon.core.types import (
AnalysisResult,
Embedder,
Formatter,
MergedBlock,
Merger,
Reader,
ScoredWindow,
Scorer,
Segmenter,
TextWindow,
Thresholder,
)

__all__ = [
Expand All @@ -15,5 +20,10 @@
"MergedBlock",
"AnalysisResult",
"Embedder",
"Formatter",
"Merger",
"Reader",
"Scorer",
"Segmenter",
"Thresholder",
]
40 changes: 40 additions & 0 deletions src/cordon/core/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,49 @@
from typing import Literal


@dataclass(frozen=True)
class WindowConfig:
"""Configuration for window segmentation.

Attributes:
window_size: Number of lines per sliding window.
max_line_length: Maximum characters per line before splitting
into virtual lines. None disables splitting.
"""

window_size: int = 4
max_line_length: int | None = None


@dataclass(frozen=True)
class ScoringConfig:
"""Configuration for anomaly scoring.

Attributes:
k_neighbors: Number of neighbors for k-NN density scoring.
anomaly_percentile: Fraction of windows to retain (e.g. 0.1 = top 10%).
anomaly_range_min: Lower bound for range mode.
anomaly_range_max: Upper bound for range mode.
device: Compute device for scoring. None for auto-detect.
scoring_batch_size: Batch size for k-NN scoring queries.
"""

k_neighbors: int = 5
anomaly_percentile: float = 0.1
anomaly_range_min: float | None = None
anomaly_range_max: float | None = None
device: Literal["cuda", "mps", "cpu"] | None = None
scoring_batch_size: int | None = None


@dataclass(frozen=True)
class AnalysisConfig:
"""Global configuration for the analysis pipeline.

Attributes:
window_size: Number of lines per sliding window.
max_line_length: Maximum characters per line before splitting into
virtual lines. None disables splitting.
k_neighbors: Number of neighbors for k-NN density scoring.
anomaly_percentile: Fraction of windows to retain (e.g. 0.1 = top 10%).
anomaly_range_min: Lower bound for range mode (e.g. 0.05 = exclude top 5%).
Expand All @@ -33,6 +70,7 @@ class AnalysisConfig:
"""

window_size: int = 4
max_line_length: int | None = None
k_neighbors: int = 5
anomaly_percentile: float = 0.1
anomaly_range_min: float | None = None
Expand Down Expand Up @@ -66,6 +104,8 @@ def _validate_core_params(self) -> None:
raise ValueError("k_neighbors must be >= 1")
if not 0.0 <= self.anomaly_percentile <= 1.0:
raise ValueError("anomaly_percentile must be between 0.0 and 1.0")
if self.max_line_length is not None and self.max_line_length < 1:
raise ValueError("max_line_length must be >= 1 if set")
if self.batch_size < 1:
raise ValueError("batch_size must be >= 1")
if self.scoring_batch_size is not None and self.scoring_batch_size < 1:
Expand Down
49 changes: 49 additions & 0 deletions src/cordon/core/types.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
from collections.abc import Iterable, Iterator, Sequence
from dataclasses import dataclass
from pathlib import Path
from typing import TYPE_CHECKING, Any, Protocol

import numpy as np
Expand Down Expand Up @@ -141,3 +142,51 @@ def score_windows(
List of scored windows
"""
...


class Reader(Protocol):
"""Protocol for log file readers."""

def read_lines(self, file_path: Path) -> Iterator[tuple[int, str]]:
"""Read lines from a file with line number tracking."""
...


class Segmenter(Protocol):
"""Protocol for line-to-window segmentation."""

def segment(
self, lines: Iterator[tuple[int, str]], config: "AnalysisConfig"
) -> Iterator[TextWindow]:
"""Segment lines into text windows."""
...


class Thresholder(Protocol):
"""Protocol for selecting significant windows."""

def select_significant(
self, scored_windows: Sequence[ScoredWindow], config: "AnalysisConfig"
) -> list[ScoredWindow]:
"""Select windows above significance threshold."""
...


class Merger(Protocol):
"""Protocol for merging overlapping windows."""

def merge_windows(self, scored_windows: Sequence[ScoredWindow]) -> list[MergedBlock]:
"""Merge overlapping scored windows into contiguous blocks."""
...


class Formatter(Protocol):
"""Protocol for formatting merged blocks into output."""

def format_blocks(
self,
merged_blocks: Sequence[MergedBlock],
lines: Sequence[tuple[int, str]],
) -> str:
"""Format merged blocks into output string."""
...
38 changes: 20 additions & 18 deletions src/cordon/embedding/__init__.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,18 @@
"""Embedding module for log analysis."""

from typing import TYPE_CHECKING
import importlib
from typing import TYPE_CHECKING, cast

if TYPE_CHECKING:
from cordon.core.config import AnalysisConfig
from cordon.core.types import Embedder

_BACKEND_REGISTRY: dict[str, str] = {
"remote": "cordon.embedding.remote.RemoteEmbedder",
"llama-cpp": "cordon.embedding.llama_cpp.LlamaCppEmbedder",
"sentence-transformers": "cordon.embedding.transformer.TransformerEmbedder",
}


def create_embedder(config: "AnalysisConfig") -> "Embedder":
"""Factory function to create the appropriate embedder for a config.
Expand All @@ -14,27 +21,22 @@ def create_embedder(config: "AnalysisConfig") -> "Embedder":
config: Analysis configuration with backend selection.

Returns:
Embedder instance matching the configured backend.
Embedder instance implementing the Embedder protocol.

Raises:
ValueError: If the backend is not recognized.
"""
if config.backend == "remote":
from cordon.embedding.remote import RemoteEmbedder

return RemoteEmbedder(config)

if config.backend == "llama-cpp":
from cordon.embedding.llama_cpp import LlamaCppEmbedder

return LlamaCppEmbedder(config)

if config.backend == "sentence-transformers":
from cordon.embedding.transformer import TransformerEmbedder

return TransformerEmbedder(config)

raise ValueError(f"Unknown backend: {config.backend}")
class_path = _BACKEND_REGISTRY.get(config.backend)
if class_path is None:
raise ValueError(
f"Unknown backend: {config.backend!r}. "
f"Available: {', '.join(sorted(_BACKEND_REGISTRY))}"
)

module_path, class_name = class_path.rsplit(".", 1)
module = importlib.import_module(module_path)
embedder_class = getattr(module, class_name)
return cast("Embedder", embedder_class(config))


__all__ = ["create_embedder"]
59 changes: 42 additions & 17 deletions src/cordon/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,21 @@
import numpy as np

from cordon.analysis.scorer import DensityAnomalyScorer
from cordon.analysis.thresholder import Thresholder
from cordon.analysis.thresholder import Thresholder as ThresholderImpl
from cordon.core.config import AnalysisConfig
from cordon.core.types import AnalysisResult, ScoredWindow
from cordon.core.types import (
AnalysisResult,
Formatter,
Merger,
Reader,
ScoredWindow,
Scorer,
Segmenter,
Thresholder,
)
from cordon.embedding import create_embedder
from cordon.ingestion.reader import LogFileReader
from cordon.postprocess.formatter import OutputFormatter
from cordon.postprocess.formatter import XmlFormatter
from cordon.postprocess.merger import IntervalMerger
from cordon.segmentation.windower import SlidingWindowSegmenter

Expand All @@ -22,14 +31,36 @@ class SemanticLogAnalyzer:
anomalies highlighted.
"""

def __init__(self, config: AnalysisConfig | None = None) -> None:
"""Initialize the analyzer with configuration.
def __init__(
self,
config: AnalysisConfig | None = None,
*,
reader: Reader | None = None,
segmenter: Segmenter | None = None,
scorer: Scorer | None = None,
thresholder: Thresholder | None = None,
merger: Merger | None = None,
formatter: Formatter | None = None,
) -> None:
"""Initialize the analyzer with configuration and optional custom components.

Args:
config: Analysis configuration (uses defaults if None).
reader: Custom reader implementation.
segmenter: Custom segmenter implementation.
scorer: Custom scorer implementation.
thresholder: Custom thresholder implementation.
merger: Custom merger implementation.
formatter: Custom formatter implementation.
"""
self.config = config if config is not None else AnalysisConfig()
self._embedder = create_embedder(self.config)
self._reader = reader if reader is not None else LogFileReader()
self._segmenter = segmenter if segmenter is not None else SlidingWindowSegmenter()
self._scorer = scorer if scorer is not None else DensityAnomalyScorer()
self._thresholder = thresholder if thresholder is not None else ThresholderImpl()
self._merger = merger if merger is not None else IntervalMerger()
self._formatter = formatter if formatter is not None else XmlFormatter()

def analyze_file(self, file_path: Path) -> str:
"""Analyze a log file and return formatted output.
Expand All @@ -55,37 +86,31 @@ def analyze_file_detailed(self, file_path: Path) -> AnalysisResult:
start_time = time.time()

# stage 1: ingestion
reader = LogFileReader()
lines_list = list(reader.read_lines(file_path))
lines_list = list(self._reader.read_lines(file_path))
total_lines = len(lines_list)

# stage 2: segmentation
segmenter = SlidingWindowSegmenter()
windows = segmenter.segment(iter(lines_list), self.config)
windows = self._segmenter.segment(iter(lines_list), self.config)

# stage 3: vectorization
embedded = list(self._embedder.embed_windows(windows))
total_windows = len(embedded)

# stage 4: scoring
scorer = DensityAnomalyScorer()
scored = scorer.score_windows(embedded, self.config)
scored = self._scorer.score_windows(embedded, self.config)
del embedded

# stage 5: thresholding
thresholder = Thresholder()
significant = thresholder.select_significant(scored, self.config)
significant = self._thresholder.select_significant(scored, self.config)
significant_windows = len(significant)

# stage 6: merging
merger = IntervalMerger()
merged = merger.merge_windows(significant)
merged = self._merger.merge_windows(significant)
merged_blocks = len(merged)
del significant

# stage 7: formatting
formatter = OutputFormatter()
output = formatter.format_blocks(merged, lines_list)
output = self._formatter.format_blocks(merged, lines_list)

# calculate statistics
processing_time = time.time() - start_time
Expand Down
4 changes: 2 additions & 2 deletions src/cordon/postprocess/__init__.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
from cordon.postprocess.formatter import OutputFormatter
from cordon.postprocess.formatter import XmlFormatter
from cordon.postprocess.merger import IntervalMerger

__all__ = ["IntervalMerger", "OutputFormatter"]
__all__ = ["IntervalMerger", "XmlFormatter"]
2 changes: 1 addition & 1 deletion src/cordon/postprocess/formatter.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
from cordon.core.types import MergedBlock


class OutputFormatter:
class XmlFormatter:
"""Generate XML-tagged output with original line content.

This formatter wraps each merged block in XML tags that specify
Expand Down
Loading
Loading