import os os.environ.setdefault("USE_TF", "0") os.environ.setdefault("USE_TORCH", "1") os.environ.setdefault("TOKENIZERS_PARALLELISM", "false") os.environ.setdefault("PYTORCH_CUDA_ALLOC_CONF", "expandable_segments:True") import json import random import shutil import time import numpy as np import pandas as pd import torch import trackio def _trackio_safe(fn): def wrapper(*a, **k): try: return fn(*a, **k) except Exception as e: print(f"trackio.{fn.__name__} failed (non-fatal): {e}") return wrapper trackio.log = _trackio_safe(trackio.log) trackio.init = _trackio_safe(trackio.init) trackio.finish = _trackio_safe(trackio.finish) from huggingface_hub import HfApi, snapshot_download from safetensors.torch import load_file, save_file from sklearn.metrics import f1_score from sklearn.model_selection import train_test_split from transformers import AutoTokenizer import laya from laya.agent import _fix_tokenizer_config from laya.common import QTYPES, build_model, build_sequence, proper_reward, render_options SMOKE = os.environ.get("SMOKE", "0") == "1" MODEL_ID = "convaiinnovations/laya" OUTPUT_DIR = os.environ.get("LAYA_OUTPUT_DIR", "/data/laya_vulnerability_groups") PARQUET_URL = "https://huggingface.co/datasets/GIZ/vulnerability_training_data_full/resolve/refs%2Fconvert%2Fparquet/default/train/0000.parquet" LABELS = [ "Agricultural communities", "Coastal communities", "Ethnic, racial or other minorities", "Fishery communities", "Informal sector workers", "Members of indigenous and local communities", "Migrants and displaced persons", "Older persons", "Other", "Persons living in poverty", "Persons with disabilities", "Persons with pre-existing health conditions", "Residents of drought-prone regions", "Rural populations", "Sexual minorities (LGBTQI+)", "Urban populations", "Women and other genders", ] REPO_ID = "peter2000/laya-vulnerability-groups" QIDS = [f"g{i}" for i in range(len(LABELS))] def make_questions(): return { qid: { "type": "noul", "instructions": f"Does this text indicate that {label} are targeted, supported, or affected as a vulnerable group? Answer true or false.", } for qid, label in zip(QIDS, LABELS) } QUESTIONS = make_questions() def load_data(): df = pd.read_parquet(PARQUET_URL) assert len(df) == 475, f"expected 475 rows, got {len(df)}" Y = df[LABELS].values.astype(np.int64) nlab = Y.sum(1) idx_tr, idx_te = train_test_split( np.arange(len(df)), test_size=0.2, random_state=42, stratify=np.minimum(nlab, 3) ) return df, Y, np.asarray(idx_tr), np.asarray(idx_te) def ece(conf, correct, n_bins=15): conf = np.asarray(conf, dtype=np.float64) corr = np.asarray(correct, dtype=np.float64) bins = np.linspace(0.0, 1.0, n_bins + 1) e = 0.0 for lo, hi in zip(bins[:-1], bins[1:]): m = (conf > lo) & (conf <= hi) if m.sum() > 0: e += m.mean() * abs(corr[m].mean() - conf[m].mean()) return float(e) def evaluate(Y_true, P_pred, threshold=0.5): pred = (P_pred >= threshold).astype(int) per_label = f1_score(Y_true, pred, average=None, zero_division=0) macro = float(f1_score(Y_true, pred, average="macro", zero_division=0)) micro = float(f1_score(Y_true, pred, average="micro", zero_division=0)) conf = np.where(pred == 1, P_pred, 1.0 - P_pred) corr = (pred == Y_true).astype(np.float64) return { "macro_f1": macro, "micro_f1": micro, "ece": ece(conf, corr), "subset_accuracy": float(((pred == Y_true).all(axis=1)).mean()), "per_label_f1": {LABELS[i]: round(float(per_label[i]), 4) for i in range(len(LABELS))}, } def probs_from_answers(res): return np.array([res["answers"][qid]["noul"] for qid in QIDS], dtype=np.float64) def eval_agent(agent, texts, Y_true): t0 = time.time() P = np.stack([probs_from_answers(agent.predict(t, QUESTIONS)) for t in texts]) metrics = evaluate(Y_true, P) metrics["eval_seconds"] = round(time.time() - t0, 1) return metrics def collate_train_batch(items, pad_id): n, L = len(items), max(len(it["ids"]) for it in items) kmax = max(len(it["markers"]) for it in items) ids = torch.full((n, L), pad_id, dtype=torch.long) att = torch.zeros((n, L), dtype=torch.long) mpos = torch.zeros((n, kmax), dtype=torch.long) mmask = torch.zeros((n, kmax), dtype=torch.bool) target = torch.zeros((n, kmax), dtype=torch.float32) for i, it in enumerate(items): ids[i, : len(it["ids"])] = torch.tensor(it["ids"]) att[i, : len(it["ids"])] = 1 k = len(it["markers"]) mpos[i, :k] = torch.tensor(it["markers"]) mmask[i, :k] = True target[i, : len(it["target"])] = torch.tensor(it["target"], dtype=torch.float32) return { "input_ids": ids, "attention_mask": att, "marker_pos": mpos, "marker_mask": mmask, "target": target, "qtype": torch.tensor([it["qtype"] for it in items]), "label": torch.tensor([it["label"] for it in items]), } def fit_one_temp(sel): if len(sel) < 10: return 1.0 kmax = max(len(z) for z, _ in sel) Z = torch.full((len(sel), kmax), -1e4) T = torch.zeros((len(sel), kmax)) for i, (z, t) in enumerate(sel): Z[i, :len(z)] = torch.tensor(z) T[i, :len(t)] = torch.tensor(t, dtype=torch.float32) log_t = torch.zeros(1, requires_grad=True) opt = torch.optim.LBFGS([log_t], lr=0.1, max_iter=100) def closure(): opt.zero_grad() loss = -(T * torch.log_softmax(Z / log_t.exp(), -1)).sum(-1).mean() loss.backward() return loss opt.step(closure) return float(torch.clamp(log_t.exp(), 0.1, 10.0).item()) def build_items(texts, yvecs, tok, cfg): items, dropped = [], 0 k_expected = len(render_options({"t": "noul", "crit": {}})) for text, yvec in zip(texts, yvecs): for j in range(len(LABELS)): p_true = float(yvec[j]) target = [1.0 - p_true, p_true] q = {"t": "noul", "ins": QUESTIONS[QIDS[j]]["instructions"], "crit": {}} seq, markers = build_sequence(tok, text, q, cfg["max_len"], cfg["head_max_len"]) if len(markers) != k_expected: dropped += 1 continue items.append({ "ids": seq, "markers": markers, "qtype": QTYPES["noul"], "target": target, "label": j, }) print(f"built {len(items)} items, dropped {dropped} (marker mismatch), k_expected={k_expected}") return items def train(items, model, tok, cfg, device): EPOCHS = 1 if SMOKE else 4 MICRO_BATCH = 8 GRAD_ACCUM = 4 GROUP_SIZE = 4 LR_ENCODER = 2.5e-5 LR_HEAD = 1.0e-4 SIGMA_START = 0.4 SIGMA_END = 0.1 enc_params = [p for n, p in model.named_parameters() if "encoder." in n] head_params = [p for n, p in model.named_parameters() if "encoder." not in n] optimizer = torch.optim.AdamW([ {"params": enc_params, "lr": LR_ENCODER}, {"params": head_params, "lr": LR_HEAD}, ], weight_decay=0.01) total_updates = max(1, (len(items) // (MICRO_BATCH * GRAD_ACCUM)) * EPOCHS) scheduler = torch.optim.lr_scheduler.CosineAnnealingLR(optimizer, T_max=total_updates, eta_min=1e-6) scaler = torch.amp.GradScaler("cuda", enabled=True) pad_id = tok.pad_token_id t0 = time.time() for epoch in range(EPOCHS): random.seed(42 + epoch) random.shuffle(items) epoch_loss, n_batches = 0.0, 0 optimizer.zero_grad(set_to_none=True) accum_step = 0 progress = epoch / max(1, EPOCHS - 1) sigma = SIGMA_START + (SIGMA_END - SIGMA_START) * progress for b_idx in range(0, len(items), MICRO_BATCH): chunk = items[b_idx:b_idx + MICRO_BATCH] if not chunk: continue batch = collate_train_batch(chunk, pad_id) with torch.autocast("cuda", dtype=torch.float16): logits, act = model( batch["input_ids"].to(device), batch["attention_mask"].to(device), batch["marker_pos"].to(device), batch["marker_mask"].to(device), batch["qtype"].to(device), ) logits = logits.float() mask = batch["marker_mask"].to(device) k = mask.sum(-1, keepdim=True).float() target = batch["target"].to(device) eps = torch.randn((GROUP_SIZE,) + logits.shape, device=device) * sigma * mask eps = (eps - eps.sum(-1, keepdim=True) / k) * mask z = logits.detach().unsqueeze(0) + eps q = torch.softmax(z.masked_fill(~mask, -1e4), -1) with torch.no_grad(): r = proper_reward(q, target.unsqueeze(0), batch["qtype"].to(device), mask, w_sph=0.75, w_rps=1.0) adv = r - r.mean(0, keepdim=True) adv = adv / (adv.std() + 1e-6) logp = -(((z - logits.unsqueeze(0)) ** 2) * mask).sum(-1) / (2 * sigma ** 2) loss_rl = -(adv * logp).mean() loss_ce = -(target * torch.log_softmax(logits.masked_fill(~mask, -1e4), -1)).sum(-1).mean() loss = (loss_rl + 1.0 * loss_ce) / GRAD_ACCUM + 0.0 * act.sum() scaler.scale(loss).backward() accum_step += 1 if accum_step % GRAD_ACCUM == 0 or (b_idx + MICRO_BATCH) >= len(items): scaler.unscale_(optimizer) torch.nn.utils.clip_grad_norm_(model.parameters(), 1.0) scaler.step(optimizer) scaler.update() scheduler.step() optimizer.zero_grad(set_to_none=True) epoch_loss += loss.item() * GRAD_ACCUM n_batches += 1 if n_batches % 50 == 0: print(f" epoch {epoch+1}/{EPOCHS} step {n_batches} loss {loss.item()*GRAD_ACCUM:.4f} reward {r.mean().item():.3f} lr {scheduler.get_last_lr()[0]:.2e}") trackio.log({"laya_loss": loss.item() * GRAD_ACCUM, "laya_reward": r.mean().item()}, step=epoch * 100000 + n_batches) print(f"=== epoch {epoch+1}/{EPOCHS} done in {time.time()-t0:.0f}s avg_loss {epoch_loss/max(1,n_batches):.4f} ===") trackio.log({"laya_epoch_loss": epoch_loss / max(1, n_batches)}, step=epoch + 1) ckpt_dir = os.path.join(OUTPUT_DIR, "checkpoint_latest") os.makedirs(ckpt_dir, exist_ok=True) ckpt_sd = {k: v.half().contiguous().cpu() for k, v in model.state_dict().items()} save_file(ckpt_sd, os.path.join(ckpt_dir, "model.safetensors")) model.encoder.config.save_pretrained(os.path.join(ckpt_dir, "encoder")) tok.save_pretrained(os.path.join(ckpt_dir, "tokenizer")) with open(os.path.join(ckpt_dir, "checkpoint_meta.json"), "w") as f: json.dump({"epoch": epoch + 1, "total_epochs": EPOCHS, "avg_loss": epoch_loss / max(1, n_batches)}, f, indent=2) del optimizer, scaler, scheduler torch.cuda.empty_cache() return model def fit_and_apply_temperature(items, model, cfg, device, pad_id): model.eval() calib_items = items[::15][:400] calib_preds = [] with torch.no_grad(): for c_idx in range(0, len(calib_items), 16): c_chunk = calib_items[c_idx:c_idx + 16] cb = collate_train_batch(c_chunk, pad_id) with torch.autocast("cuda", dtype=torch.float16): l_sub, _ = model( cb["input_ids"].to(device), cb["attention_mask"].to(device), cb["marker_pos"].to(device), cb["marker_mask"].to(device), cb["qtype"].to(device), ) l_np = l_sub.float().cpu().numpy() for rr, it in enumerate(c_chunk): k = len(it["markers"]) calib_preds.append((it["qtype"], l_np[rr, :k], it["target"])) fitted = [1.2, 1.2, 1.2] for qt in range(3): sel = [(z, t) for q_type, z, t in calib_preds if q_type == qt] if sel: fitted[qt] = fit_one_temp(sel) print("fitted calibration temperatures (choice, score, noul):", [round(t, 3) for t in fitted]) cfg["temperature"] = fitted if isinstance(cfg.get("temperature_by_options"), dict): cfg["temperature_by_options"]["noul:2"] = fitted[QTYPES["noul"]] return fitted def save_checkpoint(model, tok, cfg): os.makedirs(OUTPUT_DIR, exist_ok=True) sd = {k: v.half().contiguous().cpu() for k, v in model.state_dict().items()} save_file(sd, os.path.join(OUTPUT_DIR, "model.safetensors")) model.encoder.config.save_pretrained(os.path.join(OUTPUT_DIR, "encoder")) tok.save_pretrained(os.path.join(OUTPUT_DIR, "tokenizer")) cfg["fine_tuned"] = True cfg["model_name"] = "laya-vulnerability-groups" with open(os.path.join(OUTPUT_DIR, "rl_agent_config.json"), "w") as f: json.dump(cfg, f, indent=2) ckpt_dir = os.path.join(OUTPUT_DIR, "checkpoint_latest") if os.path.isdir(ckpt_dir): shutil.rmtree(ckpt_dir) print(f"checkpoint saved to {OUTPUT_DIR}") def write_readme(metrics): readme = f"""--- license: apache-2.0 library_name: transformers tags: [laya, text-classification, multi-label, climate, vulnerability, rlcd] pipeline_tag: text-classification --- # Laya fine-tuned for climate-vulnerability group detection (multi-label) [convaiinnovations/laya](https://huggingface.co/convaiinnovations/laya) (421M, ModernBERT-large backbone) fine-tuned with the official RLCD recipe on [GIZ/vulnerability_training_data_full](https://huggingface.co/datasets/GIZ/vulnerability_training_data_full) (380 train rows x 17 binary vulnerability-group questions; 36 all-negative rows included as negatives). Each vulnerability group is asked as one binary (`noul`) typed question; all 17 are answered in a single forward pass. Evaluate with `laya.load("peter2000/laya-vulnerability-groups")` and `agent.predict(state, questions)`. ## Test-set metrics (held-out {95 if not SMOKE else 5} rows, threshold 0.5) | metric | value | |---|---| | macro-F1 | {metrics['macro_f1']:.4f} | | micro-F1 | {metrics['micro_f1']:.4f} | | ECE | {metrics['ece']:.4f} | | subset accuracy | {metrics['subset_accuracy']:.4f} | Per-label F1: | label | F1 | |---|---| """ for label, f1 in metrics["per_label_f1"].items(): readme += f"| {label} | {f1:.4f} |\n" with open(os.path.join(OUTPUT_DIR, "README.md"), "w") as f: f.write(readme) def main(): df, Y, idx_tr, idx_te = load_data() texts = df["text"].tolist() X_tr = [texts[i] for i in idx_tr] X_te = [texts[i] for i in idx_te] Y_tr, Y_te = Y[idx_tr], Y[idx_te] if SMOKE: X_te, Y_te = X_te[:5], Y_te[:5] print(f"train={len(X_tr)} test={len(X_te)} labels={len(LABELS)}") trackio.init(project="vulnerability-multilabel-classifier", space_id="peter2000/vulnerability-multilabel-classifier-trackio") model_dir = snapshot_download(MODEL_ID, ignore_patterns=["multilingual/*", "typed-decisions/*", "assets/*", "eval/*", "*.py"]) _fix_tokenizer_config(model_dir) tok = AutoTokenizer.from_pretrained(os.path.join(model_dir, "tokenizer")) with open(os.path.join(model_dir, "rl_agent_config.json")) as f: cfg = json.load(f) cfg["gradient_checkpointing"] = True device = "cuda" print("--- zero-shot baseline eval ---") agent0 = laya.load(model_dir, device=device) m0 = eval_agent(agent0, X_te, Y_te) print("zero-shot:", json.dumps(m0, indent=2)) trackio.log({"laya_zeroshot_macro_f1": m0["macro_f1"], "laya_zeroshot_micro_f1": m0["micro_f1"], "laya_zeroshot_ece": m0["ece"]}, step=0) del agent0 torch.cuda.empty_cache() items = build_items(X_tr, Y_tr, tok, cfg) if SMOKE: items = items[:64] print("building model...") model = build_model(cfg, encoder_dir=os.path.join(model_dir, "encoder")) model.load_state_dict(load_file(os.path.join(model_dir, "model.safetensors")), strict=True) model.encoder.gradient_checkpointing_enable(gradient_checkpointing_kwargs={"use_reentrant": False}) model.head_checkpointing = True model.to(device) model.train() train(items, model, tok, cfg, device) fitted = fit_and_apply_temperature(items, model, cfg, device, tok.pad_token_id) print("fitted temperatures:", [round(t, 3) for t in fitted]) save_checkpoint(model, tok, cfg) print("--- fine-tuned eval ---") agent_ft = laya.load(OUTPUT_DIR, device=device) m1 = eval_agent(agent_ft, X_te, Y_te) print("fine-tuned:", json.dumps(m1, indent=2)) trackio.log({"laya_macro_f1": m1["macro_f1"], "laya_micro_f1": m1["micro_f1"], "laya_ece": m1["ece"]}, step=10) trackio.log({f"laya_f1/{k}": v for k, v in m1["per_label_f1"].items()}, step=10) write_readme(m1) api = HfApi(token=os.environ.get("HF_TOKEN")) api.upload_folder( folder_path=OUTPUT_DIR, repo_id=REPO_ID, repo_type="model", commit_message=f"Laya fine-tuned on GIZ vulnerability data: macro-F1 {m1['macro_f1']:.3f} (zero-shot {m0['macro_f1']:.3f})", ) api.upload_file( path_or_fileobj=json.dumps({"zero_shot": m0, "fine_tuned": m1}, indent=2).encode(), path_in_repo="metrics.json", repo_id=REPO_ID, repo_type="model", commit_message="Add evaluation metrics", ) trackio.finish() print("DONE") if __name__ == "__main__": t0 = time.time() main() print(f"elapsed {time.time()-t0:.0f}s")