Files
links/nginxmon/fetcher.py
T
2026-04-01 15:51:20 +11:00

182 lines
5.9 KiB
Python

"""
Fetch nginx ingress logs — from kubectl (production) or a local file (dev/testing).
Deduplicates by request_id so overlapping windows don't double-insert.
"""
import logging
import subprocess
from datetime import timedelta
from django.db import transaction
from django.utils import timezone as dj_tz
from .models import NginxSettings, NginxAccessLog
from .parser import parse_lines
from .geo import enrich_geo_batch
logger = logging.getLogger(__name__)
_OVERLAP = 60 # extra seconds to avoid missing entries near boundaries
def fetch_and_store() -> int:
"""
Pull new log lines, parse, deduplicate, and save.
Returns the number of new rows inserted.
"""
settings = NginxSettings.get()
if not settings.enabled:
return 0
raw = _read_file(settings) if settings.log_file_path else _kubectl_logs(settings)
if raw is None:
return 0
entries = parse_lines(raw)
if not entries:
_touch(settings)
return 0
since_seconds = settings.fetch_interval_seconds + _OVERLAP
existing_ids = set(
NginxAccessLog.objects.filter(
timestamp__gte=dj_tz.now() - timedelta(seconds=since_seconds + 10),
).values_list('request_id', flat=True)
)
new_logs, seen = [], set()
for e in entries:
key = e['request_id'] or (
f"{e['timestamp'].isoformat()}|{e['remote_addr']}|{e['request_uri']}"
)
if key in existing_ids or key in seen:
continue
seen.add(key)
new_logs.append(NginxAccessLog(
timestamp=e['timestamp'],
remote_addr=e['remote_addr'],
method=e['method'],
request_uri=e['request_uri'],
protocol=e['protocol'],
status=e['status'],
body_bytes_sent=e['body_bytes_sent'],
http_referer=e['http_referer'],
http_user_agent=e['http_user_agent'],
request_length=e['request_length'],
request_time=e['request_time'],
service=e['service'],
upstream_addr=e['upstream_addr'],
upstream_response_time=e['upstream_response_time'],
upstream_status=e['upstream_status'],
request_id=e['request_id'],
))
if new_logs:
with transaction.atomic():
NginxAccessLog.objects.bulk_create(new_logs, ignore_conflicts=True)
logger.info('nginxmon: inserted %d new log entries', len(new_logs))
enrich_geo_batch(new_logs)
_touch(settings)
return len(new_logs)
def ingest_raw(text: str) -> int:
"""
Parse and store log lines from a raw string (used by the paste-logs UI
and the management command). Returns the number of new rows inserted.
"""
entries = parse_lines(text)
if not entries:
return 0
existing_ids = set(
NginxAccessLog.objects.values_list('request_id', flat=True)
)
new_logs, seen = [], set()
for e in entries:
key = e['request_id'] or (
f"{e['timestamp'].isoformat()}|{e['remote_addr']}|{e['request_uri']}"
)
if key in existing_ids or key in seen:
continue
seen.add(key)
new_logs.append(NginxAccessLog(
timestamp=e['timestamp'],
remote_addr=e['remote_addr'],
method=e['method'],
request_uri=e['request_uri'],
protocol=e['protocol'],
status=e['status'],
body_bytes_sent=e['body_bytes_sent'],
http_referer=e['http_referer'],
http_user_agent=e['http_user_agent'],
request_length=e['request_length'],
request_time=e['request_time'],
service=e['service'],
upstream_addr=e['upstream_addr'],
upstream_response_time=e['upstream_response_time'],
upstream_status=e['upstream_status'],
request_id=e['request_id'],
))
if new_logs:
with transaction.atomic():
NginxAccessLog.objects.bulk_create(new_logs, ignore_conflicts=True)
logger.info('nginxmon: ingested %d log entries from raw text', len(new_logs))
enrich_geo_batch(new_logs)
return len(new_logs)
def _kubectl_logs(settings: NginxSettings) -> str | None:
since = settings.fetch_interval_seconds + _OVERLAP
cmd = [
'kubectl', 'logs',
'-n', settings.namespace,
'-l', settings.pod_label,
'--container', settings.container,
f'--since={since}s',
'--timestamps=false',
'--prefix=true',
'--max-log-requests=10',
]
try:
result = subprocess.run(cmd, capture_output=True, text=True, timeout=30)
if result.returncode != 0:
logger.error('kubectl logs failed: %s', result.stderr[:500])
return None
return result.stdout
except FileNotFoundError:
logger.error('kubectl not found in PATH')
return None
except subprocess.TimeoutExpired:
logger.error('kubectl logs timed out')
return None
except Exception as exc:
logger.error('kubectl logs error: %s', exc)
return None
def _read_file(settings: NginxSettings) -> str | None:
try:
with open(settings.log_file_path, 'r', encoding='utf-8', errors='replace') as f:
return f.read()
except FileNotFoundError:
logger.error('nginxmon: log file not found: %s', settings.log_file_path)
return None
except Exception as exc:
logger.error('nginxmon: error reading log file: %s', exc)
return None
def _touch(settings: NginxSettings):
NginxSettings.objects.filter(pk=settings.pk).update(last_fetch_at=dj_tz.now())
def cleanup_old_logs(days: int = 7):
from datetime import timedelta
cutoff = dj_tz.now() - timedelta(days=days)
deleted, _ = NginxAccessLog.objects.filter(timestamp__lt=cutoff).delete()
if deleted:
logger.info('nginxmon: pruned %d old log entries (>%d days)', deleted, days)