Skip to content

Latest commit

 

History

History
167 lines (122 loc) · 12.4 KB

File metadata and controls

167 lines (122 loc) · 12.4 KB

CMP architecture

English · 简体中文 · 繁體中文

This document defines the current execution and lifetime contracts. See the README for installation and a complete minimal program, and the implementation record (Chinese) for historical validation and remaining work.

See the Go GMP comparison (Chinese) for scheduling mechanisms, technical limits, and follow-up tasks.

Execution model

Component Responsibility and execution context
Task<T> Owns a coroutine frame and result; lazy and single-consumer, with no thread of its own
RunLoop run(task) blocks its caller and consumes scheduled work on that thread, returning the root result or exception
ThreadPool A fixed number of workers claim work from a shared FIFO; completion order is not guaranteed
IoContext One private I/O driver runs native TCP operations; multiple connections can share a context
Composition and synchronization Own child tasks or register waiters; the last completing child, notifier, or unlocking thread determines where execution resumes

Thread affinity does not propagate automatically. Use co_await caller.schedule() to return explicitly after resuming elsewhere. Blocking code blocks its current thread; a suspended task without a resumption source can leave run() waiting indefinitely.

The function snippets below share these declarations and use RunLoop::run() as shown in the README:

import std;
import mcpplibs.cmp;

namespace cmp = mcpplibs::cmp;

Tasks and structured concurrency

Task<T> supports values and void, but not reference or array results. It is move-constructible, not copyable or move-assignable. Consume a named task with co_await std::move(task); consuming it twice terminates the process. An unconsumed Task destroys its frame. Keep a started task alive until all resumption sources have finished; destroying a Task is not cancellation.

Awaiting when_all() starts every input in input order and waits for all of them before returning ordered results or throwing the first exception in input order. A child failure does not automatically cancel siblings. The variadic form returns a tuple with std::monostate for void; the vector form preserves indices. The parent resumes on the last completing child's thread.

This example waits for two timers concurrently; loop.run(total(loop.get_scheduler())) returns 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 incrementally accepts Task<void>:

  • spawn() starts the child before returning and runs it until its first suspension. Deep recursive spawn chains should explicitly schedule() at the start of each child.
  • join() is lazy and single-use. Active children may still spawn during join; admission closes permanently when join has started and the active count reaches zero. An empty group can still accept work before join starts.
  • Join waits for every child, selects the first exception in admission order, and releases terminal wrapper frames. Retention before join grows with cumulative admission; batch scopes can bound it.
  • Destruction requires an unused or fully joined group; otherwise the process terminates. If admission or the scope body can throw, create the join Task before the first spawn, capture errors and request stop, then await join before handling errors. See the executable cleanup patterns.
  • Pass get_stop_token() explicitly to children. Awaiting cancel_and_join() requests stop and joins; it does not preempt children.

Events, mutexes, and borrowed data must outlive all their waiters. A temporary capturing coroutine lambda may lose its closure before the lazy coroutine runs; prefer named coroutine functions with value parameters.

Scheduling and cancellation

API Contract
RunLoop::Scheduler schedule() always queues; schedule_after() / schedule_at() use steady_clock. Expired or nonpositive delays still queue; exact wake-up times are not guaranteed
ThreadPool::Scheduler schedule() always queues and resumes on any worker; no timer scheduling
schedule(token) and timer overloads Pre-cancelled waits still queue. Cancellation and consumption select one outcome; winning cancellation throws OperationCancelled, and late stop requests do not replace the selected outcome

Schedulers are copyable weak identity handles and do not extend executor lifetimes. A RunLoop must have an active run(). The same RunLoop allows sequential reuse only; nested or concurrent runs throw std::logic_error, and a root finishing with leftover queued work also fails.

ThreadPool destruction closes admission, drains accepted work, and joins workers; never destroy a pool on its own worker. The pool owns threads, not caller Tasks. An explicit worker count of zero throws std::invalid_argument.

Tokens are passed explicitly, and each API defines its own priority rules. A stop request is not completion; wait for tasks to finish before releasing resources.

Synchronization primitives

Primitive Waiting and notification
OneShotEvent co_await event; the first set() stays set and resumes registered waiters on the setter thread before returning, in unspecified order; no reset or cancellation
AsyncManualResetEvent wait(token) / co_await event; set resumes pending waiters FIFO, reset affects future waits only. Pre-cancellation wins; set/cancel chooses one outcome and resumes on the winning notification thread
AsyncMutex auto guard = co_await mutex.lock_async(); uncontended acquisition continues inline, otherwise ownership passes FIFO on the Guard-releasing thread; no try-lock, manual unlock, or cancellation

All three are immovable. Destroying an event with pending waiters, or a mutex that is held or has waiters, terminates the process. Direct nested OneShotEvent set() calls grow the native stack; explicitly schedule() before signalling the next event to introduce an asynchronous boundary.

Blocking calls

Use a separate ThreadPool for blocking work so it cannot occupy all CPU workers:

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;
    });
}

Create cmp::ThreadPool blockingWorkers { 2 }, then call loop.run(offload(blockingWorkers.get_scheduler(), loop.get_scheduler())) to obtain 42. run_blocking() owns its callable and invokes it once after worker claim. Queued cancellation can skip the call; a started call cannot be preempted. Outcomes go through a non-cancellable return hop; if the return Scheduler fails, its exception overrides the business result or exception.

TCP

TcpStream::connect() and TcpListener::bind() accept numeric IPv4/IPv6 addresses only. This example connects to a local service, writes the request, and performs one read. read_some() does not guarantee a complete response; the application protocol defines message boundaries:

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);
}

A server binds an ephemeral port with co_await cmp::TcpListener::bind(io, caller, "127.0.0.1", 0), obtains it with local_port(), and awaits listener.accept(caller). See TCP tests for complete loopback usage.

  • Streams and listeners are move-constructible only. A stream permits one read and one write concurrently; a listener permits one accept. Admission remains held through result delivery; overlapping operations in the same direction throw std::logic_error.
  • Read/write spans borrow storage until the Task completes. Zero from a nonempty read means EOF; an empty read returning zero does not detect EOF. EOF remains sticky for reads and leaves writing available.
  • Read priority is: resource/admission checks → cached EOF/error → pre-cancellation → empty buffer → native read. Bytes arriving with EOF or a native error not caused by cancellation are delivered first, with the terminal result delivered on the next read; cached non-EOF errors are delivered once.
  • write_all() writes everything or reports an error; an error does not mean no bytes were sent. Pre-cancellation starts no I/O and preserves a usable stream; cancelling an initiated write closes the stream.
  • close() is thread-safe, idempotent, and non-blocking. If close wins, accepted operations finish with OperationCancelled. Cancelling accept preserves the listener; closing it does not close accepted streams.
  • Success, error, and cancellation go through the explicit return Scheduler; return-hop failure overrides the business outcome. Task creation or scheduling admission can still throw allocation/construction exceptions.
  • IoContext closes resources and drains native handlers before reclaiming its driver; never destroy it on its own driver. Callers must still join application Tasks and keep the return executor available.

Implementation map

The root module is mcpplibs.cmp, with public namespace mcpplibs::cmp. CMP does not re-export std or Asio types.

Subsystem Source
Frames and cancellation exception task.cppm, cancellation.cppm
Scheduling and blocking isolation run_loop.cppm, thread_pool.cppm, blocking.cppm
Structured concurrency when_all.cppm, task_group.cppm
Synchronization one_shot_event.cppm, async_manual_reset_event.cppm, async_mutex.cppm
Native TCP tcp.cppm

RunLoop uses an intrusive FIFO of awaiter-owned nodes and an indexed timer min-heap. Ready enqueue and expiry promotion do not allocate; ready cancellation is O(1), timer adjustment O(log n), and future-timer admission may allocate. when_all uses a thread-local startup queue and symmetric completion transfer to bound nested join stack usage. TCP retains socket objects after closing native handles until handlers for initiated operations return.

Build and validation

mcpp.toml declares C++23, private TCP dependency Asio 1.38.1, and test dependency gtest 1.15.2. mcpp.lock is not tracked; version pins do not replace a complete supply-chain lock.

Scope Linux macOS / Windows
Root library and tests LLVM 22.1.8 LLVM 22.1.8
basic and both benchmark consumers GCC 16.1.0 LLVM 22.1.8; the readiness workload is POSIX-only

Run from the repository root in Bash, or Git Bash on Windows:

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

The command gate rejects nonzero exits and warning/error diagnostics; CI runs per platform. Use mcpp test --list for test discovery. Passing Linux checks does not qualify other platforms. Performance samples, Sanitizer limits, and platform gaps are recorded separately in the implementation report, not as stable API guarantees.