11"""Error → callback mapping tests for JobManager._run_job."""
2+
23from __future__ import annotations
34
45from contextlib import asynccontextmanager
5- from unittest .mock import MagicMock , AsyncMock , patch
6+ from unittest .mock import AsyncMock , MagicMock , patch
67
78import pytest
89
1112
1213class _FakeSlot :
1314 """Minimal async context manager that does nothing."""
14- async def __aenter__ (self ): return self
15- async def __aexit__ (self , * args ): pass
15+
16+ async def __aenter__ (self ):
17+ return self
18+
19+ async def __aexit__ (self , * args ):
20+ pass
1621
1722
1823@asynccontextmanager
@@ -44,9 +49,7 @@ def record():
4449
4550async def _run_job_with_error (mgr , rec , error ):
4651 mgr ._jobs [rec .job_id ] = rec
47- with patch (
48- "src.gateway.rate_limit.RateLimiter.concurrent_slot" , _fake_concurrent_slot
49- ):
52+ with patch ("src.gateway.rate_limit.RateLimiter.concurrent_slot" , _fake_concurrent_slot ):
5053 with patch .object (mgr , "_run_workflow_mode" ) as mock_run :
5154 mock_run .side_effect = error
5255 await mgr ._run_job (rec .job_id )
@@ -62,11 +65,15 @@ async def test_maps_to_suspended_not_error(self, manager, record):
6265 record .callback .on_error = AsyncMock ()
6366
6467 await _run_job_with_error (
65- manager , record ,
68+ manager ,
69+ record ,
6670 BatchPendingError (
67- "batch submitted" , batch_job_id = "batch_123" ,
71+ "batch submitted" ,
72+ batch_job_id = "batch_123" ,
6873 handle = MagicMock (provider = "gemini" ),
69- requests = [], items_snapshot = [], output_field = "analysis" ,
74+ requests = [],
75+ items_snapshot = [],
76+ output_field = "analysis" ,
7077 ),
7178 )
7279
@@ -84,11 +91,15 @@ async def test_notifies_with_batch_id(self, manager, record):
8491 record .callback .on_progress = AsyncMock ()
8592
8693 await _run_job_with_error (
87- manager , record ,
94+ manager ,
95+ record ,
8896 BatchPendingError (
89- "batch submitted" , batch_job_id = "batch_abc" ,
97+ "batch submitted" ,
98+ batch_job_id = "batch_abc" ,
9099 handle = MagicMock (provider = "gemini" ),
91- requests = [], items_snapshot = [], output_field = "analysis" ,
100+ requests = [],
101+ items_snapshot = [],
102+ output_field = "analysis" ,
92103 ),
93104 )
94105
@@ -104,7 +115,8 @@ async def test_maps_to_failed_with_on_error(self, manager, record):
104115 record .callback .on_error = AsyncMock ()
105116
106117 await _run_job_with_error (
107- manager , record ,
118+ manager ,
119+ record ,
108120 RetryableError ("rate limited" , http_status = 429 , provider = "test" ),
109121 )
110122
@@ -120,7 +132,8 @@ async def test_notify_for_concurrent_limit(self, manager, record):
120132 record .callback .on_error = AsyncMock ()
121133
122134 await _run_job_with_error (
123- manager , record ,
135+ manager ,
136+ record ,
124137 RuntimeError ("concurrent limit reached for entry_type=cli_workflow" ),
125138 )
126139
0 commit comments