Back to home page

EIC code displayed by LXR

 
 

    


File indexing completed on 2026-09-28 09:37:31

0001 #!/usr/bin/env python3
0002 """Index of the devcloud stage-out bucket, learned from the bucket's own events.
0003 
0004 The gateway knows what is in the bucket by being told: S3 emits an
0005 ObjectCreated event per report object (and per status write, counted by
0006 the hour rather than indexed) to the epic-stageout-events queue,
0007 this tool drains the queue into a local index, and the watch and the
0008 nightly cross-check read the index instead of listing S3. Listing a
0009 production-scale prefix every few minutes would cost about as much as
0010 the reporting itself; a queue costs a dollar a month
0011 (swf-epicprod docs/JOB_REPORTING.md).
0012 
0013 Standard library plus boto3 for the drain and the AWS CLI for the rest,
0014 both on this host's account credentials. No venv, no web-tier coupling:
0015 the tool owns its state and other things read it. The drain is the one
0016 hot path: one in-process SQS client, ten messages received and ten
0017 deleted per call. A CLI process per message cost 0.7 s of CPU each and at
0018 campaign rate stacked the five-minute runs 25 deep (2026-09-15).
0019 
0020 Modes:
0021   drain      receive events, record objects, delete the messages
0022   reconcile  authoritative count by listing, compared with the index
0023   apply      apply the sweeper's spooled pass records to the index
0024   guard      nightly: is this an explosion, and pull the plug if so
0025   summary    the state as JSON, for the Capcom tile and for a human
0026   wipe       by hand: delete the prefix as of now, recorded as a pass
0027 Cron on ec2dev, admin's crontab: drain every five minutes under
0028 `flock -n` and `nice -n 19` (a run that finds the lock held exits; runs
0029 never stack), reconcile at 03:40 UTC, guard at 03:50 UTC.
0030 """
0031 
0032 import argparse
0033 from datetime import datetime
0034 import json
0035 import os
0036 import re
0037 import resource
0038 import sqlite3
0039 import subprocess
0040 import sys
0041 import time
0042 
0043 BUCKET = os.environ.get("STAGEOUT_BUCKET", "epic-devcloud-stageout")
0044 PREFIX = os.environ.get("STAGEOUT_PREFIX", "reports/")
0045 QUEUE_URL = os.environ.get(
0046     "STAGEOUT_QUEUE_URL",
0047     "https://sqs.us-east-1.amazonaws.com/962718900486/epic-stageout-events")
0048 DB_PATH = os.environ.get(
0049     "STAGEOUT_INDEX_DB", "/home/admin/data/stageout-index.sqlite")
0050 REGION = os.environ.get("AWS_DEFAULT_REGION", "us-east-1")
0051 # Objects per hour above which the watch calls it a runaway. A defect that
0052 # posts in a loop inside one job is the unbounded risk the design names;
0053 # the fleet itself is finite.
0054 RATE_CEILING_PER_HOUR = int(os.environ.get("STAGEOUT_RATE_CEILING", "20000"))
0055 # A job writing more than this many objects is looping, whatever the fleet
0056 # total says.
0057 PER_JOB_CEILING = int(os.environ.get("STAGEOUT_PER_JOB_CEILING", "50"))
0058 
0059 # The nightly guard's ceilings. Expected production is of the order of
0060 # 100,000 objects a day at a few kilobytes each; these are the levels at
0061 # which the traffic stops being production and starts being a defect.
0062 REPORTING_USER = os.environ.get("STAGEOUT_REPORTING_USER", "epic-job-reporter")
0063 OBJECT_SOFT_CEILING_PER_DAY = int(
0064     os.environ.get("STAGEOUT_OBJECT_SOFT", "500000"))
0065 OBJECT_HARD_CEILING_PER_DAY = int(
0066     os.environ.get("STAGEOUT_OBJECT_HARD", "1000000"))
0067 BYTES_SOFT_CEILING_PER_DAY = int(os.environ.get("STAGEOUT_BYTES_SOFT",
0068                                                 str(5 * 10 ** 9)))
0069 BYTES_HARD_CEILING_PER_DAY = int(os.environ.get("STAGEOUT_BYTES_HARD",
0070                                                 str(50 * 10 ** 9)))
0071 # A day that is this many times the recent norm is an explosion whatever
0072 # the absolute ceilings say; the floor keeps a quiet week from tripping it.
0073 RELATIVE_MULTIPLE = int(os.environ.get("STAGEOUT_RELATIVE_MULTIPLE", "20"))
0074 RELATIVE_FLOOR = int(os.environ.get("STAGEOUT_RELATIVE_FLOOR", "5000"))
0075 
0076 # The sweep runs inside the perimeter, on the production operations agent,
0077 # because every decision it makes needs the inside: the payload digest,
0078 # the validation that an id names a real ePIC production job, and the
0079 # record it files into. It deletes what it takes and reports what it
0080 # deleted here, since S3 tells a bucket's owner nothing about another
0081 # party's deletions and a count cannot say which rows to retire.
0082 PASS_SPOOL = os.environ.get("STAGEOUT_PASS_SPOOL",
0083                             "/var/lib/stageout/passes")
0084 # The sweeper reports hourly. Three missed hours is comfortably clear of one
0085 # slow pass and still tight against a dead sweeper (agreed 2026-09-06).
0086 PASS_STALL_SECONDS = int(os.environ.get("STAGEOUT_PASS_STALL", str(3 * 3600)))
0087 # A listing sees an object the instant it lands; its event reaches the
0088 # queue seconds to a minute later and the drain runs every five minutes.
0089 # An unindexed object younger than this is not yet heard, not unheard
0090 # (2026-09-11: two objects written in the cross-check's minute read as
0091 # unheard for a day).
0092 RECONCILE_GRACE_SECONDS = int(os.environ.get("STAGEOUT_RECONCILE_GRACE",
0093                                              str(15 * 60)))
0094 
0095 SCHEMA = """
0096 CREATE TABLE IF NOT EXISTS objects (
0097     key       TEXT PRIMARY KEY,
0098     size      INTEGER NOT NULL,
0099     etime     TEXT NOT NULL,
0100     subject   TEXT,
0101     seen_at   REAL NOT NULL
0102 );
0103 CREATE INDEX IF NOT EXISTS objects_seen ON objects (seen_at);
0104 CREATE INDEX IF NOT EXISTS objects_subject ON objects (subject);
0105 CREATE TABLE IF NOT EXISTS state (k TEXT PRIMARY KEY, v TEXT NOT NULL);
0106 CREATE TABLE IF NOT EXISTS passes (
0107     pass_id      TEXT PRIMARY KEY,
0108     reporter     TEXT,
0109     outcome      TEXT NOT NULL,
0110     reason       TEXT,
0111     window_from  TEXT,
0112     window_to    TEXT,
0113     jobs_filed   INTEGER NOT NULL,
0114     keys_read    INTEGER NOT NULL,
0115     keys_unread  INTEGER NOT NULL,
0116     retired      INTEGER NOT NULL,
0117     received_at  TEXT,
0118     applied_at   REAL NOT NULL
0119 );
0120 CREATE INDEX IF NOT EXISTS passes_applied ON passes (applied_at);
0121 CREATE TABLE IF NOT EXISTS daily (
0122     day       TEXT PRIMARY KEY,
0123     objects   INTEGER NOT NULL,
0124     bytes     INTEGER NOT NULL,
0125     at        REAL NOT NULL
0126 );
0127 CREATE TABLE IF NOT EXISTS status_writes (
0128     hour      INTEGER PRIMARY KEY,
0129     writes    INTEGER NOT NULL,
0130     bytes     INTEGER NOT NULL
0131 );
0132 """
0133 
0134 # status/<pandaid>.json is one object per job, overwritten while the job
0135 # runs in PanDA debug mode (payload 0.22.0, call home; at most 150 writes a
0136 # job). Overwrites are the cost, so these events are counted per hour as
0137 # writes rather than indexed as objects; the guard adds them to its day.
0138 STATUS_PREFIX = "status/"
0139 
0140 # reports/<subject>/<sequence>.json — the subject is the PanDA job id once
0141 # the payload carries it; anything else is recorded as it arrives.
0142 SUBJECT_RE = re.compile(r"^reports/([^/]+)/")
0143 
0144 
0145 def aws(*args, parse=True):
0146     """Run one AWS CLI call. A failure is raised with its stderr rather
0147     than swallowed: a silent indexer is worse than none."""
0148     proc = subprocess.run(("aws",) + args, capture_output=True, text=True)
0149     if proc.returncode != 0:
0150         raise RuntimeError(
0151             f"aws {' '.join(args[:3])} failed ({proc.returncode}): "
0152             f"{proc.stderr.strip()[:400]}")
0153     if not parse:
0154         return proc.stdout
0155     out = proc.stdout.strip()
0156     return json.loads(out) if out else {}
0157 
0158 
0159 def db():
0160     os.makedirs(os.path.dirname(DB_PATH), exist_ok=True)
0161     conn = sqlite3.connect(DB_PATH, timeout=30)
0162     conn.executescript(SCHEMA)
0163     return conn
0164 
0165 
0166 def set_state(conn, key, value):
0167     conn.execute("INSERT INTO state (k, v) VALUES (?, ?) "
0168                  "ON CONFLICT(k) DO UPDATE SET v = excluded.v",
0169                  (key, json.dumps(value)))
0170 
0171 
0172 def get_state(conn, key, default=None):
0173     row = conn.execute("SELECT v FROM state WHERE k = ?", (key,)).fetchone()
0174     return json.loads(row[0]) if row else default
0175 
0176 
0177 def _epoch(iso):
0178     """Seconds since the epoch for a listing's LastModified
0179     (2026-09-12T03:40:01+00:00 or ...Z); an unreadable one counts as old."""
0180     if not iso:
0181         return 0.0
0182     try:
0183         return datetime.fromisoformat(iso.replace("Z", "+00:00")).timestamp()
0184     except ValueError:
0185         return 0.0
0186 
0187 
0188 def subject_of(key):
0189     match = SUBJECT_RE.match(key)
0190     return match.group(1) if match else None
0191 
0192 
0193 def drain(conn, max_batches=500):
0194     """Receive events and record their objects. Returns (recorded, seen).
0195 
0196     A message is deleted only after its rows are committed, so a crash
0197     repeats work rather than losing it; the key is the primary key, so a
0198     repeat is a no-op. A run is bounded at max_batches rounds of ten, so
0199     a backlog drains across runs at a few percent of one core rather than
0200     all at once; the cron line's flock keeps runs from overlapping.
0201     """
0202     import boto3  # the drain alone needs it; the other modes stay on the CLI
0203     sqs = boto3.client("sqs", region_name=REGION)
0204     recorded = seen = 0
0205     for _ in range(max_batches):
0206         result = sqs.receive_message(
0207             QueueUrl=QUEUE_URL, MaxNumberOfMessages=10, WaitTimeSeconds=1)
0208         messages = result.get("Messages") or []
0209         if not messages:
0210             break
0211         entries = []
0212         for n, message in enumerate(messages):
0213             seen += 1
0214             entries.append({"Id": str(n), "ReceiptHandle": message["ReceiptHandle"]})
0215             try:
0216                 body = json.loads(message["Body"])
0217             except (KeyError, ValueError) as exc:
0218                 print(f"undecodable message body: {exc}", file=sys.stderr)
0219                 continue
0220             # S3 sends one test event when the notification is configured.
0221             if body.get("Event") == "s3:TestEvent":
0222                 continue
0223             for record in body.get("Records", []):
0224                 obj = record.get("s3", {}).get("object", {})
0225                 key = obj.get("key")
0226                 if not key:
0227                     continue
0228                 if key.startswith(STATUS_PREFIX):
0229                     conn.execute(
0230                         "INSERT INTO status_writes (hour, writes, bytes) "
0231                         "VALUES (?, 1, ?) ON CONFLICT(hour) DO UPDATE SET "
0232                         "writes = writes + 1, bytes = bytes + excluded.bytes",
0233                         (int(time.time() // 3600), int(obj.get("size") or 0)))
0234                     recorded += 1
0235                     continue
0236                 conn.execute(
0237                     "INSERT INTO objects (key, size, etime, subject, seen_at) "
0238                     "VALUES (?, ?, ?, ?, ?) ON CONFLICT(key) DO UPDATE SET "
0239                     "size = excluded.size, etime = excluded.etime",
0240                     (key, int(obj.get("size") or 0),
0241                      record.get("eventTime", ""), subject_of(key), time.time()))
0242                 recorded += 1
0243         conn.commit()
0244         reply = sqs.delete_message_batch(QueueUrl=QUEUE_URL, Entries=entries)
0245         failed = reply.get("Failed") or []
0246         if failed:
0247             # An undeleted message returns after its visibility timeout and
0248             # is re-recorded as a no-op; say so rather than hide it.
0249             print(f"delete_message_batch: {len(failed)} of {len(entries)} "
0250                   f"failed: {failed[0].get('Code')} "
0251                   f"{failed[0].get('Message', '')[:200]}", file=sys.stderr)
0252     set_state(conn, "last_drain", {"at": time.time(), "recorded": recorded})
0253     conn.commit()
0254     return recorded, seen
0255 
0256 
0257 def sweep_state(conn, now):
0258     """(time of the last applied pass, stall text or None).
0259 
0260     Three states, not two. A sweeper reporting normally is silent here. A
0261     sweeper past its threshold is stalled. A sweeper that has never
0262     reported is neither, until someone declares that it should be
0263     running: the commissioning stamp in `sweeper_expected_from` is what
0264     turns absence into staleness, and it is set by hand when the far side
0265     says the first real pass is imminent.
0266     """
0267     # A gateway wipe is recorded as a pass but says nothing of the sweeper.
0268     last_pass = conn.execute(
0269         "SELECT MAX(applied_at) FROM passes "
0270         "WHERE COALESCE(reporter, '') != 'gateway-wipe'").fetchone()[0]
0271     if last_pass:
0272         if now - last_pass > PASS_STALL_SECONDS:
0273             return last_pass, (f"no sweep pass reported for "
0274                                f"{(now - last_pass) / 3600:.1f} hours")
0275         return last_pass, None
0276     expected = get_state(conn, "sweeper_expected_from")
0277     if expected and now - float(expected) > PASS_STALL_SECONDS:
0278         return None, ("the sweeper was expected to be running and has "
0279                       "never reported a pass")
0280     return None, None
0281 
0282 
0283 def apply_passes(conn):
0284     """Apply the sweeper's spooled pass records to the index.
0285 
0286     The sweep runs inside the perimeter and deletes what it has taken, and
0287     S3 tells the owner nothing about another party's deletions. So the
0288     party that deletes reports it, and this is where those reports land in
0289     the index (swf-epicprod docs/JOB_REPORTING.md).
0290 
0291     A record is applied once and its file removed. An unreadable record is
0292     moved aside rather than deleted or retried forever: the sweeper treats
0293     its 202 as final and will never send it again, so losing it silently
0294     would be losing it for good.
0295     """
0296     if not os.path.isdir(PASS_SPOOL):
0297         return {"applied": 0, "retired": 0, "rejected": 0}
0298     applied = retired = rejected = 0
0299     now = time.time()
0300     for name in sorted(os.listdir(PASS_SPOOL)):
0301         if not name.endswith(".json"):
0302             continue
0303         path = os.path.join(PASS_SPOOL, name)
0304         try:
0305             with open(path) as handle:
0306                 record = json.load(handle)
0307         except (OSError, ValueError) as exc:
0308             print(f"unreadable pass record {name}: {exc}", file=sys.stderr)
0309             try:
0310                 os.rename(path, path + ".rejected")
0311             except OSError:
0312                 pass
0313             rejected += 1
0314             continue
0315         keys_read = [k for k in (record.get("deleted_read") or [])
0316                      if isinstance(k, str)]
0317         keys_unread = [k for k in (record.get("deleted_unread") or [])
0318                        if isinstance(k, str)]
0319         gone = 0
0320         for key in keys_read + keys_unread:
0321             gone += conn.execute(
0322                 "DELETE FROM objects WHERE key = ?", (key,)).rowcount
0323         window = record.get("window") or {}
0324         conn.execute(
0325             "INSERT INTO passes (pass_id, reporter, outcome, reason, "
0326             "window_from, window_to, jobs_filed, keys_read, keys_unread, "
0327             "retired, received_at, applied_at) "
0328             "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) "
0329             "ON CONFLICT(pass_id) DO NOTHING",
0330             (name[:-5], record.get("reporter"), record.get("outcome", "ok"),
0331              record.get("reason"), window.get("from"), window.get("to"),
0332              len(record.get("filed") or record.get("pandaids") or []),
0333              len(keys_read), len(keys_unread), gone,
0334              record.get("received_at"), now))
0335         conn.commit()
0336         os.remove(path)
0337         applied += 1
0338         retired += gone
0339     if applied:
0340         set_state(conn, "last_pass_applied", {"at": now, "records": applied,
0341                                               "retired": retired})
0342         conn.commit()
0343     return {"applied": applied, "retired": retired, "rejected": rejected}
0344 
0345 
0346 def reconcile(conn):
0347     """The nightly cross-check: count the prefix authoritatively and
0348     compare with the index, then drop rows for objects that are gone.
0349 
0350     Listing once a night costs a fraction of a cent. Drift is the signal
0351     that events were missed or that something writes without telling us;
0352     an object younger than RECONCILE_GRACE_SECONDS is left to the drain.
0353     """
0354     sizes = {}
0355     modified = {}
0356     token = None
0357     while True:
0358         args = ["s3api", "list-objects-v2", "--bucket", BUCKET,
0359                 "--prefix", PREFIX, "--output", "json"]
0360         if token:
0361             args += ["--starting-token", token]
0362         page = aws(*args)
0363         for item in page.get("Contents", []):
0364             sizes[item["Key"]] = int(item.get("Size") or 0)
0365             modified[item["Key"]] = item.get("LastModified") or ""
0366         token = page.get("NextContinuationToken")
0367         if not (page.get("IsTruncated") and token):
0368             break
0369     keys = set(sizes)
0370 
0371     indexed = {row[0] for row in conn.execute("SELECT key FROM objects")}
0372     gone_from_bucket = indexed - keys            # expired or swept
0373     # Unindexed keys split by age: an event for a fresh object is still
0374     # in flight, so only the older ones were never heard about.
0375     horizon = time.time() - RECONCILE_GRACE_SECONDS
0376     pending = {k for k in keys - indexed if _epoch(modified[k]) > horizon}
0377     missing_from_index = keys - indexed - pending
0378 
0379     for key in gone_from_bucket:
0380         conn.execute("DELETE FROM objects WHERE key = ?", (key,))
0381     # The listing carries the size, so an object recorded here is as
0382     # complete as one announced by an event, minus the event time.
0383     for key in missing_from_index:
0384         conn.execute(
0385             "INSERT INTO objects (key, size, etime, subject, seen_at) "
0386             "VALUES (?, ?, '', ?, ?) ON CONFLICT(key) DO NOTHING",
0387             (key, sizes[key], subject_of(key), time.time()))
0388 
0389     outcome = {"at": time.time(), "bucket_count": len(keys),
0390                "index_count": len(indexed),
0391                "unheard": len(missing_from_index),
0392                "pending": len(pending),
0393                "pruned": len(gone_from_bucket)}
0394     set_state(conn, "last_reconcile", outcome)
0395     conn.commit()
0396     return outcome
0397 
0398 
0399 def guard(conn, arm=True):
0400     """The nightly guard: does the traffic look like an explosion, and
0401     pull the plug if it does.
0402 
0403     The signal is the index, not the bill: objects and bytes are what
0404     boom first and the invoice is their lagging shadow. Two absolute
0405     ceilings and one relative one, because a ceiling set against today's
0406     expectation ages badly, and a sudden multiple of the recent norm is
0407     the shape of a runaway whatever the norm was.
0408 
0409     The plug is the reporting key. Disabling it stops jobs already
0410     running, since a running job keeps the environment it started with
0411     and the payload treats a failed write as ordinary (swf-epicprod
0412     docs/JOB_REPORTING.md). Only that key is touched: the sweep and
0413     stage-out keys survive, so reading and log stage-out continue.
0414     """
0415     now = time.time()
0416     day_objects, day_bytes = conn.execute(
0417         "SELECT COUNT(*), COALESCE(SUM(size), 0) FROM objects "
0418         "WHERE seen_at > ?", (now - 86400,)).fetchone()
0419     # Status overwrites are writes like any report's, on the same key.
0420     status_writes, status_bytes = status_since(conn, now, 24)
0421     day_objects += status_writes
0422     day_bytes += status_bytes
0423     # The trailing norm, excluding the last day, over whatever history
0424     # the index holds. Expiry keeps that to about a week.
0425     prior_days = conn.execute(
0426         "SELECT COUNT(*) FROM objects WHERE seen_at <= ? AND seen_at > ?",
0427         (now - 86400, now - 7 * 86400)).fetchone()[0]
0428     norm_per_day = prior_days / 6.0 if prior_days else 0.0
0429 
0430     breaches = []
0431     if day_objects > OBJECT_HARD_CEILING_PER_DAY:
0432         breaches.append(f"{day_objects:,} objects in 24 hours, hard ceiling "
0433                         f"{OBJECT_HARD_CEILING_PER_DAY:,}")
0434     if day_bytes > BYTES_HARD_CEILING_PER_DAY:
0435         breaches.append(f"{day_bytes / 1e9:.1f} GB written in 24 hours, hard "
0436                         f"ceiling {BYTES_HARD_CEILING_PER_DAY / 1e9:.0f} GB")
0437     if (norm_per_day >= RELATIVE_FLOOR
0438             and day_objects > norm_per_day * RELATIVE_MULTIPLE):
0439         breaches.append(
0440             f"{day_objects:,} objects in 24 hours against a recent norm of "
0441             f"{norm_per_day:,.0f} a day, over {RELATIVE_MULTIPLE}x")
0442 
0443     warnings = []
0444     if not breaches:
0445         if day_objects > OBJECT_SOFT_CEILING_PER_DAY:
0446             warnings.append(f"{day_objects:,} objects in 24 hours, expected "
0447                             f"under {OBJECT_SOFT_CEILING_PER_DAY:,}")
0448         if day_bytes > BYTES_SOFT_CEILING_PER_DAY:
0449             warnings.append(f"{day_bytes / 1e9:.1f} GB in 24 hours, expected "
0450                             f"under {BYTES_SOFT_CEILING_PER_DAY / 1e9:.0f} GB")
0451 
0452     # A stalled sweeper is not an explosion. The sweep is the bucket's
0453     # normal drain, so when it stops, objects accumulate for a reason that
0454     # pulling the write credential would not fix and would make worse: the
0455     # jobs would go quiet while the backlog stayed. It gets its own verdict
0456     # and never arms the plug.
0457     #
0458     # A sweeper that has never run is a third thing again, and saying
0459     # "stalled" of something not yet built would be a false alarm standing
0460     # for days, which is how a tile teaches its reader to ignore it. The
0461     # threshold therefore means nothing until commissioning is declared.
0462     last_pass, stalled = sweep_state(conn, now)
0463     if stalled:
0464         warnings.append(stalled)
0465 
0466     stopped = key_status() == "Inactive"
0467     if breaches and arm and not stopped:
0468         disable_reporting_key()
0469         stopped = True
0470 
0471     # The day's totals are kept forever, so growth outlives the objects.
0472     today = time.strftime("%Y-%m-%d", time.gmtime(now - 43200))
0473     conn.execute(
0474         "INSERT INTO daily (day, objects, bytes, at) VALUES (?, ?, ?, ?) "
0475         "ON CONFLICT(day) DO UPDATE SET objects = excluded.objects, "
0476         "bytes = excluded.bytes, at = excluded.at",
0477         (today, day_objects, day_bytes, now))
0478     history = [
0479         {"day": row[0], "objects": row[1], "bytes": row[2]}
0480         for row in conn.execute(
0481             "SELECT day, objects, bytes FROM daily ORDER BY day DESC LIMIT 14")]
0482     this_week = sum(r["objects"] for r in history[:7])
0483     last_week = sum(r["objects"] for r in history[7:14])
0484     growth = (round((this_week - last_week) / last_week, 2)
0485               if last_week else None)
0486 
0487     outcome = {"at": now, "day_objects": day_objects,
0488                "day_bytes": day_bytes,
0489                "day_status_writes": status_writes,
0490                "history": history[:14],
0491                "week_over_week": growth,
0492                "norm_objects_per_day": round(norm_per_day),
0493                "breaches": breaches, "warnings": warnings,
0494                "sweep_stalled": stalled,
0495                "last_pass_at": last_pass,
0496                "reporting_key_stopped": stopped}
0497     set_state(conn, "last_guard", outcome)
0498     conn.commit()
0499     return outcome
0500 
0501 
0502 def status_since(conn, now, hours):
0503     """(writes, bytes) to the status prefix over the last `hours` hours."""
0504     return conn.execute(
0505         "SELECT COALESCE(SUM(writes), 0), COALESCE(SUM(bytes), 0) "
0506         "FROM status_writes WHERE hour > ?",
0507         (int(now // 3600) - hours,)).fetchone()
0508 
0509 
0510 def reporting_keys():
0511     """The reporting user's access keys, from IAM rather than from the
0512     credential file, so the guard acts on what is actually enabled. The
0513     user may hold two: a retired key left Inactive beside its successor
0514     (2026-09-24, the Event Service key), so every key is considered."""
0515     return aws("iam", "list-access-keys", "--user-name", REPORTING_USER,
0516                "--output", "json").get("AccessKeyMetadata", [])
0517 
0518 
0519 def key_status():
0520     """Active if any of the user's keys can write, else Inactive."""
0521     try:
0522         keys = reporting_keys()
0523     except RuntimeError as exc:
0524         print(f"key status unreadable: {exc}", file=sys.stderr)
0525         return "Unknown"
0526     if not keys:
0527         return "Missing"
0528     return ("Active" if any(k["Status"] == "Active" for k in keys)
0529             else "Inactive")
0530 
0531 
0532 def disable_reporting_key():
0533     """Pull the plug on every active key. Reversible with one command per
0534     key, named in the notice."""
0535     active = [k["AccessKeyId"] for k in reporting_keys()
0536               if k["Status"] == "Active"]
0537     if not active:
0538         raise RuntimeError(f"no active access key found for {REPORTING_USER}")
0539     for key_id in active:
0540         aws("iam", "update-access-key", "--user-name", REPORTING_USER,
0541             "--access-key-id", key_id, "--status", "Inactive", parse=False)
0542         print(f"DISABLED reporting key {key_id}; re-enable with: "
0543               f"aws iam update-access-key --user-name {REPORTING_USER} "
0544               f"--access-key-id {key_id} --status Active", file=sys.stderr)
0545     return active
0546 
0547 
0548 def wipe(conn, reason):
0549     """Delete every object under the prefix written before this run began,
0550     retire their rows, and record the deletion as a pass.
0551 
0552     For the gateway's own clearances (2026-09-24: the 9/15-9/18 flood,
0553     cleared on Torre's word). The record goes where the sweeper's do, so
0554     the index and the guard read a wipe as a known drain and not as
0555     missing events; objects written during the run are left alone.
0556     """
0557     import boto3
0558     s3 = boto3.client("s3", region_name=REGION)
0559     started = time.time()
0560     deleted = errors = 0
0561     paginator = s3.get_paginator("list_objects_v2")
0562     for page in paginator.paginate(Bucket=BUCKET, Prefix=PREFIX):
0563         keys = [item["Key"] for item in page.get("Contents", [])
0564                 if item["LastModified"].timestamp() < started]
0565         if not keys:
0566             continue
0567         reply = s3.delete_objects(
0568             Bucket=BUCKET,
0569             Delete={"Objects": [{"Key": k} for k in keys], "Quiet": True})
0570         failed = {e["Key"] for e in reply.get("Errors", [])}
0571         if failed:
0572             errors += len(failed)
0573             first = reply["Errors"][0]
0574             print(f"delete_objects: {len(failed)} failed: {first.get('Code')} "
0575                   f"{first.get('Message', '')[:200]}", file=sys.stderr)
0576         gone = [k for k in keys if k not in failed]
0577         conn.executemany("DELETE FROM objects WHERE key = ?",
0578                          [(k,) for k in gone])
0579         conn.commit()
0580         deleted += len(gone)
0581     pass_id = time.strftime("wipe-%Y%m%dT%H%M%SZ", time.gmtime(started))
0582     window_to = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(started))
0583     conn.execute(
0584         "INSERT INTO passes (pass_id, reporter, outcome, reason, "
0585         "window_from, window_to, jobs_filed, keys_read, keys_unread, "
0586         "retired, received_at, applied_at) "
0587         "VALUES (?, 'gateway-wipe', ?, ?, NULL, ?, 0, 0, ?, ?, ?, ?)",
0588         (pass_id, "partial" if errors else "ok", reason, window_to,
0589          deleted, deleted, window_to, time.time()))
0590     conn.commit()
0591     return {"pass_id": pass_id, "deleted": deleted, "failed": errors,
0592             "seconds": round(time.time() - started)}
0593 
0594 
0595 def summary(conn):
0596     """The state other things read: totals, the last hour's rate against
0597     the ceiling, the busiest subject, and the freshness of both passes."""
0598     now = time.time()
0599     total = conn.execute("SELECT COUNT(*) FROM objects").fetchone()[0]
0600     hour = conn.execute("SELECT COUNT(*) FROM objects WHERE seen_at > ?",
0601                         (now - 3600,)).fetchone()[0]
0602     day = conn.execute("SELECT COUNT(*) FROM objects WHERE seen_at > ?",
0603                        (now - 86400,)).fetchone()[0]
0604     status_hour = status_since(conn, now, 1)[0]
0605     status_day = status_since(conn, now, 24)[0]
0606     hour += status_hour
0607     day += status_day
0608     busiest = conn.execute(
0609         "SELECT subject, COUNT(*) c FROM objects WHERE seen_at > ? "
0610         "GROUP BY subject ORDER BY c DESC LIMIT 1", (now - 86400,)).fetchone()
0611     last_reconcile = get_state(conn, "last_reconcile", {})
0612     last_pass, stalled = sweep_state(conn, now)
0613     verdict, detail = "ok", ""
0614     if stalled:
0615         verdict, detail = "warning", "the sweeper has reported no pass: " + stalled
0616     if hour > RATE_CEILING_PER_HOUR:
0617         verdict = "alarm"
0618         detail = f"{hour} objects in the last hour, ceiling {RATE_CEILING_PER_HOUR}"
0619     elif busiest and busiest[1] > PER_JOB_CEILING:
0620         verdict = "warning"
0621         detail = (f"subject {busiest[0]} wrote {busiest[1]} objects today, "
0622                   f"per-job ceiling {PER_JOB_CEILING}")
0623     elif last_reconcile.get("unheard"):
0624         verdict = "warning"
0625         detail = (f"{last_reconcile['unheard']} objects in the bucket that no "
0626                   "event announced")
0627     last_guard = get_state(conn, "last_guard", {})
0628     if last_guard.get("reporting_key_stopped"):
0629         verdict, detail = "alarm", (
0630             "the reporting key is disabled: "
0631             + "; ".join(last_guard.get("breaches") or ["stopped by the guard"]))
0632     return {"objects": total, "last_hour": hour, "last_day": day,
0633             "status_last_hour": status_hour, "status_last_day": status_day,
0634             "guard": last_guard,
0635             "last_sweep_pass": last_pass,
0636             "busiest_subject": busiest[0] if busiest else None,
0637             "busiest_count": busiest[1] if busiest else 0,
0638             "rate_ceiling": RATE_CEILING_PER_HOUR,
0639             "last_drain": get_state(conn, "last_drain", {}),
0640             "last_reconcile": last_reconcile,
0641             "verdict": verdict, "detail": detail}
0642 
0643 
0644 def main():
0645     parser = argparse.ArgumentParser(description=__doc__.split("\n\n")[0])
0646     parser.add_argument(
0647         "mode",
0648         choices=("drain", "reconcile", "apply", "guard", "summary", "wipe"))
0649     parser.add_argument("--no-arm", action="store_true",
0650                         help="report a breach without disabling the key")
0651     parser.add_argument("--reason", default="gateway wipe",
0652                         help="wipe: why, recorded with the pass")
0653     args = parser.parse_args()
0654     conn = db()
0655     try:
0656         if args.mode == "drain":
0657             recorded, seen = drain(conn)
0658             passes = apply_passes(conn)
0659             used = resource.getrusage(resource.RUSAGE_SELF)
0660             print(f"drained {seen} messages, recorded {recorded} objects; "
0661                   f"applied {passes['applied']} sweep passes retiring "
0662                   f"{passes['retired']} objects"
0663                   + (f"; {passes['rejected']} rejected"
0664                      if passes["rejected"] else "")
0665                   + f"; cpu {used.ru_utime + used.ru_stime:.1f}s")
0666         elif args.mode == "reconcile":
0667             print(json.dumps(reconcile(conn), indent=2))
0668         elif args.mode == "apply":
0669             print(json.dumps(apply_passes(conn), indent=2))
0670         elif args.mode == "guard":
0671             print(json.dumps(guard(conn, arm=not args.no_arm), indent=2))
0672         elif args.mode == "wipe":
0673             print(json.dumps(wipe(conn, args.reason), indent=2))
0674         else:
0675             print(json.dumps(summary(conn), indent=2))
0676     finally:
0677         conn.close()
0678     return 0
0679 
0680 
0681 if __name__ == "__main__":
0682     sys.exit(main())