workflow

git clone https://git.godosa.eu/workflow

master

raw · 13864 bytes

"""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