| 88 | |
| 89 | |
| 90 | class _FakeWorkResource: |
| 91 | def __init__(self, *, heartbeat_state: str = "stopping", heartbeat_error: Exception | None = None) -> None: |
| 92 | self._heartbeat_state = heartbeat_state |
| 93 | self._heartbeat_error = heartbeat_error |
| 94 | # When set, every heartbeat waits for this event before answering, so a |
| 95 | # test decides *when* in the run the heartbeat's outcome lands. |
| 96 | self.heartbeat_gate: asyncio.Event | None = None |
| 97 | self.heartbeat_calls: list[dict[str, Any]] = [] |
| 98 | self.stop_calls: list[dict[str, Any]] = [] |
| 99 | |
| 100 | async def heartbeat( |
| 101 | self, |
| 102 | work_id: str, |
| 103 | *, |
| 104 | environment_id: str, |
| 105 | expected_last_heartbeat: str, # noqa: ARG002 |
| 106 | extra_headers: Any = None, |
| 107 | ) -> Any: |
| 108 | self.heartbeat_calls.append( |
| 109 | {"work_id": work_id, "environment_id": environment_id, "extra_headers": extra_headers} |
| 110 | ) |
| 111 | if self.heartbeat_gate is not None: |
| 112 | await self.heartbeat_gate.wait() |
| 113 | if self._heartbeat_error is not None: |
| 114 | raise self._heartbeat_error |
| 115 | return SimpleNamespace(last_heartbeat="hb-1", ttl_seconds=60, state=self._heartbeat_state, lease_extended=True) |
| 116 | |
| 117 | async def stop( |
| 118 | self, |
| 119 | work_id: str, |
| 120 | *, |
| 121 | environment_id: str, |
| 122 | force: bool = False, |
| 123 | extra_headers: Any = None, |
| 124 | betas: Any = None, |
| 125 | ) -> None: |
| 126 | self.stop_calls.append( |
| 127 | { |
| 128 | "work_id": work_id, |
| 129 | "environment_id": environment_id, |
| 130 | "force": force, |
| 131 | "extra_headers": extra_headers, |
| 132 | "betas": betas, |
| 133 | } |
| 134 | ) |
| 135 | |
| 136 | |
| 137 | class _FakeSessions: |
no outgoing calls