Skip to content

Latest commit

Β 

History

History

Folders and files

NameName
Last commit message
Last commit date

parent directory

..
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 

README.md

rocketmq-runtime

Crates.io License

rocketmq-runtime is the shared runtime substrate for the rocketmq-rust workspace. It builds on Tokio to provide runtime ownership, tracked service and operation tasks, periodic scheduling, bounded blocking execution, resource budgets, metadata persistence, and shutdown diagnostics.

δΈ­ζ–‡ζ–‡ζ‘£

Unreleased API migration and contract corrections

Runtime Model

Production entrypoints own a RuntimeOwner. Libraries receive a ChildServiceContext or a narrower capability such as TaskSpawner; they do not discover or construct an independent runtime. Tasks submitted through these capabilities belong to a task ownership tree that is tracked through shutdown. The resource-budget tree accounts for explicitly reserved resources separately.

flowchart TD
    Entry["Application entrypoint"] --> Owner["RuntimeOwner"]
    Owner --> Root["RootServiceContext<br/>root TaskGroup"]
    Root --> A["ChildServiceContext A<br/>component TaskGroup A"]
    Root --> B["ChildServiceContext B<br/>sibling TaskGroup B"]
    A --> A1["ChildServiceContext A.1<br/>child component TaskGroup"]
    A1 --> A2["ChildServiceContext A.1.1<br/>another child component TaskGroup"]
    A1 --> Tasks["Service and worker tasks<br/>TaskKind labels each task"]
    A1 --> Operations["Operation tasks<br/>same component TaskGroup"]
    A1 --> Scheduled["ScheduledTaskGroup<br/>child TaskGroup"]
    Scheduled --> Driver["ScheduledDriver task"]
    Scheduled --> Runs["ScheduledRun tasks"]
    Root -.-> Blocking["Shared blocking lanes<br/>separate task registry"]
    A1 -.->|scoped admission| Blocking
    Owner -.-> Resources["RuntimeResources<br/>separate resource-budget tree"]
Loading

Solid arrows show the context and task-group ownership path; dotted arrows show shared capabilities outside the task-group tree.

Each component(...) call creates a child TaskGroup and returns a ChildServiceContext that can create further child components of the same kind; A, A.1, and A.1.1 illustrate that recursive relationship. A group may own both child groups and tasks. TaskKind classifies tasks rather than defining another tree level. OperationContext adds cancellation and a deadline to tasks in their existing component group; scheduled_tasks(...) creates a child group for its tracked driver and runs. Blocking lanes share a process-wide capacity and have their own task registry, while component scopes govern admission. Shutdown and diagnostics account for the owned subtree. RuntimeContext is a migration and test harness for an existing Tokio runtime.

Core Architecture

See implementation boundaries and completion contracts for the state owner of admission, settlement, probes and metadata persistence.

Type Responsibility
RuntimeConfig Worker threads, managed blocking capacity, thread name and stack size, keep-alive, fallback shutdown timeout, IO/time drivers, and per-lane blocking policies.
RuntimeOwner / RuntimeOwnerPlan Validate configuration, build and own a Tokio multi-thread runtime, expose the root context, and coordinate shutdown.
RootServiceContext Non-cloneable root with no public constructor; derives component contexts and exposes shared resources and diagnostics.
ChildServiceContext / TaskSpawner Component capabilities for owned work. A spawner exposes task submission and cancellation access without raw runtime access.
TaskGroup / OperationContext Track component tasks and provide operation-local cancellation, deadlines, and bounded waits. An operation does not create a new task group.
ScheduledTaskGroup Run periodic jobs with explicit overlap behavior and schedule metrics. It is the only scheduler.
ServiceManager / ServiceTask Run a wakeup-driven service loop as a service task of an owned group, with deadline-bounded shutdown.
BlockingExecutor Admit short blocking work through a bounded lane and retain its capacity until the closure actually exits.
RuntimeResources / ResourceBudget Share the process budget and derive component limits for count, retained bytes, and optional rate control.
ResourcePermit / BudgetedQueue Carry RAII reservations through queued or in-flight work and apply explicit overload policies.
MetadataIoActor Own bounded metadata snapshots, coalesce queued generations, and report durable completion.
ServiceLifecycle / ShutdownDeadline Coordinate readiness, liveness, shutdown requests, and a shared absolute shutdown deadline.
ShutdownReport Serializable evidence of completion, cancellation, aborts, failures, panics, timeouts, and remaining work.
RuntimeDiagnosticsSnapshot / RuntimeDiagnosticsViewV1 Internal runtime details and a bounded, sanitized operational view.

RuntimeHandle is an internal implementation type, not a public integration entrypoint. Common ownership types are also available from rocketmq_runtime::prelude.

The recommended path for a service that owns its runtime is documented once, as an example that runs as a test, on rocketmq_runtime::prelude: own the runtime, derive one component context from the sealed root, register services, register bounded periodic work instead of driving a raw loop, then drain final I/O inside the shutdown budget and read the shutdown report. RuntimeContext is the migration and test harness rather than a production entry point.

Runtime Ownership And Quick Start

Use RuntimeOwner::new()? for the default profile. For a named or customized profile, call RuntimeOwner::plan(config)?.build()?. Planning validates configuration without discovering system resources or starting Tokio; building performs memory-limit discovery and runtime construction.

This finite example registers a service and immediately exercises its cooperative shutdown path:

use rocketmq_runtime::{RuntimeConfig, RuntimeOwner};

fn main() -> Result<(), Box<dyn std::error::Error>> {
    let owner = RuntimeOwner::plan(RuntimeConfig::broker_default())?.build()?;
    let broker = owner.root_context().component("broker");
    let cancellation = broker.task_group().cancellation_token();

    broker.spawn_service("heartbeat", async move {
        cancellation.cancelled().await;
        // Finish any ordered asynchronous cleanup here.
    })?;

    let report = owner.shutdown_runtime_blocking()?;
    assert!(report.is_healthy(), "{}", report.to_json());
    Ok(())
}

A real entrypoint runs its startup and service future through owner.block_on(...), then consumes the owner outside the async context to shut down the runtime. See the broker entrypoint for integration with ServiceLifecycle and a shared shutdown deadline.

The error channels are intentional:

  • RuntimeContractViolation identifies invalid caller configuration or an invariant violation, including failures from plan().
  • RuntimeResult<T> contains RuntimeError for operational failures such as runtime construction, I/O, capacity, or timeout failures. Branch on RuntimeError::kind() rather than on the operation label: a submission to a closing owner reports RuntimeErrorKind::Closed, and one to a poisoned group reports RuntimeErrorKind::Poisoned.
  • Normal outcomes such as ScheduledTaskRegistrationStatus::AlreadyPresent, BudgetRejection, and metadata target conflicts have their own types.

RuntimeOperation names what failed. Its named variants are operations this crate performs; a crate built on the runtime labels its own failures with a constant such as RuntimeOperation::external("initialize-broker") instead of adding variants here. RuntimeOperation and RuntimeContractPolicy are non-exhaustive, so a match outside this crate needs a wildcard arm.

There is no automatic conversion from RuntimeContractViolation to RuntimeError. The example's application-level error type accepts both; applications can instead define an explicit startup error enum.

RuntimeConfig::for_parallelism derives worker and blocking-lane limits from the supplied CPU parallelism. The default uses std::thread::available_parallelism(), with a fallback of four workers. with_max_blocking_threads validates an override and caps lane concurrency.

Use RuntimeContext::try_from_current only in migration or test harnesses already running inside Tokio. It shuts down registered RocketMQ work without owning or closing the host runtime. Its resource budget is a permissive test budget, not the production memory-discovery path.

Task Scopes And Cancellation

Create long-lived component scopes through component(...). Use ChildServiceContext::try_component(...) to receive a creation error during shutdown or poisoning; component(...) returns a closed scope if the parent no longer accepts children. The first closed answer per owner is logged, and every one is counted in TaskGroup::event_counts() and noted in the owner's shutdown report when it is still being assembled. Validate dynamic names with ScopeId::try_new; string literals have a static-name conversion.

Cloning a context or task group shares the same owner and cancellation token. Creating a child gives it independent cancellation: parent cancellation propagates downward, while child cancellation does not cancel its parent or siblings. Dropping a context handle is not a graceful shutdown protocol; active tasks can keep their group alive.

Submission API Cancellation behavior
spawn / spawn_service Tracks the future. The service must observe its cancellation signal and perform ordered cleanup itself.
spawn_cancellable_service Drops the service future when the owner is cancelled. Use when immediate future cancellation is safe.
spawn_operation Keeps work under a fixed component owner and observes both owner cancellation and the operation's cancellation/deadline.
spawn_draining_operation Accepted work can continue after owner cancellation; operation cancellation/deadline and task-group shutdown still bound it.
spawn_with_handle Returns a join handle for a specific task while retaining group tracking.

For bounded requests or restartable work, use OperationContext instead of creating a component group for every operation. close_admission() stops new operation tasks; wait() waits until no operation task is active and the owner has settled them; cancel_and_wait() also requests cancellation. The waits require the operation's original component owner and abort unfinished work at their deadline. wait(timeout) and cancel_and_wait(timeout) then allow up to one second to confirm that the aborted futures were dropped, so they can return up to one second after timeout. wait_until(deadline) and cancel_and_wait_until(deadline) add no wait allowance beyond ShutdownDeadline: they request aborts and return, and the owner's shutdown report accounts for any task still being dropped. An operation only counts its active tasks: the tasks themselves are tagged in the owner's registry, so an idle operation costs a few hundred bytes.

Use wait_with_policy(owner, OperationWaitPolicy::new(graceful, confirmation)) to close admission and distinguish Completed, AbortConfirmed, and Unconfirmed { remaining_tasks }. Both deadlines are absolute; confirmation also caps the grace period. Call cancel() first to request cooperative cancellation. Legacy boolean false does not confirm destruction. Runtime starvation or blocking destructors can delay timer polling; these APIs do not promise a hard wall-clock bound.

TaskGroup::cancel() only broadcasts cancellation. Use shutdown(...) or shutdown_until(...) to close task admission and wait for shutdown evidence.

TaskGroup Invariants

State Meaning
Open New tasks and child groups can be registered.
Closing Shutdown sealed admission; cancellation has not yet reached every task.
Closed Admission is closed and cancellation has been broadcast; owned work is draining.
ShutdownCompleted The shutdown report is cached for repeated calls; it need not be healthy.
Poisoned A tracked task panicked while the group was open; new registration is rejected.

A shutdown moves a group through Closing, Closed, and ShutdownCompleted in that order. A closing group reports Closed as soon as its cancellation token is cancelled, so a task woken by the shutdown already observes Closed, even on another worker thread. Poisoning is fail-stop: the transition is logged once with the group path and counted in TaskGroup::event_counts().

Task metadata is registered under a spawn gate that serializes registration with shutdown transitions; the task is handed to Tokio after the gate is released. The tracker token taken under the gate keeps a shutdown waiting for the task, and an abort requested before its handle is installed is honored. Task names are TaskName values: a &'static str is stored without allocating, and String or Arc<str> names are kept as they are. The child registry uses weak references keyed by TaskGroupId; dropping the last group reference unregisters the child. Group ids are unique within the process, including across owners. Names are labels, so multiple groups can share a name without sharing identity.

Scheduled Tasks

ScheduledTaskGroup is the only scheduler. Derive one with context.scheduled_tasks("maintenance"), or wrap an owned group with ScheduledTaskGroup::new(group), and register work:

Work Entry point Ownership
Periodic maintenance schedule(config, policy, task) task is an FnMut returning a future. The driver and its runs belong to the group.
Operation-bound maintenance schedule_operation(&operation, config, policy, task) As schedule; it also stops at the operation's cancellation or deadline.
Fixed-delay work that can end itself schedule_controlled(config, task) A run that returns ScheduledTaskControl::Stop ends the schedule and counts as one completed run.

The configuration alone decides the timing mode, and the policy must agree with it:

Configuration Timing Policy
ScheduledTaskConfig::fixed_delay(name, period) Run, then wait period after the run completes. ScheduledExecutionPolicy::serial(..)
ScheduledTaskConfig::fixed_rate_no_overlap(name, period) Tick every period; runs never overlap. ScheduledExecutionPolicy::serial(..)
ScheduledTaskConfig::fixed_rate(name, period) Tick every period; up to n runs overlap. ScheduledExecutionPolicy::bounded(n, ..)

A zero period or a policy that contradicts the configuration is rejected with an unsupported error. ScheduledExecutionPolicy::default() allows one run at a time and skips a tick that arrives while it runs.

Fixed-rate schedules keep absolute ticks. When a tick finds every run slot busy, the missed-tick policy decides what happens to it: Skip discards it, CoalesceLatest keeps one pending run, and BoundedCatchUp(n) keeps up to n pending runs. The drift metric records how late each tick fired.

  • with_initial_delay(delay) sets the delay before the first run; it defaults to zero, allowing the first run immediately.
  • max_run_time bounds an individual run by dropping its future on timeout. External side effects still need an appropriate cancellation contract.
  • Duplicate names return AlreadyPresent without replacing the driver or metrics. clear_completed() clears registrations only when the scheduler's group has no active tasks.
  • Drivers and runs belong to the scheduler's task group. Ordinary runs may finish during shutdown; operation-bound registrations also observe their operation's cancellation and deadline.

Snapshots record active runs, run completions, skips, overlaps, failures, drift, and elapsed time. Close the scheduler through shutdown(timeout), shutdown_until(deadline), or its owning group.

Blocking Work

Use the executors supplied by a ChildServiceContext:

Accessor Lane Typical work
storage_io() StorageIo Short storage and filesystem operations.
metadata_io() MetadataIo Metadata persistence.
cpu_crypto() CpuCrypto Bounded CPU or cryptographic work.

Managed lanes share one global admission budget per runtime owner, bounded by RuntimeConfig::max_blocking_threads. Each lane has its own concurrency ceiling and queue bound. Idle capacity can be borrowed; a waiting lane's reservation is protected from new borrowers. Cloning an executor or deriving a context shares this capacity rather than creating another pool.

Waiters of one lane are admitted in arrival order, and a new submission never passes a queued waiter of its own lane. A released slot is handed directly to the waiter that can use it, which wakes only that waiter; a lane still below its reservation is served first. The Tokio blocking pool holds RuntimeConfig::tokio_blocking_threads() threads: the managed capacity plus max(2, capacity / 8) threads of headroom, so direct spawn_blocking calls, DNS resolution, and tokio::fs do not take the threads admitted work relies on.

max_queue_depth rejects submissions when the admission queue is full. queue_timeout bounds waiting for execution capacity; task_timeout bounds the caller's wait after admission. spawn_until and spawn_io_until also cap both phases with one absolute deadline. These methods require an active Tokio context; call them from the owning runtime's work.

Timeout or cancellation does not stop an already-running blocking closure. The closure retains its admission permit until it exits. A completion guard inside the closure removes its task record on exit; there is no separate reaper task. Cancelling a queued submission removes its queued record, while an abandoned running submission is recorded as TimedOutStillRunning until completion. See the executor implementation.

BlockingKind::LongRunning is rejected. Long-running blocking loops need a dedicated OS-thread or domain-service owner with a stop and join protocol. BlockingExecutor::new(policy, owner_group) creates an isolated executor for tests and adapters: it has an independent budget and task table outside the managed root lanes and their diagnostics. It is registered with the group's tree, so a closure it is still running counts in the blocking_still_running of that tree owner's shutdown report.

Prefer the explicit new_isolated constructor; new delegates to it for compatibility. Call stop_admission() and shutdown_until(deadline) to close and await just that executor and its clones. Queued submissions wake with a closed error. The report distinguishes completed shutdown from still-owned submissions or closures; an incomplete wait can be retried. It does not close the supplied group or sibling executors. Managed lanes reject these methods: their service context owns shutdown.

Resource Budgets And Queues

RuntimeOwner owns RuntimeResources; child contexts share its process budget. Derive narrower limits with context.process_budget().child(...). Do not create an independent ResourceBudgetTree in each component when a shared process limit is required.

The owner detects a memory limit from ROCKETMQ_PROCESS_MEMORY_LIMIT_BYTES, the process's own cgroup membership, or host physical memory. On Linux it reads /proc/self/cgroup and /proc/self/mountinfo and takes the smallest finite hard limit visible along the membership path, ignoring the cgroup v2 max and cgroup v1 unlimited markers. The detected limit and the chargeable budget are separate numbers: RuntimeOwnerPlan::with_memory_policy derives the managed budget from the effective limit as the whole limit (the default), an explicit byte count, or a bounded fraction with reserved headroom. Supply an explicit ProcessMemoryLimit through RuntimeOwnerPlan::with_memory_limit when needed. These limits account for resources admitted through the budget APIs; they do not automatically limit every process allocation or resident-memory usage, and a constraint the process cannot see is not discovered.

ResourceBudget checks count, retained bytes, and optional rate limits along the ancestor chain. Count and byte capacity is reserved with atomic operations; a request that fails at a later level rolls back the levels it reserved, so no level ever exceeds its limit, and only a level with a rate limit takes a lock. A request racing such a rollback near capacity can be rejected even though the capacity is about to return. A ResourcePermit retains count and byte reservations until dropped. BudgetClass::Control can use configured control reserves; data work cannot consume that reserved capacity. Same-tree permit rebinding keeps common-ancestor accounting while moving ownership between components.

ResourcePermit::try_resize(bytes) changes only the byte reservation. Growth reserves the delta at each ancestor and rolls back on rejection without taking another count or rate token. Shrinking remains allowed after dynamic closure; release the removed payload before shrinking its charge. Growth observes the current budget's dynamic admission gates, including after a rebind.

BudgetedQueue supports Reject, WaitUntilDeadline, CoalesceLatest, DropStale, and CloseSlowConsumer. Select a policy matching whether work can wait, be replaced, or be discarded. push_until waits for capacity only with WaitUntilDeadline and preserves the rejected item in its outcome.

The dequeue API determines the accounting lifetime:

  • try_pop() / recv() return the item and release its permit at dequeue.
  • try_pop_budgeted() / recv_budgeted() return a BudgetedItem retaining the permit during processing. into_parts() transfers the permit explicitly; into_item() releases it.

See resource-budget tests for ancestor limits, control reserves, overload handling, and permit transfer examples.

Metadata Persistence

Start the actor with MetadataIoConfig::default().into_plan()?.start(&context)?. It owns a tracked coordinator and uses the context's shared MetadataIo blocking lane. Configure actor admission with max_pending_operations and max_pending_bytes; configure the managed lane through RuntimeConfig::blocking_lane_policies.metadata_io. The actor's compatibility blocking_* settings do not replace that shared lane policy.

By default an actor writes one resource at a time, so a slow write delays its other resources. MetadataIoPlan::with_max_concurrent_writes(n) lets independent resources proceed while one write is slow, up to the metadata lane's concurrency; generations of one resource are still written one at a time and in order. effective_profile() reports the value in force. Every queued or in-flight snapshot is also charged to the owner's process budget until its write completes, including after its observer times out, so the actors of one owner share that budget. A snapshot the process budget cannot take is refused like one above max_pending_bytes.

The actual blocking closure owns its snapshot and permit together, including after cancellation destroys the actor or coordinator. The actor's local pending counters settle afterwards and may briefly lag the shared ledger. Coalescing reserves only a positive size delta; failed growth retains the old request and charge, and equal-size replacement works at the process limit. Shrinking drops the old payload before releasing excess bytes. Caller-owned allocations before admission, and caller-retained snapshot clones, are outside this actor-owned accounting.

submit and submit_next accept immutable snapshots without waiting for durability. Match MetadataWriteSubmissionStatus: Accepted provides a receipt; TargetConflict returns the request when a resource already has pending work for a different target. Wait for persistence through the receipt's wait_until, or use submit_durable / submit_next_durable and match the durable-generation or target-conflict outcome.

Queued generations for the same logical resource can coalesce; a newer durable generation can satisfy an earlier waiter. A local write uses a temporary file, file synchronization, atomic replacement, and parent-directory synchronization on supported platforms before advancing durable generation. A wait timeout does not prove that the underlying filesystem operation has stopped.

Call stop_admission() to reject new snapshots, then shutdown_until(MetadataDeadline) to drain accepted work. Inspect the returned MetadataIoShutdownReport for unfinished generations as well as the runtime's task shutdown report. See metadata I/O tests.

Service Lifecycle And Shutdown

ServiceLifecycle exposes Starting, Ready, Draining, Stopped, and Failed states. Start it under a component context, mark readiness after startup, and publish dependency readiness separately. Maintenance can suspend readiness without marking the process dead. Liveness checks lifecycle state and progress freshness, not whether a business port is open.

A task registered as critical makes its failure observable: record it with spawn_critical_service for a service that must run until cancelled, or spawn_critical for work whose panic alone is a failure. The failure stays pending in the CriticalFailureState until a handler takes it, so a full notification channel cannot lose it, and only an owner cancellation is treated as an expected exit, which keeps an ordinary shutdown out of failure handling. ServiceLifecycle::spawn_critical_failure_monitor applies a CriticalFailureRecovery policy: revoke readiness only, fail the service, or fail it and request an ordered shutdown. Spawn that monitor under an owner outside the monitored group, because a poisoned group cannot run its own monitor. Liveness stays independent, so a process that keeps progressing still reports not-ready after a handled critical failure.

With ServiceLifecycle::from_env, ROCKETMQ_HEALTH_BIND_ADDR enables the optional probe server with /readyz, /livez, and /drainz. ROCKETMQ_SHUTDOWN_TIMEOUT_SECONDS and ROCKETMQ_LIVENESS_STALE_SECONDS configure its shutdown and progress windows. Without a probe bind address, shutdown coordination still works.

/readyz and /livez accept GET and POST. /drainz starts a shutdown, so by default it accepts only POST and answers GET with 405 and Allow: POST. A Kubernetes preStop.httpGet hook can only send GET; set ROCKETMQ_HEALTH_DRAIN_METHODS=GET,POST for such a hook, as the repository's charts do. Any other value is a configuration error. Each connection is served by its own task of the lifecycle group, at most 64 at a time, so an idle or slow client does not delay other probes. The server reads until the end of the request headers, which may arrive in pieces. An accept failure caused by resource exhaustion, such as too many open files, is retried with backoff; only a listener that stays unusable marks the service failed.

The first shutdown request freezes a ShutdownDeadline; repeated pre-stop or signal requests cannot extend it. Pass that deadline through component shutdown and owner.shutdown_runtime_blocking_until(deadline).

The shutdown timeouts relate as follows:

Timeout Default Applies to
ServiceLifecycleConfig::shutdown_timeout (ROCKETMQ_SHUTDOWN_TIMEOUT_SECONDS) 45 s The process deadline frozen by the first shutdown request. Entrypoints pass it to every component and to the owner. It is the single source for a process that runs a lifecycle.
RuntimeConfig::shutdown_timeout 30 s Owner shutdown calls without an explicit deadline, for owners that do not run a lifecycle.
ServiceManager shutdown 30 s Used only when neither the call nor the parent group supplies a deadline; otherwise the earliest of the two applies.
Component budgets Per component Cap one component's share of the process deadline; they never extend it.

Task-group shutdown closes registration and broadcasts cancellation, then starts child shutdowns concurrently with waiting for the group's own tasks. Unfinished tracked tasks are aborted when the deadline expires. Reports are cached at group level; the owner additionally merges the blocking work of its managed lanes and of the isolated executors bound to its tree.

An API that takes a ShutdownDeadline never waits past it: TaskGroup::shutdown_until, OperationContext::wait_until, and BlockingExecutor::spawn_until report unconfirmed work instead. Two relative APIs confirm aborts for at most one second past their deadline: OperationContext::wait / cancel_and_wait, and an interrupting ServiceManager shutdown with less than a second left. Budget nested shutdown steps accordingly, or use the _until forms.

API Scope and guarantee
owner.shutdown_tasks().await / shutdown_tasks_until(deadline).await Close and wait for tracked tasks, retaining the Tokio runtime.
owner.shutdown_runtime_blocking() / shutdown_runtime_blocking_until(deadline) Consume the owner, shut down tracked tasks, then release Tokio within the remaining budget. Call outside a Tokio context; a separate task-shutdown call is not required.
TaskGroup::shutdown_now() Cancel and abort immediately without awaiting asynchronous completion.
owner.shutdown_background() Return immediate task-shutdown evidence and ask Tokio to shut down in the background.
RuntimeOwner::drop Emergency cleanup if explicit shutdown was omitted; not a graceful-shutdown protocol.

ShutdownReport::is_healthy() requires zero leaked, failed, panicked, timed_out, and blocking_still_running counts, and healthy child reports. timed_out counts the tasks this shutdown aborted after its deadline plus the tasks still registered; an abort requested earlier through abort_task is not a timeout. An aborted count alone does not make the report unhealthy. remaining_tasks lists at most ShutdownReport::REMAINING_TASKS_LIMIT (64) tasks; remaining_tasks_omitted counts the rest, which leaked still includes. An immediate-shutdown report is not proof that all futures completed their cleanup. Blocking snapshots are point-in-time evidence and do not terminate closures that outlive a deadline.

A ServiceManager runs a ServiceTask loop as a service task of the group passed to new_with_task_group. Its state (ServiceTaskState) is one atomic value, and shutdown_until(deadline) waits no later than the earliest of the requested deadline and the deadline installed on the parent group. An interrupting shutdown gives the aborted loop up to one second past that deadline to confirm its destruction.

Diagnostics

diagnostics_snapshot() exposes internal details such as runtime/group identity and blocking task names. Its events field, like TaskGroup::event_counts(), counts failures a tree absorbs without an error: groups poisoned by a panicking task, and component requests answered with a closed group. For authenticated operational APIs, prefer diagnostics_view_v1(RuntimeComponent::...): its versioned view aggregates bounded task-kind and lane summaries without raw IDs, names, arguments, or configuration objects. Authentication remains the caller's responsibility.

RuntimeDiagnosticsViewOptions controls summary bounds and the long-running threshold; omitted summaries set truncated. These diagnostics do not require Tokio unstable features or a console subscriber, and do not replace application health checks or performance measurements.

diagnostics_view_v2(RuntimeComponent::..., inputs) adds the same guarantees with explicit scopes: every section states whether it covers local, subtree, or process_shared state, a section whose input the caller does not own is absent rather than empty, and scheduled work, retained metadata, and shutdown results are reported as bounded aggregates. RuntimeDiagnosticsViewOptionsV2 also carries the scan and output budgets for an on-demand task detail list, so a partially scanned list reports how many tasks it examined instead of presenting a partial sum as the whole runtime. The scan budget bounds both the tasks examined and the descendant groups visited. V1 keeps its fields and meanings.

For a detail-only request use TaskGroup::diagnostics_task_details(scan, output). It skips the full-tree aggregates and returns RuntimeTaskDetails, including tasks_scanned, group_entries_scanned, and truncated. Ordered registries bound enumeration and temporary storage even after churn leaves sparse task tables or in a tree of empty groups. Names and payloads are omitted. Maintaining the sanitized task index adds registration/settlement work; the convergence benchmark measures that tradeoff. V1/V2 aggregate collection still visits every active task and group, regardless of the detail budget.

Compatibility And Workspace Integration

RocketMQRuntime has been removed, including its root and compat exports. Migrate construction to RuntimeOwner::plan(config)?.build()?, inject ChildServiceContext, select an explicit scheduling overlap policy, and inspect shutdown reports. See the API migration guide for replacements.

RuntimeContext is a migration/test harness. The executor services (TokioExecutorService, ScheduledExecutorService, FuturesExecutorService), TaskScheduler, ScheduledTaskManager, ActorRuntime, the compat module, and the legacy ServiceManager constructors have been removed; see MIGRATION.md for their replacements.

The broker, NameServer, proxy, and controller entrypoints build runtime owners and use service lifecycle deadlines. Other consumers include rocketmq-client, rocketmq-transport, rocketmq-store, rocketmq-auth, rocketmq-observability, and admin tools. ClientRuntime requires an injected application-owned child scope and never creates a fallback runtime. Store compatibility helpers retain explicit adapter boundaries; this list does not imply that every call site uses an identical ownership path. Standalone applications follow their local host-runtime and validation guides.

Features And Validation

The crate inherits its edition and minimum Rust version from the workspace manifest. It has no optional features. The filesystem helpers in common::file_utils block; asynchronous code runs them through a blocking lane or persists through MetadataIoActor.

For task-lifecycle changes, start with package-scoped checks:

cargo fmt -p rocketmq-runtime -- --check
cargo test -p rocketmq-runtime --test runtime_model

Select additional checks for the behavior being changed, rather than running every suite for every edit:

Area Test target or command
Internal units, error channels, diagnostics, service lifecycle cargo test -p rocketmq-runtime --lib
Resource limits and queue behavior cargo test -p rocketmq-runtime --test resource_budget_tree
Shared process-budget ownership cargo test -p rocketmq-runtime --test runtime_resource_ownership
Metadata persistence and fault handling cargo test -p rocketmq-runtime --test metadata_io_actor
Public scope restrictions cargo test -p rocketmq-runtime --test service_context_scope_compile_fail
Shutdown or budget interleavings task_group_shutdown_loom or resource_budget_loom via cargo test -p rocketmq-runtime --test <target>
Migration or large-future submission runtime_migration_fixture or task_submission_stack via the same test command
Filesystem helpers cargo test -p rocketmq-runtime --lib common::file_utils

When useful, run cargo clippy -p rocketmq-runtime --no-deps -- -D warnings with the affected targets/features. Validate directly affected consumers when shared behavior changes; follow standalone projects' local guides where applicable. Feature-enabled checks do not replace feature-absence coverage.

For README edits, check local links and compile/run the fenced Rust examples. cargo test --doc covers crate Rustdoc, not standalone README code blocks; test those explicitly with rustdoc --test and the built crate's --extern and dependency search path. Keep both language versions aligned.

Full-workspace checks, runtime audits, Loom models, and Criterion benchmarks belong to changes that need that evidence or the relevant CI/integration task. Benchmarks provide measurements for a specific run, not hard-coded performance guarantees. See repository validation guidance.

Platform and scale experiments

These checks provide platform and workload evidence when a change calls for it; they are not routine gates. Run cargo test -p rocketmq-runtime --no-default-features on Windows and Linux to cover the feature-absent path. The metadata actor suite exercises real filesystem replacement and cleanup as well as injected failures and gates. It does not prove power-loss recovery or coordination between independent processes.

cargo test -p rocketmq-runtime --test runtime_scale -- --nocapture checks 5,000 nested-scope/dynamic-key generations and 1,024 retired metadata histories. The latter uses an injected successful filesystem; use the metadata actor suite and metadata benchmark for real persistence evidence. Retained old receipts and escaped budgets must not revive a retired identity.

The Linux cgroup test is ignored by default. Compile it outside the constrained unit, then run its executable under a real limit:

cargo test -p rocketmq-runtime --no-default-features --test runtime_scale \
  --no-run --message-format=json > /tmp/runtime-scale-build.json
test_bin=$(python3 -c 'import json; rows=[json.loads(x) for x in open("/tmp/runtime-scale-build.json")]; print(next(x["executable"] for x in rows if x.get("executable") and x.get("target", {}).get("name") == "runtime_scale"))')
sudo systemd-run --wait --pipe --collect \
  --property=MemoryMax=536870912 --property=MemorySwapMax=0 \
  --setenv=ROCKETMQ_TEST_CGROUP_BYTES=536870912 \
  "$test_bin" --ignored --exact \
  detects_actual_cgroup_limit_and_uses_it_in_owner_planning --nocapture

This requires Linux cgroup v2 and systemd. The test verifies the detected source and byte limit and the owner's resulting memory budget. File-view tests provide separate parser evidence. An ignored test does not count as a pass.

Benchmark measurement boundaries

Set CARGO_TARGET_DIR to the intended build disk before running benchmarks. Run a benchmark with cargo bench -p rocketmq-runtime --bench <name>.

Benchmark Population and measured region
runtime_diagnostics_bench Stable 1k/10k/100k populations with eight child groups; separate group-count and fixed-32-group depth 1/4/16 cases. Setup and shutdown excluded; Criterion includes return-value destruction.
budgeted_queue_bench Reused queues; fill/reject/drain, wait/release, or replacement-at-capacity operations. Queue and Tokio runtime construction excluded; wait cases include producer spawning and joining.
blocking_executor_bench Submission through completion of 8/32 jobs with four lane slots and 1 ms simulated blocking work. Runtime creation and shutdown excluded; timeout evidence retains the blocked closure until release.
metadata_io_bench First real filesystem write held at a gate; queued submissions through release and settled receipts timed. Coalesced writes and hot/cold ordering asserted; first gate arrival, setup, and shutdown excluded.
runtime_convergence_bench Draining-operation submission from 1/4/8 threads into one shared group or one group per thread; permit acquire/release under a shared root; allocations of an idle group, operation, and child context; start order of contended blocking work. Every call is timed for percentiles; the JSON is named by ROCKETMQ_BENCH_LABEL.
runtime_boundaries_bench Mixed blocking lanes, actor budget saturation and RSS, injected slow-target isolation, bounded detail allocations versus full aggregates, and shutdown settlement. Run with cargo bench -p rocketmq-runtime --bench runtime_boundaries_bench; JSON is written under target/runtime-measurements (or CARGO_TARGET_DIR). Injected delays do not measure disk performance; use metadata_io_bench for real filesystem writes.

Criterion reports batch-derived estimates and confidence intervals, not an individual request P99. Reject scenarios intentionally reject one extra item per full batch; wait/churn scenarios assert successful admission and resource release. Do not compare older results after changing the timed region as if they measured a code speedup.

Set ROCKETMQ_MEASURE_SAMPLING=1 for independent allocation samples and per-submission observations in the diagnostics benchmark. Its sampling-costs.json artifact stores 101 samples per population/detail case and three pairs of 5,000 submit/completion observations at a 10 ms sampler delay. Allocation counters cover successful allocation/reallocation requests on the sampling thread only, not live or peak process memory. These independent timers exclude setup, teardown, and return-value destruction.

Artifacts are written under CARGO_TARGET_DIR/runtime-measurements; Criterion retains its own raw samples. Record the OS, CPU, toolchain, background load, actual sampler count, warmup/repetitions, and quantile algorithm with results. No timing number here is a release SLO; source fingerprints and file hashes are unnecessary.

Crate Layout

rocketmq-runtime/
  src/public_api.rs        deliberate ownership and diagnostics exports
  src/prelude.rs           recommended entry path and common ownership imports
  src/config.rs            runtime and blocking-lane configuration
  src/owner.rs             validated construction and owned runtime lifecycle
  src/context.rs           borrowed runtime migration/test harness
  src/service_context.rs   sealed root and child capabilities
  src/task_spawner.rs      narrow task-submission capability
  src/task_group.rs        task tracking, cancellation, and shutdown
  src/task_group/          child registry and deadline coordination
  src/operation.rs         operation-local cancellation and bounded waits
  src/scheduled.rs         periodic drivers, runs, and metrics
  src/blocking.rs          blocking API exports
  src/blocking/            lane admission, execution, and snapshots
  src/resources.rs         shared process resource capabilities
  src/resource_budget/     resource trees, permits, queues, and memory discovery
  src/metadata_io.rs       generation-aware metadata persistence
  src/service_lifecycle.rs readiness, liveness, and shutdown requests
  src/shutdown_deadline.rs shared absolute shutdown deadline
  src/shutdown_report.rs   serializable shutdown evidence
  src/diagnostics.rs       raw snapshots and sanitized views
  src/task/                service loops run by ServiceManager
  src/common/              filesystem, time, and configuration-file helpers

License

Licensed under the Apache License, Version 2.0. See LICENSE-APACHE for details.