merge(master): resolve PR #907 allocator drain vs side_effect_free
Keep #659 maintenance-drain assignment stop and #643 side_effect_free apply guard; session registration remains gated by side_effect_free. Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
This commit is contained in:
+99
-25
@@ -23,6 +23,7 @@ import json
|
||||
import os
|
||||
import uuid
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Mapping, Sequence
|
||||
|
||||
import maintenance_drain
|
||||
@@ -739,6 +740,46 @@ def normalize_exclude_issue_numbers(
|
||||
return sorted(out)
|
||||
|
||||
|
||||
def _claim_expires_at(claim: Any) -> datetime | None:
|
||||
"""Parse a claim's ``expires_at``, or ``None`` when it is absent/malformed."""
|
||||
if not isinstance(claim, Mapping):
|
||||
return None
|
||||
text = str(claim.get("expires_at") or "").strip()
|
||||
if not text:
|
||||
return None
|
||||
if text.endswith("Z"):
|
||||
text = text[:-1] + "+00:00"
|
||||
try:
|
||||
parsed = datetime.fromisoformat(text)
|
||||
except ValueError:
|
||||
return None
|
||||
if parsed.tzinfo is None:
|
||||
parsed = parsed.replace(tzinfo=timezone.utc)
|
||||
return parsed.astimezone(timezone.utc)
|
||||
|
||||
|
||||
def _drop_expired_claims(
|
||||
claims: Mapping[tuple[str, int], dict[str, Any]],
|
||||
*,
|
||||
now: datetime | None = None,
|
||||
) -> dict[tuple[str, int], dict[str, Any]]:
|
||||
"""Claims minus those whose lease has already expired (#643).
|
||||
|
||||
The read-only mirror of ``expire_stale_leases``: the sweep marks such rows
|
||||
``expired`` so they stop being returned as claims, and this reaches the same
|
||||
view without writing. A claim with no parseable ``expires_at`` is **kept** —
|
||||
an unreadable expiry is not evidence that work is free.
|
||||
"""
|
||||
moment = now or datetime.now(timezone.utc)
|
||||
kept: dict[tuple[str, int], dict[str, Any]] = {}
|
||||
for key, claim in (claims or {}).items():
|
||||
expires_at = _claim_expires_at(claim)
|
||||
if expires_at is not None and expires_at <= moment:
|
||||
continue
|
||||
kept[key] = claim
|
||||
return kept
|
||||
|
||||
|
||||
def candidate_set_fingerprint(
|
||||
candidates: Sequence[WorkCandidate],
|
||||
*,
|
||||
@@ -827,12 +868,22 @@ def allocate_next_work(
|
||||
exclude_issue_numbers: Sequence[int] | None = None,
|
||||
expected_candidate_set_fingerprint: str | None = None,
|
||||
allocation_mode: str | None = None,
|
||||
side_effect_free: bool = False,
|
||||
) -> dict[str, Any]:
|
||||
"""Select and optionally reserve the next work unit via control-plane DB.
|
||||
|
||||
*apply=False* (default): dry-run selection only — no lease/assignment.
|
||||
*apply=True*: atomic ``assign_and_lease`` for the selected candidate.
|
||||
|
||||
*side_effect_free* (#643): a dry run that writes **nothing** to the
|
||||
control-plane DB. A plain ``apply=False`` still registered a session row and
|
||||
swept stale leases globally, so a caller advertising a read-only preview was
|
||||
mutating on every call. Under this flag both writes are suppressed and stale
|
||||
leases are instead filtered out of the claim map in memory, which yields the
|
||||
same selection the sweep would have produced without persisting anything.
|
||||
Incompatible with *apply* — the combination fails closed rather than
|
||||
silently reserving.
|
||||
|
||||
*allocation_mode* (#840): ``cross_role`` (default for controller) inspects
|
||||
the complete queue and returns one authoritative selection naming the
|
||||
required downstream role/profile/action. ``role_scoped`` keeps prior
|
||||
@@ -886,7 +937,7 @@ def allocate_next_work(
|
||||
"allocation_mode": (allocation_mode or "").strip() or None,
|
||||
}
|
||||
|
||||
# #659 AC2: while maintenance drain is active, no new work is assigned —
|
||||
# #659 AC2: while maintenance drain is active, no new work is assigned —
|
||||
# for dry-run and apply alike, so a preview can never be read as evidence
|
||||
# that work was assignable during the drain. Checked before session
|
||||
# registration so a drained allocator leaves no new state behind.
|
||||
@@ -926,40 +977,57 @@ def allocate_next_work(
|
||||
),
|
||||
}
|
||||
|
||||
session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}"
|
||||
try:
|
||||
db.upsert_session(
|
||||
session_id=session_id,
|
||||
role=role_norm,
|
||||
profile=profile_name,
|
||||
pid=os.getpid(),
|
||||
controller_instance_id=controller_instance_id,
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 — surface structured
|
||||
# A side-effect-free run may never reserve: reserving is a write, and the
|
||||
# flag is the caller's assertion that this call writes nothing (#643).
|
||||
if side_effect_free and apply:
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"apply": True,
|
||||
"reasons": [
|
||||
f"failed to register session in control-plane DB: {exc} "
|
||||
"(fail closed, #613)"
|
||||
"side_effect_free is incompatible with apply=True; an "
|
||||
"assignment is a write (fail closed, #643)"
|
||||
],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
|
||||
# Expire stale leases globally before selection.
|
||||
try:
|
||||
db.expire_stale_leases()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"reasons": [f"lease expiry failed: {exc} (fail closed)"],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}"
|
||||
if not side_effect_free:
|
||||
try:
|
||||
db.upsert_session(
|
||||
session_id=session_id,
|
||||
role=role_norm,
|
||||
profile=profile_name,
|
||||
pid=os.getpid(),
|
||||
controller_instance_id=controller_instance_id,
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 — surface structured
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"reasons": [
|
||||
f"failed to register session in control-plane DB: {exc} "
|
||||
"(fail closed, #613)"
|
||||
],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
|
||||
# Expire stale leases globally before selection.
|
||||
try:
|
||||
db.expire_stale_leases()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"reasons": [f"lease expiry failed: {exc} (fail closed)"],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
|
||||
terminal = None
|
||||
try:
|
||||
@@ -994,6 +1062,12 @@ def allocate_next_work(
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
if side_effect_free:
|
||||
# ``list_active_claims`` filters on status alone, so without the
|
||||
# global sweep an already-expired lease would still read as a live
|
||||
# claim and the preview would report work as taken that is free.
|
||||
# Drop those in memory: same view the sweep produces, no write.
|
||||
claims = _drop_expired_claims(claims)
|
||||
|
||||
try:
|
||||
exclude_nums = normalize_exclude_issue_numbers(exclude_issue_numbers)
|
||||
|
||||
Reference in New Issue
Block a user