From 68e01c7aa199d38c8910760006b6a787cf1531d1 Mon Sep 17 00:00:00 2001 From: Derek Steinkamp Date: Tue, 29 Sep 2026 15:12:34 -0500 Subject: [PATCH 1/3] Recover owned readers across observed stream resets --- CMakeLists.txt | 13 +++ RedisAdapter.cpp | 152 +++++++++++++++++++++++++++-- RedisAdapter.hpp | 25 +++++ RedisConnection.hpp | 58 ++++++++++- docs/api.md | 1 + docs/stream-subscriptions.md | 57 +++++++++++ reader_cluster_test.py | 80 ++++++++++++++++ recovery_test.cpp | 180 +++++++++++++++++++++++++++++++++++ 8 files changed, 556 insertions(+), 10 deletions(-) create mode 100644 reader_cluster_test.py create mode 100644 recovery_test.cpp diff --git a/CMakeLists.txt b/CMakeLists.txt index 13208b9..b2069ff 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -89,12 +89,25 @@ if(REDIS_ADAPTER_TEST) target_link_libraries(redis-adapter-lifecycle-test ${REDIS_ADAPTER_LIBRARIES}) add_test(NAME RedisAdapter.Lifecycle COMMAND redis-adapter-lifecycle-test) set_tests_properties(RedisAdapter.Lifecycle PROPERTIES TIMEOUT 30) + add_executable(redis-adapter-recovery-test ${REDIS_ADAPTER_SOURCES} recovery_test.cpp) + target_compile_options(redis-adapter-recovery-test PRIVATE -UNDEBUG) + target_link_libraries(redis-adapter-recovery-test ${REDIS_ADAPTER_LIBRARIES}) + add_test(NAME RedisAdapter.Recovery COMMAND redis-adapter-recovery-test) + set_tests_properties(RedisAdapter.Recovery PROPERTIES TIMEOUT 30 RUN_SERIAL TRUE) add_executable(redis-adapter-write-error-test ${REDIS_ADAPTER_SOURCES} write_error_test.cpp) target_compile_options(redis-adapter-write-error-test PRIVATE -UNDEBUG) target_link_libraries(redis-adapter-write-error-test ${REDIS_ADAPTER_LIBRARIES}) add_test(NAME RedisAdapter.WriteErrors COMMAND redis-adapter-write-error-test) set_tests_properties(RedisAdapter.WriteErrors PROPERTIES TIMEOUT 30 RUN_SERIAL TRUE) find_package(Python3 REQUIRED COMPONENTS Interpreter) + find_program(REDIS_SERVER_EXECUTABLE redis-server) + find_program(REDIS_CLI_EXECUTABLE redis-cli) + if(REDIS_SERVER_EXECUTABLE AND REDIS_CLI_EXECUTABLE) + add_test(NAME RedisAdapter.ClusterRecovery + COMMAND ${Python3_EXECUTABLE} ${CMAKE_CURRENT_SOURCE_DIR}/reader_cluster_test.py + ${REDIS_SERVER_EXECUTABLE} ${REDIS_CLI_EXECUTABLE} $) + set_tests_properties(RedisAdapter.ClusterRecovery PROPERTIES TIMEOUT 75) + endif() add_test(NAME RedisAdapter.WriteTransportFailure COMMAND ${Python3_EXECUTABLE} ${CMAKE_CURRENT_SOURCE_DIR}/write_failure_test.py $) diff --git a/RedisAdapter.cpp b/RedisAdapter.cpp index 0a1a5e6..ba1867c 100644 --- a/RedisAdapter.cpp +++ b/RedisAdapter.cpp @@ -157,6 +157,15 @@ RedisAdapter::ReaderHandle& RedisAdapter::ReaderHandle::operator=(ReaderHandle&& RedisAdapter::ReaderHandle::operator bool() const { return registration_ && registration_->active.load(); } +RedisAdapter::ReaderStatus RedisAdapter::ReaderHandle::status() const { + if (!registration_) return {}; + lock_guard lock(registration_->mutex); + auto result = registration_->status; + result.active = registration_->active && !owner_.expired(); + result.cursor = registration_->cursor; + result.observedCursor = registration_->observedCursor; + return result; +} void RedisAdapter::ReaderHandle::reset() noexcept { auto registration = std::move(registration_); if (!registration) return; @@ -404,12 +413,14 @@ RedisAdapter::register_reader(const string& key, reader_sub_fn func, const strin auto registration = make_shared(); registration->id = ++_next_reader_id; registration->cursor = afterId; + registration->inspect = resolveTail; // only the owned API requests continuity inspection registration->callback = std::move(func); if (resolveTail && registration->cursor == "$") { ItemStream latest; if (_redis.xrevrange(key, "+", "-", 1, back_inserter(latest))) registration->cursor = latest.empty() ? "0-0" : latest.front().first; } + registration->observedCursor = registration->cursor; auto cursor = info.keyids.find(key); if (cursor == info.keyids.end() || cursor->second == "$" || (registration->cursor != "$" && compareStreamIds(registration->cursor, cursor->second) < 0)) @@ -438,6 +449,7 @@ void RedisAdapter::remove_registration(uint64_t id) registrations.erase(registration); if (registrations.empty()) { info.keyids.erase(key->first); + info.boundaries.erase(key->first); info.subs.erase(key); } if (info.subs.empty()) _reader.erase(bucket); @@ -460,6 +472,7 @@ bool RedisAdapter::remove_reader_helper(const string& baseKey, const string& sub for (auto& registration : entry->second) registration->active = false; info.subs.erase(entry); info.keyids.erase(key); + info.boundaries.erase(key); found = true; if (info.subs.empty()) bucket = _reader.erase(bucket); else { start_reader(bucket->first); ++bucket; } @@ -467,6 +480,94 @@ bool RedisAdapter::remove_reader_helper(const string& baseKey, const string& sub return found; } +void RedisAdapter::reader_result(reader_info& info, RedisConnection::ReadStatus result) +{ + using Result = RedisConnection::ReadStatus; + for (const auto& key : info.subs) for (const auto& registration : key.second) { + lock_guard lock(registration->mutex); + auto& status = registration->status; + if (result == Result::TimedOut) { ++status.socketTimeouts; continue; } + if (result == Result::Rejected) { ++status.readRejections; continue; } + if (result == Result::Unavailable) { + ++status.readFailures; status.connected = false; registration->continuityCheck = true; + } else { + if (!status.connected && registration->everConnected) ++status.reconnects; + status.connected = true; registration->everConnected = true; + } + } +} + +bool RedisAdapter::inspect_readers(reader_info& info, const vector& keys, + size_t& nextProbe, bool force) +{ + if (!_options.readerProbeMs || keys.empty()) return true; + using Kind = RedisConnection::StreamKind; + using Result = RedisConnection::CommandStatus; + // Limit work per pass and stop promptly even if a backend is unavailable. + for (size_t checked = 0; checked < min(keys.size(), size_t(64)) && info.run; ++checked) { + const auto& key = keys[nextProbe++ % keys.size()]; + auto& boundary = info.boundaries[key]; + const auto now = steady_clock::now(); + if (!force && now < boundary.nextProbe) continue; + const auto bounds = _redis.streamBounds(key); + boundary.nextProbe = steady_clock::now() + milliseconds(_options.readerProbeMs); + bool rewound = false; + for (const auto& registration : info.subs.at(key)) { + lock_guard lock(registration->mutex); + auto& status = registration->status; + if (bounds.status != Result::Accepted) { + status.inspected = false; + if (bounds.status == Result::Rejected) ++status.inspectionRejections; + else { ++status.inspectionFailures; status.connected = false; registration->continuityCheck = true; } + continue; + } + if (!status.connected && registration->everConnected) ++status.reconnects; + status.connected = true; registration->everConnected = true; status.inspected = true; + const auto prior = status.stream; + status.stream = bounds.kind; + status.hasData = bounds.kind == Kind::Stream && bounds.firstId != "0-0"; + const bool hasCursor = registration->observedCursor != "$" && registration->observedCursor != "0-0"; + const bool missing = bounds.kind == Kind::Missing || bounds.kind == Kind::Invalid; + const bool reset = missing ? (hasCursor || prior == Kind::Stream) + : hasCursor && compareStreamIds(bounds.lastGeneratedId, registration->observedCursor) < 0; + if (reset) { + registration->cursor = registration->observedCursor = "0-0"; + ++status.epoch; ++status.streamResets; + if (bounds.kind == Kind::Missing) ++status.disappearances; + rewound = true; + } else if (bounds.kind == Kind::Stream && registration->continuityCheck && hasCursor && + bounds.firstId != "0-0" && compareStreamIds(bounds.firstId, registration->observedCursor) > 0) { + // The saved cursor is older than retained history. This is a possible + // continuity loss, never an exact count of missed entries. + ++status.retentionGaps; + } + registration->continuityCheck = false; + } + if (bounds.status == Result::Accepted) { + boundary.kind = bounds.kind; + reader_result(info, RedisConnection::ReadStatus::Accepted); + } + if (rewound) { + // Rewind the shared read only as far as the registrations require. + auto cursor = string("$"); + for (const auto& registration : info.subs.at(key)) { + lock_guard lock(registration->mutex); + if (registration->cursor != "$" && (cursor == "$" || compareStreamIds(registration->cursor, cursor) < 0)) + cursor = registration->cursor; + } + info.keyids[key] = cursor; + } + if (bounds.status == Result::Unavailable) { + for (const auto& other : info.subs) for (const auto& registration : other.second) { + lock_guard lock(registration->mutex); + registration->status.connected = false; registration->continuityCheck = true; + } + return false; + } + } + return true; +} + bool RedisAdapter::start_reader(uint32_t token) { if (_readers_defer) return true; @@ -498,7 +599,7 @@ bool RedisAdapter::start_reader(uint32_t token) if (info.keyids.at(subscriptions.first) == "$") info.keyids[subscriptions.first] = tail; for (const auto& registration : subscriptions.second) { lock_guard cursorLock(registration->mutex); - if (registration->cursor == "$") registration->cursor = tail; + if (registration->cursor == "$") registration->cursor = registration->observedCursor = tail; } } info.thread = thread([this, &info]() { @@ -507,9 +608,30 @@ bool RedisAdapter::start_reader(uint32_t token) info.started = true; } info.start_cv.notify_all(); + vector probeKeys; + for (const auto& key : info.subs) + if (any_of(key.second.begin(), key.second.end(), [](const auto& registration) { return registration->inspect; })) + probeKeys.push_back(key.first); + size_t nextProbe = 0; + bool forceProbe = false; for (Streams out; info.run; out.clear()) { - if (!_redis.xreadMultiBlock(info.keyids.begin(), info.keyids.end(), _options.cxn.timeout, - inserter(out, out.end()))) { + if (!inspect_readers(info, probeKeys, nextProbe, forceProbe)) { + forceProbe = true; + if (info.run) this_thread::sleep_for(milliseconds(50)); + continue; + } + if (!info.run) break; + forceProbe = false; + auto cursors = info.keyids; + for (const auto& boundary : info.boundaries) + if (boundary.second.kind == RedisConnection::StreamKind::Invalid) cursors.erase(boundary.first); + if (cursors.empty()) { this_thread::sleep_for(milliseconds(50)); continue; } + RedisConnection::ReadStatus readStatus; + const bool success = _redis.xreadMultiBlock(cursors.begin(), cursors.end(), _options.cxn.timeout, + inserter(out, out.end()), &readStatus); + reader_result(info, readStatus); + if (!success) { + forceProbe = true; // redis++ reconnects its socket on the next read. A transient error must // not permanently stop a reader while a separate health PING succeeds. if (info.run) this_thread::sleep_for(milliseconds(50)); @@ -524,20 +646,38 @@ bool RedisAdapter::start_reader(uint32_t token) const auto sub = parts.first.empty() ? item.first : parts.second; auto data = make_shared(std::move(item.second)); for (const auto& registration : subscriptions->second) { - _replier_pool.job(item.first, [registration, base, sub, data]() { + uint64_t epoch; + { + lock_guard cursorLock(registration->mutex); + epoch = registration->status.epoch; + if (registration->observedCursor == "$" || compareStreamIds(data->back().first, registration->observedCursor) > 0) + registration->observedCursor = data->back().first; + } + _replier_pool.job(item.first, [registration, base, sub, data, epoch]() { if (!registration->active.load()) return; ItemStream fresh; { lock_guard cursorLock(registration->mutex); + if (registration->status.epoch != epoch) return; for (const auto& entry : *data) { if (registration->cursor == "$" || compareStreamIds(entry.first, registration->cursor) > 0) { fresh.push_back(entry); registration->cursor = entry.first; } } + if (!fresh.empty()) { + ++registration->status.callbacks; registration->status.entries += fresh.size(); + registration->status.lastReceived = steady_clock::now(); + } + } + if (!fresh.empty() && registration->active.load()) { + try { registration->callback(base, sub, fresh); } + catch (...) { + lock_guard lock(registration->mutex); + ++registration->status.callbackErrors; + throw; + } } - if (!fresh.empty() && registration->active.load()) - registration->callback(base, sub, fresh); }); } } diff --git a/RedisAdapter.hpp b/RedisAdapter.hpp index fc29a56..444c916 100644 --- a/RedisAdapter.hpp +++ b/RedisAdapter.hpp @@ -102,6 +102,8 @@ struct RA_Options std::string dogname; uint16_t workers = 1; uint16_t readers = 1; + // Owned readers inspect continuity at this interval; zero disables XINFO. + uint32_t readerProbeMs = 1000; }; //^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ @@ -142,6 +144,17 @@ class RedisAdapter bool present() const { return !fields.empty() && id != "0-0"; } }; + struct ReaderStatus { + bool active = false, connected = false, inspected = false, hasData = false; + RedisConnection::StreamKind stream = RedisConnection::StreamKind::Unknown; + std::string cursor = "0-0", observedCursor = "0-0"; + uint64_t epoch = 1, readFailures = 0, readRejections = 0, socketTimeouts = 0; + uint64_t reconnects = 0, streamResets = 0, disappearances = 0, retentionGaps = 0; + uint64_t inspectionFailures = 0, inspectionRejections = 0; + uint64_t callbacks = 0, entries = 0, callbackErrors = 0; + std::chrono::steady_clock::time_point lastReceived{}; + }; + // Cancellation prevents queued callbacks from starting. An already executing // callback may finish; consumers must fence their own generation's mutations. class ReaderHandle { @@ -154,6 +167,7 @@ class RedisAdapter ReaderHandle& operator=(const ReaderHandle&) = delete; void reset() noexcept; explicit operator bool() const; + ReaderStatus status() const; private: friend class RedisAdapter; ReaderHandle(std::weak_ptr owner, std::shared_ptr registration); @@ -611,6 +625,9 @@ class RedisAdapter std::atomic active{true}; std::mutex mutex; std::string cursor; + std::string observedCursor; + bool inspect = false, everConnected = false, continuityCheck = true; + ReaderStatus status; reader_sub_fn callback; }; std::shared_ptr _reader_owner = std::make_shared(); @@ -622,6 +639,11 @@ class RedisAdapter std::unordered_map>> subs; std::unordered_map keyids; std::string stop; + struct Boundary { + std::chrono::steady_clock::time_point nextProbe{}; + RedisConnection::StreamKind kind = RedisConnection::StreamKind::Unknown; + }; + std::unordered_map boundaries; std::atomic run = false; // used by start_reader() to confirm the reader thread has begun its read loop - @@ -632,6 +654,9 @@ class RedisAdapter bool started = false; }; std::unordered_map _reader; + bool inspect_readers(reader_info& info, const std::vector& keys, + size_t& nextProbe, bool force); + void reader_result(reader_info& info, RedisConnection::ReadStatus result); ThreadPool _replier_pool; }; diff --git a/RedisConnection.hpp b/RedisConnection.hpp index 53017ad..ba96d81 100644 --- a/RedisConnection.hpp +++ b/RedisConnection.hpp @@ -19,6 +19,13 @@ class RedisConnection { public: enum class CommandStatus { Accepted, Rejected, Unavailable }; + enum class ReadStatus { Accepted, TimedOut, Rejected, Unavailable }; + enum class StreamKind { Unknown, Missing, Stream, Invalid }; + struct StreamBounds { + CommandStatus status = CommandStatus::Unavailable; + StreamKind kind = StreamKind::Unknown; + std::string firstId = "0-0", lastGeneratedId = "0-0"; + }; struct WriteResult { CommandStatus status = CommandStatus::Unavailable; std::string id; @@ -295,19 +302,62 @@ class RedisConnection // https://redis.io/docs/reference/cluster-spec/ // template - bool xreadMultiBlock(Input fst, Input lst, uint32_t tmo, Output out) + bool xreadMultiBlock(Input fst, Input lst, uint32_t tmo, Output out, ReadStatus* status = nullptr) { + if (status) *status = ReadStatus::Unavailable; auto [cluster, singler] = snapshot(); try { - if (cluster) { cluster->xread(fst, lst, chr::milliseconds(tmo), out); return true; } - if (singler) { singler->xread(fst, lst, chr::milliseconds(tmo), out); return true; } + if (cluster) { cluster->xread(fst, lst, chr::milliseconds(tmo), out); if (status) *status = ReadStatus::Accepted; return true; } + if (singler) { singler->xread(fst, lst, chr::milliseconds(tmo), out); if (status) *status = ReadStatus::Accepted; return true; } + } + catch (const swr::TimeoutError&) { if (status) *status = ReadStatus::TimedOut; return true; } + catch (const swr::ReplyError& e) { + if (status) *status = ReadStatus::Rejected; + syslog(LOG_ERR, "RedisConnection::%s %s", __func__, e.what()); } - catch (const swr::TimeoutError&) { return true; } catch (const swr::Error& e) { syslog(LOG_ERR, "RedisConnection::%s %s", __func__, e.what()); } return false; } + // XINFO is read-only and requires no Lua/script privilege. Its first/last + // entries are decoded only for IDs; no field payload is copied into the result. + StreamBounds streamBounds(const std::string& key) { + auto [cluster, singler] = snapshot(); + if (!cluster && !singler) return {}; + try { + // Route by the actual key, not the XINFO subcommand (STREAM). + const auto command = [](swr::Connection& connection, const swr::StringView& name) { + connection.send("XINFO STREAM %b", name.data(), name.size()); + }; + const auto reply = cluster ? cluster->command(command, key) : singler->command("XINFO", "STREAM", key); + StreamBounds result; + result.status = CommandStatus::Rejected; + if (!reply || reply->type != REDIS_REPLY_ARRAY || reply->elements % 2) return result; + const auto text = [](const redisReply* value) { + return value && value->type == REDIS_REPLY_STRING ? std::string(value->str, value->len) : std::string{}; + }; + bool hasLastId = false; + for (size_t i = 0; i < reply->elements; i += 2) { + const auto field = text(reply->element[i]); + const auto value = reply->element[i + 1]; + if (field == "last-generated-id") { + result.lastGeneratedId = text(value); hasLastId = !result.lastGeneratedId.empty(); + } else if (field == "first-entry" && value->type == REDIS_REPLY_ARRAY && value->elements) + result.firstId = text(value->element[0]); + } + if (hasLastId) { result.status = CommandStatus::Accepted; result.kind = StreamKind::Stream; } + return result; + } catch (const swr::ReplyError& error) { + const std::string message(error.what()); + if (message.find("no such key") != std::string::npos) + return {CommandStatus::Accepted, StreamKind::Missing}; + if (message.find("WRONGTYPE") != std::string::npos) + return {CommandStatus::Accepted, StreamKind::Invalid}; + return {CommandStatus::Rejected, StreamKind::Unknown}; + } catch (const swr::Error&) { return {}; } + } + //^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ // xadd : add an element to the specified stream // diff --git a/docs/api.md b/docs/api.md index 3f59b29..be945de 100644 --- a/docs/api.md +++ b/docs/api.md @@ -32,6 +32,7 @@ standalone connection. Setting `cxn.path` selects a Unix-domain socket and makes | `dogname` | `std::string` | empty | If set, maintain a one-second field-TTL watchdog for this name. | | `workers` | `uint16_t` | `1` | Worker threads used to dispatch reader callbacks. | | `readers` | `uint16_t` | `1` | Reader threads across which stream keys are deterministically sharded. | +| `readerProbeMs` | `uint32_t` | `1000` | Minimum boundary-inspection interval for owned subscriptions; `0` disables inspection and automatic reset recovery. See [continuity/status](stream-subscriptions.md#continuity-and-status). | Credentials are passed directly to redis-plus-plus. Keep them out of source control and populate `RA_Options` from the consuming application's secret or diff --git a/docs/stream-subscriptions.md b/docs/stream-subscriptions.md index 822d290..b39d8b2 100644 --- a/docs/stream-subscriptions.md +++ b/docs/stream-subscriptions.md @@ -40,3 +40,60 @@ consumer to report invalid input instead of silently replacing its last value. Redis connection, socket, and connection-pool waits use the configured timeout. Stream readers retry transient read failures; the existing health/reconnect path still handles initially disconnected adapters and backend rediscovery. + +## Continuity and status + +Owned readers periodically inspect stream boundaries with read-only `XINFO +STREAM`. `RA_Options.readerProbeMs` defaults to 1000 ms; zero disables inspection +and automatic reset recovery. Legacy unowned subscriptions do not enable probes +by themselves. A pass inspects at most 64 keys per bucket and stops on transport +failure; the interval is a minimum between checks, not a deadline for inspecting +every key in a large bucket. XINFO replies include the first/last entry payloads, +so include this extra traffic when qualifying large-frame workloads. + +Deletion, replacement with a non-stream key, or a last-generated ID below an +owned registration's observed cursor starts a new stream epoch. Its delivered +and queued cursors are rewound, and queued callbacks from the previous epoch are +discarded. An already executing callback may finish; callbacks on that key stay +serialized. Recreated streams can then deliver IDs below the previous stream's +maximum. This also means an explicit future cursor is treated as a reset when +inspection finds a lower last-generated ID; disable probing for strict +wait-until-that-future-ID behavior. Ordinary connection outages preserve cursors. +Trimming a stream to empty preserves its last-generated ID and does not rewind. + +An inspected wrong-type source is excluded from the shared XREAD until a later +inspection finds it repaired, allowing other streams in the bucket to continue. +No new Redis write or Lua permission is needed. If XINFO is denied, existing +permitted reads continue with `inspected=false` and an inspection-rejection +counter; automatic reset detection is unavailable for that source. Permit +`XINFO` on the same keys for complete continuity observation. + +`ReaderHandle::status()` returns a thread-safe snapshot of registration counters +and cursors without I/O. As with the handle itself, do not race it against moving +or resetting that same handle. `active` means the registration is retained, not +that Redis is currently reachable. `connected` reflects the latest bucket read +or inspection transport observation. `stream` and `hasData` describe the last +successful boundary inspection and are meaningful only when `inspected=true`. +`cursor` is the callback deduplication cursor; `observedCursor` also includes +queued data. `lastReceived` is a process-local monotonic callback-delivery time. +`callbacks` and `entries` count attempted callback deliveries; `callbackErrors` +counts exceptions, which do not trigger replay. + +Read rejections, transport failures, socket timeouts, inspection failures and +inspection rejections have separate counters. A socket timeout is not proof of +connection loss because Redis blocking-read and client socket deadlines can +overlap. `reconnects` counts observed unavailable-to-connected transitions. +`streamResets` and `disappearances` expose observed stream epochs and absence. +`retentionGaps` counts continuity checks whose saved cursor predates retained +history, initially or after a transport failure. It is a possible gap, not a +number of lost entries. Redis IDs do not encode an exact entry count or a stable +stream identity: delete/recreate cycles entirely between inspections may be +undetectable when the new stream has already passed the old ID. Consumer frame +IDs or an application epoch are still needed to establish exact continuity. + +`RedisAdapter.Recovery` exercises reset, absence, wrong-type isolation, trim, +connection recovery, callback fencing and denied inspection permissions. When +`redis-server` and `redis-cli` are available at configure time, +`RedisAdapter.ClusterRecovery` also runs against three private loopback nodes. +Its helper creates a temporary cluster and cleans up only those processes; it +never adds nodes to an existing cluster. diff --git a/reader_cluster_test.py b/reader_cluster_test.py new file mode 100644 index 0000000..b702a5f --- /dev/null +++ b/reader_cluster_test.py @@ -0,0 +1,80 @@ +#!/usr/bin/env python3 +"""Run reader recovery against three private Redis nodes, then stop only those processes. + +Usage: reader_cluster_test.py /path/to/redis-server /path/to/redis-cli /path/to/reader-test +""" +import os +from pathlib import Path +import socket +import subprocess +import sys +import tempfile +import time + +if len(sys.argv) != 4: + raise SystemExit(__doc__) +server, cli, executable = sys.argv[1:] +processes = [] +with tempfile.TemporaryDirectory(prefix='adapter-cluster-') as directory: + root = Path(directory) + reservations = [] + ports = [] + for i in range(6): + listener = socket.socket() + listener.bind(('127.0.0.1', 0)) + ports.append(listener.getsockname()[1]) + reservations.append(listener) + for listener in reservations: + listener.close() + client_ports = ports[:3] + try: + for index, port in enumerate(client_ports): + node = root / str(index) + node.mkdir() + log = (node / 'redis.log').open('w') + process = subprocess.Popen([server, '--bind', '127.0.0.1', '--port', str(port), + '--cluster-enabled', 'yes', '--cluster-port', str(ports[index + 3]), + '--cluster-announce-ip', '127.0.0.1', '--cluster-node-timeout', '1000', + '--cluster-config-file', 'nodes.conf', '--dir', str(node), + '--save', '', '--appendonly', 'no'], stdout=log, stderr=subprocess.STDOUT) + log.close() + processes.append(process) + for port in client_ports: + for attempt in range(100): + ping = subprocess.run([cli, '-h', '127.0.0.1', '-p', str(port), 'PING'], + capture_output=True, text=True, timeout=2) + if ping.returncode == 0 and 'PONG' in ping.stdout: + break + time.sleep(.05) + else: + raise RuntimeError('cluster node startup timed out') + creation = subprocess.run([cli, '--cluster', 'create', + *[f'127.0.0.1:{port}' for port in client_ports], + '--cluster-replicas', '0', '--cluster-yes'], capture_output=True, text=True, timeout=30) + assert creation.returncode == 0, creation.stdout + creation.stderr + print(creation.stdout) + for port in client_ports: + for attempt in range(100): + status = subprocess.check_output([cli, '-h', '127.0.0.1', '-p', str(port), 'CLUSTER', 'INFO'], + text=True, timeout=2) + if 'cluster_state:ok' in status: + break + time.sleep(.05) + else: + raise RuntimeError('cluster not ready') + test = subprocess.run([executable], env=dict(os.environ, REDIS_ADAPTER_TEST_PORT=str(client_ports[0]), + REDIS_ADAPTER_CLUSTER_TEST='1'), timeout=30) + assert test.returncode == 0, test.returncode + except BaseException: + for log in root.glob('*/redis.log'): + print(log, log.read_text(), file=sys.stderr) + raise + finally: + for process in processes: + process.terminate() + for process in processes: + try: + process.wait(timeout=5) + except subprocess.TimeoutExpired: + process.kill() + process.wait(timeout=5) diff --git a/recovery_test.cpp b/recovery_test.cpp new file mode 100644 index 0000000..78b7b3b --- /dev/null +++ b/recovery_test.cpp @@ -0,0 +1,180 @@ +#include "RedisAdapter.hpp" +#include +#include +#include +#include +#include +#include +#include +#include +#include + +using namespace std::chrono_literals; + +template void eventually(F predicate) { + const auto deadline = std::chrono::steady_clock::now() + 4s; + while (!predicate()) { + assert(std::chrono::steady_clock::now() < deadline); + std::this_thread::sleep_for(5ms); + } +} + +struct Observed { + std::mutex mutex; + std::condition_variable changed; + std::vector ids; + void append(const RedisAdapter::StreamBatch& entries) { + std::lock_guard guard(mutex); + for (const auto& entry : entries) ids.push_back(entry.first); + changed.notify_all(); + } + bool wait(const std::string& id) { + std::unique_lock lock(mutex); + return changed.wait_for(lock, 4s, [&] { return std::find(ids.begin(), ids.end(), id) != ids.end(); }); + } + size_t count(const std::string& id) { + std::lock_guard guard(mutex); + return std::count(ids.begin(), ids.end(), id); + } +}; + +int main() { + RA_Options options; + if (const auto* port = std::getenv("REDIS_ADAPTER_TEST_PORT")) options.cxn.port = std::stoi(port); + options.readerProbeMs = 50; + sw::redis::ConnectionOptions connection; + connection.host = "127.0.0.1"; connection.port = options.cxn.port; + if (std::getenv("REDIS_ADAPTER_CLUSTER_TEST")) { + sw::redis::RedisCluster control(connection); + for (const auto& base : {"a", "b", "c"}) { + const auto key = "{" + std::string(base) + "}:value"; + const RedisAdapter::Attrs fields{{"_", "cluster-value"}}; + control.xadd(key, "1000-0", fields.begin(), fields.end()); + RedisAdapter adapter(base, options); + Observed observed; + auto handle = adapter.subscribeStream("value", [&](const auto&, const auto&, const auto& entries) { + observed.append(entries); + }, "0-0"); + assert(observed.wait("1000-0") && handle.status().inspected); + control.del(key); control.xadd(key, "1-0", fields.begin(), fields.end()); + assert(observed.wait("1-0") && handle.status().streamResets == 1); + assert(handle.status().inspectionRejections == 0); + } + std::cout << "cluster-routed inspection and reset recovery passed\n"; + return 0; + } + sw::redis::Redis control(connection); + const auto base = "stream-recovery-" + std::to_string(getpid()); + const auto key = "{" + base + "}:value"; + const RedisAdapter::Attrs fields{{"_", "value"}}; + control.xadd(key, "1000-0", fields.begin(), fields.end()); + RedisAdapter adapter(base, options); + Observed observed; + auto handle = adapter.subscribeStream("value", [&](const auto&, const auto&, const auto& entries) { + observed.append(entries); + }, "0-0"); + assert(observed.wait("1000-0")); + control.del(key); + control.xadd(key, "1-0", fields.begin(), fields.end()); + assert(observed.wait("1-0") && "recreated stream must resume below the previous cursor"); + assert(handle.status().streamResets == 1 && handle.status().epoch == 2); + control.xadd(key, "2-0", fields.begin(), fields.end()); + assert(observed.wait("2-0")); + assert(handle.status().cursor == "2-0" && handle.status().observedCursor == "2-0"); + + control.del(key); + eventually([&] { return handle.status().stream == RedisConnection::StreamKind::Missing; }); + assert(handle.status().disappearances == 1 && !handle.status().hasData); + control.xadd(key, "3-0", fields.begin(), fields.end()); + assert(observed.wait("3-0")); + + // A wrong-type key must not stall other readers in the same bucket. + Observed healthy; + auto other = adapter.subscribeStream("healthy", [&](const auto&, const auto&, const auto& entries) { + healthy.append(entries); + }, "0-0"); + control.del(key); control.set(key, "wrong-type"); + eventually([&] { return handle.status().stream == RedisConnection::StreamKind::Invalid; }); + control.xadd("{" + base + "}:healthy", "10-0", fields.begin(), fields.end()); + assert(healthy.wait("10-0")); + control.del(key); control.xadd(key, "4-0", fields.begin(), fields.end()); + assert(observed.wait("4-0")); + + // Connection loss preserves cursors and cannot duplicate prior callbacks. + const auto failures = handle.status().readFailures + handle.status().inspectionFailures; + const auto resets = handle.status().streamResets; + control.command("CLIENT", "KILL", "TYPE", "normal", "SKIPME", "yes"); + eventually([&] { const auto status = handle.status(); return status.readFailures + status.inspectionFailures > failures; }); + eventually([&] { return handle.status().connected && handle.status().reconnects > 0; }); + control.xadd(key, "5-0", fields.begin(), fields.end()); + assert(observed.wait("5-0")); + assert(observed.count("4-0") == 1 && handle.status().streamResets == resets); + + // A cursor below retained history is a possible gap, not a fabricated count + // of missed records. Exact trim keeps the latest stream ID intact. + const auto trimmed = "{" + base + "}:trimmed"; + control.xadd(trimmed, "10-0", fields.begin(), fields.end()); + control.xadd(trimmed, "20-0", fields.begin(), fields.end()); + control.xtrim(trimmed, 1, false); + Observed retained; + auto trim = adapter.subscribeStream("trimmed", [&](const auto&, const auto&, const auto& entries) { + retained.append(entries); + }, "10-0"); + assert(retained.wait("20-0")); + assert(trim.status().retentionGaps == 1 && trim.status().streamResets == 0); + control.xtrim(trimmed, 0, false); + eventually([&] { return !trim.status().hasData; }); + assert(trim.status().streamResets == 0 && trim.status().cursor == "20-0"); + control.xadd(trimmed, "30-0", fields.begin(), fields.end()); + assert(retained.wait("30-0")); + + // Fence data already queued from the old stream while a worker is busy. + std::mutex mutex; + std::condition_variable changed; + bool entered = false, release = false; + auto blocker = adapter.subscribeStream("blocker", [&](const auto&, const auto&, const auto&) { + std::unique_lock lock(mutex); + entered = true; changed.notify_all(); + assert(changed.wait_for(lock, 10s, [&] { return release; })); + }, "0-0"); + control.xadd("{" + base + "}:blocker", "1-0", fields.begin(), fields.end()); + { std::unique_lock lock(mutex); assert(changed.wait_for(lock, 4s, [&] { return entered; })); } + const auto fencedKey = "{" + base + "}:fenced"; + control.xadd(fencedKey, "1000-0", fields.begin(), fields.end()); + Observed fenced; + auto queued = adapter.subscribeStream("fenced", [&](const auto&, const auto&, const auto& entries) { + fenced.append(entries); + }, "0-0"); + eventually([&] { return queued.status().observedCursor == "1000-0"; }); + control.del(fencedKey); control.xadd(fencedKey, "1-0", fields.begin(), fields.end()); + eventually([&] { return queued.status().streamResets == 1 && queued.status().observedCursor == "1-0"; }); + { std::lock_guard lock(mutex); release = true; changed.notify_all(); } + assert(fenced.wait("1-0")); + assert(fenced.count("1000-0") == 0 && fenced.count("1-0") == 1); + + auto throwing = adapter.subscribeStream("errors", [](const auto&, const auto&, const auto&) { + throw std::runtime_error("intentional callback failure"); + }, "0-0"); + control.xadd("{" + base + "}:errors", "1-0", fields.begin(), fields.end()); + eventually([&] { return throwing.status().callbackErrors == 1; }); + + // Existing read-only credentials can keep delivering when XINFO is denied; + // inspection unavailability is explicit instead of masquerading as continuity. + const auto username = "recovery-reader-" + std::to_string(getpid()); + control.command("ACL", "SETUSER", username, "reset", "on", ">test-only", "~{" + base + "}:*", + "+ping", "+xread", "+xrevrange"); + { + auto restrictedOptions = options; + restrictedOptions.cxn.user = username; restrictedOptions.cxn.password = "test-only"; + RedisAdapter restricted(base, restrictedOptions); + Observed permitted; + auto read = restricted.subscribeStream("value", [&](const auto&, const auto&, const auto& entries) { + permitted.append(entries); + }, "0-0"); + assert(permitted.wait("5-0")); + eventually([&] { return read.status().inspectionRejections > 0; }); + assert(!read.status().inspected && read.status().connected); + } + control.command("ACL", "DELUSER", username); + std::cout << "stream reset, deletion, trim, outage, queue fencing and inspection status passed\n"; +} From 0c3b0e6b7fa11910afba961d065f17ac749fd345 Mon Sep 17 00:00:00 2001 From: Derek Steinkamp Date: Tue, 29 Sep 2026 15:15:20 -0500 Subject: [PATCH 2/3] Keep reader recovery and validation standalone-only --- CMakeLists.txt | 8 ---- RedisConnection.hpp | 10 ++--- docs/stream-subscriptions.md | 7 +--- reader_cluster_test.py | 80 ------------------------------------ recovery_test.cpp | 19 --------- 5 files changed, 5 insertions(+), 119 deletions(-) delete mode 100644 reader_cluster_test.py diff --git a/CMakeLists.txt b/CMakeLists.txt index b2069ff..37f11ea 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -100,14 +100,6 @@ if(REDIS_ADAPTER_TEST) add_test(NAME RedisAdapter.WriteErrors COMMAND redis-adapter-write-error-test) set_tests_properties(RedisAdapter.WriteErrors PROPERTIES TIMEOUT 30 RUN_SERIAL TRUE) find_package(Python3 REQUIRED COMPONENTS Interpreter) - find_program(REDIS_SERVER_EXECUTABLE redis-server) - find_program(REDIS_CLI_EXECUTABLE redis-cli) - if(REDIS_SERVER_EXECUTABLE AND REDIS_CLI_EXECUTABLE) - add_test(NAME RedisAdapter.ClusterRecovery - COMMAND ${Python3_EXECUTABLE} ${CMAKE_CURRENT_SOURCE_DIR}/reader_cluster_test.py - ${REDIS_SERVER_EXECUTABLE} ${REDIS_CLI_EXECUTABLE} $) - set_tests_properties(RedisAdapter.ClusterRecovery PROPERTIES TIMEOUT 75) - endif() add_test(NAME RedisAdapter.WriteTransportFailure COMMAND ${Python3_EXECUTABLE} ${CMAKE_CURRENT_SOURCE_DIR}/write_failure_test.py $) diff --git a/RedisConnection.hpp b/RedisConnection.hpp index ba96d81..126bc20 100644 --- a/RedisConnection.hpp +++ b/RedisConnection.hpp @@ -323,14 +323,10 @@ class RedisConnection // XINFO is read-only and requires no Lua/script privilege. Its first/last // entries are decoded only for IDs; no field payload is copied into the result. StreamBounds streamBounds(const std::string& key) { - auto [cluster, singler] = snapshot(); - if (!cluster && !singler) return {}; + const auto singler = snapshot().second; + if (!singler) return {}; try { - // Route by the actual key, not the XINFO subcommand (STREAM). - const auto command = [](swr::Connection& connection, const swr::StringView& name) { - connection.send("XINFO STREAM %b", name.data(), name.size()); - }; - const auto reply = cluster ? cluster->command(command, key) : singler->command("XINFO", "STREAM", key); + const auto reply = singler->command("XINFO", "STREAM", key); StreamBounds result; result.status = CommandStatus::Rejected; if (!reply || reply->type != REDIS_REPLY_ARRAY || reply->elements % 2) return result; diff --git a/docs/stream-subscriptions.md b/docs/stream-subscriptions.md index b39d8b2..7867518 100644 --- a/docs/stream-subscriptions.md +++ b/docs/stream-subscriptions.md @@ -92,8 +92,5 @@ undetectable when the new stream has already passed the old ID. Consumer frame IDs or an application epoch are still needed to establish exact continuity. `RedisAdapter.Recovery` exercises reset, absence, wrong-type isolation, trim, -connection recovery, callback fencing and denied inspection permissions. When -`redis-server` and `redis-cli` are available at configure time, -`RedisAdapter.ClusterRecovery` also runs against three private loopback nodes. -Its helper creates a temporary cluster and cleans up only those processes; it -never adds nodes to an existing cluster. +connection recovery, callback fencing and denied inspection permissions on +standalone Redis. Redis Cluster is outside the supported portfolio. diff --git a/reader_cluster_test.py b/reader_cluster_test.py deleted file mode 100644 index b702a5f..0000000 --- a/reader_cluster_test.py +++ /dev/null @@ -1,80 +0,0 @@ -#!/usr/bin/env python3 -"""Run reader recovery against three private Redis nodes, then stop only those processes. - -Usage: reader_cluster_test.py /path/to/redis-server /path/to/redis-cli /path/to/reader-test -""" -import os -from pathlib import Path -import socket -import subprocess -import sys -import tempfile -import time - -if len(sys.argv) != 4: - raise SystemExit(__doc__) -server, cli, executable = sys.argv[1:] -processes = [] -with tempfile.TemporaryDirectory(prefix='adapter-cluster-') as directory: - root = Path(directory) - reservations = [] - ports = [] - for i in range(6): - listener = socket.socket() - listener.bind(('127.0.0.1', 0)) - ports.append(listener.getsockname()[1]) - reservations.append(listener) - for listener in reservations: - listener.close() - client_ports = ports[:3] - try: - for index, port in enumerate(client_ports): - node = root / str(index) - node.mkdir() - log = (node / 'redis.log').open('w') - process = subprocess.Popen([server, '--bind', '127.0.0.1', '--port', str(port), - '--cluster-enabled', 'yes', '--cluster-port', str(ports[index + 3]), - '--cluster-announce-ip', '127.0.0.1', '--cluster-node-timeout', '1000', - '--cluster-config-file', 'nodes.conf', '--dir', str(node), - '--save', '', '--appendonly', 'no'], stdout=log, stderr=subprocess.STDOUT) - log.close() - processes.append(process) - for port in client_ports: - for attempt in range(100): - ping = subprocess.run([cli, '-h', '127.0.0.1', '-p', str(port), 'PING'], - capture_output=True, text=True, timeout=2) - if ping.returncode == 0 and 'PONG' in ping.stdout: - break - time.sleep(.05) - else: - raise RuntimeError('cluster node startup timed out') - creation = subprocess.run([cli, '--cluster', 'create', - *[f'127.0.0.1:{port}' for port in client_ports], - '--cluster-replicas', '0', '--cluster-yes'], capture_output=True, text=True, timeout=30) - assert creation.returncode == 0, creation.stdout + creation.stderr - print(creation.stdout) - for port in client_ports: - for attempt in range(100): - status = subprocess.check_output([cli, '-h', '127.0.0.1', '-p', str(port), 'CLUSTER', 'INFO'], - text=True, timeout=2) - if 'cluster_state:ok' in status: - break - time.sleep(.05) - else: - raise RuntimeError('cluster not ready') - test = subprocess.run([executable], env=dict(os.environ, REDIS_ADAPTER_TEST_PORT=str(client_ports[0]), - REDIS_ADAPTER_CLUSTER_TEST='1'), timeout=30) - assert test.returncode == 0, test.returncode - except BaseException: - for log in root.glob('*/redis.log'): - print(log, log.read_text(), file=sys.stderr) - raise - finally: - for process in processes: - process.terminate() - for process in processes: - try: - process.wait(timeout=5) - except subprocess.TimeoutExpired: - process.kill() - process.wait(timeout=5) diff --git a/recovery_test.cpp b/recovery_test.cpp index 78b7b3b..a0fbe69 100644 --- a/recovery_test.cpp +++ b/recovery_test.cpp @@ -44,25 +44,6 @@ int main() { options.readerProbeMs = 50; sw::redis::ConnectionOptions connection; connection.host = "127.0.0.1"; connection.port = options.cxn.port; - if (std::getenv("REDIS_ADAPTER_CLUSTER_TEST")) { - sw::redis::RedisCluster control(connection); - for (const auto& base : {"a", "b", "c"}) { - const auto key = "{" + std::string(base) + "}:value"; - const RedisAdapter::Attrs fields{{"_", "cluster-value"}}; - control.xadd(key, "1000-0", fields.begin(), fields.end()); - RedisAdapter adapter(base, options); - Observed observed; - auto handle = adapter.subscribeStream("value", [&](const auto&, const auto&, const auto& entries) { - observed.append(entries); - }, "0-0"); - assert(observed.wait("1000-0") && handle.status().inspected); - control.del(key); control.xadd(key, "1-0", fields.begin(), fields.end()); - assert(observed.wait("1-0") && handle.status().streamResets == 1); - assert(handle.status().inspectionRejections == 0); - } - std::cout << "cluster-routed inspection and reset recovery passed\n"; - return 0; - } sw::redis::Redis control(connection); const auto base = "stream-recovery-" + std::to_string(getpid()); const auto key = "{" + base + "}:value"; From 94470f5918c4bb544312f32f0e4a09847a122572 Mon Sep 17 00:00:00 2001 From: Derek Steinkamp Date: Fri, 2 Oct 2026 15:09:10 -0500 Subject: [PATCH 3/3] Associate reader recovery with batch admission evidence --- RedisAdapter.cpp | 30 ++++++++++++++++++--------- RedisAdapter.hpp | 17 ++++++++++++++-- docs/stream-subscriptions.md | 10 +++++++++ recovery_test.cpp | 39 ++++++++++++++++++++++++++++++++++++ 4 files changed, 84 insertions(+), 12 deletions(-) diff --git a/RedisAdapter.cpp b/RedisAdapter.cpp index 6e4c7ff..ca26bc0 100644 --- a/RedisAdapter.cpp +++ b/RedisAdapter.cpp @@ -165,9 +165,18 @@ RedisAdapter::ReaderHandle RedisAdapter::subscribeStream(const string& subKey, S RedisAdapter::ReaderHandle RedisAdapter::subscribeStreamWithEpoch(const string& subKey, EpochStreamCallback callback, const SubscriptionOptions& options) { if (!callback) throw invalid_argument("empty stream callback"); + return subscribeStreamWithMetadata(subKey, + [callback = std::move(callback)](const auto& base, const auto& sub, const auto& batch, const auto& metadata) { + callback(base, sub, batch, metadata.epoch); + }, options); +} + +RedisAdapter::ReaderHandle RedisAdapter::subscribeStreamWithMetadata(const string& subKey, MetadataStreamCallback callback, + const SubscriptionOptions& options) { + if (!callback) throw invalid_argument("empty stream callback"); const auto base = options.baseKey.empty() ? _base_key : options.baseKey; - auto wrapped = [base, subKey, callback = std::move(callback)](const auto&, const auto&, const auto& batch, uint64_t epoch) { - callback(base, subKey, batch, epoch); + auto wrapped = [base, subKey, callback = std::move(callback)](const auto&, const auto&, const auto& batch, const auto& metadata) { + callback(base, subKey, batch, metadata); }; return ReaderHandle(_reader_owner, register_reader(build_key(subKey, options.baseKey), [](const auto&, const auto&, const auto&) {}, options.afterId, true, @@ -364,14 +373,14 @@ bool RedisAdapter::add_reader_helper(const string& baseKey, const string& subKey } shared_ptr -RedisAdapter::register_reader(const string& key, reader_sub_fn func, const string& afterId, bool resolveTail, uint32_t probeMs, EpochStreamCallback epochCallback) +RedisAdapter::register_reader(const string& key, reader_sub_fn func, const string& afterId, bool resolveTail, uint32_t probeMs, MetadataStreamCallback metadataCallback) { if (!func) throw invalid_argument("empty stream callback"); if (afterId != "$") compareStreamIds(afterId, "0-0"); auto registration = make_shared(); registration->cursor = afterId; registration->callback = std::move(func); - registration->epochCallback = std::move(epochCallback); + registration->metadataCallback = std::move(metadataCallback); registration->probeMs = resolveTail ? (probeMs == UINT32_MAX ? _options.readerProbeMs : probeMs) : 0; registration->status.epoch = 1; uint32_t token; @@ -806,21 +815,22 @@ bool RedisAdapter::start_reader(uint32_t token) auto data = make_shared(std::move(item.second)); ++info.boundaries[item.first].readVersion; for (const auto& registration : subscriptions->second) { - uint64_t epoch; + StreamBatchMetadata metadata; { lock_guard lock(registration->mutex); - epoch = registration->status.epoch; + metadata.epoch = registration->status.epoch; + metadata.readRejections = registration->status.readRejections; if (registration->observedCursor == "$" || compareStreamIds(data->back().first, registration->observedCursor) > 0) registration->observedCursor = data->back().first; registration->status.streamKind = StreamKind::Stream; registration->status.hasData = true; } - _replier_pool.job(item.first, [registration, base, sub, data, epoch]() { + _replier_pool.job(item.first, [registration, base, sub, data, metadata]() { if (!registration->active.load()) return; size_t first = 0; { lock_guard cursorLock(registration->mutex); - if (registration->status.epoch != epoch) return; + if (registration->status.epoch != metadata.epoch) return; while (first < data->size() && registration->cursor != "$" && compareStreamIds((*data)[first].first, registration->cursor) <= 0) ++first; if (first == data->size()) return; @@ -829,13 +839,13 @@ bool RedisAdapter::start_reader(uint32_t token) if (!registration->active.load()) return; { lock_guard lock(registration->mutex); - if (registration->status.epoch != epoch) return; + if (registration->status.epoch != metadata.epoch) return; ++registration->status.callbacks; registration->status.entries += data->size() - first; registration->status.lastReceived = steady_clock::now(); } const auto invoke = [&](const ItemStream& batch) { - if (registration->epochCallback) registration->epochCallback(base, sub, batch, epoch); + if (registration->metadataCallback) registration->metadataCallback(base, sub, batch, metadata); else registration->callback(base, sub, batch); }; try { diff --git a/RedisAdapter.hpp b/RedisAdapter.hpp index c4f0d28..8f697a4 100644 --- a/RedisAdapter.hpp +++ b/RedisAdapter.hpp @@ -106,6 +106,14 @@ class RedisAdapter : public RedisStreamData using StreamKind = RedisConnection::StreamKind; using EpochStreamCallback = std::function; + struct StreamBatchMetadata { + uint64_t epoch = 0; + // Captured when this batch is read, before callback queueing. A delayed + // callback must not acknowledge rejection evidence from a later read. + uint64_t readRejections = 0; + }; + using MetadataStreamCallback = std::function; struct ReaderStatus { bool active = false, connected = false, inspected = false, hasData = false; StreamKind streamKind = StreamKind::Unknown; @@ -158,6 +166,11 @@ class RedisAdapter : public RedisStreamData [[nodiscard]] ReaderHandle subscribeStreamWithEpoch(const std::string& subKey, EpochStreamCallback callback) { return subscribeStreamWithEpoch(subKey, std::move(callback), SubscriptionOptions{}); } + [[nodiscard]] ReaderHandle subscribeStreamWithMetadata(const std::string& subKey, MetadataStreamCallback callback, + const SubscriptionOptions& options); + [[nodiscard]] ReaderHandle subscribeStreamWithMetadata(const std::string& subKey, MetadataStreamCallback callback) { + return subscribeStreamWithMetadata(subKey, std::move(callback), SubscriptionOptions{}); + } //^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ // Containers for getting/setting data using RedisAdapter methods // @@ -476,7 +489,7 @@ class RedisAdapter : public RedisStreamData reader_sub_fn func, const std::string& afterId, bool resolveTail = true, uint32_t probeMs = UINT32_MAX, - EpochStreamCallback epochCallback = {}); + MetadataStreamCallback metadataCallback = {}); void remove_registration(uint64_t id); template reader_sub_fn make_reader_callback(ReaderSubFn func) const; @@ -563,7 +576,7 @@ class RedisAdapter : public RedisStreamData std::mutex mutex; std::string cursor; reader_sub_fn callback; - EpochStreamCallback epochCallback; + MetadataStreamCallback metadataCallback; uint32_t probeMs = 0; bool everConnected = false, continuityCheck = true; std::string observedCursor, gapSignature; diff --git a/docs/stream-subscriptions.md b/docs/stream-subscriptions.md index f059faf..8456887 100644 --- a/docs/stream-subscriptions.md +++ b/docs/stream-subscriptions.md @@ -114,3 +114,13 @@ Queued older-epoch batches are fenced; an already executing callback may complet Three-argument callbacks should use their own frame IDs if reset attribution is required. Callback/entry counters record batches handed to the callback, while `observedCursor` can advance while a worker is still busy. + +`subscribeStreamWithMetadata()` additionally supplies `StreamBatchMetadata`, +containing the captured epoch and this registration's `readRejections` count at +the successful read. These values are captured before callback queueing and do +not change if a later XREAD is rejected while delivery is queued or executing. +Use that batch-associated counter to acknowledge read rejection recovery; polling +the mutable status in a delayed callback can incorrectly acknowledge a newer +rejection. `ReaderStatus::lastReceived` is callback-time receipt, not evidence of +when Redis accepted the read. The existing three-argument and epoch subscription +APIs retain their behavior. diff --git a/recovery_test.cpp b/recovery_test.cpp index 727deab..8239a65 100644 --- a/recovery_test.cpp +++ b/recovery_test.cpp @@ -85,6 +85,45 @@ TEST_P(Recovery, InspectionIsOptInAndCanBeSelectedPerSubscription) { EXPECT_EQ(RA::ReaderHandle{}.status().epoch, 0u); } +TEST_P(Recovery, BatchMetadataDoesNotAcknowledgeALaterReadRejection) { + RA adapter(base, options); + std::mutex mutex; + std::condition_variable changed; + bool entered = false, released = false; + std::vector observed; + auto handle = adapter.subscribeStreamWithMetadata("value", [&](const auto&, const auto&, const auto& entries, const auto& metadata) { + std::unique_lock lock(mutex); + if (entries.front().first == "1-0") { + entered = true; changed.notify_all(); + EXPECT_TRUE(changed.wait_for(lock, 4s, [&] { return released; })); + } + observed.push_back(metadata); changed.notify_all(); + }, selection(0)); + add("value", "1-0"); + { + std::unique_lock lock(mutex); + ASSERT_TRUE(changed.wait_for(lock, 3s, [&] { return entered; })); + } + control->command("ACL", "SETUSER", user, "-xread"); + EVENTUALLY(handle.status().readRejections > 0); + { + std::unique_lock lock(mutex); + released = true; changed.notify_all(); + ASSERT_TRUE(changed.wait_for(lock, 3s, [&] { return observed.size() == 1; })); + EXPECT_EQ(observed[0].readRejections, 0u); + EXPECT_EQ(observed[0].epoch, 1u); + } + EXPECT_GT(handle.status().readRejections, 0u); + control->command("ACL", "SETUSER", user, "+xread"); + add("value", "2-0"); + { + std::unique_lock lock(mutex); + ASSERT_TRUE(changed.wait_for(lock, 3s, [&] { return observed.size() == 2; })); + EXPECT_EQ(observed[1].readRejections, handle.status().readRejections); + EXPECT_EQ(observed[1].epoch, 1u); + } +} + TEST_P(Recovery, LowerIdRecreationIsDetectedWithZeroAndLongCommandTimeouts) { for (const auto timeout : {0u, 3000u}) { options.cxn.timeout = timeout;