@@ -583,3 +583,40 @@ def test_worker_start_exception_handling():
583583 output_queue .put .assert_called_with (QueueSignals .error )
584584 # Verify that the worker was cleaned up (num_active_workers decremented)
585585 assert num_active_workers .value == 0
586+
587+
588+ class ProcessFailingWorker (Worker ):
589+ @classmethod
590+ def start (cls , ** kwargs : Any ) -> "ProcessFailingWorker" :
591+ return cls ()
592+
593+ def process (self , items : Iterable [tuple [int , Any ]]) -> Iterable [tuple [int , Any ]]:
594+ raise RuntimeError ("Processing failed" )
595+
596+
597+ def test_worker_processing_exception_handling ():
598+ """Test that _worker handles exceptions during worker.process() gracefully."""
599+ input_queue = MagicMock ()
600+ output_queue = MagicMock ()
601+ num_active_workers = MagicMock ()
602+ num_active_workers .get_lock .return_value .__enter__ = MagicMock ()
603+ num_active_workers .get_lock .return_value .__exit__ = MagicMock ()
604+ num_active_workers .value = 1
605+ worker_id = 0
606+
607+ # Mock _get_items_from_queue to return a single item then stop
608+ # This triggers worker.process() and then the exception.
609+ with patch ("qwen3_embed.parallel_processor._get_items_from_queue" , return_value = [(0 , "item" )]):
610+ _worker (
611+ worker_class = ProcessFailingWorker ,
612+ input_queue = input_queue ,
613+ output_queue = output_queue ,
614+ num_active_workers = num_active_workers ,
615+ worker_id = worker_id ,
616+ kwargs = {},
617+ )
618+
619+ # Verify that QueueSignals.error was put in the output queue
620+ output_queue .put .assert_called_with (QueueSignals .error )
621+ # Verify that the worker was cleaned up (num_active_workers decremented)
622+ assert num_active_workers .value == 0
0 commit comments