From 140f0d258d2cc5096853298c265694bde6711617 Mon Sep 17 00:00:00 2001 From: rohan Date: Fri, 29 May 2026 16:55:16 +0530 Subject: [PATCH 1/4] refactor: chunked, lock-friendly READ-log purge to keep DB live under load --- .../api/management/commands/purge_app_logs.py | 149 ++++++++++++------ 1 file changed, 104 insertions(+), 45 deletions(-) diff --git a/backend/api/management/commands/purge_app_logs.py b/backend/api/management/commands/purge_app_logs.py index e2742b593..6c77cc6d3 100644 --- a/backend/api/management/commands/purge_app_logs.py +++ b/backend/api/management/commands/purge_app_logs.py @@ -1,11 +1,15 @@ +import time +from datetime import timedelta + from django.core.management.base import BaseCommand, CommandError -from api.models import Organisation, SecretEvent +from django.db import transaction from django.utils import timezone -from datetime import timedelta + +from api.models import Organisation, SecretEvent class Command(BaseCommand): - help = "Purge logs older than a specified number of days for a specific organisation or app." + help = "Purge READ logs older than a specified number of days for an org or app." def add_arguments(self, parser): parser.add_argument("org_name", type=str, help="Name of the organisation") @@ -20,64 +24,119 @@ def add_arguments(self, parser): type=str, help="ID of a specific app to delete logs for (optional)", ) + parser.add_argument( + "--batch-size", + type=int, + default=10_000, + help="Rows deleted per transaction (default: 10000)", + ) + parser.add_argument( + "--sleep-ms", + type=int, + default=500, + help="Pause between batches in ms. Gives autovacuum and replication " + "headroom; raise if replicas lag (default: 500)", + ) + + def _purge_env(self, env_id, env_name, cutoff, batch_size, sleep_ms): + # Per-environment range scan exploits the + # (environment_id, -timestamp) composite index. ORDER BY timestamp ASC + # deletes oldest-first so we walk the cold end of the index. + base_qs = SecretEvent.objects.filter( + environment_id=env_id, + event_type=SecretEvent.READ, + timestamp__lte=cutoff, + ).order_by("timestamp") + + deleted_total = 0 + batch_num = 0 + while True: + batch_num += 1 + with transaction.atomic(): + # SKIP LOCKED so a concurrent purge or any other locker is + # routed around rather than blocking this batch. + batch_ids = list( + base_qs.select_for_update(skip_locked=True) + .values_list("id", flat=True)[:batch_size] + ) + if not batch_ids: + break + + # M2M FK is ON DELETE NO ACTION at the DB level — Django + # handles cascade in the ORM, but _raw_delete bypasses that, + # so we must clear the through-table rows ourselves. + SecretEvent.tags.through.objects.filter( + secretevent_id__in=batch_ids + )._raw_delete("default") + + deleted = SecretEvent.objects.filter( + id__in=batch_ids + )._raw_delete("default") + + deleted_total += deleted + self.stdout.write( + f" {env_name} batch {batch_num}: deleted={deleted} " + f"total={deleted_total}" + ) + if sleep_ms: + time.sleep(sleep_ms / 1000.0) + return deleted_total def handle(self, *args, **options): org_name = options["org_name"] retain_days = options["retain"] app_id = options.get("app_id") + batch_size = options["batch_size"] + sleep_ms = options["sleep_ms"] if retain_days < 0: raise CommandError("The --retain argument must be a non-negative integer.") - if retain_days == 0: - time_cutoff = timezone.now() - else: - time_cutoff = timezone.now() - timedelta(days=retain_days) - - # Only show organization-wide message if no app_id is specified - if not app_id: - self.stdout.write( - f"Deleting logs older than {time_cutoff} (retaining {retain_days} days) for organisation '{org_name}'." - ) + cutoff = ( + timezone.now() if retain_days == 0 + else timezone.now() - timedelta(days=retain_days) + ) try: org = Organisation.objects.get(name=org_name) + except Organisation.DoesNotExist: + raise CommandError(f"Organisation '{org_name}' does not exist.") - # Build the filter dynamically - app_filter = {} - if app_id: - app_filter["id"] = app_id - - apps = org.apps.filter(**app_filter) - - if not apps.exists(): - raise CommandError( - f"No apps found matching the criteria (app_id: {app_id}) in organisation '{org_name}'." - ) - - for app in apps: - logs_to_delete = SecretEvent.objects.filter( - environment__in=app.environments.all(), timestamp__lte=time_cutoff - ).exclude(event_type=SecretEvent.CREATE) - - # Get IDs to delete - log_ids = list(logs_to_delete.values_list("id", flat=True)) - - count = logs_to_delete.count() + app_filter = {"id": app_id} if app_id else {} + apps = list(org.apps.filter(**app_filter)) + if not apps: + raise CommandError( + f"No apps found matching the criteria (app_id: {app_id}) " + f"in organisation '{org_name}'." + ) - # Construct queryset on the M2M through model - m2m_qs = SecretEvent.tags.through.objects.filter( - secretevent_id__in=log_ids - ) - m2m_qs._raw_delete(using="default") # raw delete on M2M table + if not app_id: + self.stdout.write( + f"Deleting READ logs older than {cutoff} " + f"(retaining {retain_days} days) for organisation '{org_name}'." + ) - logs_to_delete._raw_delete("default") + grand_total = 0 + grand_start = time.monotonic() + for app in apps: + app_total = 0 + app_start = time.monotonic() + for env in app.environments.all(): + n = self._purge_env(env.id, env.name, cutoff, batch_size, sleep_ms) + app_total += n self.stdout.write( - f"Deleted {count} logs for app '{app.name}' (id: {app.id})" + f" env '{env.name}' ({env.id}): deleted {n} READ events" ) - + elapsed = time.monotonic() - app_start self.stdout.write( - self.style.SUCCESS("Log deletion completed successfully.") + f"Deleted {app_total} logs for app '{app.name}' (id: {app.id}) " + f"in {elapsed:.1f}s" ) - except Organisation.DoesNotExist: - raise CommandError(f"Organisation '{org_name}' does not exist.") + grand_total += app_total + + elapsed = time.monotonic() - grand_start + self.stdout.write( + self.style.SUCCESS( + f"Log deletion completed: {grand_total} rows in {elapsed:.1f}s." + ) + ) From d07881ca4b720e11ec6049140a5704d49790d122 Mon Sep 17 00:00:00 2001 From: rohan Date: Fri, 29 May 2026 17:10:29 +0530 Subject: [PATCH 2/4] feat(ee): add cloud-only log-retention culling script --- backend/ee/scripts/cull_log_retention.py | 257 +++++++++++++++++++++++ 1 file changed, 257 insertions(+) create mode 100644 backend/ee/scripts/cull_log_retention.py diff --git a/backend/ee/scripts/cull_log_retention.py b/backend/ee/scripts/cull_log_retention.py new file mode 100644 index 000000000..24e1605f5 --- /dev/null +++ b/backend/ee/scripts/cull_log_retention.py @@ -0,0 +1,257 @@ +#!/usr/bin/env python3 +""" +Cloud-only operational script: cull SecretEvent READ-log retention per plan tier. + + Free orgs → retain 24 hours (1 day) + Pro orgs → retain 90 days + Enterprise → skipped (no automated cull) + +Usage (run from /app inside the container): + python ee/scripts/cull_log_retention.py # discover only + python ee/scripts/cull_log_retention.py --count # discover + count rows + python ee/scripts/cull_log_retention.py --apply # execute the purge + python ee/scripts/cull_log_retention.py --plan free --apply + python ee/scripts/cull_log_retention.py --org-id --apply + +Refuses to run when APP_HOST != "cloud" unless --force is given. +""" +import argparse +import json +import os +import sys +import time + +# Bootstrap Django when invoked as a plain script. +_HERE = os.path.dirname(os.path.abspath(__file__)) +_BACKEND = os.path.dirname(os.path.dirname(_HERE)) +sys.path.insert(0, _BACKEND) +os.environ.setdefault("DJANGO_SETTINGS_MODULE", "backend.settings") + +import django # noqa: E402 + +django.setup() + +from django.conf import settings # noqa: E402 +from django.core.management import call_command # noqa: E402 +from django.utils import timezone # noqa: E402 +from datetime import timedelta # noqa: E402 + +from api.models import Organisation, SecretEvent # noqa: E402 + + +PLAN_LABEL = { + Organisation.FREE_PLAN: "Free", + Organisation.PRO_PLAN: "Pro", +} + + +def build_arg_parser(): + p = argparse.ArgumentParser( + description=__doc__, + formatter_class=argparse.RawDescriptionHelpFormatter, + ) + p.add_argument( + "--apply", + action="store_true", + help="Actually run the purge. Without this, the script only reports.", + ) + sizing = p.add_mutually_exclusive_group() + sizing.add_argument( + "--count", + action="store_true", + help="Show exact READ-event row count per org (COUNT(*) via index — " + "accurate, can take minutes on prod-scale data).", + ) + sizing.add_argument( + "--estimate", + action="store_true", + help="Show planner-estimated row count per org (EXPLAIN-based — " + "instant, accuracy depends on how recent ANALYZE was).", + ) + p.add_argument( + "--plan", + choices=["free", "pro", "both"], + default="both", + help="Limit to a single plan tier (default: both).", + ) + p.add_argument( + "--org-id", + help="Limit to a single org id (useful for spot runs).", + ) + p.add_argument( + "--free-retain-days", + type=int, + default=1, + help="Retention for Free orgs in days (default: 1).", + ) + p.add_argument( + "--pro-retain-days", + type=int, + default=90, + help="Retention for Pro orgs in days (default: 90).", + ) + p.add_argument( + "--batch-size", + type=int, + default=10_000, + help="Forwarded to purge_app_logs --batch-size (default: 10000).", + ) + p.add_argument( + "--sleep-ms", + type=int, + default=500, + help="Forwarded to purge_app_logs --sleep-ms (default: 500).", + ) + p.add_argument( + "--force", + action="store_true", + help="Allow running when APP_HOST != 'cloud'.", + ) + return p + + +def select_orgs(args): + plan_filter = { + "free": [Organisation.FREE_PLAN], + "pro": [Organisation.PRO_PLAN], + "both": [Organisation.FREE_PLAN, Organisation.PRO_PLAN], + }[args.plan] + qs = Organisation.objects.filter(plan__in=plan_filter) + if args.org_id: + qs = qs.filter(id=args.org_id) + return list(qs.order_by("plan", "created_at")) + + +def retention_for(org, args): + if org.plan == Organisation.FREE_PLAN: + return args.free_retain_days + if org.plan == Organisation.PRO_PLAN: + return args.pro_retain_days + # Defensive — select_orgs filters Enterprise out, but if --force lets one + # through somehow, refuse to assign retention. + return None + + +def _rows_for_org(org, cutoff): + return SecretEvent.objects.filter( + environment__app__organisation=org, + event_type=SecretEvent.READ, + timestamp__lte=cutoff, + ) + + +def count_rows(org, cutoff): + """Exact COUNT(*) — uses (environment_id, -timestamp) index; slow on big orgs.""" + return _rows_for_org(org, cutoff).count() + + +def estimate_rows(org, cutoff): + """Planner row estimate via EXPLAIN — instant, accuracy depends on ANALYZE freshness.""" + plan = json.loads(_rows_for_org(org, cutoff).explain(format="json")) + return int(plan[0]["Plan"]["Plan Rows"]) + + +def report(orgs, args): + print(f"Discovery — {len(orgs)} org(s) match (plan={args.plan})") + header = f" {'plan':<5} {'org_id':<36} {'name':<30} {'apps':>4} retain" + if args.count: + header += f" {'rows_to_purge':>14}" + elif args.estimate: + header += f" {'rows_estimated':>14}" + print(header) + print(" " + "-" * (len(header) - 2)) + total_apps = 0 + total_rows = 0 + for org in orgs: + retain = retention_for(org, args) + if retain is None: + continue + apps = list(org.apps.all()) + total_apps += len(apps) + line = ( + f" {PLAN_LABEL[org.plan]:<5} {org.id:<36} " + f"{org.name[:30]:<30} {len(apps):>4} {retain}d" + ) + if args.count or args.estimate: + cutoff = timezone.now() - timedelta(days=retain) + n = count_rows(org, cutoff) if args.count else estimate_rows(org, cutoff) + total_rows += n + line += f" {n:>14,}" + print(line) + print(f"\n Total apps: {total_apps}") + if args.count: + print(f" Total rows that would be purged: {total_rows:,}") + elif args.estimate: + print( + f" Total rows estimated: ~{total_rows:,} (planner estimate, ±20% typical)" + ) + + +def apply_purge(orgs, args): + grand_start = time.monotonic() + failures = [] + for i, org in enumerate(orgs, 1): + retain = retention_for(org, args) + if retain is None: + continue + label = PLAN_LABEL[org.plan] + print( + f"\n=== [{i}/{len(orgs)}] {label} org '{org.name}' " + f"(id={org.id}, retain={retain}d) ===", + flush=True, + ) + try: + call_command( + "purge_app_logs", + org.name, + retain=retain, + batch_size=args.batch_size, + sleep_ms=args.sleep_ms, + ) + except Exception as e: + print(f" FAILED: {type(e).__name__}: {e}", file=sys.stderr, flush=True) + failures.append((org.id, org.name, repr(e))) + continue + + elapsed = time.monotonic() - grand_start + print( + f"\nCompleted {len(orgs) - len(failures)}/{len(orgs)} orgs " + f"in {elapsed:.1f}s ({elapsed/60:.1f} min)" + ) + if failures: + print(f"Failures ({len(failures)}):") + for org_id, name, err in failures: + print(f" - {org_id} {name!r}: {err}") + + +def main(): + args = build_arg_parser().parse_args() + + if settings.APP_HOST != "cloud" and not args.force: + print( + f"Refusing to run: APP_HOST={settings.APP_HOST!r}, expected 'cloud'. " + f"Use --force to override (intentionally inconvenient — this script " + f"is for cloud SaaS only).", + file=sys.stderr, + ) + return 2 + + orgs = select_orgs(args) + if not orgs: + print("No orgs match the filter.") + return 0 + + if not args.apply: + report(orgs, args) + print("\nReport-only run. Pass --apply to execute the purge.") + return 0 + + # Brief preview before executing. + report(orgs, args) + print() + apply_purge(orgs, args) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) From f7e0401531262c2bf28b8a01563562387cda3f7c Mon Sep 17 00:00:00 2001 From: rohan Date: Fri, 29 May 2026 19:19:50 +0530 Subject: [PATCH 3/4] feat: report total rows deleted and treat empty orgs as no-op --- .../api/management/commands/purge_app_logs.py | 27 +++++++++++-------- backend/ee/scripts/cull_log_retention.py | 6 ++++- 2 files changed, 21 insertions(+), 12 deletions(-) diff --git a/backend/api/management/commands/purge_app_logs.py b/backend/api/management/commands/purge_app_logs.py index 6c77cc6d3..b5299e23f 100644 --- a/backend/api/management/commands/purge_app_logs.py +++ b/backend/api/management/commands/purge_app_logs.py @@ -35,7 +35,7 @@ def add_arguments(self, parser): type=int, default=500, help="Pause between batches in ms. Gives autovacuum and replication " - "headroom; raise if replicas lag (default: 500)", + "headroom; raise if replicas lag (default: 500)", ) def _purge_env(self, env_id, env_name, cutoff, batch_size, sleep_ms): @@ -56,8 +56,9 @@ def _purge_env(self, env_id, env_name, cutoff, batch_size, sleep_ms): # SKIP LOCKED so a concurrent purge or any other locker is # routed around rather than blocking this batch. batch_ids = list( - base_qs.select_for_update(skip_locked=True) - .values_list("id", flat=True)[:batch_size] + base_qs.select_for_update(skip_locked=True).values_list( + "id", flat=True + )[:batch_size] ) if not batch_ids: break @@ -69,9 +70,9 @@ def _purge_env(self, env_id, env_name, cutoff, batch_size, sleep_ms): secretevent_id__in=batch_ids )._raw_delete("default") - deleted = SecretEvent.objects.filter( - id__in=batch_ids - )._raw_delete("default") + deleted = SecretEvent.objects.filter(id__in=batch_ids)._raw_delete( + "default" + ) deleted_total += deleted self.stdout.write( @@ -93,7 +94,8 @@ def handle(self, *args, **options): raise CommandError("The --retain argument must be a non-negative integer.") cutoff = ( - timezone.now() if retain_days == 0 + timezone.now() + if retain_days == 0 else timezone.now() - timedelta(days=retain_days) ) @@ -105,10 +107,12 @@ def handle(self, *args, **options): app_filter = {"id": app_id} if app_id else {} apps = list(org.apps.filter(**app_filter)) if not apps: - raise CommandError( - f"No apps found matching the criteria (app_id: {app_id}) " - f"in organisation '{org_name}'." - ) + if app_id: + raise CommandError( + f"App with id '{app_id}' not found in organisation '{org_name}'." + ) + self.stdout.write(f"Organisation '{org_name}' has no apps; nothing to do.") + return 0 if not app_id: self.stdout.write( @@ -140,3 +144,4 @@ def handle(self, *args, **options): f"Log deletion completed: {grand_total} rows in {elapsed:.1f}s." ) ) + return grand_total diff --git a/backend/ee/scripts/cull_log_retention.py b/backend/ee/scripts/cull_log_retention.py index 24e1605f5..261f8f3c2 100644 --- a/backend/ee/scripts/cull_log_retention.py +++ b/backend/ee/scripts/cull_log_retention.py @@ -190,6 +190,7 @@ def report(orgs, args): def apply_purge(orgs, args): grand_start = time.monotonic() failures = [] + grand_total_rows = 0 for i, org in enumerate(orgs, 1): retain = retention_for(org, args) if retain is None: @@ -201,13 +202,15 @@ def apply_purge(orgs, args): flush=True, ) try: - call_command( + n = call_command( "purge_app_logs", org.name, retain=retain, batch_size=args.batch_size, sleep_ms=args.sleep_ms, ) + if n: + grand_total_rows += int(n) except Exception as e: print(f" FAILED: {type(e).__name__}: {e}", file=sys.stderr, flush=True) failures.append((org.id, org.name, repr(e))) @@ -218,6 +221,7 @@ def apply_purge(orgs, args): f"\nCompleted {len(orgs) - len(failures)}/{len(orgs)} orgs " f"in {elapsed:.1f}s ({elapsed/60:.1f} min)" ) + print(f"Total rows deleted: {grand_total_rows:,}") if failures: print(f"Failures ({len(failures)}):") for org_id, name, err in failures: From 4e403cb60f6311b5a119714e13b8bb8348b3ca67 Mon Sep 17 00:00:00 2001 From: rohan Date: Mon, 1 Jun 2026 13:35:59 +0530 Subject: [PATCH 4/4] fix: don't return int from purge_app_logs; parse stdout in cull script --- .../api/management/commands/purge_app_logs.py | 3 +-- backend/ee/scripts/cull_log_retention.py | 16 +++++++++++++--- 2 files changed, 14 insertions(+), 5 deletions(-) diff --git a/backend/api/management/commands/purge_app_logs.py b/backend/api/management/commands/purge_app_logs.py index b5299e23f..779b20a1c 100644 --- a/backend/api/management/commands/purge_app_logs.py +++ b/backend/api/management/commands/purge_app_logs.py @@ -112,7 +112,7 @@ def handle(self, *args, **options): f"App with id '{app_id}' not found in organisation '{org_name}'." ) self.stdout.write(f"Organisation '{org_name}' has no apps; nothing to do.") - return 0 + return if not app_id: self.stdout.write( @@ -144,4 +144,3 @@ def handle(self, *args, **options): f"Log deletion completed: {grand_total} rows in {elapsed:.1f}s." ) ) - return grand_total diff --git a/backend/ee/scripts/cull_log_retention.py b/backend/ee/scripts/cull_log_retention.py index 261f8f3c2..779a2586e 100644 --- a/backend/ee/scripts/cull_log_retention.py +++ b/backend/ee/scripts/cull_log_retention.py @@ -16,8 +16,10 @@ Refuses to run when APP_HOST != "cloud" unless --force is given. """ import argparse +import io import json import os +import re import sys import time @@ -201,17 +203,25 @@ def apply_purge(orgs, args): f"(id={org.id}, retain={retain}d) ===", flush=True, ) + captured = io.StringIO() try: - n = call_command( + call_command( "purge_app_logs", org.name, retain=retain, batch_size=args.batch_size, sleep_ms=args.sleep_ms, + stdout=captured, ) - if n: - grand_total_rows += int(n) + output = captured.getvalue() + sys.stdout.write(output) + sys.stdout.flush() + m = re.search(r"Log deletion completed:\s+(\d+)\s+rows", output) + if m: + grand_total_rows += int(m.group(1)) except Exception as e: + sys.stdout.write(captured.getvalue()) + sys.stdout.flush() print(f" FAILED: {type(e).__name__}: {e}", file=sys.stderr, flush=True) failures.append((org.id, org.name, repr(e))) continue