Skip to content

Commit 97309ba

Browse files
committed
fix: distinguish shutdown cleanup outcomes from unexpected background failures
Supervisor treats early success/cancel as errors; cleanup only logs real exceptions and skips already-handled outcomes so clean shutdown stays silent. Assign the outcome marker directly to satisfy flake8 B010.
1 parent 662fd92 commit 97309ba

2 files changed

Lines changed: 196 additions & 16 deletions

File tree

skywalking/agent/__init__.py

Lines changed: 56 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -147,25 +147,53 @@ async def _shutdown_async_queue(q: asyncio.Queue, label: str) -> None:
147147
)
148148

149149

150-
def _retrieve_background_task_outcome(task: asyncio.Task):
150+
def _retrieve_background_task_outcome(
151+
task: asyncio.Task,
152+
*,
153+
report_unexpected_completion: bool = True,
154+
report_unexpected_cancellation: bool = False,
155+
):
151156
"""
152-
Return a completed background task's exception, or a sentinel for unexpected success.
157+
Return a completed background task's exception when it should be reported.
153158
154159
Retrieves the exception so asyncio does not emit "never retrieved" warnings.
160+
Successful completion / cancellation are only converted to errors when the
161+
corresponding report_* flag is set (supervisor path before shutdown).
155162
"""
156-
if task is None or not task.done() or task.cancelled():
163+
if task is None or not task.done():
164+
return None
165+
if task.cancelled():
166+
if report_unexpected_cancellation:
167+
return RuntimeError(
168+
'Python agent asyncio background task was cancelled unexpectedly'
169+
)
157170
return None
158171
exc = task.exception()
159172
if exc is not None:
160173
return exc
161-
return RuntimeError('Python agent asyncio background task finished unexpectedly')
174+
if report_unexpected_completion:
175+
return RuntimeError('Python agent asyncio background task finished unexpectedly')
176+
return None
162177

163178

164-
def _log_background_task_outcome(task: asyncio.Task) -> bool:
179+
def _log_background_task_outcome(
180+
task: asyncio.Task,
181+
*,
182+
report_unexpected_completion: bool = True,
183+
report_unexpected_cancellation: bool = False,
184+
) -> bool:
165185
"""Log and retrieve a completed background task outcome. Returns True if logged."""
166-
exc = _retrieve_background_task_outcome(task)
186+
if getattr(task, '_sw_outcome_handled', False):
187+
return False
188+
exc = _retrieve_background_task_outcome(
189+
task,
190+
report_unexpected_completion=report_unexpected_completion,
191+
report_unexpected_cancellation=report_unexpected_cancellation,
192+
)
167193
if exc is None:
168194
return False
195+
# Mark before logging so cleanup never re-logs a supervisor-handled outcome.
196+
task._sw_outcome_handled = True
169197
logger.error('Error in Python agent asyncio event loop: %s', exc, exc_info=exc)
170198
return True
171199

@@ -177,8 +205,9 @@ async def _await_shutdown_or_background_failure(
177205
"""
178206
Wait for shutdown or the first unexpected background-task completion.
179207
180-
Background tasks are expected to run until shutdown. If one finishes early,
181-
retrieve/log its outcome and signal shutdown so the root can clean up.
208+
Background tasks are expected to run until shutdown. If one finishes or is
209+
cancelled early, retrieve/log its outcome and signal shutdown so the root
210+
can clean up.
182211
"""
183212
shutdown_waiter = asyncio.create_task(finished.wait())
184213
pending = {task for task in background_tasks if task is not None}
@@ -195,7 +224,12 @@ async def _await_shutdown_or_background_failure(
195224
return
196225
for task in done:
197226
pending.discard(task)
198-
if _log_background_task_outcome(task):
227+
# Before shutdown, both early success and unexpected cancel are errors.
228+
if _log_background_task_outcome(
229+
task,
230+
report_unexpected_completion=True,
231+
report_unexpected_cancellation=True,
232+
):
199233
finished.set()
200234
return
201235
finally:
@@ -213,13 +247,20 @@ async def _cancel_pending_tasks(tasks) -> None:
213247
214248
Cancel only the given reporter / watch tasks; never the asyncio.run root.
215249
The current task is excluded so we never await ourselves.
250+
251+
During cleanup, normal completion and intentional cancellation are silent;
252+
only real exceptions that were not already handled by the supervisor are logged.
216253
"""
217254
current = asyncio.current_task()
218255
for task in tasks:
219256
if task is None or task is current:
220257
continue
221258
if task.done():
222-
_log_background_task_outcome(task)
259+
_log_background_task_outcome(
260+
task,
261+
report_unexpected_completion=False,
262+
report_unexpected_cancellation=False,
263+
)
223264
pending = [
224265
task for task in tasks
225266
if task is not None and task is not current and not task.done()
@@ -240,7 +281,11 @@ async def _cancel_pending_tasks(tasks) -> None:
240281
sum(1 for task in pending if not task.done()),
241282
)
242283
for task in pending:
243-
_log_background_task_outcome(task)
284+
_log_background_task_outcome(
285+
task,
286+
report_unexpected_completion=False,
287+
report_unexpected_cancellation=False,
288+
)
244289

245290

246291
def _close_previous_protocol(protocol) -> None:

tests/unit/test_shutdown_queue.py

Lines changed: 140 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -179,14 +179,16 @@ async def reporter():
179179
self.assertTrue(holder.get('root_finished', False))
180180

181181
def test_background_task_failure_is_logged_not_silenced(self):
182-
"""Regression: failing background tasks must be observed and logged."""
183-
holder = {}
182+
"""Regression: failing background tasks must be observed and logged exactly once."""
183+
holder = {'errors': []}
184184
error_logged = Event()
185+
aclose_done = Event()
185186

186187
class _Handler(logging.Handler):
187188
def emit(self, record):
188-
if 'Error in Python agent asyncio event loop' in record.getMessage():
189-
holder['logged'] = record.getMessage()
189+
msg = record.getMessage()
190+
if 'Error in Python agent asyncio event loop' in msg:
191+
holder['errors'].append(msg)
190192
error_logged.set()
191193

192194
agent_logger = logging.getLogger('skywalking')
@@ -199,6 +201,10 @@ async def failing_background():
199201
await asyncio.sleep(0.02)
200202
raise ValueError('command dispatch failed')
201203

204+
async def yielding_aclose():
205+
await asyncio.sleep(0.02)
206+
aclose_done.set()
207+
202208
async def root():
203209
finished = asyncio.Event()
204210
holder['finished'] = finished
@@ -211,21 +217,150 @@ async def reporter():
211217
reporter_task = asyncio.create_task(reporter())
212218
tasks = {failing_task, reporter_task}
213219
await _await_shutdown_or_background_failure(finished, tasks)
220+
await _cancel_pending_tasks(tasks)
221+
await yielding_aclose()
214222
holder['after_wait'] = True
215223
holder['failing_task'] = failing_task
216224

217225
try:
218226
asyncio.run(root())
219227
self.assertTrue(holder.get('after_wait'))
220228
self.assertTrue(error_logged.wait(2.0))
221-
self.assertIn('command dispatch failed', holder.get('logged', ''))
229+
self.assertTrue(aclose_done.wait(2.0))
230+
self.assertEqual(len(holder['errors']), 1)
231+
self.assertIn('command dispatch failed', holder['errors'][0])
222232
self.assertTrue(holder['finished'].is_set())
223233
self.assertTrue(holder['failing_task'].done())
224234
self.assertIsInstance(holder['failing_task'].exception(), ValueError)
225235
finally:
226236
agent_logger.removeHandler(handler)
227237
agent_logger.setLevel(previous_level)
228238

239+
def test_clean_shutdown_does_not_log_normal_background_completion(self):
240+
"""Clean shutdown: reporters finish after _finished; cleanup must stay silent."""
241+
holder = {'errors': []}
242+
243+
class _Handler(logging.Handler):
244+
def emit(self, record):
245+
if 'Error in Python agent asyncio event loop' in record.getMessage():
246+
holder['errors'].append(record.getMessage())
247+
248+
agent_logger = logging.getLogger('skywalking')
249+
handler = _Handler()
250+
agent_logger.addHandler(handler)
251+
previous_level = agent_logger.level
252+
agent_logger.setLevel(logging.ERROR)
253+
254+
async def root():
255+
finished = asyncio.Event()
256+
257+
async def reporter():
258+
while not finished.is_set():
259+
await asyncio.sleep(0.01)
260+
261+
tasks = {asyncio.create_task(reporter()) for _ in range(2)}
262+
finished.set()
263+
await _await_shutdown_or_background_failure(finished, tasks)
264+
# Give reporters a turn to observe _finished and return normally.
265+
await asyncio.sleep(0.05)
266+
await _cancel_pending_tasks(tasks)
267+
holder['all_done'] = all(task.done() for task in tasks)
268+
269+
try:
270+
asyncio.run(root())
271+
self.assertTrue(holder.get('all_done'))
272+
self.assertEqual(holder['errors'], [])
273+
finally:
274+
agent_logger.removeHandler(handler)
275+
agent_logger.setLevel(previous_level)
276+
277+
def test_unexpected_pre_shutdown_cancellation_is_logged(self):
278+
"""Cancellation before shutdown must be logged and trigger orderly cleanup."""
279+
holder = {'errors': []}
280+
error_logged = Event()
281+
aclose_done = Event()
282+
283+
class _Handler(logging.Handler):
284+
def emit(self, record):
285+
msg = record.getMessage()
286+
if 'Error in Python agent asyncio event loop' in msg:
287+
holder['errors'].append(msg)
288+
error_logged.set()
289+
290+
agent_logger = logging.getLogger('skywalking')
291+
handler = _Handler()
292+
agent_logger.addHandler(handler)
293+
previous_level = agent_logger.level
294+
agent_logger.setLevel(logging.ERROR)
295+
296+
async def cancelled_background():
297+
asyncio.current_task().cancel()
298+
await asyncio.sleep(0)
299+
300+
async def yielding_aclose():
301+
await asyncio.sleep(0.02)
302+
aclose_done.set()
303+
304+
async def root():
305+
finished = asyncio.Event()
306+
holder['finished'] = finished
307+
308+
async def reporter():
309+
while not finished.is_set():
310+
await asyncio.sleep(0.05)
311+
312+
cancelled_task = asyncio.create_task(cancelled_background())
313+
reporter_task = asyncio.create_task(reporter())
314+
tasks = {cancelled_task, reporter_task}
315+
await _await_shutdown_or_background_failure(finished, tasks)
316+
await _cancel_pending_tasks(tasks)
317+
await yielding_aclose()
318+
holder['after_wait'] = True
319+
320+
try:
321+
asyncio.run(root())
322+
self.assertTrue(holder.get('after_wait'))
323+
self.assertTrue(error_logged.wait(2.0))
324+
self.assertTrue(aclose_done.wait(2.0))
325+
self.assertEqual(len(holder['errors']), 1)
326+
self.assertIn('cancelled unexpectedly', holder['errors'][0])
327+
self.assertTrue(holder['finished'].is_set())
328+
finally:
329+
agent_logger.removeHandler(handler)
330+
agent_logger.setLevel(previous_level)
331+
332+
def test_intentional_cleanup_cancellation_is_silent(self):
333+
"""Tasks cancelled by _cancel_pending_tasks during cleanup must not ERROR."""
334+
holder = {'errors': []}
335+
336+
class _Handler(logging.Handler):
337+
def emit(self, record):
338+
if 'Error in Python agent asyncio event loop' in record.getMessage():
339+
holder['errors'].append(record.getMessage())
340+
341+
agent_logger = logging.getLogger('skywalking')
342+
handler = _Handler()
343+
agent_logger.addHandler(handler)
344+
previous_level = agent_logger.level
345+
agent_logger.setLevel(logging.ERROR)
346+
347+
async def forever():
348+
while True:
349+
await asyncio.sleep(0.05)
350+
351+
async def root():
352+
tasks = {asyncio.create_task(forever()) for _ in range(2)}
353+
await _cancel_pending_tasks(tasks)
354+
holder['all_done'] = all(task.done() for task in tasks)
355+
356+
try:
357+
asyncio.run(root())
358+
self.assertTrue(holder.get('all_done'))
359+
self.assertEqual(holder['errors'], [])
360+
finally:
361+
agent_logger.removeHandler(handler)
362+
agent_logger.setLevel(previous_level)
363+
229364
def test_shutdown_async_queue_bounded(self):
230365
async def _run():
231366
q = asyncio.Queue()

0 commit comments

Comments
 (0)