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
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"]
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.
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.
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:
RuntimeContractViolationidentifies invalid caller configuration or an invariant violation, including failures fromplan().RuntimeResult<T>containsRuntimeErrorfor operational failures such as runtime construction, I/O, capacity, or timeout failures. Branch onRuntimeError::kind()rather than on the operation label: a submission to a closing owner reportsRuntimeErrorKind::Closed, and one to a poisoned group reportsRuntimeErrorKind::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.
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.
| 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.
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_timebounds an individual run by dropping its future on timeout. External side effects still need an appropriate cancellation contract.- Duplicate names return
AlreadyPresentwithout 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.
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.
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 aBudgetedItemretaining 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.
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.
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_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.
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.
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_modelSelect 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.
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 --nocaptureThis 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.
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.
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
Licensed under the Apache License, Version 2.0. See
LICENSE-APACHE for details.