retention piece 1: immediate rotations for chat-history / job-log / followups (b67084660385)
- chat-history.jsonl rotates at 10MB or 7d -> logs/archive/*.jsonl.gz + .sha256 - job-log.jsonl rotates at 2MB or 14d -> same archive layout - followups.json archives resolved>7d / escalated>30d -> followups.archive.jsonl (append-only) - verify wrapper: sha256 -c, gzip -t, clean listing, spot-extract JSON + ts bounds - driver + hourly systemd user timer (retention-rotations.timer) - archive-only: nothing is ever deleted; thread backfill excluded (needs human-reviewed preview) First run 2026-10-05 16:35Z: chat-history 29.7MB + job-log 4.7MB rotated and verified; followups 0 eligible of 163.
This commit is contained in:
Executable
+93
@@ -0,0 +1,93 @@
|
||||
#!/usr/bin/env python3
|
||||
"""retention-archive-followups.py — retention piece 1 (board bug b67084660385).
|
||||
|
||||
Archives resolved follow-ups older than 7 days and escalated follow-ups older
|
||||
than 30 days from followups.json into followups.archive.jsonl (append-only).
|
||||
|
||||
ARCHIVE-ONLY: records are moved from the live file to the archive file; the
|
||||
archive itself is never pruned or hard-deleted by this script.
|
||||
|
||||
Concurrency: the live file is written atomically (tmp + os.replace) by both
|
||||
this script and bin/followup-sweeper.py. This script snapshots the file's
|
||||
(mtime_ns, size) before deciding the archival set and aborts (rc=2) if the
|
||||
file changed in the meantime — the next run picks it up.
|
||||
"""
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
from datetime import datetime, timezone, timedelta
|
||||
|
||||
ROOT = os.environ.get("NETVM_ROOT", "/home/super/Projects/NetVM")
|
||||
LIVE = os.path.join(ROOT, "followups.json")
|
||||
ARCHIVE = os.path.join(ROOT, "followups.archive.jsonl")
|
||||
RESOLVED_AFTER_DAYS = 7
|
||||
ESCALATED_AFTER_DAYS = 30
|
||||
|
||||
|
||||
def parse_ts(ts):
|
||||
if not ts:
|
||||
return None
|
||||
try:
|
||||
dt = datetime.fromisoformat(str(ts).replace("Z", "+00:00"))
|
||||
return dt if dt.tzinfo is not None else dt.replace(tzinfo=timezone.utc)
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def main():
|
||||
if not os.path.exists(LIVE):
|
||||
print("SKIP followups: no live file at %s" % LIVE)
|
||||
return 0
|
||||
|
||||
with open(LIVE, encoding="utf-8") as f:
|
||||
raw = f.read()
|
||||
st_before = os.stat(LIVE)
|
||||
data = json.loads(raw)
|
||||
now = datetime.now(timezone.utc)
|
||||
|
||||
archive, keep = [], {}
|
||||
for key, rec in data.items():
|
||||
status = rec.get("status")
|
||||
eligible = False
|
||||
reason = ""
|
||||
if status == "resolved":
|
||||
dt = parse_ts(rec.get("resolved_at")) or parse_ts(rec.get("sent_at"))
|
||||
if dt and (now - dt) > timedelta(days=RESOLVED_AFTER_DAYS):
|
||||
eligible, reason = True, "resolved>7d"
|
||||
elif status == "escalated":
|
||||
dt = parse_ts(rec.get("escalated_at")) or parse_ts(rec.get("sent_at"))
|
||||
if dt and (now - dt) > timedelta(days=ESCALATED_AFTER_DAYS):
|
||||
eligible, reason = True, "escalated>30d"
|
||||
if eligible:
|
||||
out = dict(rec)
|
||||
out["_archived_at"] = now.isoformat()
|
||||
out["_archive_reason"] = reason
|
||||
archive.append(out)
|
||||
else:
|
||||
keep[key] = rec
|
||||
|
||||
if not archive:
|
||||
print("SKIP followups: 0 eligible of %d records (resolved>7d / escalated>30d)" % len(data))
|
||||
return 0
|
||||
|
||||
st_after = os.stat(LIVE)
|
||||
if (st_after.st_mtime_ns, st_after.st_size) != (st_before.st_mtime_ns, st_before.st_size):
|
||||
print("ABORT followups: live file changed during archival decision; retry next run",
|
||||
file=sys.stderr)
|
||||
return 2
|
||||
|
||||
with open(ARCHIVE, "a", encoding="utf-8") as f:
|
||||
for rec in archive:
|
||||
f.write(json.dumps(rec) + "\n")
|
||||
|
||||
tmp = "%s.tmp.retention.%d" % (LIVE, os.getpid())
|
||||
with open(tmp, "w", encoding="utf-8") as f:
|
||||
json.dump(keep, f, indent=2)
|
||||
os.replace(tmp, LIVE)
|
||||
|
||||
print("OK followups: archived %d records -> %s (live now %d)" % (len(archive), ARCHIVE, len(keep)))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
Executable
+45
@@ -0,0 +1,45 @@
|
||||
#!/usr/bin/env bash
|
||||
# retention-rotate-chat-history.sh — retention piece 1 (board bug b67084660385).
|
||||
# Rotates logs/chat-history.jsonl when it exceeds 10MB or is older than 7 days.
|
||||
#
|
||||
# Live-writer safety: the harvester (bin/response-harvester.py) opens the file
|
||||
# in append mode per write and closes it immediately, so an atomic `mv` of the
|
||||
# live file is safe — the next write recreates a fresh file. We also `touch`
|
||||
# the new live file for tidiness and to keep ownership/visibility stable.
|
||||
#
|
||||
# Archives land in logs/archive/chat-history-YYYYMMDD-HHMMSS.jsonl.gz with a
|
||||
# sha256 sidecar. Nothing is ever deleted by this script.
|
||||
set -euo pipefail
|
||||
|
||||
ROOT="${NETVM_ROOT:-/home/super/Projects/NetVM}"
|
||||
LIVE="$ROOT/logs/chat-history.jsonl"
|
||||
ARCHIVE_DIR="$ROOT/logs/archive"
|
||||
SIZE_LIMIT=$((10 * 1024 * 1024)) # 10MB
|
||||
AGE_LIMIT_DAYS=7
|
||||
TS="$(date -u +%Y%m%d-%H%M%S)"
|
||||
|
||||
mkdir -p "$ARCHIVE_DIR"
|
||||
|
||||
if [ ! -f "$LIVE" ]; then
|
||||
echo "SKIP chat-history: no live file at $LIVE"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
size="$(stat -c %s "$LIVE")"
|
||||
age_days=$(( ( $(date +%s) - $(stat -c %Y "$LIVE") ) / 86400 ))
|
||||
|
||||
if [ "$size" -lt "$SIZE_LIMIT" ] && [ "$age_days" -lt "$AGE_LIMIT_DAYS" ]; then
|
||||
echo "SKIP chat-history: ${size}B, ${age_days}d old (limits ${SIZE_LIMIT}B / ${AGE_LIMIT_DAYS}d)"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
lines="$(wc -l < "$LIVE")"
|
||||
base="chat-history-$TS"
|
||||
echo "ROTATE chat-history: ${size}B / ${lines} lines (${age_days}d old) -> $ARCHIVE_DIR/$base.jsonl.gz"
|
||||
|
||||
mv "$LIVE" "$ARCHIVE_DIR/$base.jsonl"
|
||||
gzip -9 "$ARCHIVE_DIR/$base.jsonl"
|
||||
( cd "$ARCHIVE_DIR" && sha256sum "$base.jsonl.gz" > "$base.jsonl.gz.sha256" )
|
||||
touch "$LIVE"
|
||||
|
||||
echo "OK chat-history: archived $base.jsonl.gz, live file recreated empty"
|
||||
Executable
+45
@@ -0,0 +1,45 @@
|
||||
#!/usr/bin/env bash
|
||||
# retention-rotate-job-log.sh — retention piece 1 (board bug b67084660385).
|
||||
# Rotates job-log.jsonl when it exceeds 2MB or is older than 14 days.
|
||||
#
|
||||
# Live-writer safety: all writers (bin/job-dispatch.py, bin/followup-sweeper.py)
|
||||
# open the file in append mode per write, so an atomic `mv` of the live file
|
||||
# is safe — the next write recreates a fresh file. We also `touch` the new
|
||||
# live file for tidiness.
|
||||
#
|
||||
# Archives land in logs/archive/job-log-YYYYMMDD-HHMMSS.jsonl.gz with a
|
||||
# sha256 sidecar. Nothing is ever deleted by this script.
|
||||
set -euo pipefail
|
||||
|
||||
ROOT="${NETVM_ROOT:-/home/super/Projects/NetVM}"
|
||||
LIVE="$ROOT/job-log.jsonl"
|
||||
ARCHIVE_DIR="$ROOT/logs/archive"
|
||||
SIZE_LIMIT=$((2 * 1024 * 1024)) # 2MB
|
||||
AGE_LIMIT_DAYS=14
|
||||
TS="$(date -u +%Y%m%d-%H%M%S)"
|
||||
|
||||
mkdir -p "$ARCHIVE_DIR"
|
||||
|
||||
if [ ! -f "$LIVE" ]; then
|
||||
echo "SKIP job-log: no live file at $LIVE"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
size="$(stat -c %s "$LIVE")"
|
||||
age_days=$(( ( $(date +%s) - $(stat -c %Y "$LIVE") ) / 86400 ))
|
||||
|
||||
if [ "$size" -lt "$SIZE_LIMIT" ] && [ "$age_days" -lt "$AGE_LIMIT_DAYS" ]; then
|
||||
echo "SKIP job-log: ${size}B, ${age_days}d old (limits ${SIZE_LIMIT}B / ${AGE_LIMIT_DAYS}d)"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
lines="$(wc -l < "$LIVE")"
|
||||
base="job-log-$TS"
|
||||
echo "ROTATE job-log: ${size}B / ${lines} lines (${age_days}d old) -> $ARCHIVE_DIR/$base.jsonl.gz"
|
||||
|
||||
mv "$LIVE" "$ARCHIVE_DIR/$base.jsonl"
|
||||
gzip -9 "$ARCHIVE_DIR/$base.jsonl"
|
||||
( cd "$ARCHIVE_DIR" && sha256sum "$base.jsonl.gz" > "$base.jsonl.gz.sha256" )
|
||||
touch "$LIVE"
|
||||
|
||||
echo "OK job-log: archived $base.jsonl.gz, live file recreated empty"
|
||||
Executable
+46
@@ -0,0 +1,46 @@
|
||||
#!/usr/bin/env bash
|
||||
# retention-run-rotations.sh — retention piece 1 driver (board bug b67084660385).
|
||||
# Runs all three piece-1 rotations (chat-history, job-log, followups archive)
|
||||
# then verifies every archive touched. Idempotent: no-op when under threshold.
|
||||
# A full run report is appended to logs/retention-runs/retention-run-TS.log.
|
||||
# Exits 2 if any rotation or verification fails (so the timer run is loud).
|
||||
set -uo pipefail
|
||||
|
||||
ROOT="${NETVM_ROOT:-/home/super/Projects/NetVM}"
|
||||
BIN="$ROOT/bin"
|
||||
ARCHIVE_DIR="$ROOT/logs/archive"
|
||||
RUN_TS="$(date -u +%Y%m%d-%H%M%S)"
|
||||
LOGDIR="$ROOT/logs/retention-runs"
|
||||
mkdir -p "$LOGDIR"
|
||||
LOG="$LOGDIR/retention-run-$RUN_TS.log"
|
||||
|
||||
{
|
||||
echo "=== retention-run $RUN_TS (UTC) ==="
|
||||
|
||||
fail=0
|
||||
|
||||
"$BIN/retention-rotate-chat-history.sh" || fail=1
|
||||
"$BIN/retention-rotate-job-log.sh" || fail=1
|
||||
"$BIN/retention-archive-followups.py" || fail=1
|
||||
|
||||
echo "--- verification ---"
|
||||
|
||||
# Verify the newest chat-history / job-log archives (skip if none exist)
|
||||
for pattern in "chat-history-*.jsonl.gz" "job-log-*.jsonl.gz"; do
|
||||
newest="$(ls -t "$ARCHIVE_DIR"/$pattern 2>/dev/null | head -1 || true)"
|
||||
if [ -n "$newest" ]; then
|
||||
"$BIN/retention-verify-archive.sh" "$newest" || fail=1
|
||||
else
|
||||
echo "SKIP verify: no $pattern in $ARCHIVE_DIR"
|
||||
fi
|
||||
done
|
||||
|
||||
# Verify the followups archive (append-only, cumulative)
|
||||
"$BIN/retention-verify-archive.sh" "$ROOT/followups.archive.jsonl" || fail=1
|
||||
|
||||
echo "--- live sizes after run ---"
|
||||
ls -la "$ROOT/logs/chat-history.jsonl" "$ROOT/job-log.jsonl" "$ROOT/followups.json" 2>/dev/null || true
|
||||
|
||||
echo "=== retention-run done rc=$fail ==="
|
||||
exit $fail
|
||||
} 2>&1 | tee -a "$LOG"
|
||||
Executable
+67
@@ -0,0 +1,67 @@
|
||||
#!/usr/bin/env bash
|
||||
# retention-verify-archive.sh — retention piece 1 (board bug b67084660385).
|
||||
# Verifies one archive file: checksum, integrity, clean listing, spot-extract.
|
||||
# Usage: retention-verify-archive.sh <path-to-.jsonl.gz-or-.jsonl>
|
||||
# Exits non-zero on any verification failure. Missing file = SKIP (rc 0).
|
||||
set -uo pipefail
|
||||
|
||||
fail() { echo "FAIL verify $1: $2"; exit 1; }
|
||||
|
||||
P="${1:?usage: $0 <archive-path>}"
|
||||
|
||||
if [ ! -f "$P" ]; then
|
||||
echo "SKIP verify: $P does not exist"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
echo "== verify $P =="
|
||||
|
||||
# 1. checksum sidecar
|
||||
if [ -f "$P.sha256" ]; then
|
||||
( cd "$(dirname "$P")" && sha256sum -c "$(basename "$P").sha256" ) \
|
||||
|| fail "$P" "sha256 mismatch"
|
||||
echo " checksum: OK"
|
||||
else
|
||||
echo " checksum: no sidecar (archiver writes one for new .gz files)"
|
||||
fi
|
||||
|
||||
# 2. integrity + clean listing for gzip archives
|
||||
if [[ "$P" == *.gz ]]; then
|
||||
gzip -t "$P" || fail "$P" "gzip integrity test failed"
|
||||
echo " integrity: gzip -t OK"
|
||||
echo " listing:"; gzip -l "$P" | sed 's/^/ /'
|
||||
reader="zcat"
|
||||
else
|
||||
reader="cat"
|
||||
fi
|
||||
|
||||
# 3. spot-extract: first 3 and last 3 lines must be valid JSON
|
||||
n=0
|
||||
while IFS= read -r line; do
|
||||
n=$((n+1))
|
||||
echo "$line" | python3 -c 'import json,sys; json.loads(sys.stdin.read())' \
|
||||
|| fail "$P" "line $n is not valid JSON"
|
||||
done < <( { $reader "$P" | head -3; $reader "$P" | tail -3; } 2>/dev/null )
|
||||
|
||||
total="$( $reader "$P" | wc -l )"
|
||||
echo " spot-extract: first/last 3 lines valid JSON ($total total lines)"
|
||||
|
||||
# 4. spot-extract timestamps: report newest/oldest ts seen in the sample
|
||||
sample="$($reader "$P" | head -2000)"
|
||||
ts_line="$(echo "$sample" | python3 -c '
|
||||
import json,sys
|
||||
keys=("ts","timestamp","time","created_at","sent_at","resolved_at")
|
||||
seen=[]
|
||||
for line in sys.stdin:
|
||||
line=line.strip()
|
||||
if not line: continue
|
||||
try: rec=json.loads(line)
|
||||
except Exception: continue
|
||||
for k in keys:
|
||||
if rec.get(k): seen.append(str(rec[k])); break
|
||||
seen=sorted(set(seen))
|
||||
print(("oldest="+seen[0]+" newest="+seen[-1]) if seen else "no-ts-field")
|
||||
')"
|
||||
echo " timestamps: $ts_line"
|
||||
|
||||
echo "OK verify $P"
|
||||
Reference in New Issue
Block a user