fix(release): preserve durable recovery capacity
This commit is contained in:
@@ -59,6 +59,8 @@ ATTEMPT_OUTCOMES = {"running", "succeeded", "failed", "interrupted"}
|
||||
RECONCILIATION_OUTCOMES = {"effect_absent", "effect_succeeded", "unresolved"}
|
||||
RECONCILIATION_CONFIRMATION = "RECONCILE"
|
||||
START_RECOVERY_RESERVE = 2
|
||||
PROJECTION_TIMESTAMP = "9999-12-31T23:59:59Z"
|
||||
PROJECTION_RESULT_CODE = "r" * 64
|
||||
|
||||
|
||||
class ReleaseRunError(RuntimeError):
|
||||
@@ -200,32 +202,32 @@ class ReleaseRunStore:
|
||||
limit = min(max(int(limit), 1), MAX_LIST_LIMIT)
|
||||
with self._locked():
|
||||
self._ensure_root()
|
||||
paths = sorted(
|
||||
(
|
||||
try:
|
||||
paths = [
|
||||
path
|
||||
for path in self.root.iterdir()
|
||||
if path.name.endswith(".json")
|
||||
and RUN_ID_PATTERN.fullmatch(path.stem)
|
||||
),
|
||||
key=lambda item: item.name,
|
||||
reverse=True,
|
||||
)
|
||||
]
|
||||
except OSError as exc:
|
||||
raise ReleaseRunCorrupt(
|
||||
"Release run storage could not be enumerated safely."
|
||||
) from exc
|
||||
summaries: list[dict[str, Any]] = []
|
||||
unavailable: list[dict[str, Any]] = []
|
||||
for path in paths:
|
||||
try:
|
||||
record = self._read(path.stem)
|
||||
except ReleaseRunWorkspaceMismatch:
|
||||
continue
|
||||
except (ReleaseRunCorrupt, ReleaseRunNotFound):
|
||||
summaries.append(
|
||||
unavailable.append(
|
||||
{
|
||||
"run_id": path.stem,
|
||||
"status": "unavailable",
|
||||
"error": "Run record is unavailable or failed integrity validation.",
|
||||
}
|
||||
)
|
||||
if len(summaries) >= limit:
|
||||
break
|
||||
continue
|
||||
view = _view(record)
|
||||
immutable_input = view["immutable"]["input"]
|
||||
@@ -240,9 +242,16 @@ class ReleaseRunStore:
|
||||
"recommended_next": view["recommended_next"],
|
||||
}
|
||||
)
|
||||
if len(summaries) >= limit:
|
||||
break
|
||||
return summaries
|
||||
summaries.sort(
|
||||
key=lambda item: (
|
||||
item["updated_at"],
|
||||
item["created_at"],
|
||||
item["run_id"],
|
||||
),
|
||||
reverse=True,
|
||||
)
|
||||
unavailable.sort(key=lambda item: item["run_id"], reverse=True)
|
||||
return (summaries + unavailable)[:limit]
|
||||
|
||||
def resume(self, run_id: str, *, request_id: str) -> dict[str, Any]:
|
||||
request_fingerprint = _request_fingerprint(request_id)
|
||||
@@ -386,6 +395,7 @@ class ReleaseRunStore:
|
||||
"An unresolved outcome is already recorded for this attempt."
|
||||
)
|
||||
_ensure_command_capacity(record, reserve_after=1)
|
||||
_ensure_unresolved_recovery_byte_capacity(record, step_id)
|
||||
else:
|
||||
_ensure_command_capacity(record)
|
||||
if outcome == "effect_absent":
|
||||
@@ -446,7 +456,7 @@ class ReleaseRunStore:
|
||||
"The attempt identifier was already used for another release step."
|
||||
)
|
||||
return _view(record)
|
||||
step, _plan_step = _step_pair(record, step_id)
|
||||
step, plan_step = _step_pair(record, step_id)
|
||||
available, reason = _step_available(record, step_id)
|
||||
if not available:
|
||||
raise ReleaseRunConflict(reason)
|
||||
@@ -471,6 +481,11 @@ class ReleaseRunStore:
|
||||
}
|
||||
)
|
||||
_append_event(record, event_type="step_started", step_id=step_id)
|
||||
_ensure_start_recovery_byte_capacity(
|
||||
record,
|
||||
step_id,
|
||||
mutating=bool(plan_step.get("mutating")),
|
||||
)
|
||||
self._persist_update(record)
|
||||
return _view(record)
|
||||
|
||||
@@ -1337,12 +1352,13 @@ def _append_event(
|
||||
event_type: str,
|
||||
step_id: str | None = None,
|
||||
result_code: str | None = None,
|
||||
at: str | None = None,
|
||||
) -> None:
|
||||
if event_type not in EVENT_TYPES:
|
||||
raise ReleaseRunConflict("Release event type is invalid.")
|
||||
event: dict[str, Any] = {
|
||||
"sequence": record["state"]["next_event_sequence"],
|
||||
"at": _timestamp(),
|
||||
"at": at if at is not None else _timestamp(),
|
||||
"type": event_type,
|
||||
}
|
||||
if step_id is not None:
|
||||
@@ -1355,6 +1371,216 @@ def _append_event(
|
||||
record["state"]["events"] = record["state"]["events"][-MAX_EVENTS:]
|
||||
|
||||
|
||||
def _ensure_start_recovery_byte_capacity(
|
||||
record: dict[str, Any], step_id: str, *, mutating: bool
|
||||
) -> None:
|
||||
"""Prove that every required completion/recovery write still fits."""
|
||||
|
||||
started = deepcopy(record)
|
||||
_project_persist(started)
|
||||
projections = [started]
|
||||
|
||||
succeeded = _project_finish(started, step_id, succeeded=True)
|
||||
failed = _project_finish(started, step_id, succeeded=False)
|
||||
projections.extend((succeeded, failed, _project_retry(failed, step_id)))
|
||||
|
||||
interrupted = _project_resume(started)
|
||||
projections.append(interrupted)
|
||||
if mutating:
|
||||
projections.extend(
|
||||
_project_reconcile(interrupted, step_id, outcome=outcome)
|
||||
for outcome in ("effect_absent", "effect_succeeded")
|
||||
)
|
||||
else:
|
||||
projections.append(_project_retry(interrupted, step_id))
|
||||
_ensure_projected_records_fit(
|
||||
projections,
|
||||
message=(
|
||||
"Release run has insufficient serialized recovery capacity for a "
|
||||
"new attempt; create a new run."
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _ensure_unresolved_recovery_byte_capacity(
|
||||
record: dict[str, Any], step_id: str
|
||||
) -> None:
|
||||
"""Accept an advisory unresolved write only if a terminal write remains."""
|
||||
|
||||
unresolved = _project_reconcile(record, step_id, outcome="unresolved")
|
||||
projections = [unresolved]
|
||||
projections.extend(
|
||||
_project_reconcile(unresolved, step_id, outcome=outcome)
|
||||
for outcome in ("effect_absent", "effect_succeeded")
|
||||
)
|
||||
_ensure_projected_records_fit(
|
||||
projections,
|
||||
message=(
|
||||
"The unresolved outcome would consume capacity reserved for a "
|
||||
"terminal reconciliation."
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _project_finish(
|
||||
record: dict[str, Any], step_id: str, *, succeeded: bool
|
||||
) -> dict[str, Any]:
|
||||
projected = deepcopy(record)
|
||||
step, _plan_step = _step_pair(projected, step_id)
|
||||
attempt = _attempt_for_fingerprint(
|
||||
projected, step["attempt_fingerprint"], step_id=step_id
|
||||
)
|
||||
target_state = "succeeded" if succeeded else "failed"
|
||||
step.update(
|
||||
{
|
||||
"state": target_state,
|
||||
"finished_at": PROJECTION_TIMESTAMP,
|
||||
"result_code": PROJECTION_RESULT_CODE,
|
||||
}
|
||||
)
|
||||
attempt["outcome"] = target_state
|
||||
attempt["result_code"] = PROJECTION_RESULT_CODE
|
||||
_append_event(
|
||||
projected,
|
||||
event_type="step_succeeded" if succeeded else "step_failed",
|
||||
step_id=step_id,
|
||||
result_code=PROJECTION_RESULT_CODE,
|
||||
at=PROJECTION_TIMESTAMP,
|
||||
)
|
||||
_project_persist(projected)
|
||||
return projected
|
||||
|
||||
|
||||
def _project_resume(record: dict[str, Any]) -> dict[str, Any]:
|
||||
projected = deepcopy(record)
|
||||
running_steps = [
|
||||
step for step in projected["state"]["steps"] if step["state"] == "running"
|
||||
]
|
||||
for step in running_steps:
|
||||
attempt = _attempt_for_fingerprint(
|
||||
projected, step["attempt_fingerprint"], step_id=step["id"]
|
||||
)
|
||||
attempt["outcome"] = "interrupted"
|
||||
attempt["result_code"] = "process_interrupted"
|
||||
step["state"] = "interrupted"
|
||||
step["finished_at"] = PROJECTION_TIMESTAMP
|
||||
step["result_code"] = "process_interrupted"
|
||||
_append_event(
|
||||
projected,
|
||||
event_type="step_interrupted",
|
||||
step_id=step["id"],
|
||||
result_code="process_interrupted",
|
||||
at=PROJECTION_TIMESTAMP,
|
||||
)
|
||||
_remember_request(
|
||||
projected,
|
||||
"0" * 64,
|
||||
"resume",
|
||||
None,
|
||||
None,
|
||||
attempt_fingerprint=running_steps[0]["attempt_fingerprint"],
|
||||
)
|
||||
_append_event(projected, event_type="run_resumed", at=PROJECTION_TIMESTAMP)
|
||||
_project_persist(projected)
|
||||
return projected
|
||||
|
||||
|
||||
def _project_retry(record: dict[str, Any], step_id: str) -> dict[str, Any]:
|
||||
projected = deepcopy(record)
|
||||
step, _plan_step = _step_pair(projected, step_id)
|
||||
attempt_fingerprint = step["attempt_fingerprint"]
|
||||
retry_reason = (
|
||||
"known_failure"
|
||||
if step["state"] == "failed"
|
||||
else "read_only_interruption"
|
||||
)
|
||||
step.update(
|
||||
{
|
||||
"state": "pending",
|
||||
"attempt_fingerprint": None,
|
||||
"started_at": None,
|
||||
"finished_at": None,
|
||||
"result_code": None,
|
||||
}
|
||||
)
|
||||
_remember_request(
|
||||
projected,
|
||||
"1" * 64,
|
||||
"retry",
|
||||
step_id,
|
||||
None,
|
||||
attempt_fingerprint=attempt_fingerprint,
|
||||
)
|
||||
_append_event(
|
||||
projected,
|
||||
event_type="step_retry_requested",
|
||||
step_id=step_id,
|
||||
result_code=retry_reason,
|
||||
at=PROJECTION_TIMESTAMP,
|
||||
)
|
||||
_project_persist(projected)
|
||||
return projected
|
||||
|
||||
|
||||
def _project_reconcile(
|
||||
record: dict[str, Any], step_id: str, *, outcome: str
|
||||
) -> dict[str, Any]:
|
||||
projected = deepcopy(record)
|
||||
step, _plan_step = _step_pair(projected, step_id)
|
||||
attempt_fingerprint = step["attempt_fingerprint"]
|
||||
if outcome == "effect_absent":
|
||||
step.update(
|
||||
{
|
||||
"state": "pending",
|
||||
"attempt_fingerprint": None,
|
||||
"started_at": None,
|
||||
"finished_at": None,
|
||||
"result_code": None,
|
||||
}
|
||||
)
|
||||
elif outcome == "effect_succeeded":
|
||||
step.update(
|
||||
{
|
||||
"state": "succeeded",
|
||||
"finished_at": PROJECTION_TIMESTAMP,
|
||||
"result_code": "reconciled_effect_succeeded",
|
||||
}
|
||||
)
|
||||
_remember_request(
|
||||
projected,
|
||||
"2" * 64 if outcome == "unresolved" else "3" * 64,
|
||||
"reconcile",
|
||||
step_id,
|
||||
outcome,
|
||||
attempt_fingerprint=attempt_fingerprint,
|
||||
)
|
||||
_append_event(
|
||||
projected,
|
||||
event_type="step_reconciled",
|
||||
step_id=step_id,
|
||||
result_code=outcome,
|
||||
at=PROJECTION_TIMESTAMP,
|
||||
)
|
||||
_project_persist(projected)
|
||||
return projected
|
||||
|
||||
|
||||
def _project_persist(record: dict[str, Any]) -> None:
|
||||
record["updated_at"] = PROJECTION_TIMESTAMP
|
||||
_refresh_status(record)
|
||||
record["record_digest"] = "f" * 64
|
||||
|
||||
|
||||
def _ensure_projected_records_fit(
|
||||
records: list[dict[str, Any]], *, message: str
|
||||
) -> None:
|
||||
if any(
|
||||
len(_canonical_json(projected)) + 1 > MAX_RECORD_BYTES
|
||||
for projected in records
|
||||
):
|
||||
raise ReleaseRunConflict(message)
|
||||
|
||||
|
||||
def _ensure_command_capacity(
|
||||
record: dict[str, Any], *, reserve_after: int = 0
|
||||
) -> None:
|
||||
|
||||
Reference in New Issue
Block a user