File indexing completed on 2026-09-28 09:37:31
0001
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
0052
0053
0054 RATE_CEILING_PER_HOUR = int(os.environ.get("STAGEOUT_RATE_CEILING", "20000"))
0055
0056
0057 PER_JOB_CEILING = int(os.environ.get("STAGEOUT_PER_JOB_CEILING", "50"))
0058
0059
0060
0061
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
0072
0073 RELATIVE_MULTIPLE = int(os.environ.get("STAGEOUT_RELATIVE_MULTIPLE", "20"))
0074 RELATIVE_FLOOR = int(os.environ.get("STAGEOUT_RELATIVE_FLOOR", "5000"))
0075
0076
0077
0078
0079
0080
0081
0082 PASS_SPOOL = os.environ.get("STAGEOUT_PASS_SPOOL",
0083 "/var/lib/stageout/passes")
0084
0085
0086 PASS_STALL_SECONDS = int(os.environ.get("STAGEOUT_PASS_STALL", str(3 * 3600)))
0087
0088
0089
0090
0091
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
0135
0136
0137
0138 STATUS_PREFIX = "status/"
0139
0140
0141
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
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
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
0248
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
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
0373
0374
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
0382
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
0420 status_writes, status_bytes = status_since(conn, now, 24)
0421 day_objects += status_writes
0422 day_bytes += status_bytes
0423
0424
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
0453
0454
0455
0456
0457
0458
0459
0460
0461
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
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())