Skip to content

[P2] Supervisor::apply_put is load-mutate-store on ArcSwap (latent race) #112

Description

@moonming

Summary

Supervisor::apply_put reads the current snapshot (ArcSwap::load), clones it, mutates the clone, then stores the result. This is a load-mutate-store pattern, not RCU. Two concurrent apply_put calls would both load the same snapshot, mutate independently, and the second store overwrites the first — losing one event silently.

The function is currently safe by accident (only the watch supervisor calls it, single-threaded). But it's pub fn, and any future caller — admin-API direct write, test fixture, replay tool — can introduce the race.

Severity

🟡 MEDIUM-HIGH — latent race; not currently exploitable, but the pub fn surface invites future regression. ArcSwap exists specifically to make this kind of update safe via rcu; using it as a load-mutate-store wastes the safety property.

Location

crates/aisix-etcd/src/supervisor.rs:161 (verified on origin/main):

pub fn apply_put(&self, entry: &RawEntry) -> bool {
    let (tiny, stats) = loader::build_snapshot(&self.prefix, std::slice::from_ref(entry));
    if stats.accepted == 0 {
        return false;
    }

    let new = clone_snapshot(&self.handle.load());   // ← LOAD

    // Move any entries from `tiny` into `new`. Must cover every
    // ResourceTable on AisixSnapshot — a missing kind here means
    // a watch event silently drops on the floor and the snapshot
    // never sees the new entry, even though the loader and the
    // proxy both know about it.
    for e in tiny.models.entries() {
        new.models.insert(clone_entry(&e));            // ← MUTATE
    }
    for e in tiny.apikeys.entries() {
        new.apikeys.insert(clone_entry(&e));
    }
    for e in tiny.provider_keys.entries() {
        new.provider_keys.insert(clone_entry(&e));
    }
    for e in tiny.guardrails.entries() {
        new.guardrails.insert(clone_entry(&e));
    }
    for e in tiny.cache_policies.entries() {
        new.cache_policies.insert(clone_entry(&e));
    }
    for e in tiny.observability_exporters.entries() {
        new.observability_exporters.insert(clone_entry(&e));
    }

    self.handle.store(new);                            // ← STORE
    ...
}

Why this is risky

Concrete race:

  1. Supervisor handles event A: load → snapshot v1 (5 models)
  2. Concurrently, supervisor handles event B: load → snapshot v1 (5 models)
  3. A clones v1, inserts model M_A → v1' has 6 models including M_A
  4. B clones v1, inserts model M_B → v1'' has 6 models including M_B
  5. A stores v1' (last write step for A)
  6. B stores v1'' — M_A is lost

Today this doesn't happen because the etcd watch sub is sequential. But:

  • A future fix that adds parallel watch handlers (e.g., to handle multiple etcd prefixes) would race
  • A test fixture that replays many events fast would race
  • An admin-API direct-write path that reuses this method would race

Recommended fix

Use ArcSwap's rcu pattern (CAS retry until success):

pub fn apply_put(&self, entry: &RawEntry) -> bool {
    let (tiny, stats) = loader::build_snapshot(&self.prefix, std::slice::from_ref(entry));
    if stats.accepted == 0 {
        return false;
    }

    self.handle.rcu(|current| {
        let new = clone_snapshot(current);
        for e in tiny.models.entries() {
            new.models.insert(clone_entry(&e));
        }
        // ... same as before for other resource kinds ...
        new
    });

    // ... cache-tracking and revision update ...
    true
}

rcu retries the closure if a concurrent store happened. The retry is correct because the closure reads current fresh each retry.

Note: the state.lock() block at the bottom of the function (cache-tracking) needs separate consideration — it's already correctly using a Mutex.

Test case

#[tokio::test]
async fn apply_put_concurrent_does_not_lose_events() {
    let sup = Supervisor::new(...);
    let n = 100;
    let mut tasks = JoinSet::new();
    for i in 0..n {
        let sup = sup.clone();
        tasks.spawn(async move {
            sup.apply_put(&entry(format!("/aisix/models/m-{}", i), VALID_MODEL, i + 1));
        });
    }
    while let Some(_) = tasks.join_next().await {}
    
    let snapshot = sup.handle.load();
    // All N models should be present, no concurrent-overwrite losses.
    assert_eq!(snapshot.models.len(), n);
}

Source

Static scan, May 2026 (verified on origin/main:crates/aisix-etcd/src/supervisor.rs:161).

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    P2Long-tail integrations — backlogbugSomething isn't working

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions