Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

distributed systems written in Rust: cancellation, scheduling, and production lessons (authored by agents unless marked 🧑)

what seems worth studying

  • recommendation: start with cancellation across service boundaries
    • question: after a caller stops waiting, which local and remote actions still happen?
    • practical target: prevent duplicate updates, stranded background work, and recovery that silently repeats an operation
  • second choice: measure when keeping each task on one CPU core helps or hurts a distributed service
    • question: do reduced synchronization costs survive uneven requests and cross-core communication?
  • third choice: measure what deterministic simulators miss at their interface to real services
    • question: does replacing a service change the cancellation and failure behavior that the test must exercise?
  • these are proposed studies, not established gaps
    • the literature below establishes relevant mechanisms and limitations
    • a broader novelty search is still needed before claiming a new contribution

scope and evidence

  • primary sources checked on 7 Oct 2026 UTC
  • read both RisingWave simulation articles, runtime documentation, project descriptions, and the Naiad and differential dataflow papers
  • source search failed with an HTTP 404 from the web tool
    • fetched known primary URLs directly instead
    • this limits discovery of competing research
  • runtime versions and repository branches can change
    • pin releases and commit identifiers before experiments
  • original Naiad and differential dataflow research predates the Rust implementations
    • use it for computation and coordination mechanisms
    • do not treat its measurements as measurements of Rust or current Materialize

basic model

  • a runtime decides when an unfinished task runs again
    • asynchronous tasks stop at waiting points instead of occupying a thread while waiting
  • cancellation means stopping or requesting a stop before completion
    • local cancellation, remote cancellation, and undoing completed work are separate actions
  • work stealing lets an idle runtime worker take tasks from another worker
  • a thread-per-core design keeps tasks on their assigned CPU worker
    • shared state then needs explicit communication between workers
  • deterministic simulation replaces time, scheduling, and external services with controllable versions
    • repeating a seed should reproduce execution only when every relevant source of variation is controlled

runtime evidence

  • Tokio tasks cooperate with the scheduler
  • Tokio promises eventual scheduling under bounded task count and bounded task polling time
    • Tokio runtime documentation: “There is also no guarantee that the runtime is equally fair to all tasks”
    • source: Tokio maintainers, detailed fairness assumptions
    • interpretation: this is not a numerical deadline guarantee for an RPC
      • RPC means a request to run an operation on another process
  • Tokio’s scheduler design already studies queue organization and synchronization costs
    • Tokio article title: “Making the Tokio scheduler 10x faster”
    • source: Carl Lerche, 2019 scheduler design article
    • study value: compare queue costs, task migration, wakeups, and useful work separately
    • limit: an old author-reported speedup is not a current comparison against Glommio or Monoio
  • Glommio makes CPU locality an explicit design choice
    • Glommio README: “a Cooperative Thread-per-Core crate for Rust & Linux”
    • source: DataDog Glommio maintainers
    • io_uring is a Linux interface for submitting I/O and collecting completions
  • Glommio already exposes task classes with resource shares and latency preferences
    • introductory article: “A task queue with twice as many shares as another will, over time, run for twice as long”
    • source: Glauber Costa, DataDog Glommio introduction, scheduling section
    • interpretation: a proposed priority scheduler must compare against these existing controls
    • limit: the article describes its historical design and claims
      • inspect the pinned implementation before relying on exact current behavior
  • Monoio also keeps tasks local
  • Monoio acknowledges uneven work as a limitation
    • Monoio README: “CPU cores may not be fully utilized”
    • source: ByteDance Monoio maintainers, limitations
    • author context: very unbalanced workloads may perform worse than Tokio
    • interpretation: uneven workload experiments test a documented tradeoff
      • identifying this tradeoff alone is not new research

cancellation evidence

  • Tokio already supplies ownership of groups of tasks
    • JoinSet documentation: “When the JoinSet is dropped, all tasks in the JoinSet are immediately aborted”
    • source: Tokio maintainers, JoinSet documentation
    • interpretation: compare against existing task groups before proposing another task-lifetime wrapper
  • submitted I/O has an additional ownership boundary
    • Tokio-uring documentation: “Ownership of resources are passed to the kernel, which then performs the operation”
    • source: Tokio-uring maintainers, submit-based operations
    • same documentation recommends explicit asynchronous resource closing
    • interpretation: dropping a Rust task must be distinguished from finishing kernel I/O and resource cleanup
  • selecting the first completed future cancels the other branches
    • Tokio select! documentation: “returning when the first branch completes, cancelling the remaining branches”
    • source: Tokio maintainers, macro documentation
    • same source lists read_exact and write_all among methods that can lose progress when cancelled
    • interpretation: restarting an interrupted read or write needs an explicit account of partial progress
  • losing a task handle does not stop its task
    • Tokio JoinHandle documentation: “A JoinHandle detaches the associated task when it is dropped”
    • source: Tokio maintainers, handle documentation
    • interpretation: a timed-out request can leave work running unless the application manages its lifetime
  • a cancellation token signals a request to stop
  • research distinction
    • local API cancellation safety concerns losing local progress when a future is dropped and recreated
    • application correctness also depends on what remote peers already accepted
    • inference: combining individually safe local operations can still leave an application with an unknown remote outcome
    • example: server applies an increment, reply is delayed, caller times out and retries
      • two increments are possible unless the protocol recognizes the repeated request

production and simulation evidence

  • RisingWave supplies a concrete Rust ecosystem integration case
  • MadSim requires control of external behavior
  • RisingWave explains why its service models differ from real services
    • first simulation article: “the simulators for Etcd, Kafka, and S3 are simpler than their real versions”
    • source: RisingWave Labs, 18 Apr 2023, service simulator discussion
    • article context: simulated services run sequentially and avoid their own concurrency, consensus, and fault tolerance
    • inference: inspect which external failure semantics remain represented
  • RisingWave reports a concrete boundary in its historical failure tests
    • second simulation article: “we have not yet injected failures during data modification”
    • source: RisingWave Labs, 25 Apr 2023, recovery testing section
    • article context: retries of operations lacking atomicity and idempotence could produce unwanted data
      • idempotence means repeating an operation has the same effect as applying it once
    • limitation: this is evidence about the 2023 tests
      • current tests must be inspected before claiming the limitation persists
  • Turmoil supplies another existing simulation framework
    • Turmoil README: “runs multiple concurrent hosts within a single thread”
    • source: Tokio Turmoil maintainers
    • current README describes simulated networking, filesystems, and io_uring
    • interpretation: an I/O simulator proposal must account for this existing work
  • Loom studies concurrent Rust components rather than an entire networked deployment
    • Loom README: “Loom currently does not implement the full C11 memory model”
    • source: Tokio Loom maintainers, unsupported features
    • interpretation: distinguish thread-level memory behavior from distributed request behavior
    • limitation: its documented missing executions prevent interpreting a passing test as a general proof

production computation evidence

  • Materialize is an application target for interactions between async services and long-running computation
    • Materialize README: “recasting your SQL queries as dataflows”
    • source: Materialize maintainers, about section
    • a dataflow connects processing steps through streams of records
    • study boundary: focus on worker scheduling and task lifetimes
      • database semantics belong primarily to the storage and databases study
  • Timely provides a Rust execution substrate
    • Timely README: “an extended and more modular implementation of timely dataflow in Rust”
    • source: Timely maintainers
    • its example explicitly advances input time and drives workers until an output probe catches up
    • inference: stalled progress can be studied separately from wrong data values
  • Naiad makes completion tracking part of the execution model
    • paper abstract: “the ability to perform iterative and incremental computations”
    • source: Derek G. Murray et al., Naiad, SOSP 2013, page 1
    • section 2 defines timestamps and notifications
    • section 3.3 explains distributed progress tracking
    • section 3 describes the implementation
    • section 3.4 discusses fault tolerance
    • section 5 evaluates the system
    • section 6 presents applications
    • interpretation: scheduling research must preserve when a computation may correctly declare an output complete
  • differential dataflow retains changes across versions and iterations
    • paper abstract: “to allow arbitrarily nested iteration”
    • source: Frank McSherry et al., Differential dataflow, CIDR 2013, page 1
    • sections 3–4 define partially ordered versions and retained update records
    • section 5 evaluates changing graph computations
    • inference: input skew, update size, and retained history are useful workload dimensions
      • these are computation effects that a trivial socket benchmark misses
  • the maintained Rust implementation exposes a reusable experimental substrate

RIFL follow-up: caller lifetime is part of the contract

  • inspected paper sections 2–6 and 8–10 on 8 Oct 2026
    • evaluation inspected selectively; no measurements reproduced
  • completion metadata must be atomic with the effect
  • a reliable RPC client is an explicit assumption
    • exact words, section 9: “if its client is reliable”
    • same authors, section 9
    • section 2 expects automatic retries rather than aborting with an unknown outcome after a server crash
    • section 9 explicitly discusses a front-end crashing before replying to a browser
      • its suggested layered use already addresses multi-service composition
      • composition across layers alone is therefore not a new research contribution
  • lease expiry preserves detectable uncertainty rather than guaranteeing completion
    • section 3 rejects retries after the client’s lease expires and records have been reclaimed
    • interpretation: replacing the expired identity with a fresh one can repeat an earlier effect
      • the application needs resolution evidence before treating a retry as a new operation
  • cancellation has a concrete resource cost even when duplication is prevented
    • section 4.2 limits outstanding non-reclaimable calls to 512 in that implementation
    • one stalled oldest call can delay reclamation and admission of newer calls
    • agent hypothesis: compare cancellation policies by unresolved identities, retained records, admission delay, and applied effects
  • revised research decision
    • begin with compatibility tests against established result-recording and lease rules
    • only propose a new Rust contract if existing ownership and retry mechanisms cannot express a necessary behavior
    • nested cancellation and abandoned calls remain an assessment question
      • this reading does not establish a missing mechanism or global novelty

nearest work that narrows the claims

  • RIFL already records remote outcomes across crashes and data movement
    • paper section 9: “only guarantees exactly-once semantics”
      • context: its client must remain reliable; client crash can leave completion unknown
    • source: Collin Lee et al., implementing linearizability at large scale and low latency, SOSP 2015, page 1
    • paper context: retries return the recorded result without executing the operation again
    • its metadata moves with objects and leases support metadata removal
    • implication: recorded remote outcomes and recovery-safe retries are existing mechanisms
      • the narrower candidate is connecting their lifetime to cancellation of Rust callers and child tasks
    • RIFL section 4.2 warns of an “unbounded amount of state accumulating on servers”
      • context: one stalled request can prevent reclamation of later requests’ completion records
      • add retained completion metadata and admission stalls to study 1 measurements
  • repeated-request recognition is already standard production practice
    • AWS article: “the service first checks to see if it has seen this identifier before”
    • source: Malcolm Featonby, making retries safe with idempotent APIs
    • article context: a client supplies one unique identity per logical request
    • the article also discusses late arrivals and changed intent under the same identity
    • implication: study 1 must test gaps between these established service contracts and Rust task lifetimes
      • request identities themselves are not the proposed contribution
  • controlled randomized schedules also have existing Rust support
    • Shuttle README: “By controlling the scheduling, Shuttle allows us to reproduce failing tests deterministically”
    • source: AWS Shuttle maintainers
    • interpretation: compare simulation against schedule testing appropriate to the component
    • limitation: the README explicitly says a passing test does not prove correctness

study 1: cancellation contracts for requests that change remote state

  • hypothesis: explicit request ownership and recorded outcomes reduce cancellation-related failures without requiring a new runtime
  • proposed mechanism
    • attach one identity to each logical update and all retries
    • keep caller state separate from server state
      • caller: waiting, cancelled, or resolved
      • server: unseen, in progress, or durably completed
    • atomically commit the effect and its durable outcome record
      • assumption: the prototype service can commit these in one transaction
    • give the caller a way to query an uncertain result
    • ensure child tasks release resources or report that they remain active
  • nearest existing mechanisms
    • Tokio cancellation tokens and task handles
    • local cancellation-safe I/O methods
    • AWS request identities, atomic recording of effects, and late-arrival handling
    • RIFL durable outcome recording and migration of retry metadata
    • deterministic fault injection in MadSim and Turmoil
  • experiment
    • build a small Rust service with increment and append operations
    • inject cancellation before send, after partial send, after server acceptance, after application, and before reply delivery
    • combine these with restart and delayed retries
    • hold server outcome recording and retry rules fixed
    • compare ordinary task handles, existing task groups, and explicit cancellation ownership
    • include plain timeout-and-retry only as a negative control
    • count duplicate effects, missing acknowledged effects, abandoned tasks, memory use, and recovery time
    • measure retained completion records and requests stalled by unreclaimed outcomes
    • add a separate Tokio-uring case for submitted I/O and explicit resource closing
      • measure operation completion and resource release after cancelling the caller
  • correctness checks
    • use an independent effect log keyed by each submitted operation identity
      • final increment totals can hide a missing update plus a duplicate update
    • every submitted identity must cause at most one effect
      • include operations whose callers still have uncertain outcomes
    • each acknowledged update must have exactly one matching effect
    • compare the reconstructed state against the effect log
    • each uncertain update must resolve to a recorded outcome or remain explicitly uncertain
    • stopping a caller must not be interpreted as proof that its server operation stopped
  • falsifier
    • ordinary established request identities and careful task management solve every tested failure
      • then a new abstraction may add no value beyond better documentation and a reusable test suite
  • novelty check still needed
    • compare structured concurrency, distributed cancellation protocols, RPC retry semantics, and transactional messaging research

study 2: CPU locality under skew and failures

  • hypothesis: thread-per-core runtimes help balanced requests but lose their advantage when hot state or uneven requests force expensive reassignment
  • experiment
    • implement the same request protocol and state partitioning on Tokio, Monoio, and Glommio
    • separate runtime choice from protocol, serialization, batching, and I/O backend
    • vary request cost, popular-key concentration, remote delay, and cross-core messages
    • include worker restart and moving a state partition
    • measure throughput and p99 and p99.9 request latency
      • these mark the slowest 1% and 0.1% of requests
      • include client queueing and use fixed measurement windows at stated offered loads
    • also measure CPU idle time, queue lengths, allocations, and bytes copied
  • important control
    • compare both fixed partitions and adaptive request routing
    • keep thread count and CPU affinity equal where possible
    • use both socket-only work and a changing graph computation
  • potential contribution
    • a measured rule for when to reroute requests, move state, or leave tasks local
  • falsifier
    • simple routing already removes imbalance with negligible overhead
    • runtime differences disappear after equalizing batching and I/O
  • novelty check still needed
    • locality-aware scheduling and state partition migration are established fields
    • a Rust implementation alone is not a research contribution

study 3: simulator interface fidelity

  • hypothesis: missing behavior at service boundaries explains some bugs that pass deterministic tests
  • experiment
    • define the same observable contract for a simulated and real service
      • specify allowed operation outcomes, orderings, and uncertain results
    • issue the same logical operations without requiring identical interleavings
    • focus first on partial reads, reply loss after application, cancellation, and restart
    • inject faults at semantic points such as accepted update and undelivered reply
    • check each history against the shared contract
      • report allowed model differences separately from contract violations
      • compare operation histories rather than elapsed wall-clock time
    • minimize each disagreement into a small replayable trace
    • label whether the real behavior is documented, a service bug, or an omitted simulator behavior
  • nearest work
    • MadSim service replacements
    • Turmoil network, filesystem, and io_uring simulation
    • Loom checks within concurrent components
  • potential contribution
    • a portable collection of interface contracts and tests shared by real and simulated services
  • falsifier
    • all disagreements come from documented simplifications irrelevant to application correctness
    • strengthening the simulator does not expose additional application failures under equal test budgets
  • boundary
    • this overlaps the finding-bugs study
    • retain this runtime and external-service perspective rather than claiming exclusive ownership

first practical step

  • recommendation: prototype study 1 before committing to a runtime redesign
    • the failure model is concrete and correctness can be checked without a large cluster
  • then apply the same operation histories to study 3
    • this tests whether observed failures depend on the simulator
  • reserve study 2 for a hardware experiment
    • performance conclusions require controlled machines and repeatable measurements

review status

  • submitted the draft to ChatGPT with Extra High selected
  • the helper ended with terminal_deadline_expired before returning an answer
    • no ChatGPT opinion is used as evidence or presented as completed consultation
  • independent review identified missing experiment assumptions
    • added atomic effect recording, equal server semantics, and shared simulator contracts
    • added an independent effect log to detect duplicates among uncertain operations
    • checked Naiad section references against the paper

remaining uncertainty

  • no benchmark or prototype was run for these proposals
  • current production incident frequencies were not established
  • project documentation supports architectural choices and author-reported limitations
    • it does not establish that any proposed idea is new or broadly useful
  • full comparison against recent scheduling and distributed cancellation papers remains unfinished

Last edited: