From fab97463c2e724e46204fe8eab92a1a9bb8aebb2 Mon Sep 17 00:00:00 2001 From: Nick Yuan Date: Mon, 13 Apr 2026 17:24:43 -0500 Subject: [PATCH] Fix unbounded growth of Broker.set_stat cluster master list (#322) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Prune stale entries from the master list on new cluster registration and only write the list when membership changes (with timeout=None). Before this fix, set_stat was append-only — dead cluster keys accumulated forever because the only cleanup path (get_stats) runs only from the monitor UIs. On DatabaseCache backends this rewrote the growing pickled list to Postgres on every heartbeat, producing measurable egress. Co-authored-by: Claude Opus 4.6 (1M context) --- django_q/brokers/__init__.py | 7 ++++- django_q/tests/test_brokers.py | 51 ++++++++++++++++++++++++++++++++++ 2 files changed, 57 insertions(+), 1 deletion(-) diff --git a/django_q/brokers/__init__.py b/django_q/brokers/__init__.py index 909bdc5..3751ebd 100644 --- a/django_q/brokers/__init__.py +++ b/django_q/brokers/__init__.py @@ -106,8 +106,13 @@ class Broker: return key_list = self.cache.get(Conf.Q_STAT, []) if key not in key_list: + # Prune stale entries whose per-stat value has expired, so the + # master list cannot grow without bound across cluster restarts. + key_list = [k for k in key_list if self.cache.get(k) is not None] key_list.append(key) - self.cache.set(Conf.Q_STAT, key_list) + # timeout=None: master list lifetime is managed by membership + # changes, not by TTL refresh on every heartbeat. + self.cache.set(Conf.Q_STAT, key_list, None) return self.cache.set(key, value, timeout) def get_stat(self, key: str): diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py index 661d0f5..00ef2aa 100644 --- a/django_q/tests/test_brokers.py +++ b/django_q/tests/test_brokers.py @@ -35,6 +35,57 @@ def test_broker(monkeypatch): assert broker.get_stats("test:*") is None +def test_broker_set_stat_prunes_stale_master_list(): + # regression: set_stat should prune stale master-list entries on register + broker = Broker() + a_key = f"{Conf.Q_STAT}:A" + b_key = f"{Conf.Q_STAT}:B" + broker.cache.delete(Conf.Q_STAT) + + broker.set_stat(a_key, "state_a", 3) + assert a_key in broker.cache.get(Conf.Q_STAT) + + # Drop A's per-stat value to simulate a dead cluster whose TTL expired, + # leaving the master-list entry as a stale reference. + broker.cache.delete(a_key) + assert broker.get_stat(a_key) is None + assert a_key in broker.cache.get(Conf.Q_STAT) + + # Registering B should prune stale A from the master list. + broker.set_stat(b_key, "state_b", 3) + key_list = broker.cache.get(Conf.Q_STAT) + assert a_key not in key_list + assert b_key in key_list + assert broker.get_stat(b_key) == "state_b" + + +def test_broker_set_stat_skips_master_list_write_on_repeat(monkeypatch): + # set_stat should only write the master list when membership changes + broker = Broker() + a_key = f"{Conf.Q_STAT}:A" + broker.cache.delete(Conf.Q_STAT) + + writes = [] + orig_set = broker.cache.set + + def counting_set(key, value, timeout=None, **kw): + if key == Conf.Q_STAT: + writes.append(key) + return orig_set(key, value, timeout, **kw) + + monkeypatch.setattr(broker.cache, "set", counting_set) + + # First call adds a new entry, one master-list write expected + broker.set_stat(a_key, "state_a", 3) + assert len(writes) == 1 + + # Subsequent calls for the same key must not rewrite the master list + broker.set_stat(a_key, "state_a", 3) + broker.set_stat(a_key, "state_a", 3) + broker.set_stat(a_key, "state_a", 3) + assert len(writes) == 1 + + def test_redis(monkeypatch): monkeypatch.setattr(Conf, "DJANGO_REDIS", None) broker = get_broker()