mirror of
https://github.com/wahyd4/links.git
synced 2026-08-09 05:06:16 +10:00
497 lines
19 KiB
Python
497 lines
19 KiB
Python
import json
|
|
import logging
|
|
import threading
|
|
import time
|
|
from datetime import timedelta
|
|
|
|
from django.contrib import messages
|
|
from django.db.models import Avg, Count, Q
|
|
from django.db.models.functions import TruncDay, TruncHour, TruncMinute
|
|
from django.http import JsonResponse, StreamingHttpResponse
|
|
from django.shortcuts import get_object_or_404, redirect, render
|
|
from django.urls import reverse_lazy
|
|
from django.utils import timezone
|
|
from django.views.generic import CreateView, DeleteView, TemplateView, UpdateView, View
|
|
|
|
from .detector import detect_threats as run_detect
|
|
from .fetcher import fetch_and_store, ingest_raw
|
|
from .forms import NginxSettingsForm, NginxAlertProfileForm, PasteLogsForm
|
|
from .models import BannedIP, IPGeoCache, NginxAccessLog, NginxAlertProfile, NginxSettings, ThreatAlert
|
|
from .notifications import send_test_telegram
|
|
from .tasks import start_fetch_job, stop_fetch_job, schedule_profile, unschedule_profile
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# ── Time range config ─────────────────────────────────────────────────────────
|
|
|
|
_RANGES = {
|
|
'5m': {'delta': timedelta(minutes=5), 'trunc': TruncMinute, 'group_n': 1, 'label': 'Last 5 min', 'fmt': '%H:%M'},
|
|
'30m': {'delta': timedelta(minutes=30), 'trunc': TruncMinute, 'group_n': 2, 'label': 'Last 30 min', 'fmt': '%H:%M'},
|
|
'6h': {'delta': timedelta(hours=6), 'trunc': TruncMinute, 'group_n': 15, 'label': 'Last 6 hours', 'fmt': '%H:%M'},
|
|
'1d': {'delta': timedelta(hours=24), 'trunc': TruncHour, 'group_n': 1, 'label': 'Last 24 hours', 'fmt': '%H:%M'},
|
|
'7d': {'delta': timedelta(days=7), 'trunc': TruncHour, 'group_n': 4, 'label': 'Last 7 days', 'fmt': '%m/%d %H:%M'},
|
|
'30d': {'delta': timedelta(days=30), 'trunc': TruncDay, 'group_n': 1, 'label': 'Last 30 days', 'fmt': '%m/%d'},
|
|
'90d': {'delta': timedelta(days=90), 'trunc': TruncDay, 'group_n': 3, 'label': 'Last 90 days', 'fmt': '%m/%d'},
|
|
}
|
|
_DEFAULT_RANGE = '30m'
|
|
|
|
|
|
def _get_range(request):
|
|
key = request.GET.get('range', _DEFAULT_RANGE)
|
|
if key not in _RANGES:
|
|
key = _DEFAULT_RANGE
|
|
cfg = _RANGES[key]
|
|
return key, cfg, timezone.now() - cfg['delta']
|
|
|
|
|
|
def _floor_bucket_key(dt, trunc_fn, group_n, fmt):
|
|
"""Return format string for the bucket boundary dt falls into."""
|
|
if trunc_fn is TruncMinute:
|
|
m = (dt.minute // group_n) * group_n
|
|
aligned = dt.replace(minute=m, second=0, microsecond=0)
|
|
elif trunc_fn is TruncHour:
|
|
h = (dt.hour // group_n) * group_n
|
|
aligned = dt.replace(hour=h, minute=0, second=0, microsecond=0)
|
|
else: # TruncDay
|
|
if group_n <= 1:
|
|
aligned = dt.replace(hour=0, minute=0, second=0, microsecond=0)
|
|
else:
|
|
from datetime import date
|
|
epoch = date(2020, 1, 1)
|
|
d = dt.date() if hasattr(dt, 'date') else dt
|
|
days = (d - epoch).days
|
|
floored_date = epoch + timedelta(days=(days // group_n) * group_n)
|
|
aligned = dt.replace(
|
|
year=floored_date.year, month=floored_date.month,
|
|
day=floored_date.day, hour=0, minute=0, second=0, microsecond=0,
|
|
)
|
|
return aligned.strftime(fmt)
|
|
|
|
|
|
# ── Dashboard — live monitoring ───────────────────────────────────────────────
|
|
|
|
class DashboardView(TemplateView):
|
|
template_name = 'nginxmon/dashboard.html'
|
|
|
|
def get_context_data(self, **kwargs):
|
|
ctx = super().get_context_data(**kwargs)
|
|
range_key, cfg, since = _get_range(self.request)
|
|
|
|
logs = NginxAccessLog.objects.filter(timestamp__gte=since)
|
|
stats = {
|
|
'total': logs.count(),
|
|
'errors': logs.filter(status__gte=400).count(),
|
|
'unique_ips': logs.values('remote_addr').distinct().count(),
|
|
'avg_resp_ms': round(
|
|
(logs.aggregate(a=Avg('request_time'))['a'] or 0) * 1000, 1
|
|
),
|
|
}
|
|
|
|
top_services = list(
|
|
logs.values('service')
|
|
.annotate(count=Count('id'))
|
|
.order_by('-count')[:10]
|
|
)
|
|
top_paths = list(
|
|
logs.values('request_uri')
|
|
.annotate(count=Count('id'), errors=Count('id', filter=Q(status__gte=400)))
|
|
.order_by('-count')[:15]
|
|
)
|
|
status_dist = list(
|
|
logs.values('status')
|
|
.annotate(count=Count('id'))
|
|
.order_by('status')
|
|
)
|
|
top_ips = list(
|
|
logs.values('remote_addr', 'country', 'city')
|
|
.annotate(count=Count('id'), errors=Count('id', filter=Q(status__gte=400)))
|
|
.order_by('-count')[:10]
|
|
)
|
|
active_alerts = ThreatAlert.objects.filter(dismissed=False).count()
|
|
profiles = NginxAlertProfile.objects.filter(enabled=True)
|
|
banned_ips = BannedIP.objects.all()[:100]
|
|
|
|
ctx.update({
|
|
'nginx_settings': NginxSettings.get(),
|
|
'stats': stats,
|
|
'top_services': top_services,
|
|
'top_paths': top_paths,
|
|
'status_dist': status_dist,
|
|
'top_ips': top_ips,
|
|
'active_alerts': active_alerts,
|
|
'profiles': profiles,
|
|
'banned_ips': banned_ips,
|
|
'status_dist_json': json.dumps(status_dist),
|
|
'top_paths_json': json.dumps(top_paths),
|
|
'current_range': range_key,
|
|
'range_label': cfg['label'],
|
|
'all_ranges': [(k, v['label']) for k, v in _RANGES.items()],
|
|
})
|
|
return ctx
|
|
|
|
|
|
# ── Settings (one-off k8s config + paste test logs) ──────────────────────────
|
|
|
|
class SettingsView(View):
|
|
template_name = 'nginxmon/settings.html'
|
|
|
|
def _get_settings(self):
|
|
return NginxSettings.get()
|
|
|
|
def get(self, request):
|
|
obj = self._get_settings()
|
|
return render(request, self.template_name, {
|
|
'form': NginxSettingsForm(instance=obj),
|
|
'paste_form': PasteLogsForm(),
|
|
'nginx_settings': obj,
|
|
})
|
|
|
|
def post(self, request):
|
|
obj = self._get_settings()
|
|
if 'save_settings' in request.POST:
|
|
form = NginxSettingsForm(request.POST, instance=obj)
|
|
if form.is_valid():
|
|
saved = form.save()
|
|
stop_fetch_job()
|
|
if saved.enabled:
|
|
start_fetch_job(saved.fetch_interval_seconds)
|
|
messages.success(request, 'Settings saved.')
|
|
return redirect('nginxmon-settings')
|
|
return render(request, self.template_name, {
|
|
'form': form,
|
|
'paste_form': PasteLogsForm(),
|
|
'nginx_settings': obj,
|
|
})
|
|
|
|
if 'paste_logs' in request.POST:
|
|
paste_form = PasteLogsForm(request.POST)
|
|
if paste_form.is_valid():
|
|
count = ingest_raw(paste_form.cleaned_data['log_text'])
|
|
messages.success(request, f'Ingested {count} new log entries.')
|
|
return redirect('nginxmon-settings')
|
|
return render(request, self.template_name, {
|
|
'form': NginxSettingsForm(instance=obj),
|
|
'paste_form': paste_form,
|
|
'nginx_settings': obj,
|
|
})
|
|
|
|
return redirect('nginxmon-settings')
|
|
|
|
|
|
# ── Alert profiles ────────────────────────────────────────────────────────────
|
|
|
|
class ProfileCreateView(CreateView):
|
|
model = NginxAlertProfile
|
|
form_class = NginxAlertProfileForm
|
|
template_name = 'nginxmon/profile_form.html'
|
|
success_url = reverse_lazy('nginxmon-dashboard')
|
|
|
|
def form_valid(self, form):
|
|
response = super().form_valid(form)
|
|
if self.object.enabled:
|
|
schedule_profile(self.object)
|
|
messages.success(self.request, f'Alert profile "{self.object.name}" created.')
|
|
return response
|
|
|
|
|
|
class ProfileUpdateView(UpdateView):
|
|
model = NginxAlertProfile
|
|
form_class = NginxAlertProfileForm
|
|
template_name = 'nginxmon/profile_form.html'
|
|
success_url = reverse_lazy('nginxmon-dashboard')
|
|
|
|
def form_valid(self, form):
|
|
response = super().form_valid(form)
|
|
unschedule_profile(self.object)
|
|
if self.object.enabled:
|
|
schedule_profile(self.object)
|
|
messages.success(self.request, f'Alert profile "{self.object.name}" updated.')
|
|
return response
|
|
|
|
|
|
class ProfileDeleteView(DeleteView):
|
|
model = NginxAlertProfile
|
|
template_name = 'nginxmon/profile_confirm_delete.html'
|
|
success_url = reverse_lazy('nginxmon-dashboard')
|
|
|
|
def form_valid(self, form):
|
|
unschedule_profile(self.object)
|
|
return super().form_valid(form)
|
|
|
|
|
|
class TestTelegramView(View):
|
|
def post(self, request, pk):
|
|
profile = get_object_or_404(NginxAlertProfile, pk=pk)
|
|
result = send_test_telegram(profile)
|
|
if result['ok']:
|
|
messages.success(request, 'Telegram test message sent!')
|
|
else:
|
|
messages.error(request, f'Telegram error: {result["error"]}')
|
|
return redirect('nginxmon-dashboard')
|
|
|
|
|
|
# ── HTMX partials ─────────────────────────────────────────────────────────────
|
|
|
|
class LiveLogsPartialView(View):
|
|
_SORT_MAP = {
|
|
'timestamp': 'timestamp',
|
|
'ip': 'remote_addr',
|
|
'status': 'status',
|
|
'method': 'method',
|
|
'service': 'service',
|
|
'time': 'request_time',
|
|
'bytes': 'body_bytes_sent',
|
|
}
|
|
|
|
def get(self, request):
|
|
qs = NginxAccessLog.objects.all()
|
|
|
|
ip_f = request.GET.get('ip', '').strip()
|
|
status_f = request.GET.get('status', '').strip()
|
|
service_f = request.GET.get('service', '').strip()
|
|
path_f = request.GET.get('path', '').strip()
|
|
sort_col = request.GET.get('sort', 'timestamp')
|
|
sort_ord = request.GET.get('order', 'desc')
|
|
|
|
if ip_f:
|
|
qs = qs.filter(remote_addr=ip_f)
|
|
if status_f:
|
|
try:
|
|
qs = qs.filter(status=int(status_f))
|
|
except ValueError:
|
|
pass
|
|
if service_f:
|
|
qs = qs.filter(service=service_f)
|
|
if path_f:
|
|
qs = qs.filter(request_uri=path_f)
|
|
|
|
db_col = self._SORT_MAP.get(sort_col, 'timestamp')
|
|
qs = qs.order_by(f'{"" if sort_ord == "asc" else "-"}{db_col}')
|
|
|
|
return render(request, 'nginxmon/_live_logs.html', {'logs': qs[:60]})
|
|
|
|
|
|
class AlertsPartialView(View):
|
|
def get(self, request):
|
|
alerts = ThreatAlert.objects.filter(dismissed=False).select_related('profile')[:20]
|
|
return render(request, 'nginxmon/_alerts.html', {'alerts': alerts})
|
|
|
|
|
|
# ── Action views ──────────────────────────────────────────────────────────────
|
|
|
|
class TriggerFetchView(View):
|
|
def post(self, request):
|
|
threading.Thread(target=fetch_and_store, daemon=True).start()
|
|
messages.success(request, 'Log fetch triggered in the background.')
|
|
return redirect('nginxmon-dashboard')
|
|
|
|
|
|
class TriggerDetectView(View):
|
|
def post(self, request):
|
|
for profile in NginxAlertProfile.objects.filter(enabled=True):
|
|
threading.Thread(target=run_detect, args=[profile.pk], daemon=True).start()
|
|
messages.success(request, 'Threat detection triggered for all profiles.')
|
|
return redirect('nginxmon-dashboard')
|
|
|
|
|
|
class DismissAlertView(View):
|
|
def post(self, request, pk):
|
|
alert = get_object_or_404(ThreatAlert, pk=pk)
|
|
alert.dismissed = True
|
|
alert.save(update_fields=['dismissed'])
|
|
if request.headers.get('HX-Request'):
|
|
return JsonResponse({'ok': True})
|
|
return redirect('nginxmon-dashboard')
|
|
|
|
|
|
class UnbanIPView(View):
|
|
"""Mark an IP as manually unbanned — removes it from the ConfigMap block list
|
|
but keeps the BannedIP record so auto-ban never re-bans it."""
|
|
|
|
def post(self, request, pk):
|
|
ban = get_object_or_404(BannedIP, pk=pk)
|
|
ip = ban.ip
|
|
try:
|
|
from links.mini_apps_views import _read_blocked_ips, _write_blocked_ips
|
|
blocked = _read_blocked_ips()
|
|
if ip in blocked:
|
|
blocked.remove(ip)
|
|
_write_blocked_ips(blocked)
|
|
except Exception as exc:
|
|
logger.warning('nginxmon unban: could not update ConfigMap for %s: %s', ip, exc)
|
|
ban.unbanned_at = timezone.now()
|
|
ban.save(update_fields=['unbanned_at'])
|
|
# Clear ban notes from log rows
|
|
NginxAccessLog.objects.filter(remote_addr=ip).update(note='')
|
|
messages.success(request, f'{ip} has been unbanned.')
|
|
return redirect('nginxmon-dashboard')
|
|
|
|
|
|
class IPBanView(View):
|
|
"""Manual IP ban form — adds IPs/CIDRs to the k8s ConfigMap and creates BannedIP records."""
|
|
|
|
def post(self, request):
|
|
raw = request.POST.get('ips', '').strip()
|
|
if not raw:
|
|
messages.error(request, 'No IPs provided.')
|
|
return redirect('nginxmon-dashboard')
|
|
|
|
new_ips = [
|
|
ip.strip()
|
|
for part in raw.replace('\n', ',').split(',')
|
|
for ip in [part.strip()]
|
|
if ip
|
|
]
|
|
|
|
try:
|
|
from links.mini_apps_views import _read_blocked_ips, _write_blocked_ips
|
|
blocked = _read_blocked_ips()
|
|
added = 0
|
|
for ip in new_ips:
|
|
if ip not in blocked:
|
|
blocked.append(ip)
|
|
added += 1
|
|
BannedIP.objects.get_or_create(
|
|
ip=ip,
|
|
defaults=dict(
|
|
ban_source='manual',
|
|
reason='Manually banned from dashboard.',
|
|
request_count=NginxAccessLog.objects.filter(remote_addr=ip).count(),
|
|
),
|
|
)
|
|
_write_blocked_ips(blocked)
|
|
messages.success(request, f'Banned {added} new IP(s). Total blocked: {len(blocked)}.')
|
|
except Exception as exc:
|
|
logger.error('nginxmon IPBanView: %s', exc)
|
|
messages.error(request, f'Failed to ban IPs: {exc}')
|
|
|
|
return redirect('nginxmon-dashboard')
|
|
|
|
|
|
# ── Chart & ingest API ────────────────────────────────────────────────────────
|
|
|
|
class ChartDataView(View):
|
|
def get(self, request):
|
|
range_key, cfg, since = _get_range(request)
|
|
trunc_fn = cfg['trunc']
|
|
group_n = cfg['group_n']
|
|
fmt = cfg['fmt']
|
|
|
|
# Aggregate in DB using the timestamp index + trunc function.
|
|
# Returns at most ~360 rows (6h/minute) — then Python merges into final buckets.
|
|
rows = list(
|
|
NginxAccessLog.objects
|
|
.filter(timestamp__gte=since)
|
|
.annotate(bucket=trunc_fn('timestamp'))
|
|
.values('bucket')
|
|
.annotate(total=Count('id'), errors=Count('id', filter=Q(status__gte=400)))
|
|
.order_by('bucket')
|
|
)
|
|
|
|
# Merge adjacent DB buckets into the desired bucket size.
|
|
# Dict insertion order is preserved (Python 3.7+) so labels stay sorted.
|
|
buckets: dict = {}
|
|
for row in rows:
|
|
bkt = row['bucket']
|
|
if bkt is None:
|
|
continue
|
|
key = _floor_bucket_key(bkt, trunc_fn, group_n, fmt)
|
|
if key not in buckets:
|
|
buckets[key] = {'total': 0, 'errors': 0}
|
|
buckets[key]['total'] += row['total']
|
|
buckets[key]['errors'] += row['errors']
|
|
|
|
return JsonResponse({
|
|
'labels': list(buckets),
|
|
'total': [v['total'] for v in buckets.values()],
|
|
'errors': [v['errors'] for v in buckets.values()],
|
|
'range': range_key,
|
|
'label': cfg['label'],
|
|
})
|
|
|
|
|
|
class GeoStatsView(View):
|
|
"""Returns country-level aggregates + IP bubble data for the geo map."""
|
|
def get(self, request):
|
|
range_key, cfg, since = _get_range(request)
|
|
countries = list(
|
|
NginxAccessLog.objects
|
|
.filter(timestamp__gte=since).exclude(country_code='')
|
|
.values('country_code', 'country')
|
|
.annotate(total=Count('id'), errors=Count('id', filter=Q(status__gte=400)))
|
|
.order_by('-total')
|
|
)
|
|
ip_agg = list(
|
|
NginxAccessLog.objects.filter(timestamp__gte=since)
|
|
.values('remote_addr')
|
|
.annotate(total=Count('id'), errors=Count('id', filter=Q(status__gte=400)))
|
|
.order_by('-total')[:300]
|
|
)
|
|
ip_set = [r['remote_addr'] for r in ip_agg]
|
|
geo_cache = {
|
|
c.ip: (c.lat, c.lon, c.country, c.city)
|
|
for c in IPGeoCache.objects.filter(ip__in=ip_set, lat__isnull=False, is_private=False)
|
|
}
|
|
bubbles = []
|
|
for row in ip_agg:
|
|
c = geo_cache.get(row['remote_addr'])
|
|
if c and c[0] is not None and c[1] is not None:
|
|
bubbles.append({
|
|
'lat': round(float(c[0]), 2),
|
|
'lon': round(float(c[1]), 2),
|
|
'total': row['total'],
|
|
'errors': row['errors'],
|
|
'label': c[3] or c[2] or row['remote_addr'],
|
|
})
|
|
return JsonResponse({'countries': countries, 'bubbles': bubbles, 'range': range_key})
|
|
|
|
|
|
class IngestLogsView(View):
|
|
"""POST endpoint for the paste-logs form (used from SettingsView)."""
|
|
def post(self, request):
|
|
text = request.POST.get('log_text', '')
|
|
count = ingest_raw(text)
|
|
return JsonResponse({'inserted': count})
|
|
|
|
|
|
# ── Ingestion toggle & SSE status stream ─────────────────────────────────────
|
|
|
|
class ToggleIngestionView(View):
|
|
"""POST — flip NginxSettings.enabled and restart/stop the fetch job."""
|
|
def post(self, request):
|
|
settings = NginxSettings.get()
|
|
settings.enabled = not settings.enabled
|
|
settings.save(update_fields=['enabled'])
|
|
stop_fetch_job()
|
|
if settings.enabled:
|
|
start_fetch_job(settings.fetch_interval_seconds)
|
|
return JsonResponse({'enabled': settings.enabled})
|
|
|
|
|
|
class StatusStreamView(View):
|
|
"""SSE endpoint — pushes ingestion status JSON every 2 seconds."""
|
|
def get(self, request):
|
|
def event_stream():
|
|
while True:
|
|
try:
|
|
settings = NginxSettings.get()
|
|
total = NginxAccessLog.objects.count()
|
|
last_fetch = (
|
|
settings.last_fetch_at.isoformat()
|
|
if settings.last_fetch_at else None
|
|
)
|
|
payload = json.dumps({
|
|
'enabled': settings.enabled,
|
|
'last_fetch_at': last_fetch,
|
|
'total_logs': total,
|
|
})
|
|
yield f'data: {payload}\n\n'
|
|
except Exception:
|
|
pass
|
|
time.sleep(2)
|
|
|
|
response = StreamingHttpResponse(event_stream(), content_type='text/event-stream')
|
|
response['Cache-Control'] = 'no-cache'
|
|
response['X-Accel-Buffering'] = 'no'
|
|
return response
|