Skip to content

Commit 3adfd99

Browse files
committed
ワーカーを、時計を待たずに起こせるようにする
枠が明いているうちに回しておきたい、が普通に起きる。次の起動まで待つと、 待っているあいだに他の依頼が枠を食う(実測で、外からの 1 回が 5 時間枠を 48 ポイント持っていった)。 押せば待ち行列の先頭が 1 本流れる。起こしたことになるので、次の起動は押した時刻から 間隔ぶん後になる。 断る理由は書き分ける。行列が空なのか枠が詰まっているのかで、次にすることが逆になる (積むのを待つ / 窓が明くのを待つ)——「起きませんでした」だけでは、押した人に どちらなのか読めない。 時計から回す道と同じ手順を通す(_flush_one)。書き写すと、片方だけ直したときに 拾い方や塊の扱いが食い違う。 Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
1 parent c236f90 commit 3adfd99

4 files changed

Lines changed: 181 additions & 26 deletions

File tree

CLAUDE.md

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -466,6 +466,12 @@ GeoNames 全世界地名辞典 = `geonames`(いずれも 348 言語版・195 か
466466
人に直させる価値が無い)
467467
- **収集を消したら行列からも外す**(`collect.remove`)。消えた収集を抱えた
468468
ままだと、そのワーカーは起こそうとして 404 を踏み続ける
469+
- **時計を待たずに起こせる**(`main.wake_worker` / 画面の「今すぐ起こす」)。
470+
**枠が明いているうちに回しておきたい、が普通に起きる** —— 次の起動まで待つと、
471+
待っているあいだに他の依頼が枠を食う(実測で、外からの 1 回が 5 時間枠を
472+
48 ポイント持っていった)。押せば行列の先頭が 1 本流れ、**そこから間隔を
473+
数え直す**(起こしたことになる)。**断る理由は書き分ける** —— 行列が空なのか
474+
枠が詰まっているのかで、次にすることが逆になる(積むのを待つ / 窓が明くのを待つ)
469475
- **節に「前回起きた・次に起きる・待ち行列」を出す**(`ai_workers._queue_html`)。
470476
**積まれているのに動かないなら**枠が詰まっているか起動待ち、**積まれていない
471477
なら**巡回の側がまだ積んでいない —— どちらなのかは、時刻と行列を並べないと

app/main.py

Lines changed: 63 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -344,34 +344,71 @@ def _run_one_from_a_worker() -> bool:
344344
defined = workers.load()
345345
except ValueError:
346346
return False
347-
for worker in defined:
348-
batch = workers.queued(worker.name)
347+
# **1 周で流すのは 1 本。** 取り込みは同時に 1 ジョブしか受けないので、
348+
# 先に流せたワーカーで打ち切る(残りは次の周)
349+
return any(_flush_one(worker, now) for worker in defined)
350+
351+
352+
def _flush_one(worker, now: datetime, wake: bool = False) -> bool:
353+
"""そのワーカーから 1 本流す。流したら True。
354+
355+
`wake` は**時計を待たずに起こす**(画面の「今すぐ起こす」)。枠が明いている
356+
うちに回しておきたい、が普通に起きる —— 次の起動まで待つと、待っているあいだに
357+
誰かが枠を食う。**起こした時刻は普通に控える**ので、そこから間隔を数え直す。
358+
"""
359+
batch = workers.queued(worker.name)
360+
if not batch:
361+
return False
362+
# **いま流している塊が先。** 無ければ、起動の時刻が来ていれば拾う
363+
if not workers.claim_ready(worker.name):
364+
if not wake and not _worker_due(worker, now):
365+
return False
366+
batch = workers.claim(worker.name, worker.per_run, _iso(now))
349367
if not batch:
350-
continue
351-
# **いま流している塊が先。** 無ければ、起動の時刻が来ていれば拾う
352-
if not workers.claim_ready(worker.name):
353-
if not _worker_due(worker, now):
354-
continue
355-
batch = workers.claim(worker.name, worker.per_run, _iso(now))
356-
if not batch:
357-
continue
358-
entry = batch[0]
359-
step = workers.pick(worker)
360-
if step is None:
361-
# **どれも詰まっていたら流さない。** 塊はそのまま残るので、
362-
# 窓が明けた周で続きから流れる
363-
continue
364-
try:
365-
start_collection_bake(entry["collection"], entry["sweep"])
366-
except HTTPException as e:
367-
# 取り込みが混んでいる・巡回が消えた。**塊からは外す** ——
368-
# 消えた巡回を抱えたままだと、そのワーカーが二度と進まない
369-
if e.status_code == 404:
370-
workers.done(worker.name, entry["collection"], entry["sweep"])
371368
return False
372-
workers.done(worker.name, entry["collection"], entry["sweep"])
373-
return True
374-
return False
369+
entry = batch[0]
370+
if workers.pick(worker) is None:
371+
# **どれも詰まっていたら流さない。** 塊はそのまま残るので、
372+
# 窓が明けた周で続きから流れる
373+
return False
374+
try:
375+
start_collection_bake(entry["collection"], entry["sweep"])
376+
except HTTPException as e:
377+
# 取り込みが混んでいる・巡回が消えた。**塊からは外す** ——
378+
# 消えた巡回を抱えたままだと、そのワーカーが二度と進まない
379+
if e.status_code == 404:
380+
workers.done(worker.name, entry["collection"], entry["sweep"])
381+
return False
382+
workers.done(worker.name, entry["collection"], entry["sweep"])
383+
return True
384+
385+
386+
def wake_worker(name: str) -> dict:
387+
"""ワーカーを**時計を待たずに起こす**(画面の「今すぐ起こす」)。
388+
389+
**断る理由は書き分ける。** 押しても何も起きないときに「起きませんでした」
390+
だけだと、行列が空なのか枠が詰まっているのかが押した人に読めない ——
391+
どちらなのかで次にすることが逆になる(積むのを待つ / 窓が明くのを待つ)。
392+
"""
393+
try:
394+
worker = workers.get(name)
395+
except ValueError as e:
396+
raise HTTPException(409, {"error": str(e)}) from None
397+
if worker is None:
398+
raise HTTPException(404, {"error": f"ワーカー「{name}」がありません"})
399+
if not workers.queued(name):
400+
raise HTTPException(409, {
401+
"error": f"ワーカー「{name}」の待ち行列は空です",
402+
"hint": "巡回の側が自分を積むまで、起こしても流すものがありません",
403+
})
404+
if workers.pick(worker) is None:
405+
raise HTTPException(429, _all_full(name))
406+
if not _flush_one(worker, datetime.now(UTC), wake=True):
407+
raise HTTPException(409, {
408+
"error": f"ワーカー「{name}」から流せませんでした",
409+
"hint": "取り込みが走っている最中かもしれません(少し置いてからもう一度)",
410+
})
411+
return {"ok": True, "worker": name}
375412

376413

377414
def _worker_due(worker, now: datetime) -> bool:

app/views/ai_workers.py

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -144,6 +144,7 @@ def _queue_html(worker: workers.Worker) -> str:
144144
<thead><tr><th>待ち行列(先に積まれた順)</th><th></th></tr></thead>
145145
<tbody>{rows}</tbody>
146146
</table>
147+
{_wake_form(worker, bool(waiting))}
147148
<p class="muted">
148149
<strong>次に起きる時刻は、前回「起きた」時刻から数えます</strong>(流し終えた時刻では
149150
ない)—— 塊を流し切るのに何周かかっても、次の起動は最初の起動から間隔ぶん後になる。<br>
@@ -153,6 +154,26 @@ def _queue_html(worker: workers.Worker) -> str:
153154
"""
154155

155156

157+
def _wake_form(worker: workers.Worker, waiting: bool) -> str:
158+
"""時計を待たずに 1 本流す口。**待っているものが無ければ出さない**。
159+
160+
**枠が明いているうちに回しておきたい、が普通に起きる** —— 次の起動まで待つと、
161+
待っているあいだに他の依頼が枠を食う(実測で、外からの 1 回が 5 時間枠を
162+
48 ポイント持っていった)。押せば行列の先頭が 1 本流れ、**そこから間隔を
163+
数え直す**(起こしたことになるので、次の起動は押した時刻からずれる)。
164+
"""
165+
if not waiting:
166+
return ('<p class="muted">待っているものが無いので、起こしても流すものが'
167+
"ありません。</p>")
168+
return (
169+
f'<form class="init-form" method="post" action="/admin/ai/workers/wake">'
170+
f'<input type="hidden" name="worker_name" value="{esc(worker.name)}">'
171+
f'<button type="submit"'
172+
f' title="間隔を待たずに、待ち行列の先頭を 1 本流します">今すぐ起こす</button>'
173+
"</form>"
174+
)
175+
176+
156177
def section_html(selects) -> str:
157178
"""節ぜんたい。`selects` は相手・モデル・考える量のセレクトを作る 3 つ。
158179
@@ -196,6 +217,23 @@ def section_html(selects) -> str:
196217
"""
197218

198219

220+
@router.post("/admin/ai/workers/wake")
221+
async def wake_worker(request: Request):
222+
"""ワーカーを**時計を待たずに起こす**(行列の先頭を 1 本流す)。
223+
224+
枠が明いているうちに回しておきたい、が普通に起きる —— 次の起動まで待つと、
225+
待っているあいだに他の依頼が枠を食う。
226+
227+
**断る理由は書き分ける**(`main.wake_worker`)。行列が空なのか枠が詰まって
228+
いるのかで、次にすることが逆になる(積むのを待つ / 窓が明くのを待つ)。
229+
"""
230+
from app.main import wake_worker as wake
231+
232+
form = await request.form()
233+
wake(str(form.get("worker_name") or "").strip())
234+
return RedirectResponse(BACK_TO_SECTION, status_code=303)
235+
236+
199237
@router.post("/admin/ai/workers")
200238
async def save_worker(request: Request):
201239
"""ワーカー 1 つぶんを保存する。**その 1 つだけ**を書き換える。

tests/test_workers.py

Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -496,6 +496,80 @@ async def broken(_asked, _messages):
496496
assert not workers.full_until("antigravity")
497497

498498

499+
class TestWakingItByHand:
500+
"""時計を待たずに 1 本流す(`main.wake_worker` / 画面の「今すぐ起こす」)。
501+
502+
**枠が明いているうちに回しておきたい、が普通に起きる** —— 次の起動まで待つと、
503+
待っているあいだに他の依頼が枠を食う(実測で、外からの 1 回が 5 時間枠を
504+
48 ポイント持っていった)。
505+
"""
506+
507+
@pytest.fixture
508+
def baking(self, enabled, monkeypatch):
509+
from app import main
510+
511+
started: list[tuple] = []
512+
monkeypatch.setattr(
513+
main, "start_collection_bake",
514+
lambda name, sweep=None: started.append((name, sweep)),
515+
)
516+
return main, started
517+
518+
def test_it_runs_without_waiting_for_the_clock(self, baking):
519+
"""**間隔が来ていなくても流す。** それが押した意味。"""
520+
main, started = baking
521+
workers.save([workers.Worker("精査", (workers.Step("codex"),), interval_minutes=600)])
522+
_quota("codex", 10.0)
523+
workers.enqueue("精査", "news", "ざっと", "2026-01-01T00:00:00+00:00")
524+
workers.claim("精査", 1, "2099-01-01T00:00:00+00:00")
525+
workers.done("精査", "news", "ざっと")
526+
workers.enqueue("精査", "news", "整理", "2026-01-01T00:00:00+00:00")
527+
528+
# 時計では起きない(前回起きたのが未来の時刻になっている)
529+
assert not main._run_one_from_a_worker()
530+
531+
main.wake_worker("精査")
532+
533+
assert started == [("news", "整理")]
534+
535+
def test_an_empty_queue_says_so(self, baking):
536+
"""**理由を書き分ける。** 行列が空なのか枠が詰まっているのかで、
537+
次にすることが逆になる(積むのを待つ / 窓が明くのを待つ)。
538+
"""
539+
import fastapi
540+
541+
main, _started = baking
542+
workers.save([workers.Worker("精査", (workers.Step("codex"),))])
543+
_quota("codex", 10.0)
544+
545+
with pytest.raises(fastapi.HTTPException) as got:
546+
main.wake_worker("精査")
547+
548+
assert got.value.status_code == 409
549+
assert "待ち行列は空" in got.value.detail["error"]
550+
551+
def test_all_crowded_says_so(self, baking):
552+
import fastapi
553+
554+
main, _started = baking
555+
workers.save([workers.Worker("精査", (workers.Step("codex"),))])
556+
_quota("codex", 95.0)
557+
workers.enqueue("精査", "news", "ざっと", "2026-01-01T00:00:00+00:00")
558+
559+
with pytest.raises(fastapi.HTTPException) as got:
560+
main.wake_worker("精査")
561+
562+
assert got.value.status_code == 429
563+
564+
def test_an_unknown_worker_is_404(self, baking):
565+
import fastapi
566+
567+
main, _started = baking
568+
with pytest.raises(fastapi.HTTPException) as got:
569+
main.wake_worker("いない")
570+
assert got.value.status_code == 404
571+
572+
499573
class TestWhoActuallyRan:
500574
"""控えに残すのは、決めた相手ではなく**頼んだ相手**。
501575

0 commit comments

Comments
 (0)