Skip to content

Latest commit

 

History

History
167 lines (122 loc) · 11.1 KB

File metadata and controls

167 lines (122 loc) · 11.1 KB

CMP 架构

English · 简体中文 · 繁體中文

本文说明当前 API 的执行与生命周期契约。安装及完整最小程序见 README;历史验证和待完成事项见实施记录。

与 Go GMP 的机制对照、技术边界和待办见 GMP 对照文档。

执行模型

组件 职责与执行位置
Task<T> 持有协程帧与结果;惰性、单消费者,本身不创建线程
RunLoop run(task) 阻塞调用线程并在该线程消费调度队列,返回根任务的值或异常
ThreadPool 固定数量的 worker 按共享 FIFO 领取工作;完成顺序不保证
IoContext 一个私有 I/O driver 驱动原生 TCP,多条连接可共享同一 context
组合与同步原语 持有子任务或注册等待者;恢复线程由最后完成的子任务、通知者或解锁者决定

线程亲和不自动传播:任务从其他线程恢复后,通过 co_await caller.schedule() 显式返回目标执行器。阻塞代码会阻塞当前线程;没有恢复来源的挂起任务会使 run() 一直等待。

以下函数片段共用这些声明;调用时使用 README 中的 RunLoop::run():

import std;
import mcpplibs.cmp;

namespace cmp = mcpplibs::cmp;

任务与结构化并发

Task<T> 支持值和 void,不支持引用或数组结果。它可以移动构造,不能复制或移动赋值;co_await std::move(task) 消费命名任务,重复消费会终止进程。未消费 Task 销毁自己的帧;已启动任务必须保留到所有恢复来源结束,不能用销毁 Task 代替取消。

when_all() 等待时按输入顺序启动所有任务,先等待全部结束,再按输入顺序返回结果或抛出第一个异常,不会因一个子任务失败而自动取消其他任务。变参形式返回 tuple,void 对应 std::monostate;vector 形式返回同索引顺序的结果。父任务在最后完成的子任务线程继续。

下面并发等待两个定时任务;loop.run(total(loop.get_scheduler())) 返回 42:

cmp::Task<int> later(cmp::RunLoop::Scheduler scheduler, int value) {
    co_await scheduler.schedule_after(std::chrono::milliseconds { 1 });
    co_return value;
}

cmp::Task<int> total(cmp::RunLoop::Scheduler scheduler) {
    auto [first, second] = co_await cmp::when_all(
        later(scheduler, 20), later(scheduler, 22));
    co_return first + second;
}

TaskGroup 用于增量接纳 Task<void>:

  • spawn() 在返回前启动子任务,执行到第一次挂起;深层递归 spawn 应在子任务开头显式 schedule()。
  • join() 惰性且只能等待一次。join 期间活跃子任务仍可 spawn;join 已开始且活动数归零时永久关闭接纳。join 前暂时为空的 group 仍可接纳。
  • join 等待全部子任务结束,按接纳顺序选择首个异常并释放终态包装帧。join 前的保留量随累计任务数增长,可用分批 group 控制。
  • 析构时必须未使用或已完成 join,否则终止进程。接纳或作用域主体可能抛异常时,在首次 spawn 前创建 join Task,捕获错误、请求停止,再等待 join 后处理错误;可执行范例见结构化清理测试。
  • get_stop_token() 必须显式传给子任务。cancel_and_join() 在被等待时请求停止并 join,不会抢占子任务。

事件、锁和借用数据必须覆盖全部等待者的生命周期。临时捕获型协程 lambda 的闭包可能先于懒协程销毁,优先使用参数按值传入的具名协程函数。

调度与取消

API 契约
RunLoop::Scheduler schedule() 始终排队;schedule_after() / schedule_at() 使用 steady_clock,已到期或非正延迟也排队,不保证精确唤醒时刻
ThreadPool::Scheduler schedule() 始终排队,在任意 worker 恢复;不提供定时调度
schedule(token) 及定时重载 预取消也排队;取消与消费选择一个结果,取消获胜抛出 OperationCancelled,晚到的 stop 不替换已选结果

Scheduler 是可复制的弱身份句柄,不延长执行器寿命。RunLoop 必须处于有效的 run() 中;同一个 RunLoop 只允许顺序复用,嵌套或并发 run 抛出 std::logic_error,根任务结束时仍有遗留队列工作也会失败。

ThreadPool 析构关闭接纳、排空已接纳工作并 join worker;不得从自己的 worker 析构。它拥有线程,不拥有上层 Task。显式传入零个 worker 会抛出 std::invalid_argument。

取消需要显式传递 token,各 API 的优先级分别定义。停止请求不等于任务结束,释放资源前仍要等待任务收尾。

同步原语

原语 等待与通知规则
OneShotEvent co_await event;首次 set() 永久置位,在 setter 线程、set 返回前恢复已注册等待者,顺序未指定;不支持 reset 或取消
AsyncManualResetEvent wait(token) / co_await event;set 按 FIFO 恢复 pending 等待者,reset 只影响未来等待;预取消优先,set/cancel 只选一个结果,在获胜的通知线程恢复
AsyncMutex auto guard = co_await mutex.lock_async();无竞争时 inline 继续,否则 FIFO 交接,在 Guard 释放线程恢复;无 try-lock、手工 unlock 或取消

三者不可移动。事件有 pending 等待者,或 mutex 仍被持有/有人排队时析构,会终止进程。OneShotEvent 的直接嵌套 set() 链会增长调用栈;通知下一事件前显式 schedule() 可形成异步边界。

阻塞调用

为阻塞工作使用独立 ThreadPool,避免占满 CPU worker:

cmp::Task<int> offload(
    cmp::ThreadPool::Scheduler blockingWorkers,
    cmp::RunLoop::Scheduler caller) {
    co_return co_await cmp::run_blocking(blockingWorkers, caller, [] {
        std::this_thread::sleep_for(std::chrono::milliseconds { 1 });
        return 42;
    });
}

调用方创建 cmp::ThreadPool blockingWorkers { 2 },再用 loop.run(offload(blockingWorkers.get_scheduler(), loop.get_scheduler())),结果为 42。run_blocking() 持有 callable,worker 领取后执行一次;排队取消可跳过执行,开始后不能抢占。结果通过不可取消的返回调度交付;返回 Scheduler 失败时,其异常覆盖业务结果或异常。

TCP

TcpStream::connect() 和 TcpListener::bind() 仅接受数值 IPv4/IPv6。下面连接本地服务、写完请求并执行一次读取;read_some() 不保证读到完整响应,消息边界由上层协议处理:

cmp::Task<std::size_t> send_and_read_some(
    cmp::IoContext& io,
    cmp::RunLoop::Scheduler caller,
    std::uint16_t port,
    std::span<const std::byte> request,
    std::span<std::byte> reply,
    std::stop_token token = {}) {
    auto stream = co_await cmp::TcpStream::connect(
        io, caller, "127.0.0.1", port, token);
    co_await stream.write_all(caller, request, token);
    co_return co_await stream.read_some(caller, reply, token);
}

服务端通过 co_await cmp::TcpListener::bind(io, caller, "127.0.0.1", 0) 绑定临时端口,用 local_port() 取得端口,再 co_await listener.accept(caller)。完整回环用法见 TCP 测试。

  • stream/listener 只能移动构造。每个 stream 同时允许一个 read 和一个 write;每个 listener 允许一个 accept。方向占用持续到结果交付,同方向重叠抛出 std::logic_error。
  • read/write 的 span 借用底层存储直到 Task 完成。非空读取返回零表示 EOF;空缓冲区返回零不检测 EOF。EOF 在后续读取保持有效,write 方向仍可用。
  • 读取优先级为:资源与方向校验 → 缓存 EOF/错误 → 预取消 → 空缓冲区 → 原生读取。字节与 EOF 或非取消导致的原生错误同时到达时,先交付字节,下次读交付终态;缓存的非 EOF 错误只交付一次。
  • write_all() 全写或报错,报错不表示没有字节发出。预取消不启动 I/O、不关闭仍可用的 stream;已发起写的取消会关闭 stream。
  • close() 线程安全、幂等、非阻塞;关闭获胜时已接纳操作以 OperationCancelled 收尾。取消 accept 保持 listener 可用;关闭 listener 不关闭已接受的 stream。
  • 成功、错误和取消都经过显式 return Scheduler;返回调度失败会覆盖业务结果。Task 创建或调度接纳仍可能抛出分配/构造异常。
  • IoContext 关闭资源并排空 native handler 后回收 driver,不可从自己的 driver 析构。调用方仍需 join 上层任务,并让返回执行器保持可用。

实现入口

根模块为 mcpplibs.cmp,公共命名空间为 mcpplibs::cmp;CMP 不重新导出 std 或 Asio 类型。

子系统 源码
帧与取消异常 task.cppm、cancellation.cppm
调度与阻塞隔离 run_loop.cppm、thread_pool.cppm、blocking.cppm
结构化并发 when_all.cppm、task_group.cppm
同步 one_shot_event.cppm、async_manual_reset_event.cppm、async_mutex.cppm
原生 TCP tcp.cppm

RunLoop 使用 awaiter 内节点组成的侵入 FIFO 与带索引的 timer 最小堆;就绪入队/到期转移不分配,就绪取消为 O(1)、timer 调整为 O(log n),未来 timer 接纳可能分配。when_all 以 thread-local 启动队列和对称完成转移控制深层 join 的栈使用。TCP 关闭 native handle 后保留 socket 对象至已发起操作的 handler 返回。

构建与验证

mcpp.toml 声明 C++23、私有 TCP 依赖 Asio 1.38.1,以及测试依赖 gtest 1.15.2。当前不跟踪 mcpp.lock;版本固定不能替代完整供应链锁定。

范围 Linux macOS / Windows
根库与测试 LLVM 22.1.8 LLVM 22.1.8
basic 与两个 benchmark consumer GCC 16.1.0 LLVM 22.1.8;readiness 负载仅适用 POSIX

在仓库根目录用 Bash 执行,Windows 使用 Git Bash:

scripts/qualify-command.sh mcpp build --profile dev --strict --cache=off
scripts/qualify-command.sh mcpp test --profile dev --strict --cache=off
scripts/qualify-command.sh mcpp build --profile release --strict --cache=off
scripts/qualify-command.sh mcpp test --profile release --strict --cache=off
cd examples/basic
../../scripts/qualify-command.sh mcpp build --profile dev --strict --cache=off
../../scripts/qualify-command.sh mcpp run
../../scripts/qualify-command.sh mcpp build --profile release --strict --cache=off

命令门禁拒绝非零退出和 warning/error 诊断;CI按平台运行。测试清单以 mcpp test --list 为准,Linux 通过不代表其他平台已验收。性能样本、Sanitizer 与平台验证缺口单独记录在实施报告,不作为稳定 API 保证。