#!/usr/bin/env python3
import argparse, ast, hashlib, json, os, re, subprocess, urllib.request, urllib.error
from datetime import datetime
from pathlib import Path
API="http://127.0.0.1:8000/v1/chat/completions"
MODELS="http://127.0.0.1:8000/v1/models"
MODEL="deepseek70b-7816"
ROOT=Path("/home/harness_user_1/dgx_ai_factory")
OUTROOT=ROOT/"data/generated/deepseek70b_7816_code_sft_2000_structured_v6"
LOGROOT=ROOT/"logs/data2000_structured_hotfix_v6"
CATEGORIES=[
("python_algorithm","advanced","python"),
("python_data_pipeline","advanced","python"),
("python_debugging","advanced","python"),
("python_concurrency","advanced","python"),
("fastapi_backend","advanced","python"),
("linux_bash","advanced","bash"),
("sql_database","advanced","sql"),
("testing_quality","advanced","python"),
("performance_optimization","advanced","python"),
("ml_tooling","advanced","python"),
("system_automation","advanced","bash"),
("api_integration","advanced","python"),
]
ANGLES=[
"실행 가능한 완성 예제와 최소한의 검증 코드를 포함하세요.",
"실무에서 발생할 수 있는 실패 조건과 예외 처리를 포함하세요.",
"메모리 사용량과 확장성을 고려한 구현을 요구하세요.",
"재현 가능한 버그 상황과 수정 전후 차이를 포함하세요.",
"테스트 가능한 함수 단위 설계와 테스트 예제를 포함하세요.",
"입력 검증과 명확한 오류 메시지를 포함하세요.",
"동시성 또는 재시도 상황에서 안전하게 동작하도록 요구하세요.",
"대용량 입력을 가정하고 스트리밍 또는 배치 처리를 고려하세요.",
]
SCHEMA={
"type":"object",
"properties":{
"instruction":{"type":"string","minLength":40},
"response":{"type":"string","minLength":220},
"category":{"type":"string"},
"difficulty":{"type":"string"},
"language":{"type":"string"}
},
"required":["instruction","response","category","difficulty","language"],
"additionalProperties":False
}
THINK_BLOCK=re.compile(r".*?",re.I|re.S)
THINK_TAG=re.compile(r"?think>",re.I)
CJK=re.compile(r"[\u4e00-\u9fff]")
BAD=re.compile(r"\bTODO\b|\bFIXME\b|\bTBD\b|placeholder|생략|중략|omitted|truncated|continue here",re.I)
FENCE=re.compile(r"```([A-Za-z0-9_+\-]*)\s*\n(.*?)```",re.S)
DANGEROUS=re.compile(r"rm\s+-rf\s+/(?:\s|$)|mkfs\.|dd\s+if=.*of=/dev/|:\(\)\s*\{\s*:\|:&\s*\};:",re.I)
def now(): return datetime.now().astimezone().isoformat()
def fsync_append(path,obj):
path.parent.mkdir(parents=True,exist_ok=True)
with path.open("a",encoding="utf-8") as f:
f.write(json.dumps(obj,ensure_ascii=False,separators=(",",":"))+"\n")
f.flush(); os.fsync(f.fileno())
def atomic_json(path,obj):
path.parent.mkdir(parents=True,exist_ok=True)
tmp=path.with_suffix(path.suffix+".tmp")
with tmp.open("w",encoding="utf-8") as f:
json.dump(obj,f,ensure_ascii=False,indent=2); f.flush(); os.fsync(f.fileno())
os.replace(tmp,path)
def normalize_text(s):
if not isinstance(s,str): return s
s=s.replace("Ġ"," ").replace("Ċ","\n").replace("ĉ","\t")
s=THINK_BLOCK.sub("",s)
s=THINK_TAG.sub("",s)
return s.replace("\r\n","\n").replace("\r","\n").strip()
def parse_content(content):
s=normalize_text(content or "")
try: return json.loads(s),None
except Exception as e:
a=s.find("{"); b=s.rfind("}")
if a>=0 and b>a:
try:return json.loads(s[a:b+1]),"extracted_json"
except Exception:pass
return None,f"json_parse:{type(e).__name__}"
def validate(obj,expected):
reasons=[]
if not isinstance(obj,dict): return ["not_object"],None
obj={k:normalize_text(v) if isinstance(v,str) else v for k,v in obj.items()}
for k in ("instruction","response","category","difficulty","language"):
if not isinstance(obj.get(k),str) or not obj[k].strip(): reasons.append(f"missing_or_empty:{k}")
if reasons: return sorted(set(reasons)),obj
if len(obj["instruction"])<40: reasons.append("instruction_too_short")
if len(obj["response"])<220: reasons.append("response_too_short")
if len(obj["response"])>14000: reasons.append("response_too_long")
if obj["category"]!=expected[0]: reasons.append("category_mismatch")
if obj["difficulty"]!=expected[1]: reasons.append("difficulty_mismatch")
if obj["language"]!=expected[2]: reasons.append("language_mismatch")
alltext=obj["instruction"]+"\n"+obj["response"]
if THINK_TAG.search(alltext) or "" in alltext.lower(): reasons.append("think_tag")
if CJK.search(alltext): reasons.append("han_cjk_char")
if BAD.search(alltext): reasons.append("todo_or_incomplete")
if DANGEROUS.search(alltext): reasons.append("dangerous_destructive_pattern")
if obj["response"].count("```")%2: reasons.append("unbalanced_fence")
for lang,code in FENCE.findall(obj["response"]):
if lang.lower() in ("python","py"):
try: ast.parse(code)
except SyntaxError as e: reasons.append("python_syntax_error:"+e.msg)
elif lang.lower()=="json":
try: json.loads(code)
except Exception: reasons.append("json_fence_parse_error")
if re.search(r"json\.dumps\(|Define the JSON data|Convert the dictionary to a JSON string",obj["response"],re.I):
reasons.append("meta_json_generation_instead_of_solution")
if re.search(r'"instruction"\s*:\s*""|"response"\s*:\s*"\+\+"',obj["response"]):
reasons.append("degenerate_embedded_schema")
return sorted(set(reasons)),obj
def http_json(url,payload=None,timeout=900):
data=None; headers={}
if payload is not None:
data=json.dumps(payload,ensure_ascii=False).encode(); headers["Content-Type"]="application/json"
req=urllib.request.Request(url,data=data,headers=headers)
try:
with urllib.request.urlopen(req,timeout=timeout) as r:
return r.status,json.loads(r.read().decode("utf-8","replace"))
except urllib.error.HTTPError as e:
raw=e.read().decode("utf-8","replace")
try: body=json.loads(raw)
except Exception: body={"raw":raw}
return e.code,body
except Exception as e:
return 0,{"error":repr(e)}
def structured_payload(category,difficulty,language,variation,mode="json_schema"):
schema=json.loads(json.dumps(SCHEMA))
schema["properties"]["category"]["enum"]=[category]
schema["properties"]["difficulty"]["enum"]=[difficulty]
schema["properties"]["language"]["enum"]=[language]
system=(
"당신은 코딩 개발 비서를 위한 고품질 SFT 데이터 생성기입니다. 일반 대화나 잡담을 만들지 마십시오. "
"하나의 독립적인 실무형 코딩 문제와 그에 대한 완성된 최종 답변을 만드십시오. "
"instruction은 실제 사용자의 구체적인 개발 요청이어야 하며 빈 문자열이면 안 됩니다. "
"response는 바로 학습에 사용할 최종 Assistant 답변입니다. 필요한 경우 실행 가능한 코드, 오류 처리, 테스트를 포함하십시오. "
" 또는 내부 추론을 출력하지 마십시오. TODO/FIXME/생략/placeholder를 사용하지 마십시오. "
"JSON을 만드는 Python 예제를 답으로 쓰지 마십시오. 지정된 JSON 객체 그 자체만 반환하십시오. "
"반드시 instruction, response, category, difficulty, language 다섯 필드만 반환하십시오."
)
user=(f"새로운 코딩 SFT 샘플 1개를 작성하세요.\ncategory={category}\ndifficulty={difficulty}\nlanguage={language}\n"
f"variation_id=V6-{variation:06d}\n추가 조건: {ANGLES[variation%len(ANGLES)]}\n"
"이전 샘플을 복사하지 말고 구체적인 문제 상황, 입력/출력 또는 실패 조건을 포함해 서로 다른 과제로 만드세요.")
p={"model":MODEL,"messages":[{"role":"system","content":system},{"role":"user","content":user}],
"temperature":0.35,"top_p":0.9,"max_tokens":900}
if mode=="json_schema":
p["response_format"]={"type":"json_schema","json_schema":{"name":"coding_sft_record","schema":schema}}
elif mode=="structured_outputs":
p["structured_outputs"]={"json":schema}
elif mode=="json_object":
p["response_format"]={"type":"json_object"}
else: node-7.example.invalid ValueError(mode)
return p
def api_probe():
code,body=http_json(MODELS,timeout=10)
if code!=200:return None,{"models_http":code,"body":body}
ids=[x.get("id") for x in body.get("data",[])]
if MODEL not in ids:return None,{"models_http":code,"ids":ids,"error":"LoRA model id missing"}
results=[]
for mode in ("json_schema","structured_outputs","json_object"):
code,b=http_json(API,structured_payload("python_algorithm","advanced","python",900001,mode),timeout=900)
rec={"mode":mode,"http":code}
if code==200:
try: content=b["choices"][0]["message"]["content"]
except Exception:
rec["error"]="missing_content"; results.append(rec); continue
obj,perr=parse_content(content)
vr,norm=validate(obj,("python_algorithm","advanced","python")) if obj is not None else ([perr],None)
rec["parse_error"]=perr; rec["validation"]=vr; rec["preview"]=normalize_text(content)[:500]
results.append(rec)
if not vr:return mode,{"selected":mode,"results":results}
else:
rec["body"]=b;results.append(rec)
return None,{"selected":None,"results":results}
def mem_available_gib_local():
vals={}
for line in Path("/proc/meminfo").read_text().splitlines():
if ":" in line:
k,v=line.split(":",1);p=v.split()
if p and p[0].isdigit(): vals[k]=int(p[0])*1024
return vals.get("MemAvailable",0)/(1024**3)
def mem_available_gib_sub():
try:
p=subprocess.run(["ssh","-o","BatchMode=yes","-o","ConnectTimeout=5","harness_user_2@192.0.2.13",
"awk '/MemAvailable:/ {print $2}' /proc/meminfo"],capture_output=True,text=True,timeout=10)
if p.returncode:return None
return float(p.stdout.strip())/1024/1024
except Exception:return None
def load_hashes(path):
out=set()
if not path.is_file():return out
for line in path.read_text(errors="replace").splitlines():
try:o=json.loads(line)
except Exception:continue
ins=normalize_text(o.get("instruction","")).lower()
if ins:out.add(hashlib.sha256(ins.encode()).hexdigest())
return out
def selftest():
good={"instruction":"Python에서 대용량 JSONL 파일을 스트리밍 방식으로 중복 제거하고 결과를 검증하는 도구를 작성하세요.",
"response":"전체 파일을 한꺼번에 읽지 않고 줄 단위로 처리할 수 있습니다.\n```python\nfrom pathlib import Path\n\ndef unique_count(path: Path) -> int:\n seen=set()\n with path.open(encoding='utf-8') as f:\n for line in f:\n seen.add(line.rstrip('\\n'))\n return len(seen)\n\nassert callable(unique_count)\n```\n실무에서는 JSON 파싱 오류를 별도 quarantine 파일로 기록하고, 대용량 데이터에서는 외부 저장소나 해시 파티셔닝으로 seen 집합의 메모리 사용량을 제한하는 방식으로 확장할 수 있습니다.",
"category":"python_data_pipeline","difficulty":"advanced","language":"python"}
vr,_=validate(good,("python_data_pipeline","advanced","python")); assert not vr,vr
bad=dict(good);bad["response"]="x TODO"
vr,_=validate(bad,("python_data_pipeline","advanced","python")); assert "todo_or_incomplete" in vr
assert normalize_text("AĠBĊC")=="A B\nC"
print("SELFTEST=PASS")
def main():
ap=argparse.ArgumentParser()
ap.add_argument("--target",type=int,default=2000)
ap.add_argument("--max-attempts",type=int,default=5000)
ap.add_argument("--mode",choices=["json_schema","structured_outputs","json_object"])
ap.add_argument("--self-test",action="store_true")
a=ap.parse_args()
if a.self_test:selftest();return 0
if not a.mode:node-7.example.invalid SystemExit("--mode required")
OUTROOT.mkdir(parents=True,exist_ok=True);LOGROOT.mkdir(parents=True,exist_ok=True)
accepted=OUTROOT/"accepted.jsonl"; rejected=OUTROOT/"rejected.jsonl"; statep=OUTROOT/"state.json"
state={"phase":"generation","accepted":0,"target":a.target,"attempts":0,"rejected":0,"api_failures":0,
"mode":a.mode,"started_at":now(),"updated_at":now()}
if statep.is_file():
try:
old=json.loads(statep.read_text())
if old.get("mode")==a.mode:state.update(old);state["phase"]="generation";state["target"]=a.target
except Exception:pass
hashes=load_hashes(accepted);state["accepted"]=len(hashes);atomic_json(statep,state)
while state["accepted"]=a.target else "max_attempts";state["updated_at"]=now();atomic_json(statep,state)
report={"accepted":state["accepted"],"target":a.target,"attempts":state["attempts"],"rejected":state["rejected"],
"api_failures":state["api_failures"],"mode":a.mode,"accepted_file":str(accepted),
"sha256":hashlib.sha256(accepted.read_bytes()).hexdigest() if accepted.is_file() else None,"finished_at":now()}
atomic_json(OUTROOT/"final_report.json",report);print(json.dumps(report,ensure_ascii=False,indent=2))
return 0 if state["accepted"]>=a.target else 2
if __name__=="__main__":node-7.example.invalid SystemExit(main())