From 81d4e80fd5aabe4e80f58e960affa795cf7d34ec Mon Sep 17 00:00:00 2001 From: godosa Date: Wed, 7 Oct 2026 07:27:17 +0200 Subject: workflow: initial public history --- wflib/lanes.py | 297 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 297 insertions(+) create mode 100644 wflib/lanes.py (limited to 'wflib/lanes.py') diff --git a/wflib/lanes.py b/wflib/lanes.py new file mode 100644 index 0000000..d9b6fe6 --- /dev/null +++ b/wflib/lanes.py @@ -0,0 +1,297 @@ +"""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 -- cgit