|
|
@@ -167,6 +167,108 @@ class FencedCommandGate:
|
|
|
await self.verify_in_session(session, token)
|
|
|
return await mutation(session)
|
|
|
|
|
|
+ async def compare_and_set_input_snapshot(
|
|
|
+ self,
|
|
|
+ token: MissionOwnerToken,
|
|
|
+ *,
|
|
|
+ expected_snapshot_id: int,
|
|
|
+ new_snapshot_id: int,
|
|
|
+ ) -> None:
|
|
|
+ """Move the binding to one immutable snapshot under the owner fence."""
|
|
|
+
|
|
|
+ async with self._sessions() as session, session.begin():
|
|
|
+ await self.verify_in_session(session, token)
|
|
|
+ result = await session.execute(
|
|
|
+ update(mission_binding_table)
|
|
|
+ .where(
|
|
|
+ mission_binding_table.c.script_build_id == token.script_build_id,
|
|
|
+ mission_binding_table.c.input_snapshot_id == expected_snapshot_id,
|
|
|
+ )
|
|
|
+ .values(input_snapshot_id=new_snapshot_id, updated_at=datetime.now(UTC))
|
|
|
+ )
|
|
|
+ if result.rowcount != 1:
|
|
|
+ current = await session.scalar(
|
|
|
+ select(mission_binding_table.c.input_snapshot_id).where(
|
|
|
+ mission_binding_table.c.script_build_id == token.script_build_id
|
|
|
+ )
|
|
|
+ )
|
|
|
+ if current != new_snapshot_id:
|
|
|
+ raise MissionFencingTokenStale()
|
|
|
+
|
|
|
+ async def begin_phase(self, token: MissionOwnerToken) -> None:
|
|
|
+ """Project the legacy build to running without a stop-overwrite window."""
|
|
|
+
|
|
|
+ await self._update_runtime_state(
|
|
|
+ token,
|
|
|
+ expected_statuses=(BuildStatus.PARTIAL.value, BuildStatus.RUNNING.value),
|
|
|
+ status=BuildStatus.RUNNING.value,
|
|
|
+ end_time=None,
|
|
|
+ error_message=None,
|
|
|
+ summary=None,
|
|
|
+ )
|
|
|
+
|
|
|
+ async def set_checkpoint(
|
|
|
+ self,
|
|
|
+ token: MissionOwnerToken,
|
|
|
+ *,
|
|
|
+ checkpoint_code: str,
|
|
|
+ summary: str,
|
|
|
+ ) -> None:
|
|
|
+ """Persist the Phase 3 checkpoint in the same transaction as fencing."""
|
|
|
+
|
|
|
+ await self._update_runtime_state(
|
|
|
+ token,
|
|
|
+ expected_statuses=(BuildStatus.RUNNING.value,),
|
|
|
+ idempotent_checkpoint=checkpoint_code[:255],
|
|
|
+ status=BuildStatus.PARTIAL.value,
|
|
|
+ error_message=checkpoint_code[:255],
|
|
|
+ summary=summary[:2000],
|
|
|
+ reson_trace_id=token.root_trace_id[:200],
|
|
|
+ end_time=datetime.now(UTC),
|
|
|
+ )
|
|
|
+
|
|
|
+ async def _update_runtime_state(
|
|
|
+ self,
|
|
|
+ token: MissionOwnerToken,
|
|
|
+ *,
|
|
|
+ expected_statuses: tuple[str, ...],
|
|
|
+ idempotent_checkpoint: str | None = None,
|
|
|
+ **values: Any,
|
|
|
+ ) -> None:
|
|
|
+ async with self._sessions() as session, session.begin():
|
|
|
+ await self.verify_in_session(session, token)
|
|
|
+ result = await session.execute(
|
|
|
+ update(self._runtime_record)
|
|
|
+ .where(
|
|
|
+ self._runtime_record.c.id == token.script_build_id,
|
|
|
+ self._runtime_record.c.is_deleted.is_(False),
|
|
|
+ self._runtime_record.c.status.in_(expected_statuses),
|
|
|
+ )
|
|
|
+ .values(**values)
|
|
|
+ )
|
|
|
+ if result.rowcount != 1:
|
|
|
+ current = (
|
|
|
+ (
|
|
|
+ await session.execute(
|
|
|
+ select(
|
|
|
+ script_build_record.c.status,
|
|
|
+ script_build_record.c.error_message,
|
|
|
+ ).where(script_build_record.c.id == token.script_build_id)
|
|
|
+ )
|
|
|
+ )
|
|
|
+ .mappings()
|
|
|
+ .one_or_none()
|
|
|
+ )
|
|
|
+ if current is None:
|
|
|
+ raise BuildNotFound()
|
|
|
+ if (
|
|
|
+ idempotent_checkpoint is not None
|
|
|
+ and current["status"] == BuildStatus.PARTIAL.value
|
|
|
+ and current["error_message"] == idempotent_checkpoint
|
|
|
+ ):
|
|
|
+ return
|
|
|
+ raise MissionFencingTokenStale()
|
|
|
+
|
|
|
async def verify_in_session(
|
|
|
self,
|
|
|
session: AsyncSession,
|
|
|
@@ -204,6 +306,8 @@ class FencedCommandGate:
|
|
|
raise BuildNotFound()
|
|
|
if status in {BuildStatus.STOPPING.value, BuildStatus.STOPPED.value}:
|
|
|
raise MissionStopRequested()
|
|
|
+ if status == BuildStatus.FAILED.value:
|
|
|
+ raise MissionFencingTokenStale()
|
|
|
if status == BuildStatus.SUCCESS.value and not allow_success:
|
|
|
raise MissionFencingTokenStale()
|
|
|
|