|
17 | 17 | sleep, |
18 | 18 | ) |
19 | 19 | from anyio.abc import TaskGroup |
| 20 | +from anyio.streams.memory import MemoryObjectReceiveStream, MemoryObjectSendStream |
20 | 21 | from google.protobuf import empty_pb2 |
21 | 22 | from jumpstarter_protocol import ( |
22 | 23 | jumpstarter_pb2, |
@@ -998,111 +999,130 @@ async def process_connections(): |
998 | 999 | if self._lease_context is lease_scope: |
999 | 1000 | self._lease_context = None |
1000 | 1001 |
|
1001 | | - async def serve(self): # noqa: C901 |
1002 | | - """ |
1003 | | - Serve the exporter. |
1004 | | - """ |
1005 | | - # initial registration |
| 1002 | + async def serve(self): |
| 1003 | + """Serve the exporter, handling leases until stopped.""" |
1006 | 1004 | async with self.session(): |
1007 | 1005 | pass |
1008 | | - # Buffer status updates to avoid blocking during short processing gaps |
1009 | 1006 | status_tx, status_rx = create_memory_object_stream[jumpstarter_pb2.StatusResponse](max_buffer_size=5) |
| 1007 | + try: |
| 1008 | + await self._run_control_plane(status_tx, status_rx) |
| 1009 | + finally: |
| 1010 | + self._tg = None |
| 1011 | + self._status_drain_active = False |
| 1012 | + clear_log_context() |
1010 | 1013 |
|
| 1014 | + async def _run_control_plane( |
| 1015 | + self, |
| 1016 | + status_tx: MemoryObjectSendStream[jumpstarter_pb2.StatusResponse], |
| 1017 | + status_rx: MemoryObjectReceiveStream[jumpstarter_pb2.StatusResponse], |
| 1018 | + ) -> None: |
| 1019 | + """Start control-plane streams and process status updates.""" |
1011 | 1020 | async with create_task_group() as tg: |
1012 | 1021 | self._tg = tg |
1013 | | - # Start background status drain (makes _report_status non-blocking) |
1014 | 1022 | self._status_rpc_event = Event() |
1015 | 1023 | self._pending_status_request = None |
1016 | 1024 | self._status_drain_active = True |
1017 | 1025 | tg.start_soon(self._drain_status_reports) |
1018 | | - # Start status stream with retry logic |
1019 | 1026 | tg.start_soon( |
1020 | 1027 | self._retry_stream, |
1021 | 1028 | "Status", |
1022 | 1029 | self._status_stream_factory(), |
1023 | 1030 | status_tx, |
1024 | 1031 | ) |
1025 | 1032 | async for status in status_rx: |
1026 | | - # Check for lease state transitions |
1027 | | - previous_leased = self._previous_leased |
1028 | | - current_leased = status.leased |
1029 | | - |
1030 | | - # Check if this is a new lease assignment (no active lease context and we have a lease name) |
1031 | | - # This handles both first lease and subsequent leases after the previous one ended |
1032 | | - if self._lease_context is None and status.lease_name != "" and current_leased: |
1033 | | - self._started = True |
1034 | | - logger.info("Starting new lease: %s", status.lease_name) |
1035 | | - # Create lease scope and start handling the lease |
1036 | | - # The session will be created inside handle_lease and stay open for the lease duration |
1037 | | - lease_scope = LeaseContext( |
1038 | | - lease_name=status.lease_name, |
1039 | | - before_lease_hook=Event(), |
1040 | | - ) |
1041 | | - self._lease_context = lease_scope |
1042 | | - log_ctx = {"lease_id": status.lease_name, "exporter": self.name} |
1043 | | - if status.context: |
1044 | | - log_ctx.update(status.context) |
1045 | | - set_log_context(**log_ctx) |
1046 | | - tg.start_soon(self.handle_lease, status.lease_name, tg, lease_scope) |
1047 | | - |
1048 | | - if current_leased: |
1049 | | - if self._lease_context: |
1050 | | - self._lease_context.update_client(status.client_name) |
1051 | | - if status.client_name: |
1052 | | - set_log_context(client=status.client_name) |
1053 | | - logger.info("Currently leased by %s under %s", status.client_name, status.lease_name) |
1054 | | - |
1055 | | - # Before-lease hook when transitioning from unleased to leased |
1056 | | - if not previous_leased: |
1057 | | - if self.hook_executor and self._lease_context: |
1058 | | - tg.start_soon( |
1059 | | - self.hook_executor.run_before_lease_hook, |
1060 | | - self._lease_context, |
1061 | | - self._report_status, |
1062 | | - self.stop, # Pass shutdown callback |
1063 | | - self._request_lease_release, # Pass lease release callback |
1064 | | - ) |
1065 | | - # else: No hook configured - LEASE_READY is set inside handle_lease() |
1066 | | - # after session and Listen stream are established |
1067 | | - else: |
1068 | | - logger.info("Currently not leased") |
1069 | | - |
1070 | | - # Lease ended: signal handle_lease() so it can exit its loop and run |
1071 | | - # cleanup/afterLease hook in its finally block (where session is still open) |
1072 | | - if previous_leased and self._lease_context: |
1073 | | - lease_ctx = self._lease_context |
1074 | | - logger.info("Lease ended, signaling handle_lease to run afterLease hook") |
1075 | | - lease_ctx.lease_ended.set() |
1076 | | - |
1077 | | - # Wait for the hook to complete |
1078 | | - with CancelScope(shield=True): |
1079 | | - await lease_ctx.after_lease_hook_done.wait() |
1080 | | - logger.info("afterLease hook completed") |
1081 | | - |
1082 | | - # Clear lease scope and log context for next lease |
1083 | | - session_was_created = ( |
1084 | | - self._lease_context is not None and self._lease_context.session is not None |
| 1033 | + if await self._apply_status(status, tg): |
| 1034 | + break |
| 1035 | + |
| 1036 | + async def _apply_status( |
| 1037 | + self, |
| 1038 | + status: jumpstarter_pb2.StatusResponse, |
| 1039 | + tg: TaskGroup, |
| 1040 | + ) -> bool: |
| 1041 | + """Process a single status update. Returns True to stop the status loop.""" |
| 1042 | + previous_leased = self._previous_leased |
| 1043 | + current_leased = status.leased |
| 1044 | + |
| 1045 | + if self._lease_context is None and status.lease_name != "" and current_leased: |
| 1046 | + self._on_lease_acquired(status, tg) |
| 1047 | + |
| 1048 | + if current_leased: |
| 1049 | + self._on_lease_update(status) |
| 1050 | + if not previous_leased: |
| 1051 | + if self.hook_executor and self._lease_context: |
| 1052 | + tg.start_soon( |
| 1053 | + self.hook_executor.run_before_lease_hook, |
| 1054 | + self._lease_context, |
| 1055 | + self._report_status, |
| 1056 | + self.stop, |
| 1057 | + self._request_lease_release, |
1085 | 1058 | ) |
1086 | | - self._lease_context = None |
1087 | | - clear_log_context() |
1088 | | - if session_was_created: |
1089 | | - # Brief delay to ensure session is fully closed before next lease |
1090 | | - # This prevents SSL corruption from overlapping connections |
1091 | | - await sleep(0.2) |
1092 | | - logger.debug("Ready for next lease") |
1093 | | - |
1094 | | - if self.exit_on_lease_end and previous_leased: |
1095 | | - logger.info("Exporter configured to exit after lease, shutting down") |
1096 | | - self._stop_requested = True |
1097 | | - |
1098 | | - if self._stop_requested: |
1099 | | - self.stop(should_unregister=self._deferred_unregister) |
1100 | | - break |
1101 | | - |
1102 | | - self._previous_leased = current_leased |
1103 | | - self._tg = None |
1104 | | - self._status_drain_active = False |
| 1059 | + else: |
| 1060 | + await self._on_lease_released(previous_leased) |
| 1061 | + |
| 1062 | + self._previous_leased = current_leased |
| 1063 | + return self._check_stop_requested() if not current_leased else False |
| 1064 | + |
| 1065 | + def _on_lease_acquired( |
| 1066 | + self, |
| 1067 | + status: jumpstarter_pb2.StatusResponse, |
| 1068 | + tg: TaskGroup, |
| 1069 | + ) -> None: |
| 1070 | + """Handle new lease assignment: create context and spawn lease handler.""" |
| 1071 | + self._started = True |
| 1072 | + logger.info("Starting new lease: %s", status.lease_name) |
| 1073 | + lease_scope = LeaseContext( |
| 1074 | + lease_name=status.lease_name, |
| 1075 | + before_lease_hook=Event(), |
| 1076 | + ) |
| 1077 | + self._lease_context = lease_scope |
| 1078 | + log_ctx: dict[str, str] = {"lease_id": status.lease_name, "exporter": self.name} |
| 1079 | + if status.context: |
| 1080 | + log_ctx.update(status.context) |
| 1081 | + set_log_context(**log_ctx) |
| 1082 | + tg.start_soon(self.handle_lease, status.lease_name, tg, lease_scope) |
| 1083 | + |
| 1084 | + def _on_lease_update(self, status: jumpstarter_pb2.StatusResponse) -> None: |
| 1085 | + """Update client info on every leased status tick.""" |
| 1086 | + if self._lease_context: |
| 1087 | + self._lease_context.update_client(status.client_name) |
| 1088 | + if status.client_name: |
| 1089 | + set_log_context(client=status.client_name) |
| 1090 | + logger.info("Currently leased by %s under %s", status.client_name, status.lease_name) |
| 1091 | + |
| 1092 | + async def _on_lease_released(self, previous_leased: bool) -> None: |
| 1093 | + """Handle not-leased status: signal handle_lease on transition, clean up context.""" |
| 1094 | + logger.info("Currently not leased") |
| 1095 | + |
| 1096 | + if previous_leased and self._lease_context: |
| 1097 | + lease_ctx = self._lease_context |
| 1098 | + logger.info("Lease ended, signaling handle_lease to run afterLease hook") |
| 1099 | + lease_ctx.lease_ended.set() |
| 1100 | + |
| 1101 | + with CancelScope(shield=True): |
| 1102 | + await lease_ctx.after_lease_hook_done.wait() |
| 1103 | + logger.info("afterLease hook completed") |
| 1104 | + |
| 1105 | + session_was_created = ( |
| 1106 | + self._lease_context is not None and self._lease_context.session is not None |
| 1107 | + ) |
| 1108 | + self._lease_context = None |
1105 | 1109 | clear_log_context() |
| 1110 | + if session_was_created: |
| 1111 | + # Brief delay to ensure session is fully closed before next lease. |
| 1112 | + # Prevents SSL corruption from overlapping connections. |
| 1113 | + await sleep(0.2) |
| 1114 | + logger.debug("Ready for next lease") |
| 1115 | + |
| 1116 | + if self.exit_on_lease_end and previous_leased: |
| 1117 | + logger.info("Exporter configured to exit after lease, shutting down") |
| 1118 | + self._stop_requested = True |
| 1119 | + |
| 1120 | + def _check_stop_requested(self) -> bool: |
| 1121 | + """Check if stop was requested and initiate shutdown. Returns True to break the status loop.""" |
| 1122 | + if self._stop_requested: |
| 1123 | + self.stop(should_unregister=self._deferred_unregister) |
| 1124 | + return True |
| 1125 | + return False |
1106 | 1126 |
|
1107 | 1127 | async def serve_standalone_tcp( |
1108 | 1128 | self, |
|
0 commit comments