From 8826219173fe1b8479c109ff4127221d67b34ee8 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 14:51:38 +0800 Subject: [PATCH] refactor: bound shared thread pool shutdown --- app/runtime/thread.py | 102 +++++++++++++++-- .../backend-architecture-next-stage.md | 21 +++- docs/rules/05-architecture.md | 2 + mypy.ini | 1 + tests/conftest.py | 4 +- tests/test_lifecycle_shutdown.py | 12 ++ tests/test_thread_helper_shutdown.py | 103 ++++++++++++++++++ 7 files changed, 229 insertions(+), 16 deletions(-) create mode 100644 tests/test_thread_helper_shutdown.py diff --git a/app/runtime/thread.py b/app/runtime/thread.py index dbcac5b69..93cfa3d1f 100644 --- a/app/runtime/thread.py +++ b/app/runtime/thread.py @@ -1,21 +1,96 @@ -from concurrent.futures import ThreadPoolExecutor +import threading +import time +from concurrent.futures import Future, ThreadPoolExecutor, wait from contextvars import copy_context +from typing import Any, Callable, TypeVar, cast +from app.foundation.singleton import Singleton from app.runtime.settings import RuntimeSettingsCompat settings = RuntimeSettingsCompat() -from app.foundation.singleton import Singleton + +_Result = TypeVar("_Result") +_THREAD_POOL_STOP_TIMEOUT_SECONDS = 10.0 -class ThreadHelper(metaclass=Singleton): +class _OwnedThreadPoolExecutor(ThreadPoolExecutor): """ - 线程池管理 + 追踪所有提交入口的共享执行器,包括旧调用方直接使用的 ``pool.submit``。 """ - def __init__(self): + + def __init__(self, max_workers: int) -> None: + """初始化线程池及其 Future owner 集合。""" + super().__init__(max_workers=max_workers) + self._ownership_lock = threading.RLock() + self._owned_futures: set[Future[Any]] = set() + + def submit( + self, + fn: Callable[..., _Result], + /, + *args: Any, + **kwargs: Any, + ) -> Future[_Result]: + """提交任务并在其达到终态前保留 owner。""" + with self._ownership_lock: + future = super().submit(fn, *args, **kwargs) + self._owned_futures.add(future) + future.add_done_callback(self._discard_future) + return future + + def _discard_future(self, future: Future[Any]) -> None: + """任务达到终态后释放 owner 记录。""" + with self._ownership_lock: + self._owned_futures.discard(future) + + def shutdown_bounded(self, timeout: float) -> bool: + """ + 封口新提交并有限等待全部已接受任务。 + + :param timeout: 等待 Future 达到终态的最长秒数 + :return: 所有任务与 worker 均已终止时返回 True,否则返回 False + """ + deadline = time.monotonic() + max(0.0, timeout) + with self._ownership_lock: + # 不取消排队工作,保持历史 shutdown(wait=True) 的完成语义。 + super().shutdown(wait=False) + owned_futures = tuple(self._owned_futures) + if owned_futures: + _, pending_futures = wait( + owned_futures, + timeout=max(0.0, deadline - time.monotonic()), + ) + if pending_futures: + return False + # 标准库只提供无界 wait=True;Future 又会先标记完成再执行 done callback, + # 因此封口后读取稳定 worker 集合,复用同一 deadline 做有限 join。 + worker_threads = tuple(cast(set[threading.Thread], self._threads)) + current_thread = threading.current_thread() + for worker_thread in worker_threads: + if worker_thread is current_thread: + continue + worker_thread.join( + timeout=max(0.0, deadline - time.monotonic()), + ) + return all(not worker_thread.is_alive() for worker_thread in worker_threads) + + +# strict mypy 跳过 foundation 实现导入,因此无法在本文件解析既有 Singleton 元类类型。 +class ThreadHelper(metaclass=Singleton): # type: ignore[metaclass] + """ + 共享后台线程池 owner,负责关联上下文传播和生命周期收敛。 + """ + + def __init__(self) -> None: """按系统配置创建共享后台线程池。""" - self.pool = ThreadPoolExecutor(max_workers=settings.CONF.threadpool) + self.pool = _OwnedThreadPoolExecutor(max_workers=settings.CONF.threadpool) - def submit(self, func, *args, **kwargs): + def submit( + self, + func: Callable[..., _Result], + *args: Any, + **kwargs: Any, + ) -> Future[_Result]: """ 提交任务 :param func: 函数 @@ -26,9 +101,14 @@ class ThreadHelper(metaclass=Singleton): context = copy_context() return self.pool.submit(context.run, func, *args, **kwargs) - def shutdown(self): + def shutdown( + self, + timeout: float = _THREAD_POOL_STOP_TIMEOUT_SECONDS, + ) -> bool: """ - 关闭线程池 - :return: + 有限等待共享线程池关闭。 + + :param timeout: 等待已接受任务达到终态的最长秒数 + :return: 全部任务和 worker 均已终止时返回 True,否则返回 False """ - self.pool.shutdown() + return self.pool.shutdown_bounded(timeout=timeout) diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 5479014d8..20dc25d57 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -2,13 +2,13 @@ > 文档性质:当前架构复核、优秀 Python 后端实践对标、AI 可执行任务手册 > 适用仓库:`MoviePilot`,分支 `v3` -> 审计基线:`415335b2`(2026-08-24) +> 审计基线:`ad45bfac`(2026-08-24) > 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本 > 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文 > 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md` -> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner;阶段 43 已统一 Discord typing 的异步 owner 和 shutdown 收尾;阶段 44 已清除 WebAgent 测试临时事件循环提前关闭产生的 CI 红注解;阶段 45 已统一影视与字幕搜索的请求级逐页任务编排;阶段 46 已收口启动性能门禁的托管 runner 假失败与诊断输出;阶段 47 已补齐 Agent 渠道流式刷新任务的重入 owner;阶段 48 已统一工件上传 action 的 Node 24 主版本;阶段 49 已统一插件安装的同步/异步代际解析事实源;阶段 50 已统一插件市场 GitHub 请求降级策略;阶段 51 已统一插件索引请求与响应三态策略;阶段 52 已统一插件 Release 分页策略;阶段 53 已统一远端插件安装模式决策;阶段 54 已补齐同步安装成功后的临时回滚备份清理;阶段 55~56 已收口官方插件观察基线与报告保留策略;阶段 57 已统一进程级运行时 Facade 门禁并补齐 ModuleManager 边界;阶段 58 已消除 AgentTask 关闭回归的跨线程零时长等待竞态;阶段 59 已统一 Feishu 多实例长连接的 SDK 循环路由;阶段 60 已清除命令服务虚假的关停 owner 声明;阶段 61 已统一 Capability Runtime 同步/异步关闭的诚实收敛结果;阶段 62 已统一消息渠道长连接的多实例关闭收敛合同;阶段 63 已补齐应用消息队列线程的关闭收敛合同。 +> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner;阶段 43 已统一 Discord typing 的异步 owner 和 shutdown 收尾;阶段 44 已清除 WebAgent 测试临时事件循环提前关闭产生的 CI 红注解;阶段 45 已统一影视与字幕搜索的请求级逐页任务编排;阶段 46 已收口启动性能门禁的托管 runner 假失败与诊断输出;阶段 47 已补齐 Agent 渠道流式刷新任务的重入 owner;阶段 48 已统一工件上传 action 的 Node 24 主版本;阶段 49 已统一插件安装的同步/异步代际解析事实源;阶段 50 已统一插件市场 GitHub 请求降级策略;阶段 51 已统一插件索引请求与响应三态策略;阶段 52 已统一插件 Release 分页策略;阶段 53 已统一远端插件安装模式决策;阶段 54 已补齐同步安装成功后的临时回滚备份清理;阶段 55~56 已收口官方插件观察基线与报告保留策略;阶段 57 已统一进程级运行时 Facade 门禁并补齐 ModuleManager 边界;阶段 58 已消除 AgentTask 关闭回归的跨线程零时长等待竞态;阶段 59 已统一 Feishu 多实例长连接的 SDK 循环路由;阶段 60 已清除命令服务虚假的关停 owner 声明;阶段 61 已统一 Capability Runtime 同步/异步关闭的诚实收敛结果;阶段 62 已统一消息渠道长连接的多实例关闭收敛合同;阶段 63 已补齐应用消息队列线程的关闭收敛合同;阶段 64 已统一共享线程池的有界关闭 owner。 > 当前 canonical 状态:API/Application 公共复杂度基线已清零,组合根外 `SystemConfigOper()` 构造和 Model/Oper 隐式事务均为 0;命名 Chain/Agent 数据端口、TaskRegistry owner、Module Contract V2、typed Event、Outbox durable intent、请求关联和插件运行时 getter 已形成当前路径。插件仓适配、未知第三方 fallback 和其它 E1/E3 副作用仍按风险持续治理。 -> 最新阶段:阶段 63 已补齐应用消息队列线程的关闭收敛合同。 +> 最新阶段:阶段 64 已统一共享线程池的有界关闭 owner。 ## 当前复核结论(2026-08-24) @@ -672,6 +672,21 @@ 消息/生命周期专项 58 项、架构与兼容专项 122 项、Pylint 10.00/10、strict mypy 39 文件、宿主与质量 ratchet 以及四分片全量 `5941 passed, 3 skipped` 均通过。 +### 长期整改阶段 64:共享线程池有界关闭 owner 统一(2026-08-24) + +- 阶段 4 的消息渠道回环和阶段 60 的命令重建都声明由 `ThreadHelper` 统一持有,但原实现没有登记 Future, + `shutdown()` 直接执行无界 `ThreadPoolExecutor.shutdown(wait=True)`。任一同步任务阻塞时会卡住 + `stop_modules()` 所在事件循环,外层生命周期预算无法取消,且无法向阶段 61 的布尔合同报告未收敛。 +- 当前共享 executor 在 Future 达到终态前保留 owner,先封口新提交,再使用 10 秒默认预算有限等待; + 超时返回 `False` 且不取消正在执行或排队的历史工作,任务完成后同一 owner 可重试到真实终态。 + 追踪同时覆盖宿主 `ThreadHelper.submit()` 和旧调用方直接使用的 `.pool.submit()`,避免形成第二套旁路。 +- `ThreadHelper()`、`submit()`、无参数 `shutdown()`、公开 `.pool` 及 `app.helper.thread` 精确映射均保留; + `.pool` 仍是 `ThreadPoolExecutor` 子类,上下文传播、任务返回和异常语义不变。未修改 SDK/Compat 清单、 + 插件 Hook 或插件仓,V1/V2/V3 插件边界保持兼容。 +- 宿主/旧兼容提交、阻塞任务、排队任务、阻塞完成回调、提交封口、重试终态和 startup 失败传播均有 + 故障注入;线程池/生命周期专项 78 项、架构与兼容调用链 164 项、Pylint 10.00/10、strict mypy + 40 文件、宿主与质量 ratchet 及四分片全量 `5946 passed, 3 skipped` 均通过。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/docs/rules/05-architecture.md b/docs/rules/05-architecture.md index 390727fb5..62737190e 100644 --- a/docs/rules/05-architecture.md +++ b/docs/rules/05-architecture.md @@ -199,6 +199,8 @@ ModuleManager 与 startup 组合根继续关闭其余资源但必须向上返回 长连接、轮询或 Socket 服务只有在真实终止后才能返回成功,超时 owner 不得清空句柄。 应用消息队列的监控线程遵守同一收敛语义:停止必须有限等待,回调阻塞导致线程仍存活时保留 owner 并向 startup 返回 `False`,不得用无界 `join()` 阻塞生命周期或把日志当作成功。 +共享 `ThreadHelper` 必须追踪通过宿主 `submit()` 和旧兼容 `.pool.submit()` 接受的全部 Future;关闭时 +先封口新任务,再有限等待且保留未终止 owner,结果由 startup 聚合,不得恢复无界 executor shutdown。 API 中允许丢失或可重建的进程内任务必须登记到 `app/runtime/tasks.py`;登记器先于其他 运行资源启动,并在资源释放前停止接收、取消和有限等待。需要崩溃恢复的 E2/E3 副作用仍应 进入 Outbox 或持久任务表,不能把 TaskRegistry 当成 durable queue。 diff --git a/mypy.ini b/mypy.ini index 4316ae754..e44180357 100644 --- a/mypy.ini +++ b/mypy.ini @@ -14,6 +14,7 @@ files = app/runtime/correlation.py, app/runtime/coalesce.py, app/runtime/observability/__init__.py, + app/runtime/thread.py, app/runtime/event/contracts.py, app/runtime/event/errors.py, app/runtime/extensions/module/contracts.py, diff --git a/tests/conftest.py b/tests/conftest.py index f29583aa4..4ec119bcd 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -478,8 +478,8 @@ def pytest_sessionfinish(session, exitstatus): from app.runtime.thread import ThreadHelper helper = ThreadHelper.get_existing_instance() - if helper: - helper.shutdown() + if helper and helper.shutdown() is False: + raise RuntimeError("shared thread pool did not converge") except Exception as err: _report_session_cleanup_error(session, "thread helper", err) diff --git a/tests/test_lifecycle_shutdown.py b/tests/test_lifecycle_shutdown.py index d0e20b1af..55d140714 100644 --- a/tests/test_lifecycle_shutdown.py +++ b/tests/test_lifecycle_shutdown.py @@ -969,6 +969,18 @@ def test_stop_modules_propagates_message_queue_nonconvergence(monkeypatch): _assert_completed_once(dependency) +def test_stop_modules_propagates_shared_thread_pool_nonconvergence(monkeypatch): + """共享线程池仍有活动 Future 时必须由模块服务关闭结果向上暴露。""" + dependencies = _patch_module_shutdown_dependencies(monkeypatch) + dependencies["thread"].return_value = False + + converged = asyncio.run(modules_initializer.stop_modules()) + + assert converged is False + for dependency in dependencies.values(): + _assert_completed_once(dependency) + + def test_stop_modules_drains_web_agent_tasks_before_persistence(monkeypatch): """关闭时先收口 Web Agent,再关闭持久化准入和数据库任务。""" order = [] diff --git a/tests/test_thread_helper_shutdown.py b/tests/test_thread_helper_shutdown.py new file mode 100644 index 000000000..f1b37c4ca --- /dev/null +++ b/tests/test_thread_helper_shutdown.py @@ -0,0 +1,103 @@ +import threading +import time +from concurrent.futures import ThreadPoolExecutor + +import pytest + +from app.runtime import thread as thread_module +from app.runtime.thread import ThreadHelper + + +def _new_thread_helper() -> ThreadHelper: + """构造不进入全局 Singleton 的隔离线程池 owner。""" + helper = object.__new__(ThreadHelper) + helper.__init__() + return helper + + +@pytest.mark.parametrize("use_legacy_pool", [False, True]) +def test_thread_helper_shutdown_is_bounded_and_retryable(use_legacy_pool): + """宿主 submit 和旧 pool.submit 都必须进入同一个可重试关闭 owner。""" + helper = _new_thread_helper() + entered = threading.Event() + release = threading.Event() + + def blocked_work() -> str: + """模拟无法由线程池强制取消的同步任务。""" + entered.set() + release.wait() + return "done" + + assert isinstance(helper.pool, ThreadPoolExecutor) + submit = helper.pool.submit if use_legacy_pool else helper.submit + future = submit(blocked_work) + try: + assert entered.wait(timeout=1) + started_at = time.monotonic() + assert helper.shutdown(timeout=0.01) is False + assert time.monotonic() - started_at < 1 + assert not future.done() + + with pytest.raises(RuntimeError): + helper.submit(lambda: None) + with pytest.raises(RuntimeError): + helper.pool.submit(lambda: None) + finally: + release.set() + + assert future.result(timeout=1) == "done" + assert helper.shutdown(timeout=1) is True + + +def test_thread_helper_shutdown_preserves_queued_work(): + """关闭封口不得取消已接受的排队任务,保持历史完成语义。""" + executor = thread_module._OwnedThreadPoolExecutor(max_workers=1) + entered = threading.Event() + release = threading.Event() + + def blocked_work() -> str: + """占用唯一 worker,确保后一任务仍在队列中。""" + entered.set() + release.wait() + return "first" + + first = executor.submit(blocked_work) + queued = executor.submit(lambda: "queued") + try: + assert entered.wait(timeout=1) + assert executor.shutdown_bounded(timeout=0.01) is False + assert not queued.cancelled() + finally: + release.set() + + assert first.result(timeout=1) == "first" + assert queued.result(timeout=1) == "queued" + assert executor.shutdown_bounded(timeout=1) is True + + +def test_thread_helper_shutdown_waits_for_worker_after_future_completion(): + """Future 已完成但用户回调仍阻塞时不得误报 worker 已收敛。""" + executor = thread_module._OwnedThreadPoolExecutor(max_workers=1) + work_release = threading.Event() + callback_entered = threading.Event() + callback_release = threading.Event() + + future = executor.submit(lambda: work_release.wait()) + + def blocked_done_callback(_future) -> None: + """模拟 Future 终态之后仍占用 worker 的第三方完成回调。""" + callback_entered.set() + callback_release.wait() + + future.add_done_callback(blocked_done_callback) + work_release.set() + try: + assert callback_entered.wait(timeout=1) + assert future.done() + started_at = time.monotonic() + assert executor.shutdown_bounded(timeout=0.01) is False + assert time.monotonic() - started_at < 1 + finally: + callback_release.set() + + assert executor.shutdown_bounded(timeout=1) is True