aboutsummaryrefslogtreecommitdiffziptar.gz
path: root/wflib/lanes.py
diff options
context:
space:
mode:
Diffstat (limited to 'wflib/lanes.py')
-rw-r--r--wflib/lanes.py297
1 files changed, 297 insertions, 0 deletions
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