File size: 2,831 Bytes
eb75e03
 
 
d58a3db
 
 
 
 
 
 
eb75e03
d58a3db
eb75e03
 
 
 
 
d58a3db
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
eb75e03
 
 
 
 
 
d58a3db
eb75e03
d58a3db
 
 
 
 
eb75e03
 
2c212f0
eb75e03
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
"""Remote chain daemon: wait for collection to finish, then auto-launch
the v3 orchestrator. Run once via nohup; it exits after handing off.

Chain: collect_selfplay (nohup) --[this daemon]--> traces autosave (bg)
                                               --> t4_real25.py (nohup)
                                               watchdog_eval_v3.py (separate)

v0.4.6 lesson: the first run lost 3.3GB of collected traces because they
were never pushed. Now the daemon fires a background traces-upload the
moment collection completes, BEFORE launching the orchestrator.
"""
import os
import subprocess
import time

LOG = "/content/collect_v3.log"
REPO = "/content/qwenjev"
WORK = "/content/qwenjev_work_v3"
HF_REPO = "tchbcb/qwenjev"

# background one-shot: tar traces (read-only) then upload with retries
TRACES_UPLOADER = f"""
import tarfile, time
from pathlib import Path
from huggingface_hub import HfApi
tok = (Path({WORK!r}) / ".hf_token").read_text().strip()
tgz = "/tmp/traces_real2.tgz"
with tarfile.open(tgz, "w:gz") as t:
    t.add({REPO!r} + "/data/traces_real2", arcname="traces_real2")
print("traces tgz MB:", Path(tgz).stat().st_size // 2**20, flush=True)
for i in range(5):
    try:
        HfApi(token=tok).upload_file(
            path_or_fileobj=tgz, path_in_repo="t4_assets_v3/data/traces_real2.tgz",
            repo_id={HF_REPO!r}, repo_type="model",
            commit_message="t4_assets_v3: remote traces autosave")
        print("traces pushed to HF", flush=True)
        break
    except Exception as e:
        print("traces upload retry", i, e, flush=True)
        time.sleep(60)
"""

print("[chain] waiting for collection to finish...", flush=True)
while True:
    try:
        txt = open(LOG, errors="replace").read()
        if "done ->" in txt:
            print("[chain] collection done -> traces autosave + orchestrator",
                  flush=True)
            subprocess.Popen(
                ["python3", "-u", "-B", "-c", TRACES_UPLOADER],
                stdout=open("/content/traces_upload.log", "w"),
                stderr=subprocess.STDOUT,
                start_new_session=True)          # detached, dies with VM only
            subprocess.Popen(
                ["bash", "-c",
                 f"cd {REPO} && nohup python3 -u -B ops/t4_real25.py "
                 "> /content/t4_v3.log 2>&1 < /dev/null &"],
                stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
                start_new_session=True)
            print("[chain] orchestrator launched, chain daemon exits",
                  flush=True)
            break
        print("[chain] still collecting...", flush=True)
    except FileNotFoundError:
        print("[chain] no collect log yet", flush=True)
    except Exception as exc:
        print(f"[chain] err: {exc}", flush=True)
    time.sleep(180)