Shuiquan-Future-Lab/future-science-algorithms
#!/usr/bin/env python3
forensicoilcomplete_reform.py
完整改良版单文件:法证级多模态油脂包浆鉴定(一次性全面改革)
目标:一次性解决之前所有失败与不足,提供工程化、可测试、可部署的单文件实现。
说明:
- 生产必须替换 TODO: REPLACE 的厂商解析器与提供真实 CS-14C 标定数据
- 请通过 KMS/Secret 管理 PROVENANCE_SECRET,不要在环境中明文存放
- 运行示例:
python forensicoilcomplete_reform.py --mode build-dummy
python forensicoilcompletereform.py --mode train --manifest trainmanifest.json --model-dir models
python forensicoilcomplete_reform.py --mode infer-file --sample-file sample.json --model-dir models
python forensicoilcomplete_reform.py --mode auto
uvicorn forensicoilcomplete_reform:app --host 0.0.0.0 --port 8000
import os import io import sys import json import time import uuid import math import hmac import yaml import base64 import hashlib import logging import datetime from typing import Dict, Any, List, Tuple, Optional from threading import Thread, Event from queue import Queue, Empty from contextlib import asynccontextmanager
numpy compatibility
try: import numpy as np if not hasattr(np, 'bool'): np.bool = np.bool_ except Exception: raise
import pandas as pd from scipy import signal, optimize, stats from scipy.ndimage import gaussianfilter from sklearn.isotonic import IsotonicRegression from sklearn.metrics import rocaucscore, precisionrecallfscoresupport from sklearn.modelselection import traintest_split import xgboost as xgb
Optional heavy deps (graceful fallback)
try: import pymzml except Exception: pymzml = None try: import tifffile except Exception: tifffile = None try: import pymc3 as pm import arviz as az except Exception: pm = None try: import shap except Exception: shap = None try: from PIL import Image, UnidentifiedImageError except Exception: Image = None UnidentifiedImageError = Exception
Web API optional
try: from fastapi import FastAPI, UploadFile, File, HTTPException from fastapi.responses import HTMLResponse, JSONResponse from pydantic import BaseModel except Exception: FastAPI = None
Logging
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s: %(message)s") logger = logging.getLogger("forensicoilcomplete_reform")
----------------------------
Default configuration
----------------------------
DEFAULTCONFIG = { "modelversion": "5.0.0-complete-reform", "provenancesecretenv": "PROVENANCESECRET", "auditdir": "auditlogs", "modeldir": "models", "incomingdir": "incomingsamples", "minsamplemassmg": 0.2, "gcms": { "binnedbins": 1024, "markerwindows": [[270,290],[300,320],[340,360],[400,420]], "topk": 256, "peakprominence": 1e3, "peakwidth": (1, 200) }, "spectrum": {"binnedbins": 512, "markerpeaks": [1700,1740,2850,2920,1650]}, "image": {"patchsize": 256, "histbins": 128}, "afm": {"features": ["maxforce","minforce","adhesion","slope","auc"]}, "fusion": {"confidencehigh": 0.92, "confidencelow": 0.35}, "training": {"xgbrounds": 200, "randomseed": 42}, "blindtest": {"nmin": 200, "auctarget": 0.95, "fprtarget": 0.05}, "bayes": {"mcmcdraws": 1000, "mcmctune": 500}, "auto": { "watchdir": "incomingsamples", "pollintervals": 5, "maxworkers": 2, "humanreviewconfidencethreshold": 0.75, "autoacceptconfidence": 0.92, "maxqueue_size": 1000 } }
----------------------------
Utilities
----------------------------
def now_iso() -> str: return datetime.datetime.utcnow().replace(microsecond=0).isoformat() + "Z"
def ensuredir(p: str): os.makedirs(p, existok=True)
def canonicaljson(obj: Any) -> str: return json.dumps(obj, sortkeys=True, ensure_ascii=False, separators=(',',':'))
def computeprovenance(meta: Dict[str, Any], filesbytes: List[bytes], secret: str) -> Dict[str,str]: h = hashlib.sha256() h.update(canonicaljson(meta).encode('utf-8')) for fb in filesbytes: if fb: h.update(fb) digest = h.hexdigest() mac = hmac.new(secret.encode('utf-8'), digest.encode('utf-8'), hashlib.sha256).hexdigest() return {"hash": digest, "hmac": mac}
----------------------------
Error codes & helpers
----------------------------
ERRORCODES = { "PARSEGCMSFAIL": "E1001", "PARSESPECTRUMFAIL": "E1002", "PARSEIMAGEFAIL": "E1003", "PARSEAFMFAIL": "E1004", "MISSINGMODALITY": "E2001", "MODELLOADFAIL": "E3001", "INFERENCE_FAIL": "E3002" }
def auditerror(auditdir: str, sampleid: str, code: str, message: str, rawbytes: Optional[bytes]=None): ensuredir(auditdir) rec = {"sampleid": sampleid, "errorcode": code, "message": message, "timestamp": nowiso()} if rawbytes: rec['rawhash'] = hashlib.sha256(rawbytes).hexdigest() fname = os.path.join(auditdir, f"error{sampleid}{int(time.time())}{uuid.uuid4().hex[:6]}.json") with open(fname, 'w', encoding='utf-8') as f: json.dump(rec, f, ensure_ascii=False, indent=2) logger.debug("Audit error written: %s", fname)
----------------------------
Robust parsing utilities (comprehensive)
- Handles placeholders, single-column CSVs, non-image bytes, and tries multiple fallbacks
----------------------------
PLACEHOLDERMARKERS = {b"mzmlplaceholder", b"spectrumplaceholder", b"tiffplaceholder", b"afmplaceholder", b"cs14cplaceholder", b"tiff", b"image", b"placeholder"}
def isplaceholder(b: Optional[bytes]) -> bool: if not b: return True if isinstance(b, bytes) and len(b) < 1024: low = b.strip().lower() for m in PLACEHOLDERMARKERS: if m in low: return True return False
def tryreadcsvbytes(b: bytes, sepguesses=[',','\t',';',' ']) -> pd.DataFrame: s = b.decode('utf-8', errors='ignore') # quick sniff: if looks like JSON, raise if s.strip().startswith('{') or s.strip().startswith('['): raise ValueError("Not CSV") # try pandas default try: df = pd.readcsv(io.StringIO(s)) if df.shape[1] >= 2: return df except Exception: pass # try guesses for sep in sepguesses: try: df = pd.read_csv(io.StringIO(s), sep=sep) if df.shape[1] >= 2: return df except Exception: continue # try single-column splitting heuristics lines = [ln.strip() for ln in s.splitlines() if ln.strip()] rows = [] for ln in lines: parts = ln.replace(',', ' ').replace('\t',' ').split() if len(parts) >= 2: rows.append(parts[:2]) if rows: df = pd.DataFrame(rows) return df # fallback: raise raise ValueError("CSV parse failed or single-column without numeric pairs")
----------------------------
Parsers (robust, with synthetic fallback)
----------------------------
def parsemzmlprecise(mzml_bytes: Optional[bytes]) -> Dict[str, Any]: """ Return {'peaks': DataFrame with columns ['mz','intensity','rt'], 'meta': {...}, 'note':...} Robust behavior:
- If placeholder or None -> return empty peaks and note 'placeholder'
- If pymzml available and valid mzML -> parse
- Else try CSV heuristics
- On failure return empty peaks and audit error """ meta = {} if isplaceholder(mzmlbytes): return {'peaks': pd.DataFrame(columns=['mz','intensity','rt']), 'meta': {'parser':'placeholder'}, 'note':'placeholder'} if pymzml is not None: try: with io.BytesIO(mzmlbytes) as bio: run = pymzml.run.Reader(bio) peaks = [] for spec in run: if getattr(spec, 'mslevel', 1) == 1: mzs = np.array(getattr(spec, 'mz', [])) ints = np.array(getattr(spec, 'i', [])) try: rt = spec.scantimeinminutes() except Exception: rt = np.nan if len(mzs) == 0: continue for m, it in zip(mzs, ints): peaks.append((float(m), float(it), float(rt))) if len(peaks) == 0: raise ValueError("No MS1 peaks found") df = pd.DataFrame(peaks, columns=['mz','intensity','rt']) return {'peaks': df, 'meta': {'parser':'pymzml'}, 'note':'pymzml'} except Exception as e: logger.warning("pymzml parse failed: %s. Will try CSV fallback.", e) # CSV fallback try: df = tryreadcsvbytes(mzmlbytes) # normalize column names cols = [c.lower() for c in df.columns] if 'mz' in cols and 'intensity' in cols: df2 = df.copy() df2.columns = cols if 'rt' not in df2.columns: df2['rt'] = np.nan return {'peaks': df2[['mz','intensity','rt']].copy(), 'meta': {'parser':'csvfallback'}, 'note':'csvfallback'} # if first two columns numeric, treat as mz,intensity numericcols = [] for c in df.columns[:2]: try: pd.tonumeric(df[c]) numericcols.append(c) except Exception: pass if len(numericcols) >= 2: df2 = df.copy() df2.columns = ['mz','intensity'] + list(df.columns[2:]) if 'rt' not in df2.columns: df2['rt'] = np.nan return {'peaks': df2[['mz','intensity','rt']].copy(), 'meta': {'parser':'csvguess'}, 'note':'csvguess'} except Exception as e: logger.exception("parsemzml_precise failed: %s", e) # final fallback: empty return {'peaks': pd.DataFrame(columns=['mz','intensity','rt']), 'meta': {'parser':'failed'}, 'note':'failed'}
def parsejcampprecise(b: Optional[bytes]) -> Dict[str, Any]: if isplaceholder(b): return {'wn': np.array([], dtype=np.float32), 'intensity': np.array([], dtype=np.float32), 'meta': {'parser':'placeholder'}, 'note':'placeholder'} s = b.decode('utf-8', errors='ignore') if '##' in s and 'JCAMP' in s.upper(): header = {} data = [] inxy = False for ln in s.splitlines(): ln = ln.strip() if ln.startswith('##'): parts = ln[2:].split('=',1) if len(parts) == 2: header[parts[0].strip()] = parts[1].strip() if ln.upper().startswith('##XYDATA'): inxy = True continue if inxy: for token in ln.split(): try: data.append(float(token)) except: pass if len(data) >= 4 and len(data) % 2 == 0: wn = np.array(data[0::2], dtype=np.float32) inten = np.array(data[1::2], dtype=np.float32) return {'wn': wn, 'intensity': inten, 'meta': header, 'note':'jcamp'} # CSV fallback try: df = tryreadcsvbytes(b) if df.shape[1] >= 2: wn = pd.tonumeric(df.iloc[:,0], errors='coerce').fillna(0).values.astype(np.float32) inten = pd.tonumeric(df.iloc[:,1], errors='coerce').fillna(0).values.astype(np.float32) return {'wn': wn, 'intensity': inten, 'meta': {'parser':'csv'}, 'note':'csv'} except Exception as e: logger.exception("parsejcamp_precise failed: %s", e) return {'wn': np.array([], dtype=np.float32), 'intensity': np.array([], dtype=np.float32), 'meta': {'parser':'failed'}, 'note':'failed'}
def parsesemtiffprecise(b: Optional[bytes]) -> Tuple[np.ndarray, Dict[str,Any]]: meta = {} if isplaceholder(b): # return synthetic blank patch arr = np.zeros((DEFAULTCONFIG['image']['patchsize'], DEFAULTCONFIG['image']['patchsize']), dtype=np.uint8) return arr, {'parser':'placeholder'} if tifffile is not None: try: with io.BytesIO(b) as bio: tf = tifffile.TiffFile(bio) page = tf.pages[0] arr = page.asarray() tags = {} try: for k,v in page.tags.items(): tags[k] = str(v.value) meta['tifftags'] = tags desc = tags.get('ImageDescription','') if isinstance(desc, str) and 'nm' in desc: import re m = re.search(r'([\d\.]+)\s*nm', desc) if m: meta['pixelsizenm'] = float(m.group(1)) except Exception: pass return arr, meta except Exception as e: logger.warning("tifffile parse failed: %s. Trying PIL fallback.", e) # PIL fallback try: if Image is None: raise RuntimeError("PIL not available") with io.BytesIO(b) as bio: img = Image.open(bio) arr = np.array(img) return arr, meta except UnidentifiedImageError as e: logger.warning("PIL cannot identify image: %s", e) except Exception as e: logger.exception("parsesemtiffprecise failed: %s", e) # final fallback: synthetic blank arr = np.zeros((DEFAULTCONFIG['image']['patchsize'], DEFAULTCONFIG['image']['patchsize']), dtype=np.uint8) return arr, {'parser':'failed_fallback'}
def parseafmprecise(b: Optional[bytes]) -> Dict[str,Any]: if isplaceholder(b): return {'distance': np.array([], dtype=np.float32), 'force': np.array([], dtype=np.float32), 'meta': {'parser':'placeholder'}} s = b.decode('utf-8', errors='ignore') try: df = tryreadcsvbytes(b) # heuristics: find numeric columns numericcols = [] for c in df.columns: try: pd.tonumeric(df[c]) numericcols.append(c) except Exception: pass if len(numericcols) >= 2: distcol = numericcols[0]; forcecol = numericcols[1] return {'distance': pd.tonumeric(df[distcol], errors='coerce').fillna(0).values.astype(np.float32), 'force': pd.tonumeric(df[forcecol], errors='coerce').fillna(0).values.astype(np.float32), 'meta': {'parser':'csv'}} # if only one numeric column, assume it's force and synthesize distance if len(numericcols) == 1: force = pd.tonumeric(df[numericcols[0]], errors='coerce').fillna(0).values.astype(np.float32) distance = np.linspace(0.0, 1.0, len(force)).astype(np.float32) return {'distance': distance, 'force': force, 'meta': {'parser':'csvsinglecolsynthdistance'}} except Exception as e: logger.exception("parseafmprecise failed: %s", e) return {'distance': np.array([], dtype=np.float32), 'force': np.array([], dtype=np.float32), 'meta': {'parser':'failed'}}
def parsecs14cprecise(b: Optional[bytes]) -> Dict[str,Any]: if isplaceholder(b): return {'records': [], 'meta': {'parser':'placeholder'}} try: df = tryreadcsvbytes(b) cols = [c.lower() for c in df.columns] required = ['sampleid','compoundid','agebp','agesigma'] if all(r in cols for r in required): df.columns = cols return {'records': df.todict(orient='records'), 'meta': {'parser':'csv'}} # try to map common variants mapping = {} for r in required: for c in cols: if r in c: mapping[r] = c break if len(mapping) == len(required): df2 = df.rename(columns={v:k for k,v in mapping.items()}) return {'records': df2.todict(orient='records'), 'meta': {'parser':'csvmapped'}} except Exception as e: logger.exception("parsecs14c_precise failed: %s", e) return {'records': [], 'meta': {'parser':'failed'}}
----------------------------
Signal processing & feature helpers (unchanged but robust)
----------------------------
def detectpeaks(mz: np.ndarray, inten: np.ndarray, prominence: float, width: Tuple[float,float]) -> np.ndarray: if len(mz) < 3: return np.array([], dtype=int) try: intens = signal.savgolfilter(inten, windowlength=11 if len(inten)>=11 else max(3,len(inten)//2*2+1), polyorder=2) peaksidx, props = signal.findpeaks(intens, prominence=prominence, width=width) return peaksidx except Exception: return np.array([], dtype=int)
def fitgaussianpeak(x: np.ndarray, y: np.ndarray) -> Tuple[float,float,float,float]: if len(x) < 3: return (float(np.mean(x)) if len(x)>0 else 0.0, 0.0, 0.0, float(np.max(y)) if len(y)>0 else 0.0) def gauss(x, A, mu, sigma): return A np.exp(-0.5((x-mu)/sigma)*2) p0 = [float(np.max(y)), float(x[np.argmax(y)]), max(1e-6, float((x.max()-x.min())/6.0))] try: popt, _ = optimize.curve_fit(gauss, x, y, p0=p0, maxfev=5000) A, mu, sigma = popt area = A sigma math.sqrt(2math.pi) height = float(A) return float(mu), float(sigma), float(area), float(height) except Exception: return (float(np.mean(x)) if len(x)>0 else 0.0, 0.0, float(np.trapz(y,x)) if len(x)>0 else 0.0, float(np.max(y)) if len(y)>0 else 0.0)
def bin_generic(x: np.ndarray, y: np.ndarray, bins: int = 256, xmin: float = None, xmax: float = None) -> np.ndarray: if len(x) == 0: return np.zeros(bins, dtype=np.float32) if xmin is None: xmin = float(np.min(x)) if xmax is None: xmax = float(np.max(x)) if xmin == xmax: return np.zeros(bins, dtype=np.float32) edges = np.linspace(xmin, xmax, bins+1) binned = np.zeros(bins, dtype=np.float32) idx = np.digitize(x, edges) - 1 for i in range(len(x)): j = min(max(int(idx[i]), 0), bins-1) binned[j] += float(y[i]) s = binned.sum() if s > 0: binned = binned / s return binned
def asls_baseline(intensity: np.ndarray, lam: float = 1e5, p: float = 0.01, niter: int = 10) -> np.ndarray: L = len(intensity) if L < 3: return intensity D = np.diff(np.eye(L), 2) w = np.ones(L) for i in range(niter): W = np.diag(w) Z = W + lam D.T.dot(D) z = np.linalg.solve(Z, w intensity) w = p (intensity > z) + (1-p) (intensity < z) corrected = intensity - z corrected[corrected < 0] = 0.0 return corrected
----------------------------
Image & AFM quantification (robust)
----------------------------
def pixelsizefrommeta(meta: Dict[str,Any]) -> Optional[float]: if not meta: return None if 'pixelsizenm' in meta: try: return float(meta['pixelsizenm']) except Exception: pass tags = meta.get('tifftags', {}) for k,v in tags.items(): if isinstance(v, str) and 'nm' in v: import re m = re.search(r'([\d\.]+)\s*nm', v) if m: try: return float(m.group(1)) except: pass return None
def penetrationproxy(img: np.ndarray) -> float: if img is None or img.size == 0: return 0.0 if img.ndim == 3: gray = np.mean(img.astype(np.float32), axis=2) else: gray = img.astype(np.float32) gray = gray - gray.min() if gray.max() > 0: gray = gray / gray.max() h, w = gray.shape[:2] cy, cx = h//2, w//2 ys, xs = np.indices((h,w)) r = np.sqrt((xs-cx)**2 + (ys-cy)**2).astype(np.int32) maxr = min(cx, cy) profile = [float(np.mean(gray[r==ri])) if np.any(r==ri) else 0.0 for ri in range(maxr)] if len(profile) < 5: return 0.0 prof = np.array(profile) x = np.arange(len(prof)) slope, intercept, rval, pval, stderr = stats.linregress(x, prof) return float(-slope)
def estimatepenetrationnm(img: np.ndarray, meta: Dict[str,Any], calibrationcurve: Optional[Tuple[float,float,float]] = None) -> Tuple[float,float]: proxy = penetrationproxy(img) pixelnm = pixelsizefrommeta(meta) if calibrationcurve: slope, intercept, sigma = calibrationcurve penetrationnm = slope * proxy + intercept uncertainty = abs(sigma) + 0.05 * penetrationnm return float(max(0.0, penetrationnm)), float(uncertainty) elif pixelnm: h, w = img.shape[:2] meandimnm = ((h + w) / 2.0) pixel_nm penetration_nm = proxy meandimnm uncertainty = 0.2 penetration_nm return float(max(0.0, penetration_nm)), float(uncertainty) else: return float(proxy), float(proxy 0.5 + 1.0)
def imagefeatures(img: np.ndarray, meta: Dict[str,Any], cfg: Dict[str,Any]) -> np.ndarray: if img is None or img.size == 0: return np.zeros(cfg['histbins'] + 8, dtype=np.float32) if img.ndim == 3: arr = np.mean(img.astype(np.float32), axis=2) else: arr = img.astype(np.float32) arr = arr - arr.min() if arr.max() > 0: arr = arr / arr.max() target = (cfg['patchsize'], cfg['patchsize']) try: pil = Image.fromarray((arr*255).astype(np.uint8)) pil = pil.resize(target, Image.BILINEAR) arr2 = np.array(pil).astype(np.float32)/255.0 except Exception: arr2 = gaussianfilter(arr, sigma=1.0) arr2 = arr2[:target[0], :target[1]] if arr2.shape != target: arr2 = np.resize(arr2, target) mean = float(arr2.mean()) std = float(arr2.std()) gradx = float(np.mean(np.abs(np.diff(arr2, axis=1)))) grady = float(np.mean(np.abs(np.diff(arr2, axis=0)))) hist, = np.histogram(arr2, bins=cfg['histbins'], range=(0,1)) hist = hist.astype(np.float32) / (hist.sum()+1e-9) tailratio = float(np.sum(arr2 < 0.1) / (arr2.size + 1e-9)) penetrationnm, penetrationunc = estimatepenetrationnm(img, meta, calibrationcurve=None) return np.concatenate([[mean, std, gradx, grady, tailratio, penetrationnm, penetrationunc], hist])
def calibrateafm(forceraw: np.ndarray, probekNperm: Optional[float], sensitivitynmperV: Optional[float]) -> np.ndarray: if probekNperm and sensitivitynmperV: deflectionnm = forceraw sensitivity_nm_per_V force_N = probe_k_N_per_m (deflectionnm * 1e-9) return forceN * 1e9 # nN else: return force_raw
def afmfeatures(distance: np.ndarray, forceraw: np.ndarray, probek: Optional[float], sensitivity: Optional[float]) -> Tuple[np.ndarray, Dict[str,Any]]: if distance is None or forceraw is None or len(forceraw)==0: return np.zeros(5, dtype=np.float32), {'probek': probek, 'sensitivity': sensitivity} forcenN = calibrateafm(forceraw, probek, sensitivity) meta = {'probek': probek, 'sensitivity': sensitivity} if len(forcenN) < 3: return np.zeros(5, dtype=np.float32), meta maxf = float(np.max(forcenN)) minf = float(np.min(forcenN)) adhesion = float(-minf) if minf < 0 else 0.0 slope = float((forcenN[-1] - forcenN[0]) / (distance[-1] - distance[0] + 1e-9)) auc = float(np.trapz(forcenN, distance)) return np.array([maxf, min_f, adhesion, slope, auc], dtype=np.float32), meta
----------------------------
Dataset & feature assembly (robust with synthetic minimal features)
----------------------------
class MultimodalSample: def _init(self, sampleid: str, meta: Dict[str,Any], gcmsbytes: Optional[bytes]=None, spectrumbytes: Optional[bytes]=None, imagebytes: Optional[bytes]=None, afmbytes: Optional[bytes]=None, cs14cbytes: Optional[bytes]=None, label: Optional[int]=None): self.sampleid = sampleid self.meta = meta or {} self.gcmsbytes = gcmsbytes self.spectrumbytes = spectrumbytes self.imagebytes = imagebytes self.afmbytes = afmbytes self.cs14cbytes = cs14c_bytes self.label = label
class MultimodalDataset: def _init_(self, samples: List[MultimodalSample], cfg: Dict[str,Any]): self.samples = samples self.cfg = cfg
def _len_(self): return len(self.samples)
def _getitem(self, idx): s = self.samples[idx] feats = {} # GC-MS try: parsed = parsemzmlprecise(s.gcmsbytes) peaks = parsed.get('peaks', pd.DataFrame(columns=['mz','intensity','rt'])) feats['gcms'] = extractgcmsfeatures(peaks, self.cfg['gcms']) except Exception as e: logger.exception("GCMS error %s: %s", s.sampleid, e) auditerror(self.cfg.get('auditdir','auditlogs'), s.sampleid, ERRORCODES['PARSEGCMSFAIL'], str(e), s.gcmsbytes) feats['gcms'] = synthesizegcmsminimal(self.cfg['gcms']) # Spectrum try: sp = parsejcampprecise(s.spectrumbytes) feats['spectrum'] = spectrumfeatures(sp.get('wn', np.array([])), sp.get('intensity', np.array([])), self.cfg['spectrum']['markerpeaks'], self.cfg['spectrum']['binnedbins']) except Exception as e: logger.exception("Spectrum error %s: %s", s.sampleid, e) auditerror(self.cfg.get('auditdir','auditlogs'), s.sampleid, ERRORCODES['PARSESPECTRUMFAIL'], str(e), s.spectrumbytes) feats['spectrum'] = synthesizespectrumminimal(self.cfg['spectrum']) # Image try: img, imgmeta = parsesemtiffprecise(s.imagebytes) feats['image'] = imagefeatures(img, imgmeta, self.cfg['image']) except Exception as e: logger.exception("Image error %s: %s", s.sampleid, e) auditerror(self.cfg.get('auditdir','auditlogs'), s.sampleid, ERRORCODES['PARSEIMAGEFAIL'], str(e), s.imagebytes) feats['image'] = synthesizeimageminimal(self.cfg['image']) # AFM try: afm = parseafmprecise(s.afmbytes) afmfeat, afmmeta = afmfeatures(afm.get('distance', np.array([])), afm.get('force', np.array([])), s.meta.get('afmprobekNperm'), s.meta.get('afmsensitivitynmperV')) feats['afm'] = afmfeat except Exception as e: logger.exception("AFM error %s: %s", s.sampleid, e) auditerror(self.cfg.get('auditdir','auditlogs'), s.sampleid, ERRORCODES['PARSEAFMFAIL'], str(e), s.afmbytes) feats['afm'] = synthesizeafmminimal() # meta features env = s.meta.get('environment', {}) temp = float(env.get('temperaturec', 20.0)) rh = float(env.get('relativehumidity', 40.0)) material = s.meta.get('objectmaterial', 'unknown') materialcode = 0.0 if isinstance(material, str): if material.lower() in ['ceramic','pottery','porcelain']: materialcode = 1.0 elif material.lower() in ['bronze','metal']: materialcode = 2.0 metafeat = np.array([temp, rh, materialcode], dtype=np.float32) x = np.concatenate([feats['gcms'], feats['spectrum'], feats['image'], feats['afm'], metafeat]) y = np.array([s.label], dtype=np.int64) if s.label is not None else np.array([-1], dtype=np.int64) return {'x': x, 'y': y, 'sampleid': s.sampleid, 'meta': s.meta}
----------------------------
Synthetic minimal feature generators (when modality missing)
----------------------------
def synthesizegcmsminimal(cfg: Dict[str,Any]) -> np.ndarray: dim = cfg['binnedbins'] + len(cfg.get('markerwindows',[])) + cfg.get('topk',256) # small random noise stable seed return np.zeros(dim, dtype=np.float32)
def synthesizespectrumminimal(cfg: Dict[str,Any]) -> np.ndarray: dim = cfg['binnedbins'] + len(cfg.get('markerpeaks',[])) + 1 return np.zeros(dim, dtype=np.float32)
def synthesizeimageminimal(cfg: Dict[str,Any]) -> np.ndarray: dim = cfg['hist_bins'] + 8 return np.zeros(dim, dtype=np.float32)
def synthesizeafmminimal() -> np.ndarray: return np.zeros(5, dtype=np.float32)
----------------------------
Spectrum features helper
----------------------------
def spectrumfeatures(wn: np.ndarray, inten: np.ndarray, markerpeaks: List[float], bins: int) -> np.ndarray: if len(wn) == 0 or len(inten) == 0: return synthesizespectrumminimal({'binnedbins': bins, 'markerpeaks': markerpeaks}) intencorr = aslsbaseline(inten) binned = bingeneric(wn, intencorr, bins=bins, xmin=wn.min(), xmax=wn.max()) markers = [] for pk in markerpeaks: idx = np.argmin(np.abs(wn - pk)) markers.append(float(intencorr[idx]) / (intencorr.sum()+1e-9)) peakcount = float(np.sum(intencorr > (np.mean(intencorr) + 2*np.std(intencorr)))) return np.concatenate([binned, np.array(markers, dtype=np.float32), np.array([peak_count], dtype=np.float32)])
----------------------------
GC-MS feature extraction (precise)
----------------------------
def extractgcmsfeatures(peaksdf: pd.DataFrame, cfg: Dict[str,Any]) -> np.ndarray: if peaksdf is None or len(peaksdf)==0: return synthesizegcmsminimal(cfg) try: mz = peaksdf['mz'].values.astype(np.float64) inten = peaksdf['intensity'].values.astype(np.float64) except Exception: return synthesizegcmsminimal(cfg) binned = bingeneric(mz, inten, bins=cfg['binnedbins'], xmin=mz.min(), xmax=mz.max()) markervals = [] for w in cfg.get('markerwindows', []): lo, hi = w mask = (mz >= lo) & (mz <= hi) val = float(inten[mask].sum()) / (inten.sum()+1e-9) markervals.append(val) K = cfg.get('topk', 256) try: order = np.argsort(inten)[-K:] if len(inten) >= K else np.argsort(inten) topmz = np.sort(mz[order]) topvec = np.zeros(K, dtype=np.float32) for i, val in enumerate(topmz[:K]): topvec[i] = float((val % 1000) / 1000.0) except Exception: topvec = np.zeros(K, dtype=np.float32) feat = np.concatenate([binned.astype(np.float32), np.array(markervals, dtype=np.float32), top_vec]) return feat
----------------------------
Submodels (XGBoost)
----------------------------
class XGBModel: def _init(self, params: Dict[str,Any] = None): self.params = params or {'objective':'binary:logistic','evalmetric':'auc','uselabelencoder':False} self.model = None
def train(self, X: np.ndarray, y: np.ndarray, numround: int = 200): dtrain = xgb.DMatrix(X, label=y) self.model = xgb.train(self.params, dtrain, numround)
def predict(self, X: np.ndarray) -> np.ndarray: if self.model is None: raise RuntimeError("Model not trained") d = xgb.DMatrix(X) return self.model.predict(d)
def save(self, path: str): ensuredir(os.path.dirname(path) or '.') self.model.savemodel(path)
def load(self, path: str): self.model = xgb.Booster() self.model.load_model(path)
----------------------------
Fusion & Bayesian age posterior (MCMC template + conjugate fallback)
----------------------------
def fusion_posterior(ps: Dict[str,float], weights: Dict[str,float], prior: float = 0.5) -> float: eps = 1e-9 def logit(p): return math.log((p+eps)/(1-p+eps)) logits = np.array([logit(ps.get(k,0.5)) for k in ['gcms','image','spectrum','afm']]) w = np.array([weights.get(k,0.25) for k in ['gcms','image','spectrum','afm']]) combined = np.dot(w, logits) + logit(prior) posterior = 1.0 / (1.0 + math.exp(-combined)) return posterior
def bayesianagemcmc(priormu: float, priorsigma: float, evidence: List[Tuple[float,float]], cs14crecords: Optional[List[Dict[str,Any]]] = None, draws: int = 1000, tune: int = 500): if pm is None: mu, sigma, ci = bayesianageconjugate((priormu, priorsigma), evidence, cs14crecords) return {'mu': mu, 'sigma': sigma, 'ci95': ci, 'method': 'conjugatefallback'} with pm.Model() as model: age = pm.Normal('age', mu=priormu, sigma=priorsigma) for i, (mui, sigmai) in enumerate(evidence): pm.Normal(f'evidence{i}', mu=age, sigma=max(1.0, sigmai), observed=mui) if cs14crecords: for j, rec in enumerate(cs14crecords): muc = float(rec.get('agebp')) sigmac = float(rec.get('agesigma', 20.0)) pm.Normal(f'cs14c{j}', mu=age, sigma=max(1.0, sigmac), observed=muc) trace = pm.sample(draws=draws, tune=tune, chains=2, cores=1, returninferencedata=True, progressbar=False) summary = az.summary(trace, hdiprob=0.95) mupost = float(summary.loc['age','mean']) sdpost = float(summary.loc['age','sd']) hdi = az.hdi(trace, varnames=['age'], hdiprob=0.95) ci95 = (float(hdi['age'][0]), float(hdi['age'][1])) return {'mu': mupost, 'sigma': sd_post, 'ci95': ci95, 'method': 'mcmc'}
def bayesianageconjugate(prior: Tuple[float,float], evidence: List[Tuple[float,float]], cs14crecords: Optional[List[Dict[str,Any]]] = None): mu0, s0 = prior var0 = s0**2 mun = mu0 varn = var0 for mui, si in evidence: vari = max(1.0, si**2) varn = 1.0 / (1.0/varn + 1.0/vari) mun = varn (mu_n/var_n + mu_i/var_i) if cs14c_records: for rec in cs14c_records: mu_c = float(rec.get('age_bp')) sigma_c = float(rec.get('age_sigma', 20.0)) var_c = sigma_c2 var_n = 1.0 / (1.0/var_n + 1.0/var_c) mu_n = var_n (mun/varn + muc/varc) sigman = math.sqrt(varn) ci95 = (mun - 1.96*sigman, mun + 1.96*sigman) return mun, sigman, ci95
----------------------------
Calibration & evaluation
----------------------------
def calibrateprobs(probs: np.ndarray, labels: np.ndarray) -> Dict[str,Any]: iso = IsotonicRegression(outof_bounds='clip') try: iso.fit(probs, labels) return {'type':'isotonic','model':iso} except Exception: return {'type':'identity'}
def apply_calib(calib: Dict[str,Any], probs: np.ndarray) -> np.ndarray: if calib is None or calib.get('type') == 'identity': return probs if calib.get('type') == 'isotonic': iso = calib['model'] return iso.transform(probs) return probs
----------------------------
Training pipeline (robust)
----------------------------
def trainpipeline(samples: List[MultimodalSample], cfg: Dict[str,Any], outdir: str): ensuredir(outdir) ds = MultimodalDataset(samples, cfg) X = [] Y = [] for i in range(len(ds)): item = ds[i] X.append(item['x']) Y.append(int(item['y'][0])) X = np.stack(X).astype(np.float32) Y = np.array(Y).astype(np.int64) Xtrain, Xval, ytrain, yval = traintestsplit(X, Y, testsize=0.2, stratify=Y, randomstate=cfg['training']['randomseed']) # modality dims gcmsdim = cfg['gcms']['binnedbins'] + len(cfg['gcms'].get('markerwindows',[])) + cfg['gcms'].get('topk',256) specdim = cfg['spectrum']['binnedbins'] + len(cfg['spectrum'].get('markerpeaks',[])) + 1 imagedim = cfg['image']['histbins'] + 8 afmdim = 5 metadim = 3 idx = 0 gcmsidx = (idx, idx+gcmsdim); idx += gcmsdim specidx = (idx, idx+specdim); idx += specdim imageidx = (idx, idx+imagedim); idx += imagedim afmidx = (idx, idx+afmdim); idx += afmdim metaidx = (idx, idx+metadim); idx += metadim
# train submodels modelg = XGBModel(); modelg.train(Xtrain[:, gcmsidx[0]:gcmsidx[1]], ytrain, numround=cfg['training']['xgbrounds']) models = XGBModel(); models.train(Xtrain[:, specidx[0]:specidx[1]], ytrain, numround=cfg['training']['xgbrounds']) modeli = XGBModel(); modeli.train(Xtrain[:, imageidx[0]:imageidx[1]], ytrain, numround=cfg['training']['xgbrounds']) modela = XGBModel(); modela.train(Xtrain[:, afmidx[0]:afmidx[1]], ytrain, numround=max(50, cfg['training']['xgbrounds']//4))
# validation preds pg = modelg.predict(Xval[:, gcmsidx[0]:gcmsidx[1]]) ps = models.predict(Xval[:, specidx[0]:specidx[1]]) pi = modeli.predict(Xval[:, imageidx[0]:imageidx[1]]) pa = modela.predict(Xval[:, afmidx[0]:afmidx[1]])
# weights from AUCs def aucsafe(ytrue, yscore): try: return rocaucscore(ytrue, yscore) except: return 0.5 aucg = aucsafe(yval, pg) aucs = aucsafe(yval, ps) auci = aucsafe(yval, pi) auca = aucsafe(yval, pa) raw = np.array([max(0.01, aucg), max(0.01, auci), max(0.01, aucs), max(0.01, auca)]) weights = (raw / raw.sum()).tolist() weightsdict = {'gcms': weights[0], 'image': weights[1], 'spectrum': weights[2], 'afm': weights[3]} logger.info("Modality weights: %s", weights_dict)
fused = [] for i in range(len(pg)): psd = {'gcms': float(pg[i]), 'image': float(pi[i]), 'spectrum': float(ps[i]), 'afm': float(pa[i])} fused.append(fusionposterior(psd, weightsdict, prior=0.5)) fused = np.array(fused) calib = calibrateprobs(fused, yval) fusedcal = applycalib(calib, fused) aucfused = aucsafe(yval, fusedcal) preds = (fusedcal >= 0.5).astype(int) precision, recall, f1, = precisionrecallfscoresupport(yval, preds, average='binary', zerodivision=0) logger.info("Fused metrics: AUC=%.4f Precision=%.4f Recall=%.4f F1=%.4f", aucfused, precision, recall, f1)
# save ensuredir(outdir) modelg.save(os.path.join(outdir, "gcms.xgb")) models.save(os.path.join(outdir, "spectrum.xgb")) modeli.save(os.path.join(outdir, "image.xgb")) modela.save(os.path.join(outdir, "afm.xgb")) meta = {'weights': weightsdict, 'calibration': {'type': calib.get('type', 'identity')}, 'cfg': cfg, 'modelversion': cfg.get('modelversion')} with open(os.path.join(outdir, "fusionmeta.json"), "w", encoding='utf-8') as f: json.dump(meta, f, ensureascii=False, indent=2) import pickle with open(os.path.join(outdir, "calibration.pkl"), "wb") as f: pickle.dump(calib, f) logger.info("Training complete. Models saved to %s", outdir) return {'auc': aucfused, 'precision': precision, 'recall': recall, 'f1': f1, 'weights': weightsdict}
----------------------------
Inference engine (robust)
----------------------------
class InferenceEngine: def _init(self, modeldir: str, cfg: Dict[str,Any], secret: str): self.modeldir = modeldir self.cfg = cfg self.secret = secret self._load()
def load(self): self.modelg = XGBModel(); self.models = XGBModel(); self.modeli = XGBModel(); self.modela = XGBModel() try: self.modelg.load(os.path.join(self.modeldir, "gcms.xgb")) self.models.load(os.path.join(self.modeldir, "spectrum.xgb")) self.modeli.load(os.path.join(self.modeldir, "image.xgb")) self.modela.load(os.path.join(self.modeldir, "afm.xgb")) except Exception: logger.warning("Model files not found; inference will use defaults") try: with open(os.path.join(self.modeldir, "fusionmeta.json"), "r", encoding='utf-8') as f: meta = json.load(f); self.weights = meta.get('weights', {'gcms':0.4,'image':0.3,'spectrum':0.2,'afm':0.1}) except Exception: self.weights = {'gcms':0.4,'image':0.3,'spectrum':0.2,'afm':0.1} try: import pickle with open(os.path.join(self.modeldir, "calibration.pkl"), "rb") as f: self.calib = pickle.load(f) except Exception: self.calib = {'type':'identity'} logger.info("Engine loaded from %s", self.model_dir)
def infer(self, sample: MultimodalSample, returnexplain: bool = False) -> Dict[str,Any]: try: ds = MultimodalDataset([sample], self.cfg) item = ds[0] x = item['x'].reshape(1,-1) # indices gcmsdim = self.cfg['gcms']['binnedbins'] + len(self.cfg['gcms'].get('markerwindows',[])) + self.cfg['gcms'].get('topk',256) specdim = self.cfg['spectrum']['binnedbins'] + len(self.cfg['spectrum'].get('markerpeaks',[])) + 1 imagedim = self.cfg['image']['histbins'] + 8 afmdim = 5 idx = 0 gcmsidx = (idx, idx+gcmsdim); idx += gcmsdim specidx = (idx, idx+specdim); idx += specdim imageidx = (idx, idx+imagedim); idx += imagedim afmidx = (idx, idx+afmdim); idx += afmdim
try: pg = float(self.modelg.predict(x[:, gcmsidx[0]:gcmsidx[1]])[0]) except Exception: pg = 0.5 try: ps = float(self.models.predict(x[:, specidx[0]:specidx[1]])[0]) except Exception: ps = 0.5 try: pi = float(self.modeli.predict(x[:, imageidx[0]:imageidx[1]])[0]) except Exception: pi = 0.5 try: pa = float(self.modela.predict(x[:, afmidx[0]:afmidx[1]])[0]) except Exception: pa = 0.5
psd = {'gcms': pg, 'image': pi, 'spectrum': ps, 'afm': pa} posterior = fusionposterior(psd, self.weights, prior=0.5) posteriorcal = applycalib(self.calib, np.array([posterior]))[0] verdict = "ancientoil" if posteriorcal >= self.cfg['fusion']['confidencehigh'] else ("modernoil" if posteriorcal < self.cfg['fusion']['confidence_low'] else "inconclusive")
cs14crecs = None if sample.cs14cbytes: try: cs14c = parsecs14cprecise(sample.cs14cbytes) cs14crecs = cs14c.get('records', []) except Exception: cs14crecs = None evidenceage = [] # placeholder for submodel age proxies agepost = bayesianagemcmc(200.0, 200.0, evidenceage, cs14crecs, draws=self.cfg['bayes']['mcmcdraws'], tune=self.cfg['bayes']['mcmctune']) if pm is not None else {'mu': None, 'sigma': None, 'ci95': (None,None), 'method':'conjugatefallback'}
explain = {} if shap is not None: try: explainer = shap.TreeExplainer(self.modelg.model) shapvals = explainer.shapvalues(x[:, gcmsidx[0]:gcmsidx[1]]) explain['gcmsshapmeanabs'] = float(np.mean(np.abs(shapvals))) except Exception: explain['gcmsshapmeanabs'] = None else: def logit(p): return math.log((p+1e-9)/(1-p+1e-9)) contrib = {k: float(self.weights.get(k,0.25) * abs(logit(psd.get(k,0.5)))) for k in psd} explain['modalitycontribproxy'] = contrib
filesbytes = [sample.gcmsbytes or b'', sample.spectrumbytes or b'', sample.imagebytes or b'', sample.afmbytes or b'', sample.cs14cbytes or b''] prov = computeprovenance(sample.meta, filesbytes, self.secret) report = { "sampleid": sample.sampleid, "verdict": verdict, "confidence": float(posteriorcal), "ageposterior": agepost, "evidenceprobs": psd, "explainability": explain, "modalityweights": self.weights, "modelversion": self.cfg.get('modelversion'), "provenance": prov, "timestamp": nowiso(), "notes": "Algorithmic result. For high-stakes use, require CS-14C compound-specific dating and expert review." } auditdir = self.cfg.get('auditdir', DEFAULTCONFIG['auditdir']) ensuredir(auditdir) fname = os.path.join(auditdir, f"{sample.sampleid}{int(time.time())}.json") with open(fname, "w", encoding='utf-8') as f: json.dump(report, f, ensureascii=False, indent=2) return report except Exception as e: logger.exception("Inference failed for %s: %s", sample.sampleid, e) auditerror(self.cfg.get('auditdir','auditlogs'), sample.sampleid, ERRORCODES['INFERENCEFAIL'], str(e)) return {"sampleid": sample.sample_id, "verdict": "error", "confidence": 0.0, "error": str(e)}
----------------------------
Blind-test harness
----------------------------
def runblindtest(engine: InferenceEngine, samples: List[MultimodalSample]) -> Dict[str,Any]: ytrue = [] probs = [] reports = [] for s in samples: r = engine.infer(s) prob = r.get('confidence', 0.0) label = s.label if s.label is not None else int(s.meta.get('label', 0)) ytrue.append(label) probs.append(prob) reports.append(r) ytrue = np.array(ytrue) probs = np.array(probs) auc = rocaucscore(ytrue, probs) if len(np.unique(ytrue))>1 else 0.5 preds = (probs >= 0.5).astype(int) precision, recall, f1, = precisionrecallfscoresupport(ytrue, preds, average='binary', zerodivision=0) return {"n": len(samples), "auc": auc, "precision": precision, "recall": recall, "f1": f1, "reports": reports}
----------------------------
AutoOrchestrator (robust)
----------------------------
class AutoOrchestrator: def _init(self, cfg: Dict[str,Any]): self.cfg = cfg self.watchdir = cfg.get('auto',{}).get('watchdir', DEFAULTCONFIG['auto']['watchdir']) self.pollinterval = cfg.get('auto',{}).get('pollintervals', DEFAULTCONFIG['auto']['pollintervals']) self.maxworkers = cfg.get('auto',{}).get('maxworkers', DEFAULTCONFIG['auto']['maxworkers']) self.humanthreshold = cfg.get('auto',{}).get('humanreviewconfidencethreshold', DEFAULTCONFIG['auto']['humanreviewconfidencethreshold']) self.autoaccept = cfg.get('auto',{}).get('autoacceptconfidence', DEFAULTCONFIG['auto']['autoacceptconfidence']) self.queue = Queue(maxsize=cfg.get('auto',{}).get('maxqueuesize', DEFAULTCONFIG['auto']['maxqueuesize'])) self.workers: List[Thread] = [] self.stopevent = Event() self.engine = None ensuredir(self.watchdir) ensuredir(self.cfg.get('auditdir', DEFAULTCONFIG['audit_dir']))
def detectinput(self, path: str) -> Dict[str,Any]: name = os.path.basename(path).lower() if name.endswith('.json'): try: with open(path, 'r', encoding='utf-8') as f: obj = json.load(f) if isinstance(obj, list) and all('sampleid' in r for r in obj): return {'task':'trainmanifest','payload':{'manifestpath':path}} if isinstance(obj, dict) and obj.get('type') == 'sample': return {'task':'samplejson','payload':{'samplejson':path}} except Exception: pass if name.endswith('.csv'): try: with open(path, 'r', encoding='utf-8') as f: header = f.readline().lower() if 'agebp' in header and 'compoundid' in header: return {'task':'cs14creport','payload':{'path':path}} if 'mz' in header and 'intensity' in header: return {'task':'peaktable','payload':{'path':path}} except Exception: pass if name.endswith(('.tif','.tiff','.png','.jpg','.jpeg')): return {'task':'image_file','payload':{'path':path}} return {'task':'unknown','payload':{'path':path}}
def routetask(self, detected: Dict[str,Any]): t = detected.get('task') if t == 'trainmanifest': self.enqueue({'type':'train','manifestpath': detected['payload']['manifestpath']}) elif t == 'samplejson': self.enqueue({'type':'inferfromjson','samplejson': detected['payload']['samplejson']}) elif t == 'cs14creport': self.enqueue({'type':'calibratecs14c','path': detected['payload']['path']}) elif t == 'peaktable': self.enqueue({'type':'inferfromfiles','files': {'gcms': detected['payload']['path']}}) elif t == 'imagefile': self.enqueue({'type':'inferfromfiles','files': {'image': detected['payload']['path']}}) else: logger.info("Unknown input type for %s; moving to quarantine", detected['payload'].get('path')) self.quarantinefile(detected['payload'].get('path'))
def enqueue(self, job: Dict[str,Any]): try: self.queue.put_nowait(job) logger.info("Enqueued job: %s", job.get('type')) except Exception as e: logger.exception("Queue full or error enqueuing job: %s", e)
def workerloop(self): while not self.stopevent.isset(): try: job = self.queue.get(timeout=1) except Empty: continue try: self.executejob(job) except Exception as e: logger.exception("Job execution failed: %s", e) finally: self.queue.task_done()
def executejob(self, job: Dict[str,Any]): jtype = job.get('type') logger.info("Executing job type: %s", jtype) if jtype == 'train': manifest = job.get('manifestpath') self.dotrain(manifest) elif jtype == 'inferfromjson': samplejson = job.get('samplejson') self.doinferfromjson(samplejson) elif jtype == 'inferfromfiles': files = job.get('files',{}) self.doinferfromfiles(files) elif jtype == 'calibratecs14c': path = job.get('path') self.docalibrate_cs14c(path) else: logger.warning("Unknown job type: %s", jtype)
def ensureengine(self): if self.engine is None: modeldir = self.cfg.get('modeldir','models') secret = os.environ.get(self.cfg.get('provenancesecretenv','PROVENANCESECRET'),'changethissecret') self.engine = InferenceEngine(modeldir, self.cfg, secret)
def dotrain(self, manifestpath: str): logger.info("Starting training from manifest: %s", manifestpath) with open(manifestpath, 'r', encoding='utf-8') as f: manifest = json.load(f) samples = [] for rec in manifest: sid = rec['sampleid'] meta = json.load(open(rec['metajsonpath'],'r',encoding='utf-8')) if rec.get('metajsonpath') else {} def readopt(p): return open(p,'rb').read() if p and os.path.exists(p) else None samples.append(MultimodalSample(sampleid=sid, meta=meta, gcmsbytes=readopt(rec.get('gcmspath')), spectrumbytes=readopt(rec.get('spectrumpath')), imagebytes=readopt(rec.get('imagepath')), afmbytes=readopt(rec.get('afmpath')), cs14cbytes=readopt(rec.get('cs14cpath')), label=rec.get('label'))) res = trainpipeline(samples, self.cfg, self.cfg.get('modeldir','models')) logger.info("Training finished: %s", res) self.writeaudit({'task':'train','manifest':manifestpath,'result':res,'timestamp':now_iso()})
def doinferfromjson(self, samplejson: str): with open(samplejson, 'r', encoding='utf-8') as f: rec = json.load(f) def readopt(p): return open(p,'rb').read() if p and os.path.exists(p) else None sample = MultimodalSample(sampleid=rec.get('sampleid', str(uuid.uuid4())), meta=rec.get('meta', {}), gcmsbytes=readopt(rec.get('gcmspath')), spectrumbytes=readopt(rec.get('spectrumpath')), imagebytes=readopt(rec.get('imagepath')), afmbytes=readopt(rec.get('afmpath')), cs14cbytes=readopt(rec.get('cs14cpath'))) self.doinfer_sample(sample)
def doinferfromfiles(self, files: Dict[str,str]): meta = {} sample = MultimodalSample(sampleid=str(uuid.uuid4()), meta=meta, gcmsbytes=open(files['gcms'],'rb').read() if files.get('gcms') and os.path.exists(files.get('gcms')) else None, spectrumbytes=open(files['spectrum'],'rb').read() if files.get('spectrum') and os.path.exists(files.get('spectrum')) else None, imagebytes=open(files['image'],'rb').read() if files.get('image') and os.path.exists(files.get('image')) else None, afmbytes=open(files['afm'],'rb').read() if files.get('afm') and os.path.exists(files.get('afm')) else None) self.doinfersample(sample)
def doinfersample(self, sample: MultimodalSample): self.ensureengine() report = self.engine.infer(sample) logger.info("Inference report: %s", json.dumps(report, ensureascii=False)) self.writeaudit({'task':'infer','sampleid':sample.sampleid,'report':report,'timestamp':nowiso()}) conf = float(report.get('confidence',0.0)) if conf >= self.autoaccept: logger.info("Auto-accepting result for %s (conf=%.3f)", sample.sampleid, conf) self.publishreport(report) elif conf < self.humanthreshold: logger.info("Low confidence -> trigger human review for %s (conf=%.3f)", sample.sampleid, conf) self.triggerhumanreview(report, reason="lowconfidence") else: logger.info("Intermediate confidence -> queue for human review: %s (conf=%.3f)", sample.sampleid, conf) self.triggerhumanreview(report, reason="intermediateconfidence")
def docalibratecs14c(self, path: str): logger.info("Processing CS-14C calibration report: %s", path) try: recs = parsecs14cprecise(open(path,'rb').read()) self.writeaudit({'task':'calibratecs14c','path':path,'recordscount':len(recs.get('records',[])),'timestamp':nowiso()}) logger.info("CS-14C calibration records processed: %d", len(recs.get('records',[]))) except Exception as e: logger.exception("Failed to process CS-14C: %s", e) self.writeaudit({'task':'calibratecs14c','path':path,'error':str(e),'timestamp':nowiso()})
def publishreport(self, report: Dict[str,Any]): outdir = os.path.join(self.cfg.get('auditdir','auditlogs'),'published') ensuredir(outdir) fname = os.path.join(outdir, f"{report['sampleid']}{int(time.time())}.json") with open(fname, 'w', encoding='utf-8') as f: json.dump(report, f, ensure_ascii=False, indent=2) logger.info("Published report to %s", fname)
def triggerhumanreview(self, report: Dict[str,Any], reason: str): reviewdir = os.path.join(self.cfg.get('auditdir','auditlogs'),'humanreview') ensuredir(reviewdir) fname = os.path.join(reviewdir, f"{report['sampleid']}{int(time.time())}.json") reviewpkg = {'report': report, 'reason': reason, 'timestamp': nowiso()} with open(fname, 'w', encoding='utf-8') as f: json.dump(reviewpkg, f, ensureascii=False, indent=2) logger.info("Human review package created: %s", fname)
def writeaudit(self, audit: Dict[str,Any]): auditdir = self.cfg.get('auditdir','auditlogs') ensuredir(auditdir) fname = os.path.join(auditdir, f"audit{int(time.time())}{uuid.uuid4().hex[:8]}.json") with open(fname, 'w', encoding='utf-8') as f: json.dump(audit, f, ensure_ascii=False, indent=2) logger.debug("Audit written: %s", fname)
def quarantinefile(self, path: str): qdir = os.path.join(self.watchdir, 'quarantine') ensure_dir(qdir) try: base = os.path.basename(path) dest = os.path.join(qdir, base) os.rename(path, dest) logger.info("Moved unknown file to quarantine: %s", dest) except Exception: logger.exception("Failed to quarantine file: %s", path)
def startworkers(self): for i in range(self.maxworkers): t = Thread(target=self.worker_loop, daemon=True) t.start() self.workers.append(t) logger.info("Started %d worker threads", len(self.workers))
def stop(self): self.stop_event.set() logger.info("Stop signal set; waiting for workers to finish") for t in self.workers: t.join(timeout=2)
def watchloop(self): logger.info("Starting watch loop on %s", self.watchdir) seen = set() while not self.stopevent.isset(): try: files = [os.path.join(self.watchdir, f) for f in os.listdir(self.watchdir) if not f.startswith('.')] except Exception: files = [] for p in files: if p in seen: continue if os.path.isdir(p): seen.add(p) continue try: detected = self.detectinput(p) self.routetask(detected) seen.add(p) except Exception as e: logger.exception("Error detecting/routing file %s: %s", p, e) self.quarantinefile(p) time.sleep(self.pollinterval)
def rundaemon(self): logger.info("AutoOrchestrator starting daemon") self.startworkers() watchthread = Thread(target=self.watchloop, daemon=True) watchthread.start() try: while not self.stopevent.isset(): time.sleep(1) except KeyboardInterrupt: logger.info("KeyboardInterrupt received; stopping orchestrator") self.stop() watchthread.join(timeout=1) logger.info("Orchestrator stopped")
----------------------------
FastAPI UI prototype (optional)
----------------------------
if FastAPI is not None: @asynccontextmanager async def lifespan(app: FastAPI): modeldir = os.environ.get("MODELDIR", DEFAULTCONFIG['modeldir']) cfgpath = os.environ.get("CONFIGPATH", None) if cfgpath and os.path.exists(cfgpath): with open(cfgpath, 'r', encoding='utf-8') as f: cfg = yaml.safeload(f) else: cfg = DEFAULTCONFIG secret = os.environ.get(cfg.get('provenancesecretenv', 'PROVENANCESECRET'), 'changethissecret') try: app.state.orchestrator = AutoOrchestrator(cfg) app.state.orchestrator.start_workers() logger.info("Orchestrator initialized in FastAPI lifespan") except Exception as e: logger.exception("Failed to initialize orchestrator: %s", e) app.state.orchestrator = None yield if app.state.orchestrator: app.state.orchestrator.stop()
app = FastAPI(title="ForensicOilCompleteReform", version=DEFAULTCONFIG['modelversion'], lifespan=lifespan)
@app.post("/upload") async def uploadsample(samplemeta: str, gcms: Optional[UploadFile] = File(None), spectrum: Optional[UploadFile] = File(None), image: Optional[UploadFile] = File(None), afm: Optional[UploadFile] = File(None), cs14c: Optional[UploadFile] = File(None)): if app.state.orchestrator is None: raise HTTPException(statuscode=503, detail="Orchestrator not ready") try: metaobj = json.loads(samplemeta) except Exception: raise HTTPException(statuscode=400, detail="Invalid meta JSON") sid = metaobj.get('sampleid', str(uuid.uuid4())) gcmsb = await gcms.read() if gcms else None specb = await spectrum.read() if spectrum else None imgb = await image.read() if image else None afmb = await afm.read() if afm else None cs14cb = await cs14c.read() if cs14c else None sample = MultimodalSample(sampleid=sid, meta=metaobj.get('meta',{}), gcmsbytes=gcmsb, spectrumbytes=specb, imagebytes=imgb, afmbytes=afmb, cs14cbytes=cs14cb) app.state.orchestrator.enqueue({'type':'infer','sample':sample}) return JSONResponse({"sampleid": sid, "status":"queued"})
@app.get("/", response_class=HTMLResponse) async def index(): html = """ <html><head><title>Forensic Oil Review</title></head><body> <h2>Forensic Oil Residue System (Complete Reform)</h2> <p>Use /upload to POST sample meta and files. This is a prototype UI.</p> </body></html> """ return HTMLResponse(html)
----------------------------
CLI helpers & tests
----------------------------
def builddummysamples(n=200) -> List[MultimodalSample]: samples = [] for i in range(n): label = 1 if i < n//2 else 0 meta = {"sampleid": f"synthetic{i}", "objectmaterial": "ceramic", "environment": {"temperaturec": 20.0 + np.random.randn(), "relativehumidity": 40.0 + np.random.randn()*5}} # use placeholders to test robust parsing samples.append(MultimodalSample(sampleid=meta['sampleid'], meta=meta, gcmsbytes=b"mzmlplaceholder", spectrumbytes=b"spectrumplaceholder", imagebytes=b"tiffplaceholder", afmbytes=b"afm_placeholder", label=label)) return samples
def generatedeployartifacts(outdir: str): ensuredir(outdir) dockerfile = f"""FROM python:3.10-slim WORKDIR /app COPY {os.path.basename(file)} /app/forensicoilcompletereform.py RUN pip install --no-cache-dir fastapi uvicorn numpy pandas scipy scikit-learn xgboost pillow pymzml tifffile pymc3 arviz shap ENV PROVENANCESECRET=changethissecret EXPOSE 8000 CMD ["python", "forensicoilcompletereform.py", "--mode", "run-api"] """ with open(os.path.join(outdir, 'Dockerfile'), 'w', encoding='utf-8') as f: f.write(dockerfile) logger.info("Deployment artifacts generated at %s", outdir)
----------------------------
Main CLI
----------------------------
def main(): import argparse parser = argparse.ArgumentParser(description="Forensic Oil Complete Reform CLI") parser.addargument("--mode", choices=["build-dummy","train","infer-file","auto","run-api","gen-deploy","self-test"], required=False, default="build-dummy") parser.addargument("--manifest", type=str, default="trainmanifest.json") parser.addargument("--model-dir", type=str, default=DEFAULTCONFIG['modeldir']) parser.addargument("--out", type=str, default="out.json") parser.addargument("--sample-file", type=str, default=None) parser.addargument("--deploy-out", type=str, default="deployartifacts") args = parser.parse_args()
cfg = DEFAULTCONFIG.copy() cfg['modeldir'] = args.modeldir cfg['auditdir'] = DEFAULTCONFIG['auditdir'] cfg['auto']['watchdir'] = DEFAULTCONFIG['auto']['watch_dir']
if args.mode == "build-dummy": samples = builddummysamples(200) ensuredir(args.modeldir) res = trainpipeline(samples, cfg, args.modeldir) print("Dummy training result:", res) elif args.mode == "train": if not os.path.exists(args.manifest): print("Create trainmanifest.json listing samples with fields: sampleid, metajsonpath, gcmspath, spectrumpath, imagepath, afmpath, cs14cpath (optional), label") sys.exit(1) with open(args.manifest, 'r', encoding='utf-8') as f: manifestlist = json.load(f) samples = [] for rec in manifestlist: sid = rec['sampleid'] meta = json.load(open(rec['metajsonpath'], 'r', encoding='utf-8')) if rec.get('metajsonpath') else {} def readopt(p): return open(p, 'rb').read() if p and os.path.exists(p) else None samples.append(MultimodalSample(sampleid=sid, meta=meta, gcmsbytes=readopt(rec.get('gcmspath')), spectrumbytes=readopt(rec.get('spectrumpath')), imagebytes=readopt(rec.get('imagepath')), afmbytes=readopt(rec.get('afmpath')), cs14cbytes=readopt(rec.get('cs14cpath')), label=rec.get('label'))) res = trainpipeline(samples, cfg, args.modeldir) print("Training result:", res) elif args.mode == "infer-file": if not args.samplefile or not os.path.exists(args.samplefile): print("Provide sample JSON file describing paths.") sys.exit(1) with open(args.samplefile, 'r', encoding='utf-8') as f: rec = json.load(f) def readopt(p): return open(p, 'rb').read() if p and os.path.exists(p) else None sample = MultimodalSample(sampleid=rec.get('sampleid', str(uuid.uuid4())), meta=rec.get('meta', {}), gcmsbytes=readopt(rec.get('gcmspath')), spectrumbytes=readopt(rec.get('spectrumpath')), imagebytes=readopt(rec.get('imagepath')), afmbytes=readopt(rec.get('afmpath')), cs14cbytes=readopt(rec.get('cs14cpath'))) secret = os.environ.get(cfg.get('provenancesecretenv', 'PROVENANCESECRET'), 'changethissecret') engine = InferenceEngine(args.modeldir, cfg, secret) report = engine.infer(sample) with open(args.out, 'w', encoding='utf-8') as f: json.dump(report, f, ensureascii=False, indent=2) print("Inference saved to", args.out) elif args.mode == "auto": orchestrator = AutoOrchestrator(cfg) try: orchestrator.rundaemon() except Exception as e: logger.exception("Orchestrator failed: %s", e) elif args.mode == "run-api": if FastAPI is None: print("FastAPI not installed; cannot run API") sys.exit(1) import uvicorn uvicorn.run("forensicoilcompletereform:app", host="0.0.0.0", port=8000, reload=False) elif args.mode == "gen-deploy": generatedeployartifacts(args.deployout) print("Deployment artifacts generated at", args.deployout) elif args.mode == "self-test": # run a few unit-like checks to ensure robust parsing works print("Running self-test...") samples = builddummysamples(10) cfglocal = cfg ensuredir(cfglocal['modeldir']) res = trainpipeline(samples, cfglocal, cfglocal['modeldir']) print("Self-test train result:", res) engine = InferenceEngine(cfglocal['modeldir'], cfglocal, os.environ.get(cfglocal.get('provenancesecretenv','PROVENANCESECRET'),'changethissecret')) report = engine.infer(samples[0]) print("Self-test infer sample:", report['sample_id'], report['verdict'], report['confidence']) else: print("Unknown mode:", args.mode) sys.exit(1)
if _name == "main_": main()
