Files
Enelix-EMS/services/netplan-v4/netplan_v4/retention.py
T

107 lines
6.1 KiB
Python

"""Bounded online tables with verified lossless archives; never prune unarchived data.
Runs at most hourly. Raw input history: 120 days (longer than 90-day training and
comparison limits). Ordinary plans: 7 days, excluding active/acknowledged plans.
Archives are retained; external backup/long-term archive lifecycle is operational.
"""
from pathlib import Path
import gzip
import hashlib
import json
import os
import sqlite3
import tempfile
def schema(con):
con.executescript('''
CREATE TABLE IF NOT EXISTS planner_maintenance(
name TEXT PRIMARY KEY, checked_at INTEGER NOT NULL, status TEXT NOT NULL, detail TEXT NOT NULL);
CREATE TABLE IF NOT EXISTS planner_archives(
id TEXT PRIMARY KEY, created_at INTEGER NOT NULL, filename TEXT NOT NULL,
sha256 TEXT NOT NULL, records INTEGER NOT NULL, kind TEXT NOT NULL);
CREATE INDEX IF NOT EXISTS planner_plan_created ON planner_plans(created_at);
''')
def archive_batch(con,directory,now,kind,limit=100):
if kind not in ('observations','plans') or type(limit) is not int or not 1<=limit<=1000:
raise ValueError('Invalid archival scope')
directory=Path(directory)
if directory.is_symlink():raise ValueError('Archive symlink refused')
directory.mkdir(mode=0o700,parents=True,exist_ok=True)
if directory.stat().st_mode & 0o007:raise ValueError('Archive directory must not be public')
from datetime import datetime,timezone
with_context=False
con.execute('BEGIN IMMEDIATE')
try:
if kind=='observations':
rows=con.execute('SELECT * FROM planner_observations WHERE captured_at<? AND received_at<? ORDER BY captured_at LIMIT ?',
(now-120*86400,now-7*86400,limit)).fetchall()
else:
before=datetime.fromtimestamp(now-7*86400,timezone.utc).isoformat()
rows=con.execute('SELECT * FROM planner_plans p WHERE created_at<? AND NOT EXISTS(SELECT 1 FROM planner_current c WHERE c.plan_id=p.plan_id) AND NOT EXISTS(SELECT 1 FROM planner_ack a WHERE a.plan_id=p.plan_id) ORDER BY created_at LIMIT ?',
(before,limit)).fetchall()
if not rows:con.commit();return {'archived':0,'kind':kind}
records=[]
for r in rows:
item={'table':kind,'row':dict(r)}
if kind=='observations':
item['origins']=[dict(x) for x in con.execute('SELECT * FROM planner_observation_origins WHERE plant=? AND dataset=? AND captured_at=?',(r['plant'],r['dataset'],r['captured_at']))]
records.append(json.dumps(item,sort_keys=True,separators=(',',':'),allow_nan=False).encode()+b'\n')
raw=b''.join(records)
if len(raw)>67108864:raise ValueError('Archive batch exceeds memory budget')
ident=hashlib.sha256(raw).hexdigest();name=kind+'-'+ident+'.jsonl.gz';path=directory/name
if path.is_symlink():raise ValueError('Archive target symlink refused')
if not path.exists():
fd,tmp=tempfile.mkstemp(prefix='.archive-',dir=directory)
try:
os.fchmod(fd,0o600)
with os.fdopen(fd,'wb') as out:
with gzip.GzipFile(fileobj=out,mode='wb',mtime=0) as zipped:zipped.write(raw)
out.flush();os.fsync(out.fileno())
os.replace(tmp,path)
dfd=os.open(directory,os.O_RDONLY)
try:os.fsync(dfd)
finally:os.close(dfd)
finally:
if os.path.exists(tmp):os.unlink(tmp)
# Read back exact bytes before removing any database row.
with gzip.open(path,'rb') as f:verified=f.read(len(raw)+1)
if verified!=raw:raise ValueError('Archive verification failed; original rows retained')
for r in rows:
if kind=='observations':
key=(r['plant'],r['dataset'],r['captured_at'])
con.execute('DELETE FROM planner_observation_origins WHERE plant=? AND dataset=? AND captured_at=?',key)
con.execute('DELETE FROM planner_observations WHERE plant=? AND dataset=? AND captured_at=? AND fingerprint=?',(*key,r['fingerprint']))
else:con.execute('DELETE FROM planner_plans WHERE plan_id=?',(r['plan_id'],))
con.execute('INSERT OR IGNORE INTO planner_archives VALUES(?,?,?,?,?,?)',(ident,now,name,ident,len(rows),kind))
con.commit()
return {'archived':len(rows),'kind':kind,'file':name,'sha256':ident,'sourceRecoverable':True}
except Exception:
con.rollback();raise
def maintain(store,now):
if os.environ.get('NETPLAN_V4_ARCHIVE_ENABLED','0')!='1':return
con=store.con;stamp=int(now.timestamp());old=con.execute("SELECT checked_at FROM planner_maintenance WHERE name='archive'").fetchone()
if old and stamp-old[0]<3600:return
db=con.execute('PRAGMA database_list').fetchone()[2]
if not db:return # In-memory test/ephemeral databases have no archival location.
try:
directory=Path(db).resolve().parent/'archives'
results=[archive_batch(con,directory,stamp,kind) for kind in ('plans','observations')]
if any(r['archived']>=100 for r in results):stamp-=3300
state='ok';detail={'results':results,'archivesRetained':True,'rawRetentionDays':120,'ordinaryPlanRetentionDays':7}
except (OSError,ValueError,sqlite3.Error) as exc:
state='archive_error';detail={'errorType':type(exc).__name__,'unverifiedDataNotDeleted':True}
with con:con.execute('INSERT INTO planner_maintenance VALUES(?,?,?,?) ON CONFLICT(name) DO UPDATE SET checked_at=excluded.checked_at,status=excluded.status,detail=excluded.detail',
('archive',stamp,state,json.dumps(detail,separators=(',',':'))))
def status(con):
if os.environ.get('NETPLAN_V4_ARCHIVE_ENABLED','0')!='1':return {'status':'disabled_requires_operator_opt_in'}
r=con.execute("SELECT checked_at,status,detail FROM planner_maintenance WHERE name='archive'").fetchone()
# Public per-plant state must not expose archival filenames/rows of other tenants.
return {'checkedAtEpoch':r[0],'status':r[1],'rawRetentionDays':120,'ordinaryPlanRetentionDays':7,'archivesRetained':True} if r else {'status':'not_yet_run'}