Files
Kavosh/app/services/log_processor.py
T
2026-08-07 21:17:17 +03:30

149 lines
5.7 KiB
Python

"""Core log-file processing pipeline (Chapter 07), factored out of
app/cli.py so it has exactly one implementation shared by:
- the optional `flask process-logs` CLI command (for anyone who still
wants cron), and
- the automatic background-thread trigger fired right after upload and
opportunistically on page load (app/services/background.py) — the
no-cron-required simplification.
process_one_batch() keeps the same bounded-batch, checkpointed-commit
contract as the original cron design (Ch03 rule 1/4/9): it still never
loads a whole file into memory, still commits progress incrementally, and
still stops after `batch_size` lines so a very large file doesn't hold one
enormous open transaction — the only thing that changed is *who* calls it
again for the next batch.
"""
from __future__ import annotations
import gzip
import itertools
from datetime import date
from pathlib import Path
from sqlalchemy import insert
from app.extensions import db
from app.models.log_entry import LogEntry
from app.models.log_file import LogFile
from app.services import aggregator
from app.services.classification import classify_entry, write_batch_side_effects
from app.services.log_parser import ParsedEntry, get_parser_for
from app.utils.upload_paths import upload_path_for
def open_log_stream(path: Path):
"""Open a log file, transparently decompressing .gz, one line at a time."""
if path.suffix == ".gz":
return gzip.open(path, mode="rt", encoding="utf-8", errors="replace")
return path.open("r", encoding="utf-8", errors="replace")
def count_total_lines(path: Path) -> int:
"""One streaming pass to count lines (bounded memory — never the whole
file at once, Ch03 rule 1), so the UI can show a real
processed/total percentage instead of an indeterminate spinner.
This is an extra sequential read of the file beyond the parse pass
itself. That cost is now worth paying: processing used to be silently
triggered by cron with no live audience watching, but now it runs
automatically right after upload while the admin is looking at a
progress bar — the UX value of a real percentage justifies the extra
I/O pass.
"""
with open_log_stream(path) as fh:
return sum(1 for _ in fh)
def process_one_batch(log_file: LogFile, batch_size: int) -> None:
"""Parse up to `batch_size` lines from `log_file`'s current checkpoint.
Leaves status as "processing" (with progress already committed) if
the file isn't finished yet — the caller decides whether to invoke
this again: the CLI calls it once per pending file per invocation
(unchanged cron-tick semantics); the background thread loops it until
the file reaches "done"/"error".
"""
path = upload_path_for(log_file)
if not path.exists():
log_file.status = "error"
log_file.error_message = f"Upload file missing on disk: {path}"
db.session.commit()
return
if log_file.total_lines is None:
# First pickup of this file — count once, not on every batch.
log_file.total_lines = count_total_lines(path)
log_file.status = "processing"
db.session.commit()
parser = get_parser_for(log_file)
lines_seen_this_run = 0
skipped_this_run = 0
dates_touched: set[date] = set()
with open_log_stream(path) as fh:
remainder = itertools.islice(fh, log_file.processed_lines, None)
batch_entries: list[ParsedEntry] = []
for line in remainder:
entry = parser.parse_line(line)
if entry is not None:
batch_entries.append(entry)
dates_touched.add(entry.timestamp.date())
else:
skipped_this_run += 1 # malformed line — skip, don't abort the batch (Ch07)
log_file.processed_lines += 1
lines_seen_this_run += 1
if len(batch_entries) >= 500:
_bulk_insert_entries(log_file.id, batch_entries)
batch_entries = []
if lines_seen_this_run >= batch_size:
_record_skip_count(log_file, skipped_this_run)
db.session.commit() # persist checkpoint; resumable if interrupted
return
if batch_entries:
_bulk_insert_entries(log_file.id, batch_entries)
if dates_touched:
aggregator.compute_rollups_for_range(min(dates_touched), max(dates_touched))
_record_skip_count(log_file, skipped_this_run)
log_file.status = "done"
db.session.commit()
def _record_skip_count(log_file: LogFile, skipped_this_run: int) -> None:
"""Informational only — doesn't touch `status`.
STOPGAP (flagged in Ch07): log_files has no dedicated skip-counter
column, so this overwrites error_message with the latest run's count
rather than accumulating across runs.
"""
if skipped_this_run:
log_file.error_message = f"{skipped_this_run} unparsable line(s) skipped in the most recent parse run."
def _bulk_insert_entries(log_file_id: int, entries: list[ParsedEntry]) -> None:
"""Bulk-insert log_entries (Ch03 rule 4), classifying is_bot/flagged
per line, then derive bot_hits/suspicious_events/ip_registry (Ch07).
"""
if not entries:
return
payload = []
for e in entries:
classification = classify_entry(e)
payload.append({
"log_file_id": log_file_id, "timestamp": e.timestamp, "ip": e.ip,
"method": e.method, "path": e.path, "status_code": e.status_code,
"bytes_sent": e.bytes_sent, "referrer": e.referrer, "user_agent": e.user_agent,
"is_bot": classification.is_bot, "flagged": classification.flagged,
})
db.session.execute(insert(LogEntry.__table__), payload)
db.session.commit()
write_batch_side_effects(log_file_id, entries)