diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index 52f3272..58416a6 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -51,6 +51,9 @@ jobs: REDIS_ADAPTER_TEST_SOCKET: /tmp/redis.sock REDIS_ADAPTER_ISOLATED_TEST: 1 + - name: Verify owned and legacy reads on Redis Cluster + run: python3 scripts/run-test-redis-cluster.py -- ctest --test-dir build -R ClusterRecovery --output-on-failure + - name: Show Redis logs on failure if: failure() run: docker compose -f docker-compose.test.yml logs diff --git a/CHANGELOG.md b/CHANGELOG.md index 0eb62f3..097ca52 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,8 @@ versions are independent of the RedisAdapter wire-protocol version. ### Added +- Read-side status, observed/delivered cursors and captured-epoch callbacks. +- Opt-in per-subscription continuity inspection with a bounded, persistent scheduler. - `RA_REJECTED` and structured stream write/trim results with error text and a topology-refresh decision, available to real and mock consumers. - Owned subscription handles, exact stream snapshots and shared safe wire decoders. @@ -26,6 +28,11 @@ versions are independent of the RedisAdapter wire-protocol version. ### Fixed +- Wrong-type/read-denied quarantine uses XREAD permissions and keeps healthy + neighbours flowing, including legacy registrations after an owned peer is removed. +- Unsupported/denied inspection cannot suppress normal Cluster or standalone reads. +- Idle replies, socket failures, empty retention loss and inactive-handle status + have distinct diagnostics; queued older-epoch batches are fenced. - READONLY refreshes future connections while preserving its known-rejection status. - Callback capture cleanup outside worker and reader locks, including exceptions and canceled jobs; peer registrations survive handle removal. diff --git a/CMakeLists.txt b/CMakeLists.txt index 5683c56..4eb3147 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -112,6 +112,10 @@ if(REDIS_ADAPTER_TEST) message(WARNING "Python3 not found: write proxy tests disabled; native tests remain enabled") endif() + add_executable(redis-adapter-recovery-test ${REDIS_ADAPTER_SOURCES} recovery_test.cpp) + target_link_libraries(redis-adapter-recovery-test ${REDIS_ADAPTER_LIBRARIES} GTest::gtest_main) + gtest_discover_tests(redis-adapter-recovery-test PROPERTIES TIMEOUT 20 RUN_SERIAL TRUE) + endif() # Build redis-adapter-benchmark and googlebenchmark if REDIS_ADAPTER_BENCHMARK is true diff --git a/RedisAdapter.cpp b/RedisAdapter.cpp index d7d3129..ca26bc0 100644 --- a/RedisAdapter.cpp +++ b/RedisAdapter.cpp @@ -112,6 +112,17 @@ 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 == "$" ? "" : registration_->cursor; + result.observedCursor = registration_->observedCursor == "$" ? "" : registration_->observedCursor; + if (!result.active) { result.connected = result.inspected = result.hasData = false; } + return result; +} + void RedisAdapter::ReaderHandle::reset() noexcept { auto registration = std::move(registration_); if (!registration) return; @@ -142,13 +153,34 @@ RedisAdapter::StreamSnapshot RedisAdapter::getStreamSnapshot(const string& subKe } RedisAdapter::ReaderHandle RedisAdapter::subscribeStream(const string& subKey, StreamCallback callback, - const string& afterId, const string& baseKey) { + const string& afterId, const string& baseKey, uint32_t probeMs) { if (!callback) throw invalid_argument("empty stream callback"); const auto base = baseKey.empty() ? _base_key : baseKey; return ReaderHandle(_reader_owner, register_reader(build_key(subKey, baseKey), [base, subKey, callback = std::move(callback)](const auto&, const auto&, const auto& batch) { callback(base, subKey, batch); - }, afterId)); + }, afterId, true, probeMs)); +} + +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, 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, + options.probeMs.value_or(UINT32_MAX), std::move(wrapped))); } //^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ @@ -262,6 +294,8 @@ bool RedisAdapter::removeGenericReader(const string& key) retired.push_back(std::move(entry->second)); info.subs.erase(entry); info.keyids.erase(key); + info.boundaries.erase(key); + info.quarantined.erase(key); found = true; if (info.subs.empty()) it = _reader.erase(it); else { start_reader(it->first); ++it; } @@ -339,13 +373,16 @@ 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) +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->metadataCallback = std::move(metadataCallback); + registration->probeMs = resolveTail ? (probeMs == UINT32_MAX ? _options.readerProbeMs : probeMs) : 0; + registration->status.epoch = 1; uint32_t token; { lock_guard lock(_reader_mtx); @@ -376,6 +413,9 @@ RedisAdapter::register_reader(const string& key, reader_sub_fn func, const strin if (!hadCursor || oldCursor == "$" || (registration->cursor != "$" && compareStreamIds(registration->cursor, oldCursor) < 0)) info.keyids[key] = registration->cursor; + registration->observedCursor = registration->cursor; + registration->status.connected = info.readConnected; + registration->everConnected = info.readConnected; info.subs[key].push_back(registration); if (info.stop.empty()) { info.stop = stopStreamKey(key); @@ -420,6 +460,8 @@ void RedisAdapter::remove_registration(uint64_t id) registrations.erase(registration); if (registrations.empty()) { info.keyids.erase(key->first); + info.boundaries.erase(key->first); + info.quarantined.erase(key->first); info.subs.erase(key); } if (info.subs.empty()) _reader.erase(bucket); @@ -444,6 +486,8 @@ bool RedisAdapter::remove_reader_helper(const string& baseKey, const string& sub retired.push_back(std::move(entry->second)); info.subs.erase(entry); info.keyids.erase(key); + info.boundaries.erase(key); + info.quarantined.erase(key); found = true; if (info.subs.empty()) bucket = _reader.erase(bucket); else { start_reader(bucket->first); ++bucket; } @@ -476,12 +520,225 @@ size_t RedisAdapter::resolve_reader_tails(reader_info& info) 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; } } return unresolved; } +void RedisAdapter::reader_result(reader_info& info, RedisConnection::ReadStatus result, bool socketTimedOut) +{ + using Result = RedisConnection::ReadStatus; + const bool connected = result == Result::Accepted || result == Result::TimedOut; + if (connected && info.readConnected) return; + if (result == Result::Rejected) return; // identify the rejected key on the read path + for (const auto& key : info.subs) for (const auto& registration : key.second) { + lock_guard lock(registration->mutex); + auto& status = registration->status; + if (connected) { + if (!status.connected && registration->everConnected) ++status.reconnects; + status.connected = true; registration->everConnected = true; + } else { + ++status.readFailures; + if (socketTimedOut) ++status.socketTimeouts; + status.connected = false; registration->continuityCheck = true; + } + } + info.readConnected = connected; +} + +void RedisAdapter::reset_stream(reader_info& info, const string& key, StreamKind kind, bool allReaders) +{ + string minimum = "$"; + for (const auto& registration : info.subs.at(key)) { + lock_guard lock(registration->mutex); + auto& status = registration->status; + if (allReaders || registration->probeMs) { + const bool hadCursor = registration->observedCursor != "$" && registration->observedCursor != "0-0"; + if (hadCursor || status.streamKind == StreamKind::Stream) { + ++status.epoch; ++status.streamResets; + if (kind == StreamKind::Missing) ++status.disappearances; + } + registration->cursor = registration->observedCursor = "0-0"; + registration->gapSignature.clear(); + status.hasData = false; status.lastReceived = {}; + } + status.streamKind = kind; + if (registration->cursor != "$" && (minimum == "$" || compareStreamIds(registration->cursor, minimum) < 0)) + minimum = registration->cursor; + } + info.keyids[key] = minimum; +} + +void RedisAdapter::check_readable(reader_info& info, const vector& keys) +{ + if (!info.run || keys.empty()) return; + const auto checked = _redis.probeReadable(keys); + if (checked.status == RedisConnection::ReadStatus::Unavailable) { + reader_result(info, checked.status); return; + } + const auto now = steady_clock::now(); + if (checked.status == RedisConnection::ReadStatus::Accepted) { + for (const auto& key : keys) { + const auto bad = info.quarantined.find(key); + if (bad != info.quarantined.end()) { + if (bad->second.wrongType) reset_stream(info, key, StreamKind::Unknown, true); + info.quarantined.erase(bad); + } + for (const auto& registration : info.subs.at(key)) { + lock_guard lock(registration->mutex); + if (!registration->status.connected && registration->everConnected) ++registration->status.reconnects; + registration->status.connected = true; registration->everConnected = true; + } + } + return; + } + // A command-wide ACL refusal needs no per-key search. Otherwise split the + // group to isolate bad keys without K serial round trips for one bad key. + const bool commandDenied = checked.error.rfind("NOPERM", 0) == 0 && + checked.error.find("to run the 'xread' command") != string::npos; + if (keys.size() > 1 && !commandDenied) { + const auto middle = keys.begin() + keys.size() / 2; + check_readable(info, vector(keys.begin(), middle)); + check_readable(info, vector(middle, keys.end())); + return; + } + for (const auto& key : keys) { + const auto inserted = info.quarantined.emplace(key, reader_info::Quarantine{}); + auto& bad = inserted.first->second; + if (checked.wrongType && (inserted.second || !bad.wrongType)) reset_stream(info, key, StreamKind::Invalid, true); + bad.wrongType |= checked.wrongType; + bad.nextCheck = now + milliseconds(checked.wrongType ? 100 : 1000); + if (inserted.second) syslog(LOG_WARNING, "stream key quarantined after read rejection: %s", key.c_str()); + for (const auto& registration : info.subs.at(key)) { + lock_guard lock(registration->mutex); + ++registration->status.readRejections; + registration->status.connected = true; registration->everConnected = true; + } + } +} + +void RedisAdapter::prepare_probes(reader_info& info) +{ + info.probes = {}; + const auto now = steady_clock::now(); + for (const auto& key : info.subs) { + uint32_t interval = UINT32_MAX; + for (const auto& registration : key.second) if (registration->probeMs) interval = min(interval, registration->probeMs); + auto& boundary = info.boundaries[key.first]; + if (interval == UINT32_MAX) { boundary.intervalMs = 0; continue; } + if (!boundary.intervalMs || boundary.denied || boundary.nextProbe == steady_clock::time_point{}) { + boundary.nextProbe = now; boundary.backoff = 1; boundary.denied = false; + } + boundary.intervalMs = interval; + info.probes.push({boundary.nextProbe, ++info.probeOrder, key.first}); + } +} + +void RedisAdapter::apply_bounds(reader_info& info, const string& key, const RedisConnection::StreamBounds& bounds) +{ + using Result = RedisConnection::CommandStatus; + if (bounds.status != Result::Accepted) { + for (const auto& registration : info.subs.at(key)) if (registration->probeMs) { + lock_guard lock(registration->mutex); + registration->status.inspected = false; + if (bounds.status == Result::Rejected) ++registration->status.inspectionRejections; + } + return; + } + bool reset = bounds.kind == StreamKind::Missing || bounds.kind == StreamKind::Invalid; + if (bounds.kind == StreamKind::Stream) { + for (const auto& registration : info.subs.at(key)) if (registration->probeMs) { + lock_guard lock(registration->mutex); + if (registration->observedCursor != "$" && registration->observedCursor != "0-0" && + compareStreamIds(bounds.lastGeneratedId, registration->observedCursor) < 0) reset = true; + } + } + if (reset) reset_stream(info, key, bounds.kind, bounds.kind == StreamKind::Invalid); + if (bounds.kind == StreamKind::Invalid) info.quarantined[key] = {steady_clock::now() + milliseconds(100), true}; + for (const auto& registration : info.subs.at(key)) if (registration->probeMs) { + lock_guard lock(registration->mutex); + auto& status = registration->status; + status.inspected = true; status.streamKind = bounds.kind; + status.hasData = bounds.kind == StreamKind::Stream && bounds.firstId != "0-0"; + const auto& cursor = registration->observedCursor; + const bool hasCursor = cursor != "$" && cursor != "0-0"; + const bool gap = !reset && bounds.kind == StreamKind::Stream && hasCursor && + compareStreamIds(bounds.lastGeneratedId, cursor) > 0 && + (bounds.firstId == "0-0" || compareStreamIds(bounds.firstId, cursor) > 0); + const auto signature = gap ? cursor + ":" + bounds.firstId + ":" + bounds.lastGeneratedId : string{}; + if (gap && signature != registration->gapSignature) ++status.retentionGaps; + registration->gapSignature = signature; + registration->continuityCheck = false; + } +} + +void RedisAdapter::inspect_readers(reader_info& info) +{ + constexpr size_t maximum = 16; + vector due; + const auto now = steady_clock::now(); + while (!info.probes.empty() && info.probes.top().due <= now && due.size() < maximum && info.run) { + auto ticket = info.probes.top(); info.probes.pop(); + auto& boundary = info.boundaries.at(ticket.key); + if (!boundary.intervalMs) continue; + bool continuity = false; + for (const auto& registration : info.subs.at(ticket.key)) if (registration->probeMs) { + lock_guard lock(registration->mutex); + continuity |= registration->continuityCheck; + } + if (!continuity && boundary.readVersion != boundary.lastProbeVersion) { + boundary.lastProbeVersion = boundary.readVersion; boundary.backoff = 1; + boundary.nextProbe = now + milliseconds(boundary.intervalMs); + info.probes.push({boundary.nextProbe, ++info.probeOrder, ticket.key}); + } else due.push_back(std::move(ticket.key)); + } + if (due.empty() || !info.run) return; + const auto bounds = _redis.streamBoundsBatch(due); + bool transportFailure = false; + for (size_t index = 0; index < due.size(); ++index) { + const auto& key = due[index]; + auto& boundary = info.boundaries.at(key); + const auto& result = bounds[index]; + apply_bounds(info, key, result); + boundary.lastProbeVersion = boundary.readVersion; + if (result.status == RedisConnection::CommandStatus::Rejected) { + boundary.denied = true; + boundary.nextProbe = steady_clock::now() + seconds(60); + } else { + transportFailure |= result.status == RedisConnection::CommandStatus::Unavailable; + boundary.backoff = min(8, boundary.backoff * 2); + boundary.nextProbe = steady_clock::now() + milliseconds(uint64_t(boundary.intervalMs) * boundary.backoff); + } + info.probes.push({boundary.nextProbe, ++info.probeOrder, key}); + } + if (transportFailure) { + // Inspection transport failure affects every registration in the bucket, + // independently of which pipelined command first observed it. + for (const auto& key : info.subs) for (const auto& registration : key.second) { + lock_guard lock(registration->mutex); + ++registration->status.inspectionFailures; + registration->status.inspected = false; registration->status.connected = false; + registration->continuityCheck = true; + } + info.readConnected = false; + } +} + +uint32_t RedisAdapter::read_interval(const reader_info& info, bool& blocking) const +{ + uint32_t interval = _options.cxn.timeout ? min(1000, _options.cxn.timeout) : 1000; + const auto now = steady_clock::now(); + blocking = true; + const auto shorten = [&](steady_clock::time_point due) { + if (due <= now) { blocking = false; interval = 1; } + else interval = min(interval, max(1, duration_cast(due - now).count())); + }; + if (!info.probes.empty()) shorten(info.probes.top().due); + for (const auto& bad : info.quarantined) shorten(bad.second.nextCheck); + return interval; +} + bool RedisAdapter::start_reader(uint32_t token) { if (_shutdown) return false; @@ -493,6 +750,7 @@ bool RedisAdapter::start_reader(uint32_t token) info.started = false; info.run = true; size_t unresolved = resolve_reader_tails(info); + prepare_probes(info); try { info.thread = thread([this, &info, unresolved]() mutable { { @@ -502,25 +760,49 @@ bool RedisAdapter::start_reader(uint32_t token) info.start_cv.notify_all(); uint32_t delay = 50; bool reported = false; + vector allKeys; + allKeys.reserve(info.subs.size()); + for (const auto& key : info.subs) allKeys.push_back(key.first); for (Streams out; info.run; out.clear()) { if (unresolved) unresolved = resolve_reader_tails(info); + vector recheck; + const auto now = steady_clock::now(); + for (const auto& bad : info.quarantined) if (bad.second.nextCheck <= now) recheck.push_back(bad.first); + if (!recheck.empty()) check_readable(info, recheck); + inspect_readers(info); + if (!info.run) break; // Never send '$' repeatedly. An unresolved tail is excluded until the // snapshot succeeds; healthy keys in the bucket continue to flow. const auto* cursors = &info.keyids; unordered_map resolved; - if (unresolved) { - for (const auto& key : info.keyids) if (key.second != "$") resolved.insert(key); + if (unresolved || !info.quarantined.empty() || !info.controlReadable) { + for (const auto& key : info.keyids) + if (key.second != "$" && !info.quarantined.count(key.first) && + (info.controlReadable || key.first != info.stop)) resolved.insert(key); cursors = &resolved; } - RedisConnection::ReadStatus status; - const bool accepted = !cursors->empty() && _redis.xreadMultiBlock(cursors->begin(), cursors->end(), - _options.cxn.timeout, inserter(out, out.end()), &status, _options.readerBatchCount); + RedisConnection::ReadStatus status = RedisConnection::ReadStatus::TimedOut; + bool socketTimedOut = false, blocking; + const auto interval = read_interval(info, blocking); + const bool accepted = cursors->empty() || _redis.xreadMultiBlock(cursors->begin(), cursors->end(), + interval, inserter(out, out.end()), &status, _options.readerBatchCount, &socketTimedOut, blocking); + reader_result(info, status, socketTimedOut); if (!accepted) { + if (status == RedisConnection::ReadStatus::Rejected) { + const auto previous = info.quarantined.size(); + check_readable(info, allKeys); + if (previous == info.quarantined.size() && !info.stop.empty() && info.controlReadable) { + const auto control = _redis.probeReadable({info.stop}); + if (control.status == RedisConnection::ReadStatus::Rejected) info.controlReadable = false; + } + } if (!reported) { syslog(LOG_WARNING, "stream reader paused after read rejection or transport failure"); reported = true; } if (info.run) this_thread::sleep_for(milliseconds(delay)); delay = min(1000, delay * 2); continue; } + if (cursors->empty()) this_thread::sleep_for(milliseconds(min(interval, 50))); + if (reported) syslog(LOG_INFO, "stream reader recovered"); reported = false; delay = 50; for (auto& item : out) { @@ -531,22 +813,48 @@ bool RedisAdapter::start_reader(uint32_t token) const auto base = parts.first.empty() ? item.first : parts.first; const auto sub = parts.first.empty() ? item.first : parts.second; auto data = make_shared(std::move(item.second)); + ++info.boundaries[item.first].readVersion; for (const auto& registration : subscriptions->second) { - _replier_pool.job(item.first, [registration, base, sub, data]() { + StreamBatchMetadata metadata; + { + lock_guard lock(registration->mutex); + 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, metadata]() { if (!registration->active.load()) return; size_t first = 0; { lock_guard cursorLock(registration->mutex); + 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; registration->cursor = data->back().first; } if (!registration->active.load()) return; - if (first == 0) registration->callback(base, sub, *data); - else { - const ItemStream fresh(data->begin() + first, data->end()); - registration->callback(base, sub, fresh); + { + lock_guard lock(registration->mutex); + 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->metadataCallback) registration->metadataCallback(base, sub, batch, metadata); + else registration->callback(base, sub, batch); + }; + try { + if (first == 0) invoke(*data); + else { const ItemStream fresh(data->begin() + first, data->end()); invoke(fresh); } + } catch (...) { + lock_guard lock(registration->mutex); + ++registration->status.callbackErrors; + throw; } }); } diff --git a/RedisAdapter.hpp b/RedisAdapter.hpp index abf945c..8f697a4 100644 --- a/RedisAdapter.hpp +++ b/RedisAdapter.hpp @@ -21,6 +21,9 @@ using RedisAdapter = MockRedisAdapter; #include #include #include +#include +#include +#include //^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ // define RA_VERSION @@ -72,6 +75,7 @@ struct RA_Options std::string dogname; uint16_t workers = 1; uint16_t readers = 1; + uint32_t readerProbeMs = 0; // optional continuity inspection; per-subscription override uint32_t readerBatchCount = 64; // per stream; zero is normalized to one }; @@ -100,6 +104,28 @@ class RedisAdapter : public RedisStreamData bool present() const { return connected && !fields.empty() && id != "0-0"; } }; + 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; + // Empty means an unresolved future-only cursor, never a comparable ID. + std::string cursor, observedCursor; + uint64_t epoch = 0, 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 fences callbacks that have not passed their active check. A // callback past that check may still enter; consumers must fence mutations. class ReaderHandle { @@ -112,6 +138,7 @@ class RedisAdapter : public RedisStreamData ReaderHandle& operator=(const ReaderHandle&) = delete; void reset() noexcept; explicit operator bool() const; + [[nodiscard]] ReaderStatus status() const; private: friend class RedisAdapter; ReaderHandle(std::weak_ptr owner, std::shared_ptr registration); @@ -121,16 +148,28 @@ class RedisAdapter : public RedisStreamData [[nodiscard]] StreamSnapshot getStreamSnapshot(const std::string& subKey, const std::string& baseKey = ""); [[nodiscard]] ReaderHandle subscribeStream(const std::string& subKey, StreamCallback callback, - const std::string& afterId = "$", const std::string& baseKey = ""); + const std::string& afterId = "$", const std::string& baseKey = "", + uint32_t probeMs = UINT32_MAX); struct SubscriptionOptions { std::string baseKey; std::string afterId = "$"; + std::optional probeMs; }; // Throws invalid_argument for empty callbacks/invalid IDs and runtime_error // during shutdown. Retain the returned handle for the registration lifetime. [[nodiscard]] ReaderHandle subscribeStream(const std::string& subKey, StreamCallback callback, const SubscriptionOptions& options) { - return subscribeStream(subKey, std::move(callback), options.afterId, options.baseKey); + return subscribeStream(subKey, std::move(callback), options.afterId, options.baseKey, options.probeMs.value_or(UINT32_MAX)); + } + [[nodiscard]] ReaderHandle subscribeStreamWithEpoch(const std::string& subKey, EpochStreamCallback callback, + const SubscriptionOptions& options); + [[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 @@ -449,7 +488,8 @@ class RedisAdapter : public RedisStreamData std::shared_ptr register_reader(const std::string& key, reader_sub_fn func, const std::string& afterId, - bool resolveTail = true); + bool resolveTail = true, uint32_t probeMs = UINT32_MAX, + MetadataStreamCallback metadataCallback = {}); void remove_registration(uint64_t id); template reader_sub_fn make_reader_callback(ReaderSubFn func) const; @@ -536,6 +576,11 @@ class RedisAdapter : public RedisStreamData std::mutex mutex; std::string cursor; reader_sub_fn callback; + MetadataStreamCallback metadataCallback; + uint32_t probeMs = 0; + bool everConnected = false, continuityCheck = true; + std::string observedCursor, gapSignature; + ReaderStatus status; }; std::shared_ptr _reader_owner = std::make_shared(); uint64_t _next_reader_id = 0; @@ -546,6 +591,22 @@ class RedisAdapter : public RedisStreamData std::unordered_map>> subs; std::unordered_map keyids; std::string stop; + struct Boundary { + std::chrono::steady_clock::time_point nextProbe{}; + uint32_t intervalMs = 0, backoff = 1; + uint64_t readVersion = 0, lastProbeVersion = 0; + bool denied = false; + }; + struct Quarantine { std::chrono::steady_clock::time_point nextCheck{}; bool wrongType = false; }; + struct ProbeTicket { std::chrono::steady_clock::time_point due; uint64_t order; std::string key; }; + struct ProbeLater { + bool operator()(const ProbeTicket& a, const ProbeTicket& b) const { return std::tie(a.due, a.order) > std::tie(b.due, b.order); } + }; + std::unordered_map boundaries; + std::unordered_map quarantined; + std::priority_queue, ProbeLater> probes; + uint64_t probeOrder = 0; + bool readConnected = false, controlReadable = true; std::atomic run = false; // used by start_reader() to confirm the reader thread has begun its read loop - @@ -557,6 +618,13 @@ class RedisAdapter : public RedisStreamData }; std::unordered_map _reader; + void prepare_probes(reader_info& info); + void inspect_readers(reader_info& info); + void reader_result(reader_info& info, RedisConnection::ReadStatus result, bool socketTimedOut = false); + void check_readable(reader_info& info, const std::vector& keys); + void reset_stream(reader_info& info, const std::string& key, StreamKind kind, bool allReaders); + void apply_bounds(reader_info& info, const std::string& key, const RedisConnection::StreamBounds& bounds); + uint32_t read_interval(const reader_info& info, bool& blocking) const; ThreadPool _replier_pool; }; diff --git a/RedisConnection.hpp b/RedisConnection.hpp index 178dd3b..c80d3ef 100644 --- a/RedisConnection.hpp +++ b/RedisConnection.hpp @@ -19,7 +19,19 @@ namespace chr = std::chrono; class RedisConnection { public: + enum class ReadStatus { Accepted, TimedOut, Rejected, Unavailable }; enum class CommandStatus { Accepted, 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", error; + }; + struct ReadProbe { + ReadStatus status = ReadStatus::Unavailable; + bool wrongType = false; + std::string error; + }; struct WriteResult { CommandStatus status = CommandStatus::Unavailable; std::string id; @@ -32,7 +44,6 @@ class RedisConnection std::string error; bool refreshConnection = false; }; - enum class ReadStatus { Accepted, TimedOut, Rejected, Unavailable }; //^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ // struct RedisConnection::Options // @@ -352,9 +363,11 @@ class RedisConnection // template bool xreadMultiBlock(Input fst, Input lst, uint32_t tmo, Output out, - ReadStatus* status = nullptr, uint32_t count = 64) + ReadStatus* status = nullptr, uint32_t count = 64, + bool* socketTimedOut = nullptr, bool blocking = true) { if (status) *status = ReadStatus::Unavailable; + if (socketTimedOut) *socketTimedOut = false; auto [cluster, singler] = snapshot(true); const auto block = chr::milliseconds(tmo == 0 ? 1000 : std::min(tmo, 1000)); using Fields = std::unordered_map; @@ -362,14 +375,23 @@ class RedisConnection std::unordered_map result; try { const auto destination = std::inserter(result, result.end()); - if (cluster) cluster->xread(fst, lst, block, std::max(1, count), destination); - else if (singler) singler->xread(fst, lst, block, std::max(1, count), destination); + if (cluster) { + if (blocking) cluster->xread(fst, lst, block, std::max(1, count), destination); + else cluster->xread(fst, lst, std::max(1, count), destination); + } else if (singler) { + if (blocking) singler->xread(fst, lst, block, std::max(1, count), destination); + else singler->xread(fst, lst, std::max(1, count), destination); + } else return false; if (status) *status = result.empty() ? ReadStatus::TimedOut : ReadStatus::Accepted; for (auto& item : result) *out++ = std::move(item); return true; - } catch (const swr::ReplyError&) { - if (status) *status = ReadStatus::Rejected; + } catch (const swr::ReplyError& error) { + const std::string message = error.what(); + if (status) *status = message.rfind("WRONGPASS", 0) == 0 || message.rfind("NOAUTH", 0) == 0 + ? ReadStatus::Unavailable : ReadStatus::Rejected; + } catch (const swr::TimeoutError&) { + if (socketTimedOut) *socketTimedOut = true; } catch (const swr::Error&) {} return false; } @@ -384,6 +406,60 @@ class RedisConnection // return : the id of the new element if successful // empty string if unsuccsessful or not connected // + // A nonblocking '$' XREAD returns no payload and tests exactly the permission + // and type contract that a shared blocking read needs. + ReadProbe probeReadable(const std::vector& keys) { + if (keys.empty()) return {ReadStatus::Accepted}; + auto [cluster, singler] = snapshot(); + std::vector> cursors; + cursors.reserve(keys.size()); + for (const auto& key : keys) cursors.emplace_back(key, "$" ); + using Fields = std::unordered_map; + using Entries = std::vector>; + std::unordered_map ignored; + try { + auto out = std::inserter(ignored, ignored.end()); + if (cluster) cluster->xread(cursors.begin(), cursors.end(), 1, out); + else if (singler) singler->xread(cursors.begin(), cursors.end(), 1, out); + else return {}; + return {ReadStatus::Accepted}; + } catch (const swr::ReplyError& error) { + const std::string message = error.what(); + const bool auth = message.rfind("WRONGPASS", 0) == 0 || message.rfind("NOAUTH", 0) == 0; + return {auth ? ReadStatus::Unavailable : ReadStatus::Rejected, + message.rfind("WRONGTYPE", 0) == 0, message}; + } catch (const swr::Error& error) { return {ReadStatus::Unavailable, false, error.what()}; } + } + + // Bounded pipeline depth is chosen by the scheduler. FULL COUNT 1 transfers + // one retained payload rather than both first/last payloads. Cluster metadata + // inspection is unsupported, independently of its working XREAD path. + std::vector streamBoundsBatch(const std::vector& keys) { + std::vector result(keys.size()); + auto [cluster, singler] = snapshot(); + if (cluster) { + for (auto& item : result) item.status = CommandStatus::Rejected; + return result; + } + if (!singler) return result; + try { + auto pipeline = singler->pipeline(false); + for (const auto& key : keys) pipeline.command("XINFO", "STREAM", key, "FULL", "COUNT", 1); + auto replies = pipeline.exec(); + for (size_t index = 0; index < keys.size(); ++index) { + try { result[index] = parseBounds(replies.get(index)); } + catch (const swr::ReplyError& error) { result[index] = rejectedBounds(error.what()); } + } + } catch (const swr::Error& error) { + for (auto& item : result) { item.status = CommandStatus::Unavailable; item.error = error.what(); } + } + return result; + } + + StreamBounds streamBounds(const std::string& key) { + return streamBoundsBatch({key}).front(); + } + template std::string xadd(const std::string& key, const std::string& id, Input fst, Input lst) { return xaddResult(key, id, fst, lst).id; @@ -721,6 +797,31 @@ class RedisConnection } private: + static StreamBounds rejectedBounds(const std::string& message) { + if (message.rfind("ERR no such key", 0) == 0) return {CommandStatus::Accepted, StreamKind::Missing}; + if (message.rfind("WRONGTYPE", 0) == 0) return {CommandStatus::Accepted, StreamKind::Invalid}; + return {CommandStatus::Rejected, StreamKind::Unknown, "0-0", "0-0", message}; + } + static StreamBounds parseBounds(const redisReply& reply) { + StreamBounds result; + result.status = CommandStatus::Rejected; + if (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 hasLast = false; + for (size_t index = 0; index < reply.elements; index += 2) { + const auto field = text(reply.element[index]); + const auto* value = reply.element[index + 1]; + if (field == "last-generated-id") { result.lastGeneratedId = text(value); hasLast = !result.lastGeneratedId.empty(); } + else if (field == "entries" && value && value->type == REDIS_REPLY_ARRAY && value->elements) { + const auto* first = value->element[0]; + if (first && first->type == REDIS_REPLY_ARRAY && first->elements) result.firstId = text(first->element[0]); + } + } + if (hasLast) { result.status = CommandStatus::Accepted; result.kind = StreamKind::Stream; } + return result; + } static WriteResult replyFailure(const swr::ReplyError& error) { const std::string message = error.what(); const bool refresh = message == "READONLY" || message.rfind("READONLY ", 0) == 0; diff --git a/docs/api.md b/docs/api.md index ff0af4b..9368655 100644 --- a/docs/api.md +++ b/docs/api.md @@ -33,6 +33,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 and dedicated blocking-read pool capacity. | +| `readerProbeMs` | `uint32_t` | `0` | Optional continuity inspection; per-subscription override available. | | `readerBatchCount` | `uint32_t` | `64` | XREAD entries per stream; zero is normalized to one. | Blocking reads have their own pool, sized to `readers`, so they do not consume diff --git a/docs/building.md b/docs/building.md index 82f465a..ddfc6be 100644 --- a/docs/building.md +++ b/docs/building.md @@ -147,3 +147,12 @@ Python 3 enables the write proxy/fault tests. Native GoogleTest cases remain available when Python is absent. Use `scripts/run-test-redis.py -- ctest --test-dir build --output-on-failure` for a private native Redis fixture, or the pinned Redis 7.4 CI fixture for field-TTL and XREAD-tail coverage. + +Recovery tests require an isolated Redis fixture with ACL administration and +CLIENT KILL privileges. They create unique users/keys, terminate only the test +user's clients, and clean up their namespace. Run through the private test helper +or the dedicated CI fixture; do not point fault tests at a shared Redis service. + +On Linux with Docker, run the private three-master fixture with +`scripts/run-test-redis-cluster.py -- ctest --test-dir build -R ClusterRecovery +--output-on-failure`. It binds only loopback ports and removes its own containers. diff --git a/docs/stream-subscriptions.md b/docs/stream-subscriptions.md index 39a312b..8456887 100644 --- a/docs/stream-subscriptions.md +++ b/docs/stream-subscriptions.md @@ -70,3 +70,57 @@ Legacy single-item typed getters distinguish malformed data with `RA_INVALID_PAYLOAD` (`err() == 3`), preserve the destination, and reserve zero for an accepted empty result. Range and callback wrappers skip malformed entries; they do not fabricate zero values or call back with a fully rejected empty batch. + +## Read recovery and inspection + +A rejected grouped XREAD is checked with nonblocking `XREAD COUNT 1 STREAMS key +$`. Only failing keys are quarantined; valid neighbours and legacy registrations +continue. Excluded keys are rechecked without requiring XINFO. A repaired +wrong-type key rewinds to the start of its replacement stream. Read transport +failures preserve cursors and are retried with bounded backoff. + +Continuity inspection is opt-in. `RA_Options::readerProbeMs` defaults to zero; +`SubscriptionOptions::probeMs` overrides it for one owned registration. If any +registration enables inspection of a shared key, that key is inspected. Periodic +metadata lookup detects deletion, lower-ID recreation and possible retention +loss. It is diagnostic evidence, not an exact count of missed entries or a +promise to observe every delete/recreate between checks. + +The scheduler batches at most 16 due XINFO requests through an existing pooled +connection. `XINFO STREAM FULL COUNT 1` transfers one retained payload, which can +still be large. Active keys skip unnecessary probes. Idle intervals back off to +at most eight times the configured interval. Denied or unsupported inspection +backs off for 60 seconds, or is reconsidered after a subscription/reconnect +restart. Deadlines and ordering survive reader restarts. A backlogged scheduler +uses nonblocking read cycles; normal cycles shorten for upcoming checks. Probe +intervals are minimum intervals, not hard completion deadlines. + +Metadata inspection is standalone-only. On Redis Cluster, inspection is reported +as unavailable evidence while normal owned and legacy XREAD delivery continues. +Cluster routing/retry behavior remains upstream redis-plus-plus behavior. + +`ReaderHandle::status()` separates observed and delivered cursors. An empty +cursor means unresolved future-only registration; compare IDs only when nonempty. +`epoch` is a fencing token: zero means no registration, and actual registrations +start at one. Resetting a handle or destroying its adapter clears its active, +connected and inspected flags. Transport failures are bucket observations; key +rejection counters identify the quarantined registration. Idle NIL replies keep +connection state healthy and do not count as socket timeouts. Inspection failures +and rejections are separate from read failures. + +`subscribeStreamWithEpoch()` supplies the batch's captured epoch as its fourth +callback argument. It does not infer that token by polling mutable status. +Queued older-epoch batches are fenced; an already executing callback may complete. +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 new file mode 100644 index 0000000..8239a65 --- /dev/null +++ b/recovery_test.cpp @@ -0,0 +1,376 @@ +#include "RedisAdapter.hpp" +#include +#include +#include +#include +#include +#include +#include +#include + +using namespace std::chrono_literals; +using RA = RedisAdapter; + +static bool waitUntil(const std::function& predicate) { + const auto end = std::chrono::steady_clock::now() + 4s; + do { if (predicate()) return true; std::this_thread::sleep_for(5ms); } while (std::chrono::steady_clock::now() < end); + return predicate(); +} +#define EVENTUALLY(...) ASSERT_TRUE(waitUntil([&] { return (__VA_ARGS__); })) << #__VA_ARGS__ + +struct Seen { + std::mutex mutex; + std::condition_variable changed; + std::vector> entries; + void append(const RA::StreamBatch& batch, uint64_t epoch = 0) { + std::lock_guard lock(mutex); + for (const auto& entry : batch) entries.emplace_back(entry.first, epoch); + changed.notify_all(); + } + bool wait(const std::string& id) { + std::unique_lock lock(mutex); + return changed.wait_for(lock, 4s, [&] { for (const auto& entry : entries) if (entry.first == id) return true; return false; }); + } + size_t count(const std::string& id) { std::lock_guard lock(mutex); size_t result = 0; for (const auto& entry : entries) result += entry.first == id; return result; } + size_t size() { std::lock_guard lock(mutex); return entries.size(); } + uint64_t epoch(const std::string& id) { std::lock_guard lock(mutex); for (const auto& entry : entries) if (entry.first == id) return entry.second; return 0; } +}; + +class Recovery : public testing::TestWithParam { +protected: + RA_Options options; + std::string base, user; + std::unique_ptr control; + RA::Attrs fields{{"_", "value"}}; + void SetUp() override { + if (!std::getenv("REDIS_ADAPTER_ISOLATED_TEST")) GTEST_SKIP() << "Use the private Redis fixture"; + options.cxn.port = std::stoi(std::getenv("REDIS_ADAPTER_TEST_PORT")); + options.cxn.timeout = 100; + options.workers = GetParam(); + base = "recovery-" + std::to_string(getpid()) + "-" + std::to_string(std::chrono::steady_clock::now().time_since_epoch().count()); + user = "WRONGTYPE-reader-" + base; + sw::redis::ConnectionOptions co; co.host = "127.0.0.1"; co.port = options.cxn.port; + control = std::make_unique(co); + control->command("ACL", "SETUSER", user, "reset", "on", "nopass", "~{" + base + "}:*", "+@all"); + options.cxn.user = user; + } + void TearDown() override { + if (!control) return; + try { control->command("ACL", "DELUSER", user); } catch (...) {} + try { std::vector keys; control->keys("*" + base + "*", std::back_inserter(keys)); if (!keys.empty()) control->del(keys.begin(), keys.end()); } catch (...) {} + } + std::string key(const std::string& sub) const { return "{" + base + "}:" + sub; } + void add(const std::string& sub, const std::string& id) { EXPECT_EQ(control->xadd(key(sub), id, fields.begin(), fields.end()), id); } + void recreate(const std::string& sub, const std::string& id) { + auto transaction = control->transaction(); + transaction.del(key(sub)).xadd(key(sub), id, fields.begin(), fields.end()).exec(); + } + RA::SubscriptionOptions selection(uint32_t probeMs = 20) { + RA::SubscriptionOptions result; result.afterId = "0-0"; result.probeMs = probeMs; return result; + } + RA::ReaderHandle subscribe(RA& adapter, const std::string& sub, Seen& seen, uint32_t probeMs = 20) { + return adapter.subscribeStreamWithEpoch(sub, [&](const auto&, const auto&, const auto& batch, uint64_t epoch) { seen.append(batch, epoch); }, selection(probeMs)); + } +}; + +TEST_P(Recovery, InspectionIsOptInAndCanBeSelectedPerSubscription) { + EXPECT_EQ(options.readerProbeMs, 0u); + RA adapter(base, options); + Seen scalar, imaging; + auto first = subscribe(adapter, "scalar", scalar, 20); + auto second = subscribe(adapter, "imaging", imaging, 0); + EVENTUALLY(first.status().inspected && second.status().connected); + EXPECT_FALSE(second.status().inspected); + EXPECT_EQ(first.status().epoch, 1u); + 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; + const auto sub = "value-" + std::to_string(timeout); + Seen seen; + RA adapter(base, options); + add(sub, "1000-0"); + auto handle = subscribe(adapter, sub, seen); + ASSERT_TRUE(seen.wait("1000-0")); + recreate(sub, "1-0"); + ASSERT_TRUE(seen.wait("1-0")); + EXPECT_EQ(handle.status().streamResets, 1u); + EXPECT_EQ(handle.status().epoch, 2u); + EXPECT_EQ(seen.epoch("1-0"), 2u); + } +} + +TEST_P(Recovery, DeletionIsObservedAndLaterDataResumes) { + Seen seen; + RA adapter(base, options); + add("value", "1000-0"); + auto handle = subscribe(adapter, "value", seen); + ASSERT_TRUE(seen.wait("1000-0")); + control->del(key("value")); + EVENTUALLY(handle.status().streamKind == RA::StreamKind::Missing); + EXPECT_EQ(handle.status().disappearances, 1u); + EXPECT_FALSE(handle.status().hasData); + add("value", "1-0"); + ASSERT_TRUE(seen.wait("1-0")); +} + +TEST_P(Recovery, RepairAfterXinfoRevocationStillDelivers) { + Seen seen; + RA adapter(base, options); + control->set(key("value"), "wrong type"); + auto handle = subscribe(adapter, "value", seen); + EVENTUALLY(handle.status().streamKind == RA::StreamKind::Invalid); + control->command("ACL", "SETUSER", user, "-xinfo"); + recreate("value", "1-0"); + ASSERT_TRUE(seen.wait("1-0")); + EXPECT_EQ(seen.count("1-0"), 1u); +} + +TEST_P(Recovery, StillInvalidAfterXinfoRevocationDoesNotBlockItsNeighbour) { + Seen bad, healthy; + RA adapter(base, options); + control->set(key("bad"), "wrong type"); + auto a = subscribe(adapter, "bad", bad); + auto b = subscribe(adapter, "healthy", healthy); + EVENTUALLY(a.status().streamKind == RA::StreamKind::Invalid); + control->command("ACL", "SETUSER", user, "-xinfo"); + for (unsigned i = 1; i <= 6; ++i) add("healthy", std::to_string(i) + "-0"); + ASSERT_TRUE(healthy.wait("6-0")); + EXPECT_EQ(healthy.size(), 6u); + EXPECT_EQ(b.status().readRejections, 0u); +} + +TEST_P(Recovery, LegacyRegistrationSurvivesRemovalOfItsOwnedInspectingPeer) { + Seen legacy; + RA adapter(base, options); + control->set(key("value"), "wrong type"); + auto owned = adapter.subscribeStream("value", [](const auto&, const auto&, const auto&) {}, selection()); + ASSERT_TRUE(adapter.addValuesReader("value", [&](const auto&, const auto&, const auto& batch) { + RA::StreamBatch raw; for (const auto& value : batch) raw.emplace_back(value.first.id(), value.second); legacy.append(raw); + })); + EVENTUALLY(owned.status().streamKind == RA::StreamKind::Invalid); + owned.reset(); + control->del(key("value")); + // These IDs are valid legacy nanosecond encodings, unlike arbitrary sequence IDs. + for (unsigned i = 1; i <= 20; ++i) add("value", std::to_string(i) + "-0"); + ASSERT_TRUE(legacy.wait("20-0")); + EXPECT_EQ(legacy.size(), 20u); +} + +TEST_P(Recovery, LegacyWrongTypeIsIsolatedWithoutAnyInspectionPermission) { + Seen healthy; + control->command("ACL", "SETUSER", user, "-xinfo"); + RA adapter(base, options); + control->set(key("legacy"), "wrong type"); + ASSERT_TRUE(adapter.addValuesReader("legacy", [](const auto&, const auto&, const auto&) {})); + auto owned = subscribe(adapter, "healthy", healthy, 0); + for (unsigned i = 1; i <= 6; ++i) add("healthy", std::to_string(i) + "-0"); + ASSERT_TRUE(healthy.wait("6-0")); + EXPECT_EQ(healthy.size(), 6u); + EXPECT_EQ(owned.status().readRejections, 0u); +} + +TEST_P(Recovery, DeniedMetadataAndIdleReadsKeepAccurateConnectionState) { + Seen seen; + control->command("ACL", "SETUSER", user, "-xinfo"); + RA adapter(base, options); + auto handle = subscribe(adapter, "value", seen); + EVENTUALLY(handle.status().inspectionRejections > 0 && handle.status().connected); + std::this_thread::sleep_for(300ms); + const auto status = handle.status(); + EXPECT_FALSE(status.inspected); + EXPECT_NE(status.streamKind, RA::StreamKind::Invalid); // username contains WRONGTYPE + EXPECT_EQ(status.socketTimeouts, 0u); + EXPECT_EQ(status.inspectionRejections, 1u); // denied probes back off + add("value", "1-0"); + ASSERT_TRUE(seen.wait("1-0")); +} + +TEST_P(Recovery, ScopedConnectionKillCountsFailureForEveryRegistrationAndPreservesCursors) { + Seen first, second; + RA adapter(base, options); + auto a = subscribe(adapter, "first", first); + auto b = subscribe(adapter, "second", second); + add("first", "10-0"); add("second", "10-0"); + ASSERT_TRUE(first.wait("10-0")); ASSERT_TRUE(second.wait("10-0")); + const auto failuresA = a.status().readFailures + a.status().inspectionFailures; + const auto failuresB = b.status().readFailures + b.status().inspectionFailures; + EXPECT_GT(control->command("CLIENT", "KILL", "USER", user), 0); + EXPECT_TRUE(control->ping() == "PONG"); // bystander admin connection survives + EVENTUALLY(a.status().readFailures + a.status().inspectionFailures > failuresA); + EVENTUALLY(b.status().readFailures + b.status().inspectionFailures > failuresB); + EVENTUALLY(a.status().connected && b.status().connected); + add("first", "11-0"); add("second", "11-0"); + ASSERT_TRUE(first.wait("11-0")); ASSERT_TRUE(second.wait("11-0")); + EXPECT_EQ(first.count("10-0"), 1u); EXPECT_EQ(second.count("10-0"), 1u); +} + +TEST_P(Recovery, EmptyRetainedHistoryCountsACertainGapOnce) { + Seen seen; + RA adapter(base, options); + add("value", "10-0"); + auto handle = subscribe(adapter, "value", seen); + ASSERT_TRUE(seen.wait("10-0")); + adapter.setDeferReaders(true); + add("value", "20-0"); control->xtrim(key("value"), 0, false); + adapter.setDeferReaders(false); + EVENTUALLY(handle.status().retentionGaps == 1); + EXPECT_FALSE(handle.status().hasData); + EXPECT_EQ(handle.status().cursor, "10-0"); + std::this_thread::sleep_for(200ms); + EXPECT_EQ(handle.status().retentionGaps, 1u); +} + +TEST_P(Recovery, QueuedOldEpochIsFencedAndCallbacksReceiveTheirCapturedEpoch) { + Seen seen; + RA adapter(base, options); + std::promise entered, release; + auto started = entered.get_future(); auto gate = release.get_future().share(); + struct Release { std::promise& promise; ~Release() { try { promise.set_value(); } catch (...) {} } } finally{release}; + const auto blocked = std::string("blocker"); + std::string target; + for (unsigned i = 0;; ++i) { + target = "fenced-" + std::to_string(i); + if (std::hash{}(key(target)) % options.workers == std::hash{}(key(blocked)) % options.workers) break; + } + auto blocker = adapter.subscribeStream(blocked, [&](const auto&, const auto&, const auto&) { entered.set_value(); gate.wait(); }, "0-0"); + add(blocked, "1-0"); ASSERT_EQ(started.wait_for(4s), std::future_status::ready); + add(target, "1000-0"); + auto handle = subscribe(adapter, target, seen); + EVENTUALLY(handle.status().observedCursor == "1000-0"); + EXPECT_EQ(handle.status().entries, 0u); + recreate(target, "1-0"); + EVENTUALLY(handle.status().epoch == 2 && handle.status().observedCursor == "1-0"); + release.set_value(); + ASSERT_TRUE(seen.wait("1-0")); + EXPECT_EQ(seen.count("1000-0"), 0u); + EXPECT_EQ(seen.epoch("1-0"), 2u); + EXPECT_EQ(handle.status().entries, 1u); +} + +TEST_P(Recovery, ProbeSchedulingRemainsFairAcrossMoreThan64KeysAndChurn) { + RA adapter(base, options); + std::vector handles; + adapter.setDeferReaders(true); + for (unsigned i = 0; i < 80; ++i) handles.push_back(adapter.subscribeStream("key-" + std::to_string(i), [](const auto&, const auto&, const auto&) {}, selection())); + adapter.setDeferReaders(false); + for (unsigned i = 0; i < 6; ++i) { + auto churn = adapter.subscribeStream("churn", [](const auto&, const auto&, const auto&) {}, selection()); + std::this_thread::sleep_for(20ms); churn.reset(); + } + EVENTUALLY(std::all_of(handles.begin(), handles.end(), [](const auto& handle) { return handle.status().inspected; })); +} + +TEST_P(Recovery, InactiveHandlesClearHealthAndUnresolvedCursorsAreExplicit) { + RA::ReaderHandle orphan; + { RA adapter(base, options); orphan = adapter.subscribeStream("key", [](const auto&, const auto&, const auto&) {}, selection()); EVENTUALLY(orphan.status().connected); } + EXPECT_FALSE(orphan.status().active); + EXPECT_FALSE(orphan.status().connected); + EXPECT_FALSE(orphan.status().inspected); + orphan.reset(); EXPECT_EQ(orphan.status().epoch, 0u); + auto unavailable = options; unavailable.cxn.port = 0; + RA adapter(base, unavailable); + auto pending = adapter.subscribeStream("pending", [](const auto&, const auto&, const auto&) {}); + EXPECT_TRUE(pending.status().cursor.empty()); + EXPECT_TRUE(pending.status().observedCursor.empty()); +} + +TEST_P(Recovery, CallbackErrorsAreCountedWithoutStoppingDelivery) { + Seen seen; + RA adapter(base, options); + auto bad = adapter.subscribeStream("key", [](const auto&, const auto&, const auto&) { throw std::runtime_error("intentional"); }, selection()); + auto good = subscribe(adapter, "key", seen); + add("key", "1-0"); ASSERT_TRUE(seen.wait("1-0")); + EVENTUALLY(bad.status().callbackErrors == 1); +} + +TEST_P(Recovery, RetainedHistoryAboveAnExplicitCursorReportsAPossibleGap) { + Seen seen; + RA adapter(base, options); + add("trimmed", "10-0"); add("trimmed", "20-0"); + control->xtrim(key("trimmed"), 1, false); + auto selected = selection(); selected.afterId = "10-0"; + auto handle = adapter.subscribeStreamWithEpoch("trimmed", [&](const auto&, const auto&, const auto& entries, uint64_t epoch) { seen.append(entries, epoch); }, selected); + ASSERT_TRUE(seen.wait("20-0")); + EXPECT_EQ(handle.status().retentionGaps, 1u); + EXPECT_EQ(handle.status().streamResets, 0u); +} + +TEST_P(Recovery, AuthenticationFailureClearsConnectedEvenWithoutInspection) { + Seen seen; + RA adapter(base, options); + auto handle = subscribe(adapter, "value", seen, 0); + EVENTUALLY(handle.status().connected); + control->command("ACL", "SETUSER", user, "resetpass", ">new-test-password"); + EXPECT_GT(control->command("CLIENT", "KILL", "USER", user), 0); + EVENTUALLY(handle.status().readFailures > 0 && !handle.status().connected); + EXPECT_FALSE(handle.status().inspected); +} + +INSTANTIATE_TEST_SUITE_P(Workers, Recovery, testing::Values(1u, 4u)); + +TEST(ClusterRecovery, UnsupportedInspectionKeepsOwnedAndLegacyReadsFlowing) { + const auto* port = std::getenv("REDIS_ADAPTER_CLUSTER_TEST_PORT"); + if (!port || !std::getenv("REDIS_ADAPTER_ISOLATED_TEST")) GTEST_SKIP() << "Use an isolated Redis Cluster fixture"; + RA_Options options; options.cxn.port = std::stoi(port); options.cxn.timeout = 100; options.readerProbeMs = 20; + const auto base = "cluster-review-" + std::to_string(getpid()) + "-" + std::to_string(std::chrono::steady_clock::now().time_since_epoch().count()); + Seen owned, legacy; + RA adapter(base, options), producer(base, options); + RA::SubscriptionOptions selection; selection.afterId = "0-0"; + auto handle = adapter.subscribeStream("owned", [&](const auto&, const auto&, const auto& entries) { owned.append(entries); }, selection); + ASSERT_TRUE(adapter.addValuesReader("legacy", [&](const auto&, const auto&, const auto& entries) { + RA::StreamBatch batch; for (const auto& entry : entries) batch.emplace_back(entry.first.id(), entry.second); legacy.append(batch); + })); + EVENTUALLY(handle.status().inspectionRejections > 0 && handle.status().connected); + for (unsigned i = 1; i <= 6; ++i) { + RA_ArgsAdd args; args.time = RA_Time(int64_t(i) * 1000000); args.trim = 0; + ASSERT_TRUE(producer.addSingleDouble("owned", 1., args).ok()); + ASSERT_TRUE(producer.addSingleDouble("legacy", 1., args).ok()); + } + ASSERT_TRUE(owned.wait("6-0")); ASSERT_TRUE(legacy.wait("6-0")); + EXPECT_EQ(owned.size(), 6u); EXPECT_EQ(legacy.size(), 6u); + EXPECT_FALSE(handle.status().inspected); + handle.reset(); adapter.removeReader("legacy"); + producer.del("owned"); producer.del("legacy"); +} diff --git a/scripts/run-test-redis-cluster.py b/scripts/run-test-redis-cluster.py new file mode 100644 index 0000000..0bd8b2a --- /dev/null +++ b/scripts/run-test-redis-cluster.py @@ -0,0 +1,92 @@ +#!/usr/bin/env python3 +"""Run a command against three private loopback Redis Cluster nodes on Linux.""" +import argparse +import os +from pathlib import Path +import random +import socket +import subprocess +import tempfile +import time +import uuid + + +def run(*args, **kwargs): + return subprocess.run(args, check=True, text=True, **kwargs) + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument('--image', default='redis@sha256:02419de7eddf55aa5bcf49efb74e88fa8d931b4d77c07eff8a6b2144472b6952') + parser.add_argument('command', nargs=argparse.REMAINDER) + args = parser.parse_args() + command = args.command[1:] if args.command[:1] == ['--'] else args.command + if not command: + parser.error('a command is required') + ports = [] + for _ in range(3): + for _attempt in range(100): + candidate = random.randint(20000, 30000) + if candidate in ports: + continue + try: + with socket.socket() as client, socket.socket() as bus: + client.bind(('127.0.0.1', candidate)) + bus.bind(('127.0.0.1', candidate + 10000)) + ports.append(candidate) + break + except OSError: + continue + else: + raise RuntimeError('no free loopback Cluster ports') + with tempfile.TemporaryDirectory(prefix='redis-adapter-cluster-') as temporary: + containers = [] + identity = uuid.uuid4().hex[:12] + try: + for index, port in enumerate(ports): + name = f'redis-adapter-cluster-{identity}-{index}' + folder = Path(temporary) / str(index) + folder.mkdir(mode=0o700) + run('docker', 'run', '--detach', '--rm', '--name', name, + '--user', f'{os.getuid()}:{os.getgid()}', '--network', 'host', + '--mount', f'type=bind,source={folder},target=/data', args.image, + 'redis-server', '--bind', '127.0.0.1', '--port', str(port), + '--save', '', '--appendonly', 'no', '--cluster-enabled', 'yes', + '--cluster-config-file', '/data/nodes.conf', '--cluster-node-timeout', '2000', + '--cluster-announce-ip', '127.0.0.1', '--cluster-announce-port', str(port), + '--cluster-announce-bus-port', str(port + 10000), stdout=subprocess.DEVNULL) + containers.append(name) + deadline = time.monotonic() + 5 + while True: + try: + with socket.create_connection(('127.0.0.1', port), timeout=0.2) as connection: + connection.sendall(b'*1\r\n$4\r\nPING\r\n') + if connection.recv(64) == b'+PONG\r\n': + break + except OSError: + pass + if time.monotonic() > deadline: + raise RuntimeError('private Cluster node did not start') + time.sleep(0.05) + admin = ['docker', 'exec', containers[0], 'redis-cli'] + run(*admin, '--cluster', 'create', *(f'127.0.0.1:{port}' for port in ports), + '--cluster-replicas', '0', '--cluster-yes') + deadline = time.monotonic() + 10 + while True: + info = run(*admin, '-p', str(ports[0]), 'cluster', 'info', capture_output=True).stdout + if 'cluster_state:ok' in info and 'cluster_slots_assigned:16384' in info: + break + if time.monotonic() > deadline: + raise RuntimeError('private Cluster did not become ready: ' + info) + time.sleep(0.1) + environment = dict(os.environ, REDIS_ADAPTER_ISOLATED_TEST='1', + REDIS_ADAPTER_CLUSTER_TEST_PORT=str(ports[0])) + return subprocess.run(command, env=environment).returncode + finally: + for name in containers: + subprocess.run(['docker', 'rm', '--force', name], + stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) + + +if __name__ == '__main__': + raise SystemExit(main())