From 21a95a285eb579cd13884fd9300d3df70dba40bf Mon Sep 17 00:00:00 2001 From: devtechedge Date: Wed, 9 Sep 2026 23:08:17 +0000 Subject: [PATCH 1/2] fix(test): create ModelTests before concurrent pool workers ModelTest.create_test() can call to_datetime()/ttl_cache (time.time()) while another worker freezes time via time_machine for execution_time. Building tests on the main thread removes that race. Fixes #6039 Signed-off-by: devtechedge --- sqlmesh/core/test/runner.py | 19 +++++++++---------- 1 file changed, 9 insertions(+), 10 deletions(-) diff --git a/sqlmesh/core/test/runner.py b/sqlmesh/core/test/runner.py index 284558e1c8..f17693fc39 100644 --- a/sqlmesh/core/test/runner.py +++ b/sqlmesh/core/test/runner.py @@ -125,9 +125,11 @@ def run_tests( # Ensure workers are not greater than the number of tests num_workers = min(len(model_test_metadata) or 1, default_test_connection.concurrent_tasks) - def _run_single_test( - metadata: ModelTestMetadata, engine_adapter: EngineAdapter - ) -> t.Optional[ModelTextTestResult]: + # Build ModelTest instances on the calling thread before workers start. create_test() + # can call to_datetime() / ttl_cache (time.time()), which races with another worker's + # time_machine freeze when execution_time is set under concurrent_tasks > 1. + tests: list[ModelTest] = [] + for metadata, engine_adapter in metadata_to_adapter.items(): test = ModelTest.create_test( body=metadata.body, test_name=metadata.test_name, @@ -140,10 +142,10 @@ def _run_single_test( concurrency=num_workers > 1, verbosity=verbosity, ) + if test: + tests.append(test) - if not test: - return None - + def _run_single_test(test: ModelTest) -> ModelTextTestResult: result = t.cast( ModelTextTestResult, ModelTextTestRunner().run(t.cast(unittest.TestCase, test)), @@ -159,10 +161,7 @@ def _run_single_test( start_time = time.perf_counter() try: with ThreadPoolExecutor(max_workers=num_workers) as pool: - futures = [ - pool.submit(_run_single_test, metadata=metadata, engine_adapter=engine_adapter) - for metadata, engine_adapter in metadata_to_adapter.items() - ] + futures = [pool.submit(_run_single_test, test) for test in tests] for future in concurrent.futures.as_completed(futures): test_results.append(future.result()) From 4d515a72feba60b4801e81d8d966be58cbc73670 Mon Sep 17 00:00:00 2001 From: Dev M Date: Thu, 10 Sep 2026 00:00:08 +0000 Subject: [PATCH 2/2] fix(test): create ModelTests inside connection cleanup try Move calling-thread create_test into the same try/finally that closes engine adapters so invalid tests still clean up connections. Add a NOTE about a possible future parallel create stage that must not overlap runs. Signed-off-by: devtechedge Signed-off-by: Dev M --- sqlmesh/core/test/runner.py | 42 +++++++++++++++++++------------------ 1 file changed, 22 insertions(+), 20 deletions(-) diff --git a/sqlmesh/core/test/runner.py b/sqlmesh/core/test/runner.py index f17693fc39..17a4f55d69 100644 --- a/sqlmesh/core/test/runner.py +++ b/sqlmesh/core/test/runner.py @@ -125,26 +125,6 @@ def run_tests( # Ensure workers are not greater than the number of tests num_workers = min(len(model_test_metadata) or 1, default_test_connection.concurrent_tasks) - # Build ModelTest instances on the calling thread before workers start. create_test() - # can call to_datetime() / ttl_cache (time.time()), which races with another worker's - # time_machine freeze when execution_time is set under concurrent_tasks > 1. - tests: list[ModelTest] = [] - for metadata, engine_adapter in metadata_to_adapter.items(): - test = ModelTest.create_test( - body=metadata.body, - test_name=metadata.test_name, - models=models, - engine_adapter=engine_adapter, - dialect=dialect, - path=metadata.path, - default_catalog=default_catalog, - preserve_fixtures=preserve_fixtures, - concurrency=num_workers > 1, - verbosity=verbosity, - ) - if test: - tests.append(test) - def _run_single_test(test: ModelTest) -> ModelTextTestResult: result = t.cast( ModelTextTestResult, @@ -160,6 +140,28 @@ def _run_single_test(test: ModelTest) -> ModelTextTestResult: start_time = time.perf_counter() try: + # Build ModelTest instances on the calling thread before workers start. create_test() + # can call to_datetime() / ttl_cache (time.time()), which races with another worker's + # time_machine freeze when execution_time is set under concurrent_tasks > 1. + # NOTE: We can run create_tests in a separate parallel stage for a future optimization. + # We just can't overlap runs/creations. + tests: list[ModelTest] = [] + for metadata, engine_adapter in metadata_to_adapter.items(): + test = ModelTest.create_test( + body=metadata.body, + test_name=metadata.test_name, + models=models, + engine_adapter=engine_adapter, + dialect=dialect, + path=metadata.path, + default_catalog=default_catalog, + preserve_fixtures=preserve_fixtures, + concurrency=num_workers > 1, + verbosity=verbosity, + ) + if test: + tests.append(test) + with ThreadPoolExecutor(max_workers=num_workers) as pool: futures = [pool.submit(_run_single_test, test) for test in tests]