2026-08-22 07:23:37 +07:00 | | | # -*- coding: utf-8 -*-
|
2026-08-22 07:46:00 +07:00 | | | """MACHINE ANNOTATOR v2 — annotation machine utk golden cases (Liquidity Sweep).
|
| | |
|
| | | v2 (P3-S.0): f7 memakai EVENT LIFECYCLE corrected:
|
| | | f7_lifecycle(sweep_dir, sweep_bar) -> NO_SWEEP .. SWEEP_ONSET .. EXPIRED .. NO_SWEEP
|
| | | expiration = SEQ_WINDOW (40 bar) = InpSeqWindow v4.4/v4.5 (semantics existing).
|
| | | f10/f11 (DetectEQ) TIDAK diubah (SEMANTIC REVIEW REQUIRED).
|
| | |
|
| | | Output:
|
| | | machine_annotations_f7_v2.csv (corrected — primary utk human annotation nanti)
|
| | | f7_v1_vs_v2_comparison.json (legacy vs corrected: changed cases, events removed/added)
|
| | |
|
| | | Metadata per baris: source commit, source hash, dataset hash, contract hash, annotation version.
|
| | | Parity: f10/f11 harus tetap identik dgn F cache P2.6; f7 diharapkan BERBEDA utk kasus stale
|
| | | (itu efek fix). Semua angka dari kode existing + fix lifecycle; tidak ada alasan baru.
|
| | |
|
2026-08-22 07:23:37 +07:00 | | | Usage: python machine_annotator.py
|
| | | """
|
| | | import os
|
| | | import sys
|
2026-08-22 07:46:00 +07:00 | | | import csv
|
2026-08-22 07:23:37 +07:00 | | | import json
|
2026-08-22 07:46:00 +07:00 | | | import hashlib
|
2026-08-22 07:23:37 +07:00 | | | import datetime as dt
|
| | |
|
| | | import numpy as np
|
| | |
|
| | | HERE = os.path.dirname(os.path.abspath(__file__))
|
| | | sys.path.insert(0, HERE)
|
| | | import smc_semantic_common as SC
|
2026-08-22 07:46:00 +07:00 | | | import p3_common as P3
|
| | |
|
| | | ANNOTATION_VERSION = "f7_v2"
|
2026-08-22 08:10:48 +07:00 | | | SOURCE_COMMIT = "b519a34" # commit yang mem-freeze f7 lifecycle fix (P3-S.0)
|
2026-08-22 07:46:00 +07:00 | | |
|
| | |
|
| | | def sha256_file(path):
|
| | | h = hashlib.sha256()
|
| | | with open(path, "rb") as f:
|
| | | for blk in iter(lambda: f.read(1 << 20), b""):
|
| | | h.update(blk)
|
| | | return h.hexdigest()
|
2026-08-22 07:23:37 +07:00 | | |
|
| | |
|
| | | def iso(ts):
|
| | | return dt.datetime.fromtimestamp(int(ts), tz=dt.timezone.utc).strftime("%Y-%m-%d %H:%M:%S")
|
| | |
|
| | |
|
| | | def main():
|
| | | t, o, h, l, c, v, htf = SC.load_data()
|
| | | F = SC.load_F()
|
| | | A = SC.atr_series(h, l, c)
|
| | | n = len(c)
|
| | |
|
| | | sw_full, inn_full = SC.build_structures(o, h, l, c)
|
| | | pairs_h, pairs_l = SC.eq_pairs(h, l, sw_full["pivots"])
|
| | | sweep_dir, sweep_bar = SC.sweep_state(h, l, c, inn_full["pivots"])
|
2026-08-22 07:46:00 +07:00 | | | f7_corr = SC.f7_lifecycle(sweep_dir, sweep_bar)
|
2026-08-22 07:23:37 +07:00 | | |
|
2026-08-22 07:46:00 +07:00 | | | # level grab legacy (sweep_state tidak simpan level)
|
2026-08-22 07:23:37 +07:00 | | | sweep_level = np.full(n, np.nan)
|
| | | cur_sd, cur_sb, cur_sl = 0, -1, 0.0
|
| | | for (p, lvl, is_high) in inn_full["pivots"]:
|
| | | last = min(n - 1, p + SC.GRAB_WINDOW)
|
| | | for b in range(p + 1, last + 1):
|
| | | if is_high and h[b] > lvl and c[b] < lvl:
|
| | | if b > cur_sb:
|
| | | cur_sd, cur_sb, cur_sl = -1, b, lvl
|
| | | break
|
| | | if (not is_high) and l[b] < lvl and c[b] > lvl:
|
| | | if b > cur_sb:
|
| | | cur_sd, cur_sb, cur_sl = 1, b, lvl
|
| | | break
|
| | | if cur_sb >= 0:
|
| | | sweep_level[cur_sb:] = cur_sl
|
| | |
|
| | | meta_path = os.path.join(HERE, "output", "cases_meta.json")
|
| | | if not os.path.exists(meta_path):
|
| | | raise SystemExit("cases_meta.json tidak ada — jalankan sample_cases.py dulu")
|
| | | meta = json.load(open(meta_path, encoding="utf-8"))
|
| | |
|
2026-08-22 07:46:00 +07:00 | | | # hash provenance
|
| | | prov = {
|
| | | "annotation_version": ANNOTATION_VERSION,
|
| | | "source_commit": SOURCE_COMMIT,
|
| | | "ea_source_sha": sha256_file(r"D:\TradingTerminal\HFM Metatrader 5\MQL5\Experts\AlgoForge_Backtest_Baseline.mq5"),
|
| | | "annotator_sha": sha256_file(os.path.abspath(__file__)),
|
| | | "common_sha": sha256_file(os.path.join(HERE, "smc_semantic_common.py")),
|
| | | "dataset_hash_prefix": "e85a0861", # P2.6 dataset (F cache)
|
| | | "feature_contract_hash": "C44CC6F2B740C32D06F776BD7C3E669DC5A8A6DE0484230544EBFFCF517D38DD",
|
| | | "SEQ_WINDOW": SC.SEQ_WINDOW,
|
| | | "note": "f7 corrected (event lifecycle); f10/f11 UNCHANGED",
|
| | | }
|
| | |
|
2026-08-22 07:23:37 +07:00 | | | rows = []
|
2026-08-22 07:46:00 +07:00 | | | n_f7_mismatch_cache = 0
|
| | | n_f1011_mismatch = 0
|
2026-08-22 07:23:37 +07:00 | | | for case in meta["cases"]:
|
| | | cid = case["case_id"]
|
| | | r = int(case["bar_idx"])
|
| | | ts = case["decision_timestamp"]
|
| | |
|
2026-08-22 07:46:00 +07:00 | | | f7_legacy = int(sweep_dir[r])
|
| | | f7_new = int(f7_corr[r])
|
| | | sb = int(sweep_bar[r])
|
| | | age = (r - sb) if sb >= 0 else None
|
2026-08-22 07:23:37 +07:00 | | |
|
2026-08-22 07:46:00 +07:00 | | | eq_st = SC.eq_state_at(r, pairs_h, pairs_l, A)
|
| | | f10 = eq_st["eqh"]
|
| | | f11 = eq_st["eql"]
|
| | |
|
| | | # parity thd F cache
|
| | | pf7, pf10, pf11 = int(F[r, 7]), int(F[r, 10]), int(F[r, 11])
|
| | | if f7_new != pf7:
|
| | | n_f7_mismatch_cache += 1
|
| | | if (f10, f11) != (pf10, pf11):
|
| | | n_f1011_mismatch += 1
|
2026-08-22 07:23:37 +07:00 | | |
|
| | | grab_detail = None
|
2026-08-22 07:46:00 +07:00 | | | if sb >= 0 and not np.isnan(sweep_level[r]):
|
2026-08-22 07:23:37 +07:00 | | | lvl = float(sweep_level[r])
|
2026-08-22 07:46:00 +07:00 | | | extreme = float(h[sb]) if f7_legacy < 0 else float(l[sb])
|
2026-08-22 07:23:37 +07:00 | | | grab_detail = {
|
2026-08-22 07:46:00 +07:00 | | | "dir_legacy": f7_legacy,
|
2026-08-22 07:23:37 +07:00 | | | "sweep_bar": sb,
|
2026-08-22 07:46:00 +07:00 | | | "state_age_bars": age,
|
| | | "expired_by_seq_window": bool(age is not None and age > SC.SEQ_WINDOW),
|
2026-08-22 07:23:37 +07:00 | | | "reference_level": round(lvl, 2),
|
| | | "sweep_price": round(extreme, 2),
|
| | | "excess": round(abs(extreme - lvl), 2),
|
| | | "excess_atr": round(abs(extreme - lvl) / A[r], 3),
|
| | | }
|
2026-08-22 07:46:00 +07:00 | | | eqh_detail = None
|
2026-08-22 07:23:37 +07:00 | | | if eq_st["pair_h"]:
|
| | | p1, p2, fc, dp = eq_st["pair_h"]
|
2026-08-22 07:46:00 +07:00 | | | eqh_detail = {"p1": p1, "p2": p2, "first_cross": fc,
|
| | | "reference_level": round(float(h[p2]), 2),
|
| | | "sweep_price": round(float(h[fc]), 2),
|
| | | "excess": round(float(h[fc] - h[p2]), 2),
|
| | | "excess_atr": round(float(h[fc] - h[p2]) / A[r], 3),
|
| | | "dp": round(dp, 2), "tol": round(eq_st["tol_atr"], 2)}
|
| | | eql_detail = None
|
2026-08-22 07:23:37 +07:00 | | | if eq_st["pair_l"]:
|
| | | p1, p2, fc, dp = eq_st["pair_l"]
|
2026-08-22 07:46:00 +07:00 | | | eql_detail = {"p1": p1, "p2": p2, "first_cross": fc,
|
| | | "reference_level": round(float(l[p2]), 2),
|
| | | "sweep_price": round(float(l[fc]), 2),
|
| | | "excess": round(float(l[p2] - l[fc]), 2),
|
| | | "excess_atr": round(float(l[p2] - l[fc]) / A[r], 3),
|
| | | "dp": round(dp, 2), "tol": round(eq_st["tol_atr"], 2)}
|
2026-08-22 07:23:37 +07:00 | | |
|
2026-08-22 07:46:00 +07:00 | | | decision = "YES" if (f7_new != 0 or f10 == 1 or f11 == 1) else "NO"
|
2026-08-22 07:23:37 +07:00 | | | active = []
|
2026-08-22 07:46:00 +07:00 | | | if f7_new != 0:
|
2026-08-22 07:23:37 +07:00 | | | active.append("f7_grab")
|
| | | if f10 == 1:
|
| | | active.append("f10_eqh")
|
| | | if f11 == 1:
|
| | | active.append("f11_eql")
|
| | | if not active:
|
| | | active.append("NONE")
|
| | | primary = "f10_eqh" if f10 == 1 else ("f11_eql" if f11 == 1 else
|
2026-08-22 07:46:00 +07:00 | | | ("f7_grab" if f7_new != 0 else "NONE"))
|
2026-08-22 07:23:37 +07:00 | | |
|
| | | if primary == "f10_eqh" and eqh_detail:
|
2026-08-22 07:46:00 +07:00 | | | ref_type, ref_level = "EQH (swing len=50)", eqh_detail["reference_level"]
|
| | | sweep_price, excess = eqh_detail["sweep_price"], eqh_detail["excess"]
|
| | | rejection = "N/A (kode TIDAK pakai close-back utk EQH — tidak diubah sesi ini)"
|
2026-08-22 07:23:37 +07:00 | | | elif primary == "f11_eql" and eql_detail:
|
2026-08-22 07:46:00 +07:00 | | | ref_type, ref_level = "EQL (swing len=50)", eql_detail["reference_level"]
|
| | | sweep_price, excess = eql_detail["sweep_price"], eql_detail["excess"]
|
| | | rejection = "N/A (kode TIDAK pakai close-back utk EQL — tidak diubah sesi ini)"
|
2026-08-22 07:23:37 +07:00 | | | elif primary == "f7_grab" and grab_detail:
|
2026-08-22 07:46:00 +07:00 | | | ref_type, ref_level = "Swing (internal len=5)", grab_detail["reference_level"]
|
| | | sweep_price, excess = grab_detail["sweep_price"], grab_detail["excess"]
|
2026-08-22 07:23:37 +07:00 | | | rejection = "YES (close-back ada di kondisi grab)"
|
| | | else:
|
2026-08-22 07:46:00 +07:00 | | | ref_type = ref_level = sweep_price = excess = "N/A"
|
2026-08-22 07:23:37 +07:00 | | | rejection = "N/A"
|
| | |
|
| | | reasons = []
|
2026-08-22 07:46:00 +07:00 | | | if f7_new != 0 and grab_detail:
|
| | | d = "bearish (buy-side swept)" if f7_new < 0 else "bullish (sell-side swept)"
|
2026-08-22 07:23:37 +07:00 | | | reasons.append(f"f7 grab {d} @bar{grab_detail['sweep_bar']} "
|
| | | f"lvl={grab_detail['reference_level']} px={grab_detail['sweep_price']} "
|
| | | f"excess={grab_detail['excess_atr']}ATR close-back=YES")
|
2026-08-22 07:46:00 +07:00 | | | if sb >= 0 and f7_new == 0:
|
| | | reasons.append(f"f7 grab @bar{sb} EXPIRED (age={age}>SEQ_WINDOW={SC.SEQ_WINDOW}) "
|
| | | f"-> f7=NO_SWEEP")
|
2026-08-22 07:23:37 +07:00 | | | if eqh_detail:
|
| | | reasons.append(f"f10 EQH swept @bar{eqh_detail['first_cross']} "
|
| | | f"lvl={eqh_detail['reference_level']} px={eqh_detail['sweep_price']} "
|
2026-08-22 07:46:00 +07:00 | | | f"excess={eqh_detail['excess_atr']}ATR")
|
2026-08-22 07:23:37 +07:00 | | | if eql_detail:
|
| | | reasons.append(f"f11 EQL swept @bar{eql_detail['first_cross']} "
|
| | | f"lvl={eql_detail['reference_level']} px={eql_detail['sweep_price']} "
|
2026-08-22 07:46:00 +07:00 | | | f"excess={eql_detail['excess_atr']}ATR")
|
2026-08-22 07:23:37 +07:00 | | | if not reasons:
|
2026-08-22 07:46:00 +07:00 | | | reasons.append("tidak ada primitif sweep aktif di bar keputusan (f7=0, f10=0, f11=0)")
|
2026-08-22 07:23:37 +07:00 | | |
|
| | | rows.append({
|
| | | "case_id": cid,
|
| | | "decision_timestamp": ts,
|
| | | "symbol": SC.SYMBOL,
|
| | | "decision_tf": SC.DECISION_TF,
|
| | | "machine_decision": decision,
|
| | | "primary_primitive": primary,
|
2026-08-22 07:46:00 +07:00 | | | "active_primitives": "+".join(active),
|
| | | "f7_legacy": f7_legacy,
|
| | | "f7_corrected": f7_new,
|
| | | "f7_state_age": age if age is not None else "N/A",
|
| | | "f7_expired": "YES" if (sb >= 0 and age > SC.SEQ_WINDOW) else "NO",
|
| | | "f10": f10, "f11": f11,
|
2026-08-22 07:23:37 +07:00 | | | "reference_type": ref_type,
|
| | | "reference_level": ref_level,
|
| | | "sweep_price": sweep_price,
|
| | | "excess": excess,
|
| | | "close_price": round(float(c[r]), 2),
|
| | | "rejection": rejection,
|
| | | "timeframe": SC.DECISION_TF,
|
| | | "machine_reason": "; ".join(reasons),
|
| | | "grab_detail": json.dumps(grab_detail) if grab_detail else "N/A",
|
| | | "eqh_detail": json.dumps(eqh_detail) if eqh_detail else "N/A",
|
| | | "eql_detail": json.dumps(eql_detail) if eql_detail else "N/A",
|
2026-08-22 07:46:00 +07:00 | | | "parity_f7_vs_cache": "EXPECTED_DIFF" if f7_new != pf7 else "OK",
|
| | | "parity_f10_f11_vs_cache": "OK" if (f10, f11) == (pf10, pf11) else "MISMATCH",
|
| | | "annotation_version": ANNOTATION_VERSION,
|
2026-08-22 07:23:37 +07:00 | | | })
|
| | |
|
| | | outdir = os.path.join(HERE, "output")
|
| | | os.makedirs(outdir, exist_ok=True)
|
2026-08-22 07:46:00 +07:00 | | | out_csv = os.path.join(outdir, "machine_annotations_f7_v2.csv")
|
| | | with open(out_csv, "w", newline="", encoding="utf-8") as f:
|
| | | w = csv.DictWriter(f, fieldnames=list(rows[0].keys()))
|
2026-08-22 07:23:37 +07:00 | | | w.writeheader()
|
| | | w.writerows(rows)
|
2026-08-22 07:46:00 +07:00 | | | print(f"[saved] {out_csv} ({len(rows)} rows)")
|
| | | print(f"f7 corrected vs F cache (P2.6 legacy) : {n_f7_mismatch_cache}/{len(rows)} differ (expected utk stale)")
|
| | | print(f"f10/f11 vs F cache : {n_f1011_mismatch}/{len(rows)} differ (should be 0)")
|
| | |
|
| | | # ---- comparison legacy vs corrected ----
|
| | | legacy_path = os.path.join(outdir, "machine_annotations_f7_v1_LEGACY.csv")
|
| | | legacy_rows = {}
|
| | | if os.path.exists(legacy_path):
|
| | | with open(legacy_path, encoding="utf-8-sig") as f:
|
| | | for row in csv.DictReader(f):
|
| | | legacy_rows[row["case_id"]] = row
|
| | |
|
| | | cmp = {"annotation_version": ANNOTATION_VERSION,
|
| | | "legacy_source": "machine_annotations_f7_v1_LEGACY.csv",
|
| | | "generated_utc": dt.datetime.now(dt.timezone.utc).isoformat(),
|
| | | "provenance": prov,
|
| | | "per_case": []}
|
| | | n_changed = 0
|
| | | n_unchanged = 0
|
| | | n_events_removed = 0
|
| | | n_events_added = 0
|
| | | n_persistent_corrected = 0
|
| | | for row in rows:
|
| | | cid = row["case_id"]
|
| | | old = legacy_rows.get(cid)
|
| | | old_f7 = int(old["f7_legacy"]) if old and "f7_legacy" in old else int(old["parity_f7"])
|
| | | new_f7 = int(row["f7_corrected"])
|
| | | old_dec = old["machine_decision"] if old else "?"
|
| | | changed = (old_f7 != new_f7)
|
| | | if changed:
|
| | | n_changed += 1
|
| | | if old_f7 != 0 and new_f7 == 0:
|
| | | n_events_removed += 1
|
| | | if int(row["f7_state_age"]) > SC.SEQ_WINDOW:
|
| | | n_persistent_corrected += 1
|
| | | elif old_f7 == 0 and new_f7 != 0:
|
| | | n_events_added += 1
|
| | | else:
|
| | | n_unchanged += 1
|
| | | cmp["per_case"].append({
|
| | | "case_id": cid,
|
| | | "stratum": next((cc["sampling_stratum"] for cc in meta["cases"]
|
| | | if cc["case_id"] == cid), "?"),
|
| | | "f7_legacy": old_f7,
|
| | | "f7_corrected": new_f7,
|
| | | "f7_state_age": row["f7_state_age"],
|
| | | "f7_expired": row["f7_expired"],
|
| | | "machine_decision_legacy": old_dec,
|
| | | "machine_decision_corrected": row["machine_decision"],
|
| | | "changed": changed,
|
| | | "events_removed": 1 if (old_f7 != 0 and new_f7 == 0) else 0,
|
| | | "events_added": 1 if (old_f7 == 0 and new_f7 != 0) else 0,
|
| | | "persistent_state_corrected": 1 if (old_f7 != 0 and new_f7 == 0 and
|
| | | int(row["f7_state_age"]) > SC.SEQ_WINDOW) else 0,
|
| | | })
|
| | | cmp["summary"] = {
|
| | | "n_cases": len(rows),
|
| | | "cases_changed": n_changed,
|
| | | "cases_unchanged": n_unchanged,
|
| | | "events_removed": n_events_removed,
|
| | | "events_added": n_events_added,
|
| | | "persistent_state_cases_corrected": n_persistent_corrected,
|
| | | "f10_f11_unchanged": (n_f1011_mismatch == 0),
|
| | | }
|
| | | cmp_path = os.path.join(outdir, "f7_v1_vs_v2_comparison.json")
|
| | | with open(cmp_path, "w", encoding="utf-8") as f:
|
| | | json.dump(cmp, f, indent=2, default=str)
|
| | | print(f"[saved] {cmp_path}")
|
| | | print("summary:", json.dumps(cmp["summary"], indent=2))
|
2026-08-22 07:23:37 +07:00 | | |
|
| | |
|
| | | if __name__ == "__main__":
|
| | | main()
|