raw · 13864 bytes
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 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 |