"""Lanes (by task size): config records, membership, pick order, output blocks. Sessions are dicts {"lane", "model", "socket", "alive", …} keyed by lane name or "all" (registry I/O lives in wf.py).""" import re from dataclasses import dataclass from . import tasks LANE_KEYS = {"efforts", "order", "fallback", "slices"} ORDERS = ("unblock", "priority") LANE_NAME_RE = re.compile(r"^[a-z][a-z0-9-]*$") @dataclass(frozen=True) class Lane: name: str efforts: tuple[str, ...] order: str = "priority" fallback: str | None = None slices: bool = False DEFAULT_LANES = (Lane("fast", ("<1h",), "unblock", None, True), Lane("slow", ("1h",), "priority", "fast", False)) def parse_lanes(table: dict | None, slice_above: str) -> tuple[Lane, ...]: """[lanes.*] of workflow.toml → lanes in file order; None → DEFAULT_LANES. Raises ValueError.""" if slice_above not in tasks.EFFORTS: raise ValueError(f"slice_above '{slice_above}' (want {', '.join(tasks.EFFORTS)})") allowed = tasks.EFFORTS[:tasks.EFFORTS.index(slice_above) + 1] if table is None: out = DEFAULT_LANES else: if not isinstance(table, dict): raise ValueError("'lanes' must be a table of tables") out = [] for name, t in table.items(): if not LANE_NAME_RE.match(name): raise ValueError(f"bad lane name '{name}' (want a-z 0-9 -)") if name == "all": raise ValueError("lane name 'all' is reserved") if not isinstance(t, dict): raise ValueError(f"'lanes.{name}' must be a table") for k in t: if k not in LANE_KEYS: raise ValueError(f"unknown key 'lanes.{name}.{k}'") efforts = t.get("efforts") if not isinstance(efforts, list) or not efforts or not all(isinstance(e, str) for e in efforts): raise ValueError(f"lane {name}: 'efforts' must be a non-empty list of strings") order = t.get("order", "priority") if order not in ORDERS: raise ValueError(f"lane {name}: order '{order}' (want {', '.join(ORDERS)})") fallback, slices = t.get("fallback"), t.get("slices", False) if fallback is not None and not isinstance(fallback, str) or not isinstance(slices, bool): raise ValueError(f"lane {name}: 'fallback' must be a string, 'slices' true/false") out.append(Lane(name, tuple(efforts), order, fallback, slices)) names = [l.name for l in out] owner: dict[str, str] = {} for l in out: for e in l.efforts: if e not in tasks.EFFORTS: raise ValueError(f"lane {l.name}: effort '{e}' (want {', '.join(tasks.EFFORTS)})") if e not in allowed: raise ValueError(f"lane {l.name}: effort '{e}' above slice_above '{slice_above}'") if e in owner: raise ValueError(f"effort '{e}' in lanes {owner[e]} and {l.name}") owner[e] = l.name if l.fallback is not None and (l.fallback not in names or l.fallback == l.name): raise ValueError(f"lane {l.name}: fallback '{l.fallback}' is not another lane") for e in allowed: if e not in owner: raise ValueError(f"effort '{e}' (≤ slice_above) in no lane") by = {l.name: l for l in out} for l in out: if l.fallback and by[l.fallback].fallback == l.name: a, b = sorted((l.name, l.fallback)) raise ValueError(f"fallback cycle: {a} ↔ {b}") flagged = [l for l in out if l.slices] if len(flagged) > 1: raise ValueError("slices = true on more than one lane") if not flagged: out = [Lane(l.name, l.efforts, l.order, l.fallback, owner["<1h"] == l.name) for l in out] return tuple(out) def _pending(doc: tasks.Doc) -> list[tasks.Item]: return [i for s in doc.sections if s.key == "pending" for i in s.items if not i.error] def counts(doc: tasks.Doc, archived: set[str], lanes, slice_above: str, held: dict[str, str] | None = None, runner: bool = False) -> dict[str, tuple[int, int, int]]: """lane → (pickable, waiting, pickable slice jobs) over Pending, in lane order. held: claimed ids count waiting. runner: pickable = what `list --runner` shows (has a usable Done); the rest is not_ready().""" out = {l.name: [0, 0, 0] for l in lanes} for item in _pending(doc): c = out[lane_of(item, lanes, slice_above)] if _why(doc, item, archived, slice_above, held, 0, True) is None: if runner and not item.runner_ready: continue c[0] += 1 c[2] += is_slice_job(item, slice_above) else: c[1] += 1 return {name: tuple(c) for name, c in out.items()} def not_ready(doc: tasks.Doc, archived: set[str], lanes, slice_above: str, held: dict[str, str] | None = None) -> dict[str, int]: """lane → pickable items without a Done that a prep worker can fill (owner-bound ones not counted).""" out = {l.name: 0 for l in lanes} for item in _pending(doc): if _why(doc, item, archived, slice_above, held, 0, True) is None \ and not (item.done_text and item.done_text.strip()) and item.sessions != "owner": out[lane_of(item, lanes, slice_above)] += 1 return out def prep_targets(doc: tasks.Doc, archived: set[str], lanes, slice_above: str, only: set[str] | None = None) -> list[tasks.Item]: """Pending items a prep worker can make runner-ready: no Done, not owner-bound, no status, no open slices; only: lane names. Order: pickable first, then priority, then file order.""" out = [] for n, item in enumerate(_pending(doc)): if item.done_text and item.done_text.strip() or item.sessions == "owner" or item.status \ or tasks.open_slices(doc, item.id): continue if only is not None and lane_of(item, lanes, slice_above) not in only: continue waits = _why(doc, item, archived, slice_above, None, 0, True) is not None out.append(((waits, 9 if item.prio is None else item.prio, n), item)) return [item for _, item in sorted(out, key=lambda x: x[0])] def cross_waits(doc: tasks.Doc, archived: set[str], lane: str, lanes, slice_above: str) -> list[tuple[tasks.Item, tasks.Item]]: """(my waiting item, pickable item of another lane it waits on, transitively).""" items = _pending(doc) wait = tasks.waiters(doc) name = {i.id: lane_of(i, lanes, slice_above) for i in items} out = [] for mine in items: if name[mine.id] != lane or tasks.why_not(doc, mine, archived) is None: continue for other in items: if name[other.id] != lane and mine.id in wait.get(other.id, ()) \ and tasks.why_not(doc, other, archived) is None: out.append((mine, other)) return out def lanes_block(counts: dict[str, tuple[int, int, int]], sessions: dict[str, dict], me: str | None, nready: dict[str, int] | None = None) -> list[str]: """One line per lane (always all: lanes are few). nready: lane → not runner-ready count.""" out = [] for name, (p, w, sl) in counts.items(): line = (f"{name}{' (you)' if name == me else ''}: {p} pickable" + (f" ({sl} slice{'s' if sl > 1 else ''})" if sl else "") + f" · {w} waiting") if (nready or {}).get(name): line += f" · {nready[name]} not runner-ready (no Done)" if name != me: s = sessions.get(name) if s and s.get("alive"): line += f" · session uds:{s['socket']} (alive, {s.get('model') or '?'})" else: line += " · no session" + (f" → orchestrator or owner: start one (wf next --lane {name})" if p else "") out.append(line) return out def waiting_block(waits: list[tuple[tasks.Item, tasks.Item]], sessions: dict[str, dict], lanes, slice_above: str) -> list[str]: out = [] for mine, other in waits: name = lane_of(other, lanes, slice_above) s = sessions.get(name) head = f"- {mine.id} waits on {other.id} ({name} lane) → " if s and s.get("alive"): out.append(head + f'message uds:{s["socket"]}: "{other.id} blocks my {mine.id}, please take it"') else: out.append(head + f"no {name} session: tell the owner") return out def pickable_ids(doc: tasks.Doc, archived: set[str]) -> set[str]: return {i.id for i in _pending(doc) if tasks.why_not(doc, i, archived) is None} def notify_block(freed: list[tasks.Item], sessions: dict[str, dict], lanes, slice_above: str) -> list[str]: """One line per lane with newly pickable tasks, in lane order.""" out = [] for lane in lanes: ids = ", ".join(i.id for i in freed if lane_of(i, lanes, slice_above) == lane.name) if not ids: continue s = sessions.get(lane.name) if s and s.get("alive"): out.append(f"notify {lane.name} uds:{s['socket']}: now pickable {ids}") else: out.append(f"{lane.name} work now pickable: {ids} (no session: tell the owner)") return out def solo_done_block(solo: list[str], sessions: dict[str, dict], me_pid: str | None) -> list[str]: """After solo task(s) are done: one line per other live session, so it can pick again.""" if not solo: return [] ids = ", ".join(solo) return [f"notify {name} uds:{s['socket']}: solo {ids} done, run wf next" for name, s in sessions.items() if s.get("alive") and str(s.get("pid")) != me_pid] MODEL_RANK = {m: n for n, m in enumerate(tasks.MODELS)} def over(effort: str | None, slice_above: str) -> bool: return effort in tasks.EFFORTS and tasks.EFFORTS.index(effort) > tasks.EFFORTS.index(slice_above) def is_slice_job(item: tasks.Item, slice_above: str) -> bool: """Too big to implement and not sliced yet: the slices lane splits it.""" return over(item.effort, slice_above) and not item.slices def sliced_out(item: tasks.Item, slice_above: str) -> bool: """Too big and already sliced once: never picked again (wf done it or add slices by hand).""" return over(item.effort, slice_above) and bool(item.slices) def lane_of(item: tasks.Item, lanes: tuple[Lane, ...], slice_above: str) -> str: if over(item.effort, slice_above): return next(l.name for l in lanes if l.slices) effort = item.effort if item.effort in tasks.EFFORTS else "<1h" return next(l.name for l in lanes if effort in l.efforts) def _why(doc, item, archived, slice_above, held, others, owner) -> str | None: why = tasks.why_not(doc, item, archived, held, others, owner) if why is None and not item.error and sliced_out(item, slice_above): return f"sliced: all slices done → wf done {item.id} or add slices" return why def rank(doc: tasks.Doc, items: list[tasks.Item], order: str, lane: dict[str, str]) -> dict[str, tuple]: """Sort key per item. lane: item id → lane name (for 'waited on by another lane').""" wait = tasks.waiters(doc) by_id = {i.id: i for s in doc.sections if s.key in ("pending", "human") for i in s.items if not i.error} out = {} for n, item in enumerate(items): ws = [by_id[w] for w in wait.get(item.id, ())] prio = min([p for p in [item.prio, *(w.prio for w in ws)] if p is not None], default=9) cross = any(lane.get(w.id) != lane.get(item.id) for w in ws) out[item.id] = ((not ws, -len(ws), prio, n) if order == "unblock" else (prio, not cross, -len(ws), n)) return out def _candidates(doc, lanes, slice_above, lane: str | None, model: str | None) -> list[tasks.Item]: items = doc.section("pending").items return [i for i in items if i.error or ( (lane is None or lane_of(i, lanes, slice_above) == lane) and (model is None or MODEL_RANK[i.model] <= MODEL_RANK[model]))] def _ordered(doc, lanes, slice_above, lane: str | None, model: str | None) -> list[tasks.Item]: items = _candidates(doc, lanes, slice_above, lane, model) order = next((l.order for l in lanes if l.name == lane), "priority") names = {i.id: lane_of(i, lanes, slice_above) for i in doc.section("pending").items if not i.error} keys = rank(doc, items, order, names) return sorted(items, key=lambda i: keys.get(i.id, (9, True, 0, items.index(i)))) def pick(doc, archived, lanes, slice_above, lane=None, model=None, held=None, others=0, owner=False): """Best pickable Pending item of the lane (None = all lanes, priority order), Model ≤ model; the lane empty → its fallback lane (one hop). Plus unpickable items ranked before it, with the reason.""" skipped = [] hops = [lane] fb = next((l.fallback for l in lanes if l.name == lane), None) if fb: hops.append(fb) for name in hops: for item in _ordered(doc, lanes, slice_above, name, model): why = _why(doc, item, archived, slice_above, held, others, owner) if why is None: return item, skipped if (item, why) not in skipped: skipped.append((item, why)) return None, skipped def ranked(doc, archived, lanes, slice_above, lane: str | None) -> list[tasks.Item]: """Every pickable item of the lane in pick order, then its fallback lane's (runner lists).""" hops = [lane] + [l.fallback for l in lanes if l.name == lane and l.fallback] out = [] for name in hops: out += [i for i in _ordered(doc, lanes, slice_above, name, None) if not i.error and _why(doc, i, archived, slice_above, None, 0, False) is None and i not in out] return out