Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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.
Expand Down
4 changes: 4 additions & 0 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
336 changes: 322 additions & 14 deletions RedisAdapter.cpp

Large diffs are not rendered by default.

74 changes: 71 additions & 3 deletions RedisAdapter.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@ using RedisAdapter = MockRedisAdapter;
#include <limits>
#include <memory>
#include <type_traits>
#include <queue>
#include <optional>
#include <tuple>

//^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
// define RA_VERSION
Expand Down Expand Up @@ -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
};

Expand Down Expand Up @@ -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<void(const std::string&, const std::string&, const StreamBatch&, uint64_t)>;
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<void(const std::string&, const std::string&,
const StreamBatch&, const StreamBatchMetadata&)>;
struct ReaderStatus {
Comment thread
derekste marked this conversation as resolved.
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 {
Expand All @@ -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<ReaderOwner> owner, std::shared_ptr<ReaderRegistration> registration);
Expand All @@ -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<uint32_t> 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
Expand Down Expand Up @@ -449,7 +488,8 @@ class RedisAdapter : public RedisStreamData
std::shared_ptr<ReaderRegistration> 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<typename T> reader_sub_fn make_reader_callback(ReaderSubFn<T> func) const;
Expand Down Expand Up @@ -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<ReaderOwner> _reader_owner = std::make_shared<ReaderOwner>();
uint64_t _next_reader_id = 0;
Expand All @@ -546,6 +591,22 @@ class RedisAdapter : public RedisStreamData
std::unordered_map<std::string, std::vector<std::shared_ptr<ReaderRegistration>>> subs;
std::unordered_map<std::string, std::string> 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<std::string, Boundary> boundaries;
std::unordered_map<std::string, Quarantine> quarantined;
std::priority_queue<ProbeTicket, std::vector<ProbeTicket>, ProbeLater> probes;
uint64_t probeOrder = 0;
bool readConnected = false, controlReadable = true;
std::atomic<bool> run = false;

// used by start_reader() to confirm the reader thread has begun its read loop -
Expand All @@ -557,6 +618,13 @@ class RedisAdapter : public RedisStreamData
};
std::unordered_map<uint32_t, reader_info> _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<std::string>& 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;
};

Expand Down
113 changes: 107 additions & 6 deletions RedisConnection.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -32,7 +44,6 @@ class RedisConnection
std::string error;
bool refreshConnection = false;
};
enum class ReadStatus { Accepted, TimedOut, Rejected, Unavailable };
//^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
// struct RedisConnection::Options
//
Expand Down Expand Up @@ -352,24 +363,35 @@ class RedisConnection
//
template<typename Input, typename Output>
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<uint32_t>(tmo, 1000));
using Fields = std::unordered_map<std::string, std::string>;
using Entries = std::vector<std::pair<std::string, Fields>>;
std::unordered_map<std::string, Entries> result;
try {
const auto destination = std::inserter(result, result.end());
if (cluster) cluster->xread(fst, lst, block, std::max<uint32_t>(1, count), destination);
else if (singler) singler->xread(fst, lst, block, std::max<uint32_t>(1, count), destination);
if (cluster) {
if (blocking) cluster->xread(fst, lst, block, std::max<uint32_t>(1, count), destination);
else cluster->xread(fst, lst, std::max<uint32_t>(1, count), destination);
} else if (singler) {
if (blocking) singler->xread(fst, lst, block, std::max<uint32_t>(1, count), destination);
else singler->xread(fst, lst, std::max<uint32_t>(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;
}
Expand All @@ -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<std::string>& keys) {
if (keys.empty()) return {ReadStatus::Accepted};
auto [cluster, singler] = snapshot();
std::vector<std::pair<std::string, std::string>> cursors;
cursors.reserve(keys.size());
for (const auto& key : keys) cursors.emplace_back(key, "$" );
using Fields = std::unordered_map<std::string, std::string>;
using Entries = std::vector<std::pair<std::string, Fields>>;
std::unordered_map<std::string, Entries> 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<StreamBounds> streamBoundsBatch(const std::vector<std::string>& keys) {
std::vector<StreamBounds> 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<typename Input>
std::string xadd(const std::string& key, const std::string& id, Input fst, Input lst) {
return xaddResult(key, id, fst, lst).id;
Expand Down Expand Up @@ -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;
Expand Down
1 change: 1 addition & 0 deletions docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 9 additions & 0 deletions docs/building.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Loading
Loading