| 27 | |
| 28 | |
| 29 | class BatchRunTests(unittest.IsolatedAsyncioTestCase): |
| 30 | async def test_happy_path(self) -> None: |
| 31 | async def flow(input: RawText, *, cfg: Any = None) -> str: |
| 32 | return f"ok:{input.text}" |
| 33 | |
| 34 | inputs = [RawText(text=str(i)) for i in range(5)] |
| 35 | result = await batch_run(flow, inputs, concurrency=3) |
| 36 | self.assertEqual(result.total, 5) |
| 37 | self.assertEqual(result.success_count, 5) |
| 38 | self.assertEqual(result.failure_count, 0) |
| 39 | self.assertEqual(result.results, [f"ok:{i}" for i in range(5)]) |
| 40 | self.assertEqual(result.errors, []) |
| 41 | self.assertGreaterEqual(result.duration_seconds, 0) |
| 42 | self.assertEqual(result.tokens_total, {}) |
| 43 | self.assertEqual(result.cost_estimate_usd, 0.0) |
| 44 | |
| 45 | async def test_empty_inputs(self) -> None: |
| 46 | async def flow(input: RawText, *, cfg: Any = None) -> str: |
| 47 | return "x" |
| 48 | |
| 49 | result = await batch_run(flow, [], concurrency=4) |
| 50 | self.assertEqual(result.total, 0) |
| 51 | self.assertEqual(result.success_count, 0) |
| 52 | self.assertEqual(result.results, []) |
| 53 | |
| 54 | async def test_on_error_skip_collects(self) -> None: |
| 55 | async def flow(input: RawText, *, cfg: Any = None) -> str: |
| 56 | if input.text in ("2", "4"): |
| 57 | raise ValueError(f"bad:{input.text}") |
| 58 | return f"ok:{input.text}" |
| 59 | |
| 60 | inputs = [RawText(text=str(i)) for i in range(5)] |
| 61 | result = await batch_run(flow, inputs, on_error="skip") |
| 62 | self.assertEqual(result.success_count, 3) |
| 63 | self.assertEqual(result.failure_count, 2) |
| 64 | self.assertEqual(result.results[0], "ok:0") |
| 65 | self.assertIsNone(result.results[2]) |
| 66 | self.assertIsNone(result.results[4]) |
| 67 | self.assertEqual([i for i, _ in result.errors], [2, 4]) |
| 68 | |
| 69 | async def test_on_error_raise_propagates(self) -> None: |
| 70 | boom = RuntimeError("boom") |
| 71 | |
| 72 | async def flow(input: RawText, *, cfg: Any = None) -> str: |
| 73 | raise boom |
| 74 | |
| 75 | with self.assertRaises(RuntimeError) as ctx: |
| 76 | await batch_run( |
| 77 | flow, |
| 78 | [RawText(text="x")], |
| 79 | concurrency=1, |
| 80 | on_error="raise", |
| 81 | ) |
| 82 | self.assertIs(ctx.exception, boom) |
| 83 | |
| 84 | async def test_memory_kwarg_rejected(self) -> None: |
| 85 | async def flow(input: RawText, *, cfg: Any = None) -> str: |
| 86 | return "x" |
nothing calls this directly
no outgoing calls
no test coverage detected