feat(data): Ctx 加 token 预算断路——BudgetStop/tokens_used/add_tokens, stop_now 合并检查, rc=3 checkpoint(P4-3 F4) [nas] [no-doc]
This commit is contained in:
@@ -122,6 +122,11 @@ class WallClockStop(Exception):
|
||||
"""--until 墙钟到: 完成当前 unit 后抛出, 顶层收拾落盘并 rc=3。"""
|
||||
|
||||
|
||||
class BudgetStop(Exception):
|
||||
"""token/调用预算触顶(完成当前 unit 后抛, checkpoint 语义同 WallClockStop;
|
||||
P4-3 F4 漏斗夜批专用, 既有 lane 不设 token_budget 零感知)。"""
|
||||
|
||||
|
||||
# ---------- 单实例锁(flock, 跨容器) ----------
|
||||
|
||||
def acquire_lock():
|
||||
@@ -719,9 +724,12 @@ def _save_missing(path_key, ids):
|
||||
# ---------- 运行上下文 ----------
|
||||
|
||||
class Ctx:
|
||||
def __init__(self, lane, until=None, limit=None):
|
||||
def __init__(self, lane, until=None, limit=None, token_budget=None):
|
||||
self.lane = lane
|
||||
self.limit = limit
|
||||
self.token_budget = token_budget # P4-3 F4: LLM 夜批 token 上限
|
||||
self.tokens_used = 0
|
||||
self.budget_stopped = False
|
||||
self.deadline = None
|
||||
if until:
|
||||
hh, mm = until.split(":")
|
||||
@@ -774,18 +782,26 @@ class Ctx:
|
||||
mark_done(lane, stage, unit)
|
||||
self.pending_marks.clear()
|
||||
|
||||
def add_tokens(self, prompt_tokens, completion_tokens):
|
||||
"""LLM token 计量(P4-3 F4): 每次调用后累加, 预算判定在 stop_now。"""
|
||||
self.tokens_used += int(prompt_tokens or 0) + int(completion_tokens or 0)
|
||||
|
||||
def stop_now(self):
|
||||
"""完成当前 unit 后调用: 过墙钟 → 抛 WallClockStop(checkpoint 语义)。"""
|
||||
"""完成当前 unit 后调用: 过墙钟→WallClockStop; token 预算触顶→
|
||||
BudgetStop(均 checkpoint 语义, rc=3 次夜续)。"""
|
||||
if self.deadline and dt.datetime.now() >= self.deadline:
|
||||
self.wallclock = True
|
||||
raise WallClockStop()
|
||||
if self.token_budget is not None and self.tokens_used >= self.token_budget:
|
||||
self.budget_stopped = True
|
||||
raise BudgetStop()
|
||||
|
||||
def rc(self):
|
||||
if self.failed or self.hard_cool:
|
||||
return 1
|
||||
if self.rate_limited:
|
||||
return 2
|
||||
if self.wallclock:
|
||||
if self.wallclock or self.budget_stopped:
|
||||
return 3
|
||||
return 0
|
||||
|
||||
|
||||
@@ -1245,3 +1245,28 @@ def test_fetch_em_kuaixun_stops_at_boundary():
|
||||
rows2 = cd.fetch_em_kuaixun(cl2, until)
|
||||
assert len(rows2) == 1
|
||||
assert cl2.calls == 1
|
||||
|
||||
|
||||
# ---------- Ctx token 预算断路(P4-3 F4, 2026-10-02) ----------
|
||||
|
||||
def test_ctx_token_budget_stop():
|
||||
ctx = cd.Ctx("funnel", token_budget=100)
|
||||
ctx.add_tokens(60, 50) # 110 >= 100
|
||||
with pytest.raises(cd.BudgetStop):
|
||||
ctx.stop_now()
|
||||
assert ctx.rc() == 3 # 预算停摆=checkpoint 让路
|
||||
|
||||
|
||||
def test_ctx_no_budget_no_behavior_change():
|
||||
ctx = cd.Ctx("daily") # 默认 token_budget=None
|
||||
ctx.add_tokens(10**9, 10**9)
|
||||
ctx.stop_now() # 不抛(墙钟未设)
|
||||
assert ctx.rc() == 0
|
||||
|
||||
|
||||
def test_ctx_budget_under_threshold_passes():
|
||||
ctx = cd.Ctx("funnel", token_budget=1000)
|
||||
ctx.add_tokens(100, 200)
|
||||
ctx.stop_now() # 300 < 1000 不抛
|
||||
assert ctx.tokens_used == 300
|
||||
assert ctx.budget_stopped is False
|
||||
|
||||
Reference in New Issue
Block a user