start project
This commit is contained in:
@@ -0,0 +1,148 @@
|
||||
"""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)
|
||||
Reference in New Issue
Block a user