# AsyncFramework **Repository Path**: yixiao666/async-framework ## Basic Information - **Project Name**: AsyncFramework - **Description**: No description available - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-04-10 - **Last Updated**: 2026-04-11 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # BaseAsyncTask - 可扩展异步任务框架 基于 C++17 的线程安全异步任务框架,提供线程池、定时调度、周期任务、线程健康监控、多线程同步组和异步状态机能力。零外部依赖,仅依赖 C++17 标准库和 pthread。 ## 架构总览 ``` ┌──────────────────────────────────────────────────────────────────┐ │ Application Layer │ │ (用户自定义 IPeriodicTask / IState 实现) │ ├──────────────────┬──────────────────┬────────────────────────────┤ │ PeriodicTask │ AsyncStateMachine│ Future Modules... │ │ Framework │ Framework │ (OneShot, EventBus) │ ├──────────────────┤ │ │ │ TaskSupervisor │ │ │ │ (健康监控/同步组)│ │ │ ├──────────────────┴──────────────────┴────────────────────────────┤ │ TaskScheduler │ │ (定时调度, 小顶堆优先级队列, 单调度线程) │ ├──────────────────────────────────────────────────────────────────┤ │ ThreadPool │ │ (固定大小工作线程, 任务队列, future 支持) │ ├──────────────────────────────────────────────────────────────────┤ │ Core Primitives │ │ (ThreadSafeQueue, CancellationToken, TaskHandle, HqLogger) │ └──────────────────────────────────────────────────────────────────┘ ``` ## 核心特性 ### 1. 线程池 (ThreadPool) - 固定大小线程池,默认 `hardware_concurrency()` 个工作线程 - `submit()` 接受任意可调用对象,返回 `std::future` - `submit_cancellable()` 返回 `TaskHandle`,支持协作式取消和等待 - 基于 `ThreadSafeQueue` 的 MPMC 任务分发 ### 2. 任务调度器 (TaskScheduler) - 独立调度线程,基于小顶堆管理定时任务 - `schedule_periodic()` — 周期性调度,返回 `CancellationToken` 用于取消 - `schedule_once()` — 一次性延迟调度 - 调度线程通过 `condition_variable::wait_until` 精确等待下一个截止时间 ### 3. 周期任务框架 (PeriodicTask) 用户继承 `IPeriodicTask` 实现自定义周期任务: ```cpp class MyTask : public IPeriodicTask { void execute(const CancellationToken& token) override { if (token.is_cancelled()) return; // 任务逻辑 } std::string task_name() const override { return "my_task"; } std::chrono::milliseconds period() const override { return std::chrono::seconds(1); } // 可选回调 void on_start() override {} void on_stop() override {} bool on_error(const std::exception& e) override { return true; } }; ``` ### 4. 线程健康监控 (TaskSupervisor) - 独立监控线程,周期性轮询已注册任务的 `is_healthy()` 状态 - 连续失败计数达到阈值后标记为不健康 - 健康状态变化时触发 `on_health_change` 回调 - 通过 `query_health()` 查询任意任务的实时健康状态 ```cpp // 启用健康监控 class MyTask : public IPeriodicTask { bool supports_health_check() const override { return true; } bool is_healthy() const override { return /* 自定义健康判断 */; } int max_consecutive_failures() const override { return 3; } }; // 注册健康变化回调 fw.task_supervisor().on_health_change([](const std::string& name, bool healthy) { // 处理健康状态变化 }); ``` ### 5. 多线程同步组 (Sync Group) - 同一同步组内的任务由 `TaskSupervisor` 统一调度,保证同时派发 - 使用组内最短周期作为调度间隔 - 适用于生产者-消费者等需要协同执行的场景 ```cpp class ProducerTask : public IPeriodicTask { bool supports_sync() const override { return true; } std::string sync_group() const override { return "pipeline"; } // ... }; ``` ### 6. 异步状态机框架 (AsyncStateMachine) - 串行化事件处理(专用事件线程),彻底消除并发状态切换 bug - 状态切换流程:cancel 当前任务 → on_exit → 切换 → on_enter - 支持转换表 (transition table) 和状态内部事件处理(fallback) - 状态内派发的任务自动跟踪,状态退出时自动取消 - 事件支持类型擦除的 payload(`std::any`) ```cpp class IdleState : public IState { StateId id() const override { return "idle"; } void on_enter(IStateMachine& machine, const Event& trigger) override { machine.dispatch_task([](const CancellationToken& ct) { while (!ct.is_cancelled()) { /* 后台任务 */ } }); } void on_exit(IStateMachine& machine, const Event& trigger) override {} std::optional handle_event( IStateMachine& machine, const Event& event) override { if (event.type == "start") return StateId("processing"); return std::nullopt; } }; ``` ### 7. 日志系统 (HqLogger) - 线程安全的流式日志输出,基于 mutex + RAII LogStream - 支持毫秒级时间戳 - 五级日志:ERROR / WARN / INFO / DEBUG / VERBOSE - 宏接口:`HQLOGE` / `HQLOGW` / `HQLOGI` / `HQLOGD` / `HQLOGV` ## 目录结构 ``` BaseAsyncTask/ ├── CMakeLists.txt # 顶层构建配置 ├── src/ │ ├── core/ # 核心基础设施 │ │ ├── cancellation_token.h # 协作式取消令牌 (atomic) │ │ ├── task_handle.h # 任务句柄 (future + token) │ │ ├── thread_safe_queue.h # 线程安全 MPMC 队列 │ │ ├── thread_pool.h # 固定大小线程池 │ │ ├── task_scheduler.h/.cpp # 定时调度器 (小顶堆) │ │ ├── async_framework.h/.cpp # 框架生命周期管理 │ │ └── HqLogger.h # 线程安全日志 │ ├── periodic/ # 周期任务框架 │ │ ├── i_periodic_task.h # 周期任务接口 │ │ ├── periodic_task_manager.h/.cpp # 周期任务管理器 │ │ └── task_supervisor.h/.cpp # 健康监控 + 同步组 │ └── state_machine/ # 异步状态机框架 │ ├── event.h # 事件定义 (type + any payload) │ ├── i_state.h # 状态接口 │ ├── i_state_machine.h # 状态机接口 (dispatch_task/post_event) │ └── async_state_machine.h/.cpp # 状态机实现 └── examples/ ├── periodic_example.cpp # 周期任务示例 ├── supervisor_example.cpp # 健康监控 + 同步组示例 ├── state_machine_example.cpp # 状态机基础示例 ├── context_sharing_example.cpp # 状态间上下文共享示例 └── light_control_example.cpp # 灯光控制状态机示例 ``` ## 编译 ```bash mkdir build && cd build cmake .. make -j$(nproc) ``` 生成的可执行文件: - `examples/periodic_example` — 周期任务(心跳 + 倒计时) - `examples/supervisor_example` — 健康监控与同步组 - `examples/state_machine_example` — 三状态机(Idle → Processing → Done) - `examples/context_sharing_example` — 状态间上下文传递 - `examples/light_control_example` — 灯光控制(渐亮/渐灭/中断) ## 使用示例 ### 初始化框架 ```cpp #include "src/core/async_framework.h" int main() { AsyncFramework framework; AsyncFramework::Config config; config.thread_count = 4; framework.init(config); // 使用框架... framework.shutdown(); // 按正确顺序关闭所有模块 return 0; } ``` ### 注册周期任务 ```cpp auto& mgr = framework.periodic_task_manager(); mgr.register_task(std::make_unique()); mgr.stop_task("task_name"); // 停止特定任务 mgr.stop_all(); // 停止所有任务 ``` ### 创建状态机 ```cpp auto& sm = framework.create_state_machine("main_sm"); sm.add_state(std::make_unique()); sm.add_state(std::make_unique()); sm.add_state(std::make_unique()); // 可选:声明式转换表 sm.add_transition("idle", "start", "processing"); sm.add_transition("processing", "done", "done"); sm.set_initial_state("idle"); sm.start(); // 发送事件(带 payload) sm.post_event(Event{"start", Color{255, 0, 0}}); sm.stop(); ``` ### 优雅关闭 ```cpp // 方式一:主线程直接关闭 framework.shutdown(); // 方式二:等待外部信号 // 在信号处理器中调用 framework.request_shutdown() framework.wait_for_shutdown(); ``` ## 线程安全保证 | 组件 | 同步机制 | 说明 | |------|---------|------| | ThreadSafeQueue | mutex + condition_variable | MPMC 队列,支持优雅关闭 | | ThreadPool | 委托给 ThreadSafeQueue | 工作线程只通过队列交互 | | TaskScheduler | mutex + condition_variable | 单调度线程 + 堆保护 | | PeriodicTaskManager | mutex | 注册/移除操作保护 | | TaskSupervisor | mutex | 健康状态 + 同步组保护 | | AsyncStateMachine (状态) | state_mutex_ | 保护当前状态指针 | | AsyncStateMachine (任务) | tasks_mutex_ (独立) | 避免持有状态锁时派发任务导致死锁 | | AsyncStateMachine (事件) | 串行化事件循环 | 单线程处理,消除并发状态切换 | | CancellationToken | atomic\ | 无锁,任务热循环中检查 | | HqLogger | mutex + RAII LogStream | 保证单条日志原子输出 | ## 关键设计决策 1. **共享线程池** — 所有模块共享单一 ThreadPool,避免线程泛滥,统一背压控制 2. **串行化事件循环** — 状态机事件单线程处理,彻底消除并发状态切换 bug 3. **协作式取消** — 任务定期检查 `CancellationToken`,支持优雅退出 4. **单调度线程** — 定时逻辑集中管理,工作线程只做实际工作 5. **双 mutex 策略** — 状态机使用独立的 state_mutex_ 和 tasks_mutex_,防止 on_enter 中 dispatch_task 导致死锁 6. **虚函数扩展** — 用户运行时继承 IPeriodicTask / IState,支持异构容器 7. **监控与调度分离** — TaskSupervisor 独立于 TaskScheduler,健康监控不影响任务调度 ## 扩展性 - **新任务类型**:继承 `IPeriodicTask` 或直接使用 `ThreadPool::submit()` - **新模块**:只需依赖 `std::shared_ptr` 或 `TaskScheduler` - **事件总线**:可基于 `ThreadPool` 实现发布-订阅模式 - **一次性任务**:已有 `TaskScheduler::schedule_once()` - **健康监控**:实现 `supports_health_check()` + `is_healthy()` 即可接入 ## 依赖 - C++17 标准库(``, ``, ``, ``, ``, ``) - pthread(通过 CMake `Threads::Threads`) - 无外部第三方库 ## 注意事项 1. **取消检查**:长时间运行的任务必须定期检查 `token.is_cancelled()`,否则会阻塞状态切换 2. **关闭顺序**:`AsyncFramework::shutdown()` 按 状态机 → 监控器 → 周期管理器 → 调度器 → 线程池 的顺序关闭 3. **异常处理**:周期任务的 `execute()` 抛出异常会触发 `on_error()`,返回 `false` 停止任务 4. **同步组**:同步组任务由 TaskSupervisor 调度,不经过 TaskScheduler 5. **状态锁**:不要在持有 state_mutex_ 时手动调用 dispatch_task(),使用 IStateMachine 接口