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
|