forked from b.peck/BAT
755 lines
33 KiB
Python
755 lines
33 KiB
Python
#!/usr/bin/env python3
|
|
"""seed_journal.py - backdate synthetic alarm activity into the alarm journal.
|
|
|
|
The alarm simulator (SimHarness/alarmsim) generates activity in real time, so
|
|
"give me 8 hours of alarms" would take 8 hours. This writes the same kind of
|
|
activity straight into the journal tables with backdated timestamps, shaped to
|
|
land on a target Alarm Health grade or an exact Alarm Health score.
|
|
|
|
python3 tools/seed_journal.py --hours 8 --score 85 --dry-run # plan only
|
|
python3 tools/seed_journal.py --hours 8 --score 85 # insert
|
|
python3 tools/seed_journal.py --hours 8 --grade B # canned profile
|
|
|
|
How the grade/score is hit (weights from PrimeControls.calc.DEFAULTS):
|
|
|
|
rate 0.30 activations/hr vs the 6/hr ISA target
|
|
flood 0.25 % of window inside a flood episode (>10 activations/10 min)
|
|
chatter 0.20 sources re-activating >10x/hr with a median gap <=120 s
|
|
standing 0.15 alarms active >24 h (LIVE queryStatus - not seeded, see below)
|
|
priority 0.10 deviation from the ISA 80/15/5 low/medium/high mix
|
|
|
|
Only `rate`, `flood` and `priority` are seedable. `standing` reads live gateway
|
|
alarm state via system.alarm.queryStatus, so whatever is actually active now is
|
|
measured as-is and the profile budgets for it (--standing-now). `chatter` is
|
|
structurally 0 (SOURCE_CAP < calc's chatter_min_count).
|
|
|
|
--score solves for the three seedable dials instead of using a canned profile:
|
|
total activations (rate), whether exactly one 10-min cell floods, and the exact
|
|
low/medium/high split (priority deviation). Rows already in the window - a live
|
|
active alarm, an earlier seed - are read back from the journal and counted as
|
|
part of the target, so the solve describes the window the gateway will actually
|
|
score. Two shapes are available:
|
|
|
|
--shape stable (default) rate held at/below the ISA target, so the rate
|
|
sub-score sits pinned at 100 on the flat part of its curve
|
|
and the score does NOT drift as the dashboard's rolling
|
|
window slides forward and old activations fall out. The
|
|
whole deficit lands on flood + priority mix.
|
|
--shape balanced deficit spread across rate, flood and priority - a more
|
|
typical-looking plant, but the score climbs a few tenths
|
|
per 10 minutes as the window slides past the seeded data.
|
|
|
|
The plan is verified offline against the real PrimeControls.calc before anything
|
|
is written - the same module the gateway runs - and the run aborts if the
|
|
predicted grade/score misses the target.
|
|
|
|
Journal conventions this mirrors (harvested from existing rows):
|
|
eventtype 0=active 1=clear 2=ack eventflags 0 for tag alarms
|
|
eventtime DATETIME in GATEWAY LOCAL time (America/Chicago), not UTC
|
|
event_data active/clear -> eventValue (dtype 0, intvalue 1/0)
|
|
ack -> ackUser + ackUserName (dtype 2, strvalue)
|
|
|
|
Flood episodes are stable under window drift: calc bins on `origin =
|
|
(start_ms // bin_ms) * bin_ms`, always a multiple of bin_ms, so bin edges sit on
|
|
absolute 10-minute epoch marks no matter when the dashboard re-anchors "now". A
|
|
burst placed inside one such cell stays inside one bin.
|
|
"""
|
|
import argparse
|
|
import collections
|
|
import importlib.util
|
|
import json
|
|
import os
|
|
import random
|
|
import subprocess
|
|
import sys
|
|
import types
|
|
import uuid
|
|
from datetime import datetime
|
|
|
|
try:
|
|
from zoneinfo import ZoneInfo
|
|
except ImportError:
|
|
ZoneInfo = None
|
|
|
|
ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
|
PROJECT = os.path.join(ROOT, "ignition", "gateway", "projects", "PrimeBAT",
|
|
"ignition", "script-python", "PrimeControls")
|
|
TAGS_FILE = os.path.join(ROOT, "test-data", "simulation_tags.json")
|
|
|
|
DB_CONTAINER = "buildathon-db"
|
|
DB_NAME = "ignition"
|
|
DB_USER = "ignition"
|
|
DB_PASS = "ignition"
|
|
EVENTS_TABLE = "PrimeControls_alarm_events"
|
|
DATA_TABLE = "PrimeControls_alarm_event_data"
|
|
|
|
GATEWAY_TZ = "America/Chicago"
|
|
TAG_PROVIDER = "default"
|
|
TAG_ROOT = "BuildathonSim"
|
|
|
|
PRIORITY_LEVEL = {"Diagnostic": 0, "Low": 1, "Medium": 2, "High": 3, "Critical": 4}
|
|
BUCKET_OF = {0: "low", 1: "low", 2: "medium", 3: "high", 4: "high"}
|
|
ACK_USERS = [("Prime", 953), ("operator", 284), ("admin", 122)]
|
|
|
|
# Sources whose names advertise a pathology - excluded from a healthy seed.
|
|
BAD_ACTOR_MARKERS = ("_Chatter", "_Standing", "_Fleeting")
|
|
|
|
BIN_MS = 600000 # calc DEFAULTS bin_ms (10 min)
|
|
HOUR_MS = 3600000
|
|
|
|
# Per-window profiles. `counts` is the exact low/medium/high activation split;
|
|
# `burst` is (hours_before_end, activation_count) placed inside one 10-min cell.
|
|
PROFILES = {
|
|
"B": {
|
|
"current": {"counts": {"low": 38, "medium": 9, "high": 3},
|
|
"burst": (2.2, 11), "burst_area": "Packaging"},
|
|
# Prior period only feeds the period-over-period delta chips; its grade
|
|
# is never scored. Made deliberately worse so deltas read as improving.
|
|
# Kept within the low-bucket capacity (13 sources x SOURCE_CAP).
|
|
"prior": {"counts": {"low": 46, "medium": 14, "high": 6},
|
|
"burst": (2.6, 13), "burst_area": "BoilerHouse"},
|
|
},
|
|
}
|
|
|
|
BASELINE_BIN_CAP = 4 # keeps non-burst cells under the flood-start (>10) and
|
|
# flood-continue (>=5) thresholds
|
|
SOURCE_CAP = 4 # < calc chatter_min_count (5) => structurally no chatter
|
|
MIN_DURATION_S = 45 # > calc fleeting_s (10) => structurally no fleeting
|
|
MAX_DURATION_S = 420
|
|
ACK_RATE = 0.85
|
|
CLEAR_RATE = 0.94
|
|
|
|
# --score solver
|
|
BURST_COUNT = 11 # > calc flood_per_10min (10) => one flooding cell
|
|
SOLVE_N_MAX = 260 # activation ceiling searched (32/hr over 8 h)
|
|
PRIOR_WORSE = 1.25 # prior-window activations, as a multiple of current
|
|
HIGH_SHARE_PRIOR = 0.18 # tie-break: how much of the mix should be High/Critical
|
|
|
|
|
|
# ---------------- the gateway's own calc module ----------------
|
|
|
|
|
|
def load_calc():
|
|
"""Import PrimeControls.calc from the project exactly as tests/conftest.py
|
|
does, so the prediction uses the code the gateway actually runs."""
|
|
pkg = types.ModuleType("PrimeControls")
|
|
pkg.__path__ = [PROJECT]
|
|
sys.modules["PrimeControls"] = pkg
|
|
out = {}
|
|
for mod in ("calc", "fmt"):
|
|
spec = importlib.util.spec_from_file_location(
|
|
"PrimeControls." + mod, os.path.join(PROJECT, mod, "code.py"))
|
|
m = importlib.util.module_from_spec(spec)
|
|
sys.modules["PrimeControls." + mod] = m
|
|
spec.loader.exec_module(m)
|
|
setattr(pkg, mod, m)
|
|
out[mod] = m
|
|
return out["calc"]
|
|
|
|
|
|
# ---------------- sources ----------------
|
|
|
|
|
|
def load_sources():
|
|
"""[(rel_path, priority_level)] for every real (non-bad-actor) sim alarm."""
|
|
with open(TAGS_FILE) as fh:
|
|
doc = json.load(fh)
|
|
rows = []
|
|
|
|
def walk(node, path):
|
|
name = node.get("name")
|
|
p = path + [name] if name else path
|
|
for alarm in node.get("alarms") or []:
|
|
rel = "/".join(p[1:])
|
|
if any(m in rel for m in BAD_ACTOR_MARKERS):
|
|
continue
|
|
rows.append((rel, PRIORITY_LEVEL[alarm.get("priority")]))
|
|
for child in node.get("tags") or []:
|
|
walk(child, p)
|
|
|
|
walk(doc, [])
|
|
return sorted(set(rows))
|
|
|
|
|
|
def source_path(rel):
|
|
name = rel.split("/")[-1]
|
|
return "prov:%s:/tag:%s/%s:/alm:%s" % (TAG_PROVIDER, TAG_ROOT, rel, name)
|
|
|
|
|
|
# ---------------- planning ----------------
|
|
|
|
|
|
def cell_start(ms):
|
|
"""Absolute 10-minute epoch cell containing ms."""
|
|
return (ms // BIN_MS) * BIN_MS
|
|
|
|
|
|
def bucket_capacity(sources):
|
|
"""Activations each bucket can absorb without any source reaching
|
|
SOURCE_CAP - i.e. without manufacturing a chattering alarm."""
|
|
cap = {"low": 0, "medium": 0, "high": 0}
|
|
for rel, lvl in sources:
|
|
cap[BUCKET_OF[lvl]] += SOURCE_CAP
|
|
return cap
|
|
|
|
|
|
def solve_spec(calc, target, hours, standing_now, existing, shape, tol, capacity):
|
|
"""Search the seedable dial space for a window that scores `target`.
|
|
|
|
Dials: total activations N (rate sub-score), whether exactly one 10-min cell
|
|
floods (flood sub-score), and the integer low/medium/high split (priority
|
|
deviation). chatter is structurally 0 and standing is live gateway state, so
|
|
both are inputs. Every candidate's score comes from the gateway's own
|
|
calc.health_score, and `existing` (activations already in the window, per
|
|
bucket) is part of the total - only the remainder gets seeded.
|
|
|
|
Returns a PROFILES-style spec for the seeded remainder, plus the totals and
|
|
the predicted sub-scores for reporting.
|
|
"""
|
|
span_ms = int(hours * HOUR_MS)
|
|
ex_total = sum(existing.values())
|
|
weight_p = calc.DEFAULTS["health_weights"]["priority"]
|
|
ptarget = calc.DEFAULTS["priority_target"]
|
|
cands = []
|
|
|
|
for flood_bins in (1, 0):
|
|
pct_flood = 100.0 * flood_bins * BIN_MS / span_ms
|
|
burst = BURST_COUNT if flood_bins else 0
|
|
for n in range(max(ex_total, burst, 1), SOLVE_N_MAX + 1):
|
|
rate = n / hours
|
|
# score(dev) is linear in the priority deviation while priority_s > 0,
|
|
# so two probes give the deviation the target needs - no duplicated
|
|
# copy of calc's weighting here.
|
|
base = calc.health_score(rate, pct_flood, 0, standing_now, 0.0, True)["score"]
|
|
slope = base - calc.health_score(rate, pct_flood, 0, standing_now,
|
|
1.0, True)["score"]
|
|
if slope <= 0:
|
|
continue
|
|
dev_needed = (base - target) / slope
|
|
if dev_needed < 0 or dev_needed > 100.0 / max(weight_p, 1e-9):
|
|
continue
|
|
best = None
|
|
hi_max = min(int(0.30 * n), existing["high"] + capacity["high"])
|
|
med_max = min(int(0.40 * n), existing["medium"] + capacity["medium"])
|
|
for h in range(existing["high"], hi_max + 1):
|
|
for m in range(existing["medium"], med_max + 1):
|
|
lo = n - m - h
|
|
if lo < existing["low"] or lo < m or lo < h:
|
|
continue # low stays the biggest bucket
|
|
if lo - existing["low"] > capacity["low"]:
|
|
continue # more low activations than sources can carry
|
|
dev = (abs(100.0 * lo / n - ptarget["low"]) +
|
|
abs(100.0 * m / n - ptarget["medium"]) +
|
|
abs(100.0 * h / n - ptarget["high"]))
|
|
key = (abs(dev - dev_needed), abs(h / float(n) - HIGH_SHARE_PRIOR))
|
|
if best is None or key < best[0]:
|
|
best = (key, lo, m, h, dev)
|
|
if best is None:
|
|
continue
|
|
_key, lo, m, h, dev = best
|
|
hs = calc.health_score(rate, pct_flood, 0, standing_now, dev, True)
|
|
miss = abs(hs["score"] - target)
|
|
if miss > tol:
|
|
continue
|
|
subs = dict((s["key"], s["score"]) for s in hs["subs"])
|
|
cands.append({"flood_bins": flood_bins, "burst": burst, "n": n,
|
|
"total": {"low": lo, "medium": m, "high": h},
|
|
"dev": dev, "score": hs["score"], "grade": hs["grade"],
|
|
"subs": subs, "miss": miss})
|
|
if not cands:
|
|
raise SystemExit(
|
|
"no seedable window scores %.2f +/-%.2f with %d standing alarm(s) and "
|
|
"%d activation(s) already in the window (score reachable range with "
|
|
"these inputs is roughly 60-100; wipe the window or adjust --hours)"
|
|
% (target, tol, standing_now, ex_total))
|
|
|
|
def rank(c):
|
|
# accuracy first, in 0.02-score buckets (the UI shows one decimal), then
|
|
# the shape preference, then more data over less.
|
|
bucket = int(c["miss"] / 0.02)
|
|
if shape == "stable":
|
|
shape_cost = (0 if c["subs"]["rate"] >= 100.0 else 1,
|
|
abs(c["total"]["high"] / float(c["n"]) - HIGH_SHARE_PRIOR))
|
|
else:
|
|
# spread the deficit: penalise deep single-sub holes quadratically
|
|
spread = sum((100.0 - c["subs"][k]) ** 2 * w for k, w in
|
|
calc.DEFAULTS["health_weights"].items()
|
|
if c["subs"].get(k) is not None)
|
|
shape_cost = (round(spread / 100.0), 0.0)
|
|
return (bucket, shape_cost, -c["n"])
|
|
|
|
best = sorted(cands, key=rank)[0]
|
|
best["counts"] = dict((b, best["total"][b] - existing[b]) for b in best["total"])
|
|
return best
|
|
|
|
|
|
def allocate(sources, counts, burst_count, burst_area, rng):
|
|
"""Pick (rel, level) for every activation, honouring the exact bucket split
|
|
and SOURCE_CAP. Returns (burst_picks, baseline_picks)."""
|
|
by_bucket = collections.defaultdict(list)
|
|
for rel, lvl in sources:
|
|
by_bucket[BUCKET_OF[lvl]].append((rel, lvl))
|
|
|
|
used = collections.Counter()
|
|
|
|
def take(bucket, area=None):
|
|
pool = [s for s in by_bucket[bucket] if used[s[0]] < SOURCE_CAP
|
|
and (area is None or s[0].split("/")[0] == area)]
|
|
if not pool:
|
|
if area is not None:
|
|
return take(bucket, None) # area can't serve it; go global
|
|
raise SystemExit("exhausted %s sources (raise SOURCE_CAP)" % bucket)
|
|
pool.sort(key=lambda s: (used[s[0]], s[0]))
|
|
least = [s for s in pool if used[s[0]] == used[pool[0][0]]]
|
|
pick = rng.choice(least)
|
|
used[pick[0]] += 1
|
|
return pick
|
|
|
|
slots = []
|
|
for bucket, n in counts.items():
|
|
slots += [bucket] * n
|
|
rng.shuffle(slots)
|
|
|
|
# The burst is an area cascade: draw its slots from one area where possible.
|
|
burst_slots, baseline_slots = slots[:burst_count], slots[burst_count:]
|
|
burst = [take(b, burst_area) for b in burst_slots]
|
|
baseline = [take(b) for b in baseline_slots]
|
|
return burst, baseline
|
|
|
|
|
|
def place_times(start_ms, end_ms, burst, baseline, burst_at_ms, rng, pre_cells=None):
|
|
"""Assign an activation timestamp to every pick.
|
|
|
|
Burst picks land inside the single 10-minute cell containing burst_at_ms
|
|
(with an edge margin so they cannot spill into the neighbouring cell).
|
|
Baseline picks spread over the remaining cells, capped at BASELINE_BIN_CAP.
|
|
`pre_cells` seeds the per-cell counter with activations already in the
|
|
journal, so the cap accounts for them too.
|
|
"""
|
|
burst_cell = cell_start(burst_at_ms)
|
|
events = []
|
|
|
|
margin = 45000
|
|
for pick in burst:
|
|
t = rng.randint(burst_cell + margin, burst_cell + BIN_MS - margin - 1)
|
|
events.append((t, pick))
|
|
|
|
cells = []
|
|
c = cell_start(start_ms)
|
|
if c < start_ms:
|
|
c += BIN_MS # first whole cell inside the window
|
|
while c + BIN_MS <= end_ms:
|
|
if c != burst_cell:
|
|
cells.append(c)
|
|
c += BIN_MS
|
|
if not cells:
|
|
raise SystemExit("window too short to place baseline activations")
|
|
|
|
capacity = len(cells) * BASELINE_BIN_CAP
|
|
if len(baseline) > capacity:
|
|
raise SystemExit("baseline %d exceeds capacity %d (cap %d/cell over %d cells)"
|
|
% (len(baseline), capacity, BASELINE_BIN_CAP, len(cells)))
|
|
|
|
per_cell = collections.Counter(pre_cells or {})
|
|
for pick in baseline:
|
|
# Free choice among under-cap cells, so activity clumps the way real
|
|
# activity does. BASELINE_BIN_CAP (4) is below both the flood-start
|
|
# (>10) and flood-continue (>=5) thresholds, so no clump can start or
|
|
# extend a flood episode however the draws fall.
|
|
choices = [c for c in cells if per_cell[c] < BASELINE_BIN_CAP]
|
|
cell = rng.choice(choices)
|
|
per_cell[cell] += 1
|
|
t = rng.randint(cell + 5000, cell + BIN_MS - 5000)
|
|
events.append((t, pick))
|
|
|
|
events.sort(key=lambda e: e[0])
|
|
return events, burst_cell, per_cell
|
|
|
|
|
|
def build_instances(events, end_ms, rng):
|
|
"""Expand each activation into its active / ack / clear lifecycle."""
|
|
instances = []
|
|
for active_ms, (rel, lvl) in events:
|
|
dur_s = rng.randint(MIN_DURATION_S, MAX_DURATION_S)
|
|
clear_ms = active_ms + dur_s * 1000
|
|
# A clear past the window end would not be fetched; leave it open, which
|
|
# is what a still-active alarm looks like anyway.
|
|
if rng.random() > CLEAR_RATE or clear_ms > end_ms - 30000:
|
|
clear_ms = None
|
|
ack_ms = None
|
|
ack_user = None
|
|
if rng.random() < ACK_RATE:
|
|
ceiling = (clear_ms if clear_ms else end_ms) - 5000
|
|
floor = active_ms + 15000
|
|
if ceiling > floor:
|
|
ack_ms = rng.randint(floor, ceiling)
|
|
ack_user = rng.choices([u for u, _ in ACK_USERS],
|
|
weights=[w for _, w in ACK_USERS])[0]
|
|
instances.append({
|
|
"event_id": str(uuid.uuid4()), "rel": rel, "source": source_path(rel),
|
|
"priority": lvl, "active_ms": active_ms, "ack_ms": ack_ms,
|
|
"clear_ms": clear_ms, "ack_user": ack_user,
|
|
})
|
|
return instances
|
|
|
|
|
|
def plan_window(sources, spec, start_ms, end_ms, rng, pre_cells=None):
|
|
counts = spec["counts"]
|
|
burst_hours, burst_count = spec["burst"]
|
|
burst, baseline = allocate(sources, counts, burst_count,
|
|
spec.get("burst_area"), rng)
|
|
burst_at = end_ms - int(burst_hours * HOUR_MS)
|
|
events, burst_cell, per_cell = place_times(start_ms, end_ms, burst, baseline,
|
|
burst_at, rng, pre_cells)
|
|
instances = build_instances(events, end_ms, rng)
|
|
return instances, burst_cell, per_cell
|
|
|
|
|
|
# ---------------- verification against the real calc ----------------
|
|
|
|
|
|
def to_journal_rows(instances, calc):
|
|
"""Flatten instances into the row dicts alarms._norm_event would produce."""
|
|
rows = []
|
|
for inst in instances:
|
|
lvl, name = calc.normalize_priority(inst["priority"])
|
|
base = {"event_id": inst["event_id"], "source": inst["source"],
|
|
"display_path": "", "priority": lvl, "priority_name": name,
|
|
"is_system": False, "ack_user": None}
|
|
rows.append(dict(base, state="active", ts=inst["active_ms"]))
|
|
if inst["ack_ms"]:
|
|
rows.append(dict(base, state="ack", ts=inst["ack_ms"],
|
|
ack_user=inst["ack_user"]))
|
|
if inst["clear_ms"]:
|
|
rows.append(dict(base, state="clear", ts=inst["clear_ms"]))
|
|
return rows
|
|
|
|
|
|
def predict(calc, rows, active_now, start_ms, end_ms, now_ms):
|
|
bundle = calc.build_bundle(rows, active_now, 0, start_ms, end_ms, now_ms,
|
|
{"filters": {"priorities": [], "areas": [],
|
|
"states": [], "search": ""}})
|
|
return bundle
|
|
|
|
|
|
def synthetic_standing(calc, count, now_ms):
|
|
"""Stand-ins for the live queryStatus rows, so the offline prediction sees
|
|
the same standing-alarm penalty the gateway will apply."""
|
|
rows = []
|
|
for i in range(count):
|
|
lvl, name = calc.normalize_priority(1)
|
|
rows.append({"source": source_path("TankFarm/ManifoldLeakDetect"),
|
|
"display_path": "", "priority": lvl, "priority_name": name,
|
|
"active_ms": now_ms - (48 + i) * HOUR_MS, "unacked": True})
|
|
return rows
|
|
|
|
|
|
# ---------------- SQL ----------------
|
|
|
|
|
|
def local_dt(ms):
|
|
if ZoneInfo is not None:
|
|
return datetime.fromtimestamp(ms / 1000.0, ZoneInfo(GATEWAY_TZ)).strftime(
|
|
"%Y-%m-%d %H:%M:%S")
|
|
return datetime.fromtimestamp(ms / 1000.0).strftime("%Y-%m-%d %H:%M:%S")
|
|
|
|
|
|
def sql_rows(instances, first_id):
|
|
"""(event_values, data_values) as SQL literal tuples."""
|
|
events, data = [], []
|
|
next_id = first_id
|
|
for inst in instances:
|
|
stages = [(0, inst["active_ms"])]
|
|
if inst["ack_ms"]:
|
|
stages.append((2, inst["ack_ms"]))
|
|
if inst["clear_ms"]:
|
|
stages.append((1, inst["clear_ms"]))
|
|
for etype, ts in stages:
|
|
rid = next_id
|
|
next_id += 1
|
|
events.append("(%d,'%s','%s','',%d,%d,0,'%s')" % (
|
|
rid, inst["event_id"], inst["source"], inst["priority"],
|
|
etype, local_dt(ts)))
|
|
if etype == 0:
|
|
data.append("(%d,'eventValue',0,1,NULL,NULL)" % rid)
|
|
elif etype == 1:
|
|
data.append("(%d,'eventValue',0,0,NULL,NULL)" % rid)
|
|
else:
|
|
user = inst["ack_user"]
|
|
data.append("(%d,'ackUser',2,NULL,NULL,'usr-prov:%s:/usr:%s')"
|
|
% (rid, TAG_PROVIDER, user))
|
|
data.append("(%d,'ackUserName',2,NULL,NULL,'%s')" % (rid, user))
|
|
return events, data, next_id
|
|
|
|
|
|
def mariadb(sql, capture=True):
|
|
cmd = ["docker", "exec", "-i", DB_CONTAINER, "mariadb",
|
|
"-u" + DB_USER, "-p" + DB_PASS, DB_NAME, "-N", "-B", "-e", sql]
|
|
res = subprocess.run(cmd, capture_output=capture, text=True)
|
|
if res.returncode != 0:
|
|
raise SystemExit("mariadb failed: %s" % (res.stderr or "").strip())
|
|
return (res.stdout or "").strip()
|
|
|
|
|
|
def max_id():
|
|
out = mariadb("SELECT COALESCE(MAX(id),0) FROM %s;" % EVENTS_TABLE)
|
|
return int(out.split("\n")[0])
|
|
|
|
|
|
def parse_local(text):
|
|
"""'YYYY-MM-DD HH:MM:SS' in GATEWAY_TZ (how the journal stores it) -> ms."""
|
|
dt = datetime.strptime(text.strip(), "%Y-%m-%d %H:%M:%S")
|
|
if ZoneInfo is not None:
|
|
dt = dt.replace(tzinfo=ZoneInfo(GATEWAY_TZ))
|
|
return int(dt.timestamp() * 1000)
|
|
|
|
|
|
STATE_OF_EVENTTYPE = {0: "active", 1: "clear", 2: "ack"}
|
|
|
|
|
|
def read_existing(calc, start_ms, end_ms):
|
|
"""Journal rows already in [start_ms, end_ms], shaped like the dicts
|
|
alarms._norm_event hands calc - so the offline prediction sees the window the
|
|
gateway will actually score, not just the rows this run adds.
|
|
|
|
Ack attribution is left off: it feeds MTTA/ack-user lists, never health.
|
|
"""
|
|
out = mariadb("SELECT eventid,source,priority,eventtype,eventtime FROM %s "
|
|
"WHERE eventtime >= '%s' AND eventtime <= '%s' ORDER BY eventtime;"
|
|
% (EVENTS_TABLE, local_dt(start_ms), local_dt(end_ms)))
|
|
rows = []
|
|
for line in (out or "").split("\n"):
|
|
if not line.strip():
|
|
continue
|
|
eventid, source, priority, etype, eventtime = line.split("\t")
|
|
state = STATE_OF_EVENTTYPE.get(int(etype))
|
|
if state is None:
|
|
continue
|
|
lvl, name = calc.normalize_priority(int(priority))
|
|
rows.append({"event_id": eventid, "source": source, "display_path": "",
|
|
"priority": lvl, "priority_name": name, "state": state,
|
|
"ts": parse_local(eventtime), "ack_user": None,
|
|
"is_system": source.startswith("evt:") or ":/alm:" not in source})
|
|
return rows
|
|
|
|
|
|
def existing_activations(rows, start_ms, end_ms):
|
|
"""(bucket counts, per-10-min-cell counts) for the real activations already
|
|
inside the scored window - the head start the solver has to plan around."""
|
|
buckets = {"low": 0, "medium": 0, "high": 0}
|
|
cells = collections.Counter()
|
|
for r in rows:
|
|
if r["is_system"] or r["state"] != "active":
|
|
continue
|
|
if not (start_ms <= r["ts"] <= end_ms):
|
|
continue
|
|
bucket = BUCKET_OF.get(r["priority"])
|
|
if bucket is None:
|
|
continue # unbucketed priority: calc ignores it
|
|
buckets[bucket] += 1
|
|
cells[cell_start(r["ts"])] += 1
|
|
return buckets, cells
|
|
|
|
|
|
# ---------------- reporting ----------------
|
|
|
|
|
|
def show_health(label, bundle):
|
|
h = bundle["health"]
|
|
print(" %s: grade %s score %.2f" % (label, h["grade"], h["score"]))
|
|
for sub in h["subs"]:
|
|
score = "none" if sub["score"] is None else "%6.2f" % sub["score"]
|
|
print(" %-9s w=%.2f score=%s %s"
|
|
% (sub["key"], sub["weight"], score, sub["detail"]))
|
|
|
|
|
|
def prior_spec(total_counts, burst_count, burst_hours, existing_prior, capacity):
|
|
"""A deliberately worse preceding window, so the period-over-period delta
|
|
chips read as improving. Never scored - only the chips consume it, so its
|
|
counts just get clamped to what the source pool can carry."""
|
|
counts = {}
|
|
for bucket in total_counts:
|
|
want = int(round(total_counts[bucket] * PRIOR_WORSE))
|
|
counts[bucket] = max(0, min(want - existing_prior.get(bucket, 0),
|
|
capacity[bucket]))
|
|
burst = burst_count + 2 if burst_count else 0
|
|
if sum(counts.values()) < burst:
|
|
burst = 0
|
|
return {"counts": counts, "burst": (burst_hours + 0.4, burst),
|
|
"burst_area": "BoilerHouse"}
|
|
|
|
|
|
def main():
|
|
ap = argparse.ArgumentParser(description=__doc__,
|
|
formatter_class=argparse.RawDescriptionHelpFormatter)
|
|
ap.add_argument("--hours", type=float, default=8.0,
|
|
help="window length; must match the dashboard preset (default 8)")
|
|
ap.add_argument("--score", type=float, default=None,
|
|
help="target Alarm Health score (e.g. 85). Solves for the "
|
|
"seedable dials instead of using a canned profile")
|
|
ap.add_argument("--grade", default=None, choices=sorted(PROFILES),
|
|
help="target Alarm Health grade via a canned profile "
|
|
"(assumed to be B when --score is omitted)")
|
|
ap.add_argument("--shape", default="stable", choices=("stable", "balanced"),
|
|
help="--score only: stable holds the rate sub-score pinned at 100 so "
|
|
"the score does not drift as the rolling window slides; balanced "
|
|
"spreads the deficit across rate/flood/priority (default stable)")
|
|
ap.add_argument("--tolerance", type=float, default=0.05,
|
|
help="--score only: allowed miss on the predicted score; the UI shows "
|
|
"one decimal (default 0.05)")
|
|
ap.add_argument("--burst-hours", type=float, default=2.2,
|
|
help="--score only: put the flood burst this many hours before the "
|
|
"window end - it survives window drift until it ages out "
|
|
"(default 2.2)")
|
|
ap.add_argument("--standing-now", type=int, default=1,
|
|
help="live alarms active >24h, which cost 20 pts each on the "
|
|
"standing sub-score. NOT seedable - standing() reads "
|
|
"system.alarm.queryStatus. Confirm the real number with "
|
|
"the SimHarness ProbeTick probe (default 1)")
|
|
ap.add_argument("--no-prior", action="store_true",
|
|
help="skip the preceding comparison window (delta chips go 'new')")
|
|
ap.add_argument("--attempts", type=int, default=6,
|
|
help="re-plan with a fresh rng draw this many times before giving up "
|
|
"(placement can interact with rows already in the window)")
|
|
ap.add_argument("--seed", type=int, default=20260731)
|
|
ap.add_argument("--dry-run", action="store_true",
|
|
help="plan and predict only; write nothing")
|
|
args = ap.parse_args()
|
|
if args.score is None and args.grade is None:
|
|
args.grade = "B"
|
|
|
|
calc = load_calc()
|
|
sources = load_sources()
|
|
|
|
now_ms = int(datetime.now().timestamp() * 1000)
|
|
span = int(args.hours * HOUR_MS)
|
|
end_ms, start_ms = now_ms, now_ms - span
|
|
prior_start = start_ms - span
|
|
standing_n = args.standing_now
|
|
|
|
print("window %s -> %s (%.1f h)"
|
|
% (local_dt(start_ms), local_dt(end_ms), args.hours))
|
|
print("sources %d real sim alarms (bad-actor tags excluded)" % len(sources))
|
|
|
|
# Rows already in the window (a live active alarm, an earlier seed) are part
|
|
# of what the gateway will score, so they are part of the target.
|
|
existing_rows = read_existing(calc, prior_start, end_ms)
|
|
ex_buckets, ex_cells = existing_activations(existing_rows, start_ms, end_ms)
|
|
ex_prior, _ = existing_activations(existing_rows, prior_start, start_ms - 1)
|
|
print("existing %d journal row(s) in the window -> %d activation(s) %d/%d/%d "
|
|
"low/medium/high (folded into the target)"
|
|
% (len([r for r in existing_rows if start_ms <= r["ts"] <= end_ms]),
|
|
sum(ex_buckets.values()), ex_buckets["low"], ex_buckets["medium"],
|
|
ex_buckets["high"]))
|
|
print("standing %d live alarm(s) >24h (measured, not seeded)" % standing_n)
|
|
|
|
if args.score is not None:
|
|
capacity = bucket_capacity(sources)
|
|
solved = solve_spec(calc, args.score, args.hours, standing_n, ex_buckets,
|
|
args.shape, args.tolerance, capacity)
|
|
spec_current = {"counts": solved["counts"],
|
|
"burst": (args.burst_hours, solved["burst"]),
|
|
"burst_area": "Packaging"}
|
|
spec_prior = prior_spec(solved["total"], solved["burst"], args.burst_hours,
|
|
ex_prior, capacity)
|
|
t = solved["total"]
|
|
print()
|
|
print("solved (%s shape) for score %.2f:" % (args.shape, args.score))
|
|
print(" %d activations in window (%.2f/hr), %d/%d/%d low/medium/high, "
|
|
"dev %.2f, %d flooding 10-min cell(s)"
|
|
% (solved["n"], solved["n"] / args.hours, t["low"], t["medium"],
|
|
t["high"], solved["dev"], solved["flood_bins"]))
|
|
print(" sub-scores rate %.1f flood %.1f chatter %.1f standing %.1f "
|
|
"priority %.1f -> %.4f (%s)"
|
|
% (solved["subs"]["rate"], solved["subs"]["flood"],
|
|
solved["subs"]["chatter"], solved["subs"]["standing"],
|
|
solved["subs"]["priority"], solved["score"], solved["grade"]))
|
|
print(" to seed: %d/%d/%d low/medium/high (%d activations)"
|
|
% (solved["counts"]["low"], solved["counts"]["medium"],
|
|
solved["counts"]["high"], sum(solved["counts"].values())))
|
|
else:
|
|
profile = PROFILES[args.grade]
|
|
spec_current = profile["current"]
|
|
spec_prior = profile["prior"]
|
|
|
|
def attempt(seed):
|
|
rng = random.Random(seed)
|
|
cur, _burst_cell, per_cell = plan_window(sources, spec_current, start_ms,
|
|
end_ms, rng, ex_cells)
|
|
instances = list(cur)
|
|
if not args.no_prior:
|
|
pri, _, _ = plan_window(sources, spec_prior, prior_start, start_ms, rng)
|
|
instances = pri + instances
|
|
rows = existing_rows + to_journal_rows(instances, calc)
|
|
active_now = synthetic_standing(calc, standing_n, now_ms)
|
|
bundle = predict(calc, rows, active_now, start_ms, end_ms, now_ms)
|
|
return cur, instances, rows, per_cell, bundle
|
|
|
|
def missed(bundle):
|
|
if args.score is not None:
|
|
s = bundle["health"]["score"]
|
|
if s is None:
|
|
return "no score (empty window)"
|
|
if abs(s - args.score) > args.tolerance:
|
|
return "score %.4f != %.2f +/-%.2f" % (s, args.score, args.tolerance)
|
|
return None
|
|
g = bundle["health"]["grade"]
|
|
return None if g == args.grade else "grade %s != %s" % (g, args.grade)
|
|
|
|
why = None
|
|
for i in range(max(1, args.attempts)):
|
|
cur, instances, rows, per_cell, bundle = attempt(args.seed + i)
|
|
why = missed(bundle)
|
|
if why is None:
|
|
break
|
|
print(" attempt %d missed (%s); re-planning" % (i + 1, why))
|
|
|
|
m = bundle["meta"]
|
|
print()
|
|
print("predicted (offline, via the gateway's own PrimeControls.calc):")
|
|
print(" activations %d over %.1f h = %.2f/hr"
|
|
% (m["activation_count"], m["window_hours"], bundle["kpis"]["rate_per_hr"]["value"]))
|
|
print(" floods %d episode(s), %.3f%% of window"
|
|
% (len(bundle["floods"]["episodes"]), bundle["floods"]["pct_time_in_flood"]))
|
|
print(" chattering %d fleeting %d standing %d"
|
|
% (len(bundle["chattering"]), bundle["fleeting"]["total"],
|
|
bundle["standing"]["count"]))
|
|
pct = bundle["priority"]["pct"]
|
|
print(" priority %.1f/%.1f/%.1f vs 80/15/5 (dev %.1f)"
|
|
% (pct["low"], pct["medium"], pct["high"], bundle["priority"]["sum_abs_dev"]))
|
|
print(" busiest 10-min cell (burst cell excluded): %d burst cell: %d"
|
|
% (max([c for k, c in per_cell.items()
|
|
if c < spec_current["burst"][1] or not spec_current["burst"][1]] or [0]),
|
|
spec_current["burst"][1]))
|
|
show_health("health", bundle)
|
|
|
|
if why is not None:
|
|
raise SystemExit("\nABORT after %d attempt(s): %s; nothing written."
|
|
% (max(1, args.attempts), why))
|
|
if args.score is not None:
|
|
print("\nprediction matches target score %.2f (%.4f, grade %s)"
|
|
% (args.score, bundle["health"]["score"], bundle["health"]["grade"]))
|
|
else:
|
|
print("\nprediction matches target grade %s" % args.grade)
|
|
|
|
event_rows = len(rows) - len(existing_rows)
|
|
if args.dry_run:
|
|
print("dry run: would insert %d event rows for %d activations "
|
|
"(%d in the scored window)"
|
|
% (event_rows, len(instances), len(cur)))
|
|
return
|
|
|
|
first = max_id() + 1
|
|
events, data, last = sql_rows(instances, first)
|
|
stmt = ("INSERT INTO %s (id,eventid,source,displaypath,priority,eventtype,"
|
|
"eventflags,eventtime) VALUES %s;\nINSERT INTO %s "
|
|
"(id,propname,dtype,intvalue,floatvalue,strvalue) VALUES %s;"
|
|
% (EVENTS_TABLE, ",".join(events), DATA_TABLE, ",".join(data)))
|
|
mariadb(stmt)
|
|
print("inserted %d event rows + %d data rows (ids %d..%d)"
|
|
% (len(events), len(data), first, last - 1))
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|