fix(eval): allow awaiting Python workflow helpers
This commit is contained in:
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed Python Eval `parallel()` and `pipeline()` rejecting `await` after their work completed, allowing synchronous and awaited calls without repeating operations.
|
||||
|
||||
## [17.2.15] - 2026-08-12
|
||||
|
||||
### Added
|
||||
|
||||
@@ -569,6 +569,14 @@ if "__omp_prelude_loaded__" not in globals():
|
||||
return 0
|
||||
return n if n > 0 else 0
|
||||
|
||||
class _AwaitableList(list):
|
||||
"""Completed list result accepted by both sync and ``await`` syntax."""
|
||||
|
||||
def __await__(self):
|
||||
yield from ()
|
||||
return self
|
||||
|
||||
|
||||
def _pool_map(items, fn):
|
||||
"""Run ``fn`` over ``items`` through a bounded thread pool.
|
||||
|
||||
@@ -582,10 +590,10 @@ if "__omp_prelude_loaded__" not in globals():
|
||||
|
||||
items = list(items)
|
||||
if not items:
|
||||
return []
|
||||
return _AwaitableList()
|
||||
limit = _concurrency_limit()
|
||||
workers = min(limit, len(items)) if limit > 0 else len(items)
|
||||
results = [None] * len(items)
|
||||
results = _AwaitableList(None for _ in items)
|
||||
errors = {}
|
||||
with concurrent.futures.ThreadPoolExecutor(max_workers=workers) as pool:
|
||||
futures = {}
|
||||
@@ -621,7 +629,7 @@ if "__omp_prelude_loaded__" not in globals():
|
||||
stage). Stage 1 receives the original item; later stages receive the
|
||||
previous stage's result. Pool width tracks ``task.maxConcurrency``.
|
||||
"""
|
||||
current = list(items)
|
||||
current = _AwaitableList(items)
|
||||
for stage in stages:
|
||||
if not callable(stage):
|
||||
raise TypeError("pipeline() stages must be callables")
|
||||
|
||||
@@ -31,6 +31,31 @@ describe.skipIf(!SHOULD_RUN)("python eval workflow helpers", () => {
|
||||
}
|
||||
});
|
||||
|
||||
it("parallel and pipeline results may be awaited without repeating work", async () => {
|
||||
using tempDir = TempDir.createSync("@eval-workflow-awaitable-results-");
|
||||
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
|
||||
try {
|
||||
const code = [
|
||||
"calls = []",
|
||||
"def mark(value):",
|
||||
" calls.append(value)",
|
||||
" return value",
|
||||
"sync_parallel = parallel([lambda: mark('sync')])",
|
||||
"awaited_parallel = await parallel([lambda: mark('awaited')])",
|
||||
"sync_pipeline = pipeline([1], lambda value: value + 1)",
|
||||
"awaited_pipeline = await pipeline([1], lambda value: value + 2)",
|
||||
"empty_parallel = await parallel([])",
|
||||
"empty_pipeline = await pipeline([])",
|
||||
"print(sync_parallel, awaited_parallel, sync_pipeline, awaited_pipeline, empty_parallel, empty_pipeline, calls)",
|
||||
].join("\n");
|
||||
const result = await executePythonWithKernel(kernel, code);
|
||||
expect(result.exitCode).toBe(0);
|
||||
expect(result.output).toContain("['sync'] ['awaited'] [2] [3] [] [] ['sync', 'awaited']");
|
||||
} finally {
|
||||
await kernel.shutdown();
|
||||
}
|
||||
});
|
||||
|
||||
it("parallel runs thunks concurrently", async () => {
|
||||
using tempDir = TempDir.createSync("@eval-workflow-parallel-concurrent-");
|
||||
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
|
||||
|
||||
Reference in New Issue
Block a user