Files
2026-08-07 21:17:17 +03:30

200 lines
8.5 KiB
Python

"""Rollup computation (Chapter 06 + Method A extensions from Ch08/09/10).
Called once per parsed file (app/cli.py::process_logs) for the dates it
touched, ad hoc via `flask rollup` for a manual recompute, and now also
after a file deletion (app/services/file_deletion.py) for whatever dates
the deleted file touched. Every _upsert_* function scans log_entries
inside this background batch job, never at request time — that's what
makes Ch03 rule 6 compliance possible.
CORRECTNESS FIX: every rollup writer below now deletes a day's existing
rows before writing whatever the fresh scan finds (including writing
nothing, if a day now has zero data). Three of the five writers
previously only ever upserted-when-present and silently left stale rows
behind when a day's data disappeared — unreachable before file deletion
existed (rollups only ever grew), but a real correctness bug once
deletion makes "this day now has less data than before" possible. Only
the two per-IP writers already had this right (Ch10 follow-up); the
other three are fixed here to match.
"""
from __future__ import annotations
from collections import defaultdict
from datetime import date, datetime, time, timedelta
from sqlalchemy import case, func
from app.extensions import db
from app.models.browser_stats import BrowserStatsDaily
from app.models.human_path_stats import HumanPathStatsDaily
from app.models.ip_traffic_stats import IpPathStatsDaily, IpStatusStatsDaily
from app.models.log_entry import LogEntry
from app.models.referrer_stats import ReferrerStatsDaily
from app.models.request_stats import RequestStatsDaily, RequestStatsHourly
from app.services.blocklist import refresh_blocklist_suggestions
from app.services.referrer import referrer_domain
from app.services.ua_classifier import classify_browser, classify_os
from app.utils.http_status import status_bucket
TOP_PATHS_PER_IP_PER_DAY = 15 # bounds ip_path_stats_daily row growth (Ch10 follow-up)
def compute_rollups_for_range(start: date, end: date) -> None:
"""Recompute every rollup for each day in [start, end]. Site-wide,
not per-file (Ch01: single site) — recomputing from scratch per day
avoids double-counting when two uploads cover the same period, and
correctly shrinks a day's numbers back down when a file covering
that day is deleted.
"""
current = start
while current <= end:
_upsert_hourly(current)
_upsert_daily(current)
_upsert_per_line_derived_stats(current)
refresh_blocklist_suggestions(current)
current += timedelta(days=1)
def _day_bounds(day: date) -> tuple[datetime, datetime]:
start = datetime.combine(day, time.min)
return start, start + timedelta(days=1)
def _upsert_hourly(day: date) -> None:
start, end = _day_bounds(day)
rows = (
db.session.query(
func.strftime("%Y-%m-%d %H:00:00", LogEntry.timestamp).label("date_hour"),
LogEntry.path,
LogEntry.status_code,
func.count().label("count"),
func.coalesce(func.sum(LogEntry.bytes_sent), 0).label("bytes_sent_sum"),
)
.filter(LogEntry.timestamp >= start, LogEntry.timestamp < end)
.group_by("date_hour", LogEntry.path, LogEntry.status_code)
.all()
)
# Delete-then-insert: replaces the day's hourly rows entirely,
# including leaving none behind if `rows` is now empty (e.g. the
# only file covering this day was just deleted).
db.session.query(RequestStatsHourly).filter(
RequestStatsHourly.date_hour >= start, RequestStatsHourly.date_hour < end
).delete()
if rows:
payload = [
{
"date_hour": datetime.strptime(r.date_hour, "%Y-%m-%d %H:%M:%S"),
"path": r.path,
"status_code": r.status_code,
"count": r.count,
"bytes_sent_sum": r.bytes_sent_sum,
}
for r in rows
]
db.session.execute(RequestStatsHourly.__table__.insert(), payload)
db.session.commit()
def _upsert_daily(day: date) -> None:
start, end = _day_bounds(day)
result = (
db.session.query(
func.count().label("count"),
func.count(func.distinct(LogEntry.ip)).label("unique_ips"),
func.coalesce(func.sum(LogEntry.bytes_sent), 0).label("bytes_sum"),
func.coalesce(func.sum(case((LogEntry.status_code >= 400, 1), else_=0)), 0).label("error_count"),
)
.filter(LogEntry.timestamp >= start, LogEntry.timestamp < end)
.one()
)
# Delete-then-insert: if this day now has zero entries (its only
# contributing file was deleted), the stale row is removed rather
# than left behind — no rollup row is better than a wrong one.
db.session.query(RequestStatsDaily).filter(RequestStatsDaily.date == day).delete()
if result.count > 0:
db.session.execute(
RequestStatsDaily.__table__.insert(),
{
"date": day,
"count": result.count,
"unique_ips": result.unique_ips,
"bytes_sum": result.bytes_sum,
"error_count": result.error_count,
},
)
db.session.commit()
def _upsert_per_line_derived_stats(day: date) -> None:
"""Referrer domain, browser/OS, human-only path counts (Ch08/09), and
per-IP path/status counts (Ch10 follow-up) — one streamed pass over
log_entries (Ch03 rule 1: bounded per-chunk memory via yield_per,
never the whole day loaded at once). Every table here uses the same
delete-then-insert pattern so a day's rows are fully replaced by
whatever the fresh scan finds, including nothing.
"""
start, end = _day_bounds(day)
referrer_counts: dict[str, int] = defaultdict(int)
browser_counts: dict[tuple[str, str], int] = defaultdict(int)
human_path_counts: dict[str, int] = defaultdict(int)
ip_path_counts: dict[str, dict[str, int]] = defaultdict(lambda: defaultdict(int))
ip_status_counts: dict[str, dict[str, int]] = defaultdict(lambda: defaultdict(int))
query = (
db.session.query(
LogEntry.referrer, LogEntry.user_agent, LogEntry.is_bot,
LogEntry.path, LogEntry.ip, LogEntry.status_code,
)
.filter(LogEntry.timestamp >= start, LogEntry.timestamp < end)
)
for referrer, user_agent, is_bot, path, ip, status_code in query.yield_per(1000):
domain = referrer_domain(referrer)
if domain:
referrer_counts[domain] += 1
if not is_bot: # Ch08: bot traffic excluded from human browser/OS breakdown
browser_counts[(classify_browser(user_agent), classify_os(user_agent))] += 1
human_path_counts[path] += 1
ip_path_counts[ip][path] += 1
ip_status_counts[ip][status_bucket(status_code)] += 1
db.session.query(ReferrerStatsDaily).filter(ReferrerStatsDaily.date == day).delete()
if referrer_counts:
payload = [{"date": day, "referrer_domain": d, "count": c} for d, c in referrer_counts.items()]
db.session.execute(ReferrerStatsDaily.__table__.insert(), payload)
db.session.query(BrowserStatsDaily).filter(BrowserStatsDaily.date == day).delete()
if browser_counts:
payload = [{"date": day, "browser": b, "os": o, "count": c} for (b, o), c in browser_counts.items()]
db.session.execute(BrowserStatsDaily.__table__.insert(), payload)
db.session.query(HumanPathStatsDaily).filter(HumanPathStatsDaily.date == day).delete()
if human_path_counts:
payload = [{"date": day, "path": p, "count": c} for p, c in human_path_counts.items()]
db.session.execute(HumanPathStatsDaily.__table__.insert(), payload)
db.session.query(IpPathStatsDaily).filter(IpPathStatsDaily.date == day).delete()
if ip_path_counts:
payload = []
for ip, paths in ip_path_counts.items():
top_paths = sorted(paths.items(), key=lambda kv: kv[1], reverse=True)[:TOP_PATHS_PER_IP_PER_DAY]
payload.extend({"date": day, "ip": ip, "path": p, "count": c} for p, c in top_paths)
if payload:
db.session.execute(IpPathStatsDaily.__table__.insert(), payload)
db.session.query(IpStatusStatsDaily).filter(IpStatusStatsDaily.date == day).delete()
if ip_status_counts:
payload = [
{"date": day, "ip": ip, "status_bucket": bucket, "count": c}
for ip, buckets in ip_status_counts.items()
for bucket, c in buckets.items()
]
db.session.execute(IpStatusStatsDaily.__table__.insert(), payload)
db.session.commit()