Files
playbook/skills/cook-it-through/scripts/main_loop_scheduler.py
T

388 lines
14 KiB
Python

"""Pure identities and dependency scheduling for the ticket main loop.
The module deliberately has no filesystem, Git, clock, or locking dependency.
Callers provide an immutable snapshot and receive deterministic domain results.
"""
from dataclasses import dataclass
import re
from typing import TypeAlias
FEATURE_SLUG_RE = re.compile(r"^[a-z0-9][a-z0-9-]*$")
TICKET_NUMBER_RE = re.compile(r"^\d{2,}$")
TICKET_ID_RE = re.compile(
r"^(?P<feature>[a-z0-9][a-z0-9-]*)/(?P<number>\d{2,})$"
)
FEATURE_INTEGRATION_ID_RE = re.compile(
r"^(?P<feature>[a-z0-9][a-z0-9-]*)@integrated$"
)
SATISFIED_TICKET_STATUSES = frozenset({"resolved", "skipped"})
VALID_TICKET_STATUSES = frozenset(
{
"ready-for-agent",
"claimed",
"blocked",
"resolved",
"skipped",
}
)
class SchedulerError(ValueError):
"""The queued dependency graph violates the scheduler contract."""
@dataclass(frozen=True, order=True)
class FeatureId:
value: str
def __post_init__(self) -> None:
if not isinstance(self.value, str) or not FEATURE_SLUG_RE.fullmatch(self.value):
raise SchedulerError(f"invalid feature identity: {self.value}")
@classmethod
def parse(cls, raw: str) -> "FeatureId":
if not isinstance(raw, str):
raise SchedulerError(f"invalid feature identity: {raw}")
return cls(raw)
def __str__(self) -> str:
return self.value
@dataclass(frozen=True)
class TicketId:
feature: FeatureId
number: str
def __post_init__(self) -> None:
if not isinstance(self.feature, FeatureId) or not isinstance(
self.number, str
) or not TICKET_NUMBER_RE.fullmatch(self.number):
raise SchedulerError(
f"invalid ticket identity: {self.feature}/{self.number}"
)
@classmethod
def parse(cls, raw: str) -> "TicketId":
if not isinstance(raw, str):
raise SchedulerError(
f"invalid ticket identity '{raw}'; expected <feature>/<number>"
)
match = TICKET_ID_RE.fullmatch(raw)
if not match:
raise SchedulerError(
f"invalid ticket identity '{raw}'; expected <feature>/<number>"
)
return cls(
FeatureId.parse(match.group("feature")),
match.group("number"),
)
def __str__(self) -> str:
return f"{self.feature}/{self.number}"
@dataclass(frozen=True)
class FeatureIntegrationId:
feature: FeatureId
@classmethod
def parse(cls, raw: str) -> "FeatureIntegrationId":
if not isinstance(raw, str):
raise SchedulerError(
"invalid feature integration identity "
f"'{raw}'; expected <feature>@integrated"
)
match = FEATURE_INTEGRATION_ID_RE.fullmatch(raw)
if not match:
raise SchedulerError(
"invalid feature integration identity "
f"'{raw}'; expected <feature>@integrated"
)
return cls(FeatureId.parse(match.group("feature")))
def __str__(self) -> str:
return f"{self.feature}@integrated"
Dependency: TypeAlias = TicketId | FeatureIntegrationId
@dataclass(frozen=True)
class TicketRecord:
id: TicketId
slug: str
status: str
dependencies: tuple[Dependency, ...]
@dataclass(frozen=True)
class FeatureRecord:
id: FeatureId
tickets: tuple[TicketRecord, ...]
integrated: bool = False
integration_blocked: bool = False
class Scheduler:
"""Validate and query one immutable, queue-ordered global graph snapshot."""
def __init__(self, features: tuple[FeatureRecord, ...]) -> None:
self._features = features
feature_ids = [feature.id for feature in features]
for feature_id in feature_ids:
if not isinstance(feature_id, FeatureId):
raise SchedulerError(f"invalid feature identity: {feature_id}")
if len(feature_ids) != len(set(feature_ids)):
duplicate = next(
feature_id
for feature_id in feature_ids
if feature_ids.count(feature_id) > 1
)
raise SchedulerError(f"duplicate queued feature: {duplicate}")
self._feature_by_id = {feature.id: feature for feature in features}
self._ticket_by_id: dict[TicketId, TicketRecord] = {}
for feature in features:
for ticket in feature.tickets:
if not isinstance(ticket.id, TicketId):
raise SchedulerError(f"invalid ticket identity: {ticket.id}")
if ticket.id.feature != feature.id:
raise SchedulerError(
f"ticket {ticket.id} is stored under feature {feature.id}"
)
if ticket.status not in VALID_TICKET_STATUSES:
raise SchedulerError(
f"{ticket.id}: invalid status {ticket.status}"
)
if ticket.id in self._ticket_by_id:
raise SchedulerError(f"duplicate ticket identity: {ticket.id}")
self._ticket_by_id[ticket.id] = ticket
self._validate_integration_state()
self._queue_index = {
feature.id: index for index, feature in enumerate(features)
}
self._validate_dependency_targets()
self._edges = self._build_edges()
self._validate_acyclic()
def _validate_integration_state(self) -> None:
first_pending: FeatureIntegrationId | None = None
for feature in self._features:
integration_id = FeatureIntegrationId(feature.id)
if feature.integrated:
if first_pending is not None:
raise SchedulerError(
f"{integration_id}: earlier integration is pending: "
f"{first_pending}"
)
for ticket in sorted(feature.tickets, key=lambda item: str(item.id)):
if ticket.status not in SATISFIED_TICKET_STATUSES:
raise SchedulerError(
f"{integration_id}: unsatisfied ticket {ticket.id} "
f"has status {ticket.status}"
)
elif first_pending is None:
first_pending = integration_id
def _validate_dependency_targets(self) -> None:
for ticket_id in sorted(self._ticket_by_id, key=str):
ticket = self._ticket_by_id[ticket_id]
for dependency in ticket.dependencies:
if not isinstance(dependency, (TicketId, FeatureIntegrationId)):
raise SchedulerError(
f"{ticket.id}: invalid dependency {dependency}"
)
if len(ticket.dependencies) != len(set(ticket.dependencies)):
duplicate = next(
dependency
for dependency in ticket.dependencies
if ticket.dependencies.count(dependency) > 1
)
raise SchedulerError(
f"{ticket.id}: duplicate dependency {duplicate}"
)
for dependency in sorted(ticket.dependencies, key=str):
if dependency == ticket.id:
raise SchedulerError(
f"{ticket.id}: self dependency {dependency}"
)
if dependency.feature not in self._feature_by_id:
raise SchedulerError(
f"{ticket.id}: dependency feature not queued: "
f"{dependency.feature}"
)
if (
isinstance(dependency, TicketId)
and dependency not in self._ticket_by_id
):
raise SchedulerError(
f"{ticket.id}: dependency ticket not found: {dependency}"
)
def _build_edges(self) -> dict[Dependency, tuple[Dependency, ...]]:
edges: dict[Dependency, tuple[Dependency, ...]] = {
ticket_id: tuple(sorted(ticket.dependencies, key=str))
for ticket_id, ticket in self._ticket_by_id.items()
}
previous: FeatureIntegrationId | None = None
for feature in self._features:
integration_id = FeatureIntegrationId(feature.id)
dependencies: list[Dependency] = [
ticket.id for ticket in feature.tickets
]
if previous is not None:
dependencies.append(previous)
edges[integration_id] = tuple(sorted(dependencies, key=str))
previous = integration_id
return edges
def _validate_acyclic(self) -> None:
visited: set[Dependency] = set()
active: set[Dependency] = set()
stack: list[Dependency] = []
def visit(node: Dependency) -> None:
if node in visited:
return
if node in active:
start = stack.index(node)
cycle = (*stack[start:], node)
raise SchedulerError(
"dependency cycle: " + " -> ".join(map(str, cycle))
)
active.add(node)
stack.append(node)
for dependency in self._edges[node]:
visit(dependency)
stack.pop()
active.remove(node)
visited.add(node)
for node in sorted(self._edges, key=str):
visit(node)
def _dependency_is_satisfied(self, dependency: Dependency) -> bool:
if isinstance(dependency, TicketId):
return (
self._ticket_by_id[dependency].status
in SATISFIED_TICKET_STATUSES
)
return self._feature_by_id[dependency.feature].integrated
def unsatisfied_dependencies(
self, node: Dependency
) -> tuple[Dependency, ...]:
return tuple(
dependency
for dependency in self._edges[node]
if not self._dependency_is_satisfied(dependency)
)
def integration_dependencies(
self,
ticket_id: TicketId,
) -> tuple[FeatureIntegrationId, ...]:
return tuple(
dependency
for dependency in self._edges[ticket_id]
if isinstance(dependency, FeatureIntegrationId)
)
@property
def integration_frontier(self) -> FeatureIntegrationId | None:
for feature in self._features:
if not feature.integrated:
return FeatureIntegrationId(feature.id)
return None
@property
def integration_frontier_ready(self) -> bool:
frontier = self.integration_frontier
if frontier is None:
return False
feature = self._feature_by_id[frontier.feature]
return (
not feature.integration_blocked
and not self.unsatisfied_dependencies(frontier)
)
def feature_state(self, feature_id: FeatureId) -> str:
feature = self._feature_by_id[feature_id]
if feature.integrated:
return "integrated"
if feature.integration_blocked:
return "blocked"
statuses = {ticket.status for ticket in feature.tickets}
if statuses <= SATISFIED_TICKET_STATUSES:
return "ready-to-integrate"
if "claimed" in statuses or statuses & SATISFIED_TICKET_STATUSES:
return "active"
if any(ticket_id.feature == feature_id for ticket_id in self.ticket_frontier):
return "queued"
return "blocked"
@property
def ticket_frontier(self) -> tuple[TicketId, ...]:
claimable = (
ticket
for feature in self._features
if not feature.integrated
for ticket in feature.tickets
if ticket.status == "ready-for-agent"
and not self.unsatisfied_dependencies(ticket.id)
)
return tuple(
ticket.id
for ticket in sorted(
claimable,
key=lambda ticket: (
self._queue_index[ticket.id.feature],
int(ticket.id.number),
ticket.id.number,
ticket.slug,
),
)
)
def parse_dependencies(raw: str, owner: TicketId) -> tuple[Dependency, ...]:
"""Parse the one canonical ``Blocked by`` representation."""
if not isinstance(raw, str):
raise SchedulerError(
f"{owner}: invalid dependency value; expected a Markdown string"
)
if not isinstance(owner, TicketId):
raise SchedulerError(f"invalid dependency owner: {owner}")
value = raw.strip()
if value == "None":
return ()
entries = value.split(";")
dependencies: list[Dependency] = []
for raw_entry in entries:
entry = raw_entry.strip()
ticket_match = TICKET_ID_RE.fullmatch(entry)
integration_match = FEATURE_INTEGRATION_ID_RE.fullmatch(entry)
if ticket_match:
dependency: Dependency = TicketId(
FeatureId.parse(ticket_match.group("feature")),
ticket_match.group("number"),
)
elif integration_match:
dependency = FeatureIntegrationId(
FeatureId.parse(integration_match.group("feature"))
)
else:
raise SchedulerError(
f"{owner}: invalid dependency '{entry}'; expected "
"<feature>/<number> or <feature>@integrated separated by ';'"
)
if dependency in dependencies:
raise SchedulerError(f"{owner}: duplicate dependency {dependency}")
if dependency == owner:
raise SchedulerError(f"{owner}: self dependency {dependency}")
dependencies.append(dependency)
return tuple(dependencies)