Skip to content

Add owned stream subscriptions and safe payload decoding - #108

Merged
derekste merged 7 commits into
mainfrom
dev/subscription-lifecycle
Oct 7, 2026
Merged

derekste merged 7 commits into
mainfrom
dev/subscription-lifecycle

Conversation

@derekste

@derekste derekste commented Sep 7, 2026 •

Copy link
Copy Markdown
Member

Add owned stream subscriptions so a consumer can remove one registration, retain exact Redis cursors across snapshot and registration, and safely validate raw payloads. This PR depends on the separate connection-policy foundation in #127; #109 and #111 remain stacked above it.

Handle reset fences jobs that have not passed their active check. A callback past that check may still enter, so consumers must retain captured state and fence their own generation mutations. Callback captures and canceled queues are released outside worker locks, and legacy removal releases callbacks outside the reader lock. Shutdown prevents reader restarts.

Snapshots distinguish server rejection from transport failure and keep a future-only cursor on failure. Pending tails retry without repeatedly reading $; Redis 7.4 consumers with XREAD-only access can resolve their tail through XREAD +. XREAD uses a configurable entry count, and complete fresh batches share storage. Named subscription options, input exceptions, strict decoding and RA_INVALID_PAYLOAD are documented and covered by real and mock builds.

Validation: 70/70 Linux Release CTest cases passed against pinned Redis 7.4.2 with C++20. macOS ASan+UBSan passed 67 cases; the two Redis 7.4-only ACL cases passed on Linux. Tests cover normal and throwing capture cleanup at one/four workers, last-owner destruction, legacy removal owning sibling handles, queued cancellation, shared rewind filtering, deferral, initially unavailable registration, wrong-type repair, opaque IDs, malformed payloads, batch bounds and zero-timeout cancellation.

@derekste
derekste requested a review from a team as a code owner September 7, 2026 22:02
Copilot AI lite review requested due to automatic review settings September 7, 2026 22:02

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

The per-bucket stop-stream key generation can land in the wrong Redis Cluster hash slot for already-hash-tagged keys (e.g. baseKey override / fully-qualified keys), which can break reader blocking/unblocking and multi-key reads.

Once you've addressed the issues Copilot identified, you can request another Copilot review.

Pull request overview

This PR adds an owned stream-subscription API to RedisAdapter (movable ReaderHandle), introduces safe scalar/array payload decoding helpers, and hardens reader/worker execution (retry behavior, cancellation semantics, and exception containment) while preserving existing reader APIs and deferred-start behavior.

Changes:

  • Added ReaderHandle + subscribeStream() / getStreamSnapshot() / compareStreamIds() for per-registration stream subscriptions with exact IDs.
  • Reworked reader registration/storage to support independent cursors per registration and avoid callback batching issues.
  • Added safe decodeScalar() / decodeArray() with byte-length checks and payload budgeting; added lifecycle test + documentation.
File summaries
File Description
ThreadPool.hpp Makes shutdown flag updates thread-safe and catches/logs callback exceptions in workers.
RedisConnection.hpp Applies configured timeout to connect and pool wait operations.
RedisAdapterTempl.hpp Switches legacy field decoding to safe decode helpers and avoids reference captures in callbacks.
RedisAdapter.hpp Adds owned stream subscription types/API plus safe decoding utilities and new reader-registration plumbing.
RedisAdapter.cpp Implements owned subscription lifecycle, exact stream-ID parsing/ordering, and updated reader loop semantics.
lifecycle_test.cpp Adds a standalone lifecycle/decoding/callback-survival executable test.
docs/stream-subscriptions.md Documents owned subscriptions, exact cursors, and safe decoding behavior.
CMakeLists.txt Builds and registers the new lifecycle test under REDIS_ADAPTER_TEST.
Review details

Suppressed comments (1)

RedisAdapter.cpp:567

  • In reconnect(), NO_TOKEN reader buckets are re-homed and a new info.stop is generated using split_key(key). For keys that don't contain this adapter's _base_key (e.g. subscriptions created with a different baseKey, or any fully-qualified key), split_key() returns empty and the stop key is built in the wrong cluster slot, so the restarted reader can fail to block/unblock correctly.
              auto prior = info.keyids.find(key);
              const auto cursor = cursors.at(key);
              if (prior == info.keyids.end() || prior->second == "$" ||
                  (cursor != "$" && compareStreamIds(cursor, prior->second) < 0))
                info.keyids[key] = cursor;
  • Files reviewed: 8/8 changed files
  • Comments generated: 1
  • Review effort level: Lite

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread RedisAdapter.cpp Outdated

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Needs a closer look

It introduces substantial concurrency and lifecycle changes in the reader/worker pipeline, and one confirmed correctness issue remains in boolean array decoding.

Review details

Suppressed comments (1)

Previously missed (1) — in code that hasn't changed since the last review.

RedisAdapter.hpp:191

  • decodeArray<bool> indexes bytes[index] even though the payload length is validated using sizeof(bool). On platforms where sizeof(bool) > 1, this will read the wrong byte for all but the first element and can mis-decode (or incorrectly accept/reject) boolean arrays. Use index * sizeof(bool) when reading each element’s first byte (mirroring the scalar bool check).
  • Files reviewed: 8/8 changed files
  • Comments generated: 0 new
  • Review effort level: Lite

Copilot AI review requested due to automatic review settings September 8, 2026 01:35
@derekste

derekste commented Sep 8, 2026

Copy link
Copy Markdown
Member Author

Addressed the boolean-decoding finding in the latest Copilot review in 0bb9c35. Array elements now advance by sizeof(bool), and both scalar and array decoding validate the complete native true/false representation before assigning a boolean. Invalid bytes are never copied into a bool object. Regression coverage checks mixed true/false arrays, invalid elements without modifying the prior output, scalar true/false and invalid values, empty arrays, and byte limits.

Validation: native macOS lifecycle test passed; the full isolated Linux adapter suite passed 15/15.

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Needs a closer look

It introduces substantial concurrency, lifecycle, and public-API behavior changes where correctness depends on subtle multi-threaded interactions best validated by a final human review.

Review details
  • Files reviewed: 8/8 changed files
  • Comments generated: 0 new
  • Review effort level: Lite

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot encountered an error and was unable to review this pull request. You can try again by re-requesting a review.

@derekste

Copy link
Copy Markdown
Member Author

Fresh release-prerequisite validation on 2026-09-29: hosted build/test run https://github.com/fermi-ad/redis-adapter/actions/runs/36600172667 passed all 15 tests at 0bb9c35. The native macOS lifecycle executable also passed against a private loopback Redis.

The earlier boolean-array review finding is addressed at this head: decoding advances by sizeof(bool) and validates each complete native boolean representation; regression cases cover malformed scalar/array encodings and unchanged output on failure. The previously resolved control-stream slot finding remains covered.

This is the first dependency for the approved IOC 0.8.2 → 0.9.0 sequence. Independent Instrumentation code-owner approval remains required before merge; downstream IOC #100 will then pin the merged revision. No approval or merge bypass is being used.

@derekste derekste left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Follow-up review, 30 September 2026: hold this PR for the reproduced lifetime defect in #112.

The owned subscription, exact-cursor and safe-decoding changes address concrete reload and payload-safety needs. The callback-owned destruction path still has a blocking cleanup case, however. Existing lifecycle/rejection/transport tests pass, while the new pool-only and full-adapter reproductions hang when a completed callback releases its last owner.

A narrow temporary-header experiment destroying the completed job before reacquiring the worker mutex makes the pool-only reproducer finish. The PR source has not been changed. Add the capture-cleanup regression and validate the fix before merging or updating the downstream IOC pin. Independent Instrumentation approval remains required.

Comment thread ThreadPool.hpp

@bigsamich bigsamich left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed at 0bb9c35, and re-checked at the #109/#111 heads. All of this was reproduced against pristine builds in an isolated Redis 7.4.11 (gcc 13 / clang 18, plain, ASan+UBSan and TSan) with targeted reproducers; numbers are medians of repeated, interleaved runs. Reproducers and the validated patches are available if useful.

Summary. The owned-subscription design is sound, and it fixes real bugs on main (listed below). I agree with holding for #112, and the #112 fix needs to cover more than the issue describes (see my reply on ThreadPool.hpp:93). Beyond that, a handful of should-fix items should land before merge. Most are small and each has a validated fix.

What this gets right (verified)

  • Shared-batch bug fixed. On main, a second reader on a key got an empty batch (and the full key as both base and sub). Now every registration gets the intact batch with the correct base and sub.
  • Data race / heap-use-after-free fixed. On main, add_reader_helper mutated info.subs/keyids before stop_reader while the reader thread iterated them (ASan/TSan reports on main, clean here). Under churn (2,818 subscribe/reset cycles against 36,788 writes), 570,153 checked deliveries had 0 gaps, duplicates or reorders.
  • Add-to-write gap closed. main missed up to 151 of the first 200 writes issued right after addValuesReader; this PR misses 0.
  • Resilience. Readers survive outages; main's reader dies on the first failed XREAD and never resumes. Callback exceptions no longer terminate the process. The ThreadPool _go race and shutdown hang are gone (main hung 3/600 in ~ThreadPool). The info.started predicate closes main's late-start race.
  • Safer by construction. Safe decoding removes real UB: over-reads, misaligned casts, and arbitrary bytes copied into bool. ReaderOwner makes handle-outlives-adapter and reset-versus-destructor safe (TSan/ASan clean). connect_timeout bounds the constructor against a SYN-dropping host (main was still inside the constructor after 20 s).
  • Tests. 15/15 at this head. Lifecycle passed 100/100 (40 plain, 30 ASan+UBSan, 30 TSan). No new -Wall -Wextra -Wpedantic warnings with gcc or clang, and it builds as real C++20. The only TSan reports in the whole suite come from test.cpp's bool waiting flags; with those made atomic this PR is completely TSan-clean, which makes a sanitizer CI job possible (#125).

Blocking

  • #112, broadened. The same line also deadlocks with no adapter destruction at all (captured state owning a sibling ReaderHandle wedges worker and reader against each other), and on the throw path. The validated fix and a regression-test list are in my reply on ThreadPool.hpp:93.

Should fix before merge

Details, measurements and validated fixes are in the inline comments.

  1. removeReader()/removeGenericReader() destroy callbacks under _reader_mtx: a self-deadlock when the captured state owns a ReaderHandle. The #112 fix does not cover it (RedisAdapter.cpp:439).
  2. Legacy typed readers call back with an empty list on a width mismatch, and RedisCache<T> segfaults on it (RedisAdapterTempl.hpp:544).
  3. The count-less XREAD plus explicit cursors delivers the whole backlog in one reply (+600 MiB RSS for 200 × 1 MiB), and livelocks once building the reply takes longer than the socket timeout (RedisAdapter.cpp:489).
  4. An unresolvable $ tail causes silent per-key loss (fast=50 slow=0) (RedisAdapter.cpp:474).
  5. Reading the stop stream from $ causes ~500 ms stalls under _reader_mtx, in chains of up to 85 (RedisAdapter.cpp:398).
  6. The retry loop logs LOG_ERR 20 times a second for the whole outage, with no backoff (RedisAdapter.cpp:493).
  7. cpo.wait_timeout produces false transport failures and reconnect storms when readers >= cxn.size (RedisConnection.hpp:98).
  8. getStreamSnapshot() treats a rejection as a disconnect and reconnects, and a failed snapshot returns a replay cursor (RedisAdapter.cpp:194).

Worth fixing / follow-ups

  • Cancellation wording, and a test for the fence (RedisAdapter.hpp:144).

  • [[nodiscard]], parameter order and documented throws (RedisAdapter.hpp:164).

  • A _shutdown check in start_reader() (RedisAdapter.cpp:378).

  • A zero-copy fast path per registration (RedisAdapter.cpp:512).

  • lifecycle_test hygiene and untested claims (lifecycle_test.cpp:111).

  • register_reader isn't exception-safe after stop_reader(). If creating the std::thread throws (EAGAIN), the caller gets an exception and no handle. The registration still stays active and later receives callbacks, and the bucket stays stopped until an unrelated reader call.

  • Generic keys containing braces without a usable hash tag (a{b, {}abc) get no stop stream, even on standalone Redis. Each stop then waits up to timeout under _reader_mtx; with timeout=0, removeGenericReader() never returns (main created a stop key for every shape).

  • Docs and CHANGELOG (CONTRIBUTING asks for both). Behaviour changes not recorded:

    • exact-width typed reads; empty arrays are now valid items; strict bool decoding;
    • RA_Time(string) rejects malformed IDs;
    • what removeReader returns now;
    • corrected base/sub for foreign base keys;
    • readers retry instead of stopping;
    • cxn.timeout now also bounds connect and pool waits.

    docs/stream-subscriptions.md isn't linked from README.md, docs/README.md or docs/api.md. Its deferral sentence is inaccurate: an owned default $ subscription is resolved at subscribe time, so it already keeps updates that arrive during deferral.

  • Mock. Under MOCK_REDIS_ADAPTER the new connection-free helpers and types (decodeScalar, decodeArray, compareStreamIds, StreamBatch) are unavailable. A small shared header would fix that at no cost (#124).

  • Nits.

    • STOP_STUB is dead code now.
    • "_" duplicates DEFAULT_FIELD.
    • StreamEntry/StreamBatch duplicate Item/ItemStream.
    • Four comments still say "true if reader started".
    • add_reader_helper/addGenericReader call reader_token() a second time outside the lock, so the returned bool can disagree with where the registration was placed.

Scope question

Besides the owned API, this PR bundles several independent changes:

  • strict decoding for every legacy getter and reader;
  • the RA_Time parser;
  • ThreadPool semantics;
  • reader retry and $ resolution;
  • connect/pool timeouts;
  • rewritten legacy remove semantics.

When closing #91 you asked for each of its fixes to "land as a focused change on current main with its own regression coverage", and the same reasoning applies here. Could you split out at least the timeout change, which carries its own risk (item 7)? Otherwise, please list each change with its compatibility impact in the template's Compatibility section.

Issues

  • Please add Fixes #84. Removal now scans every bucket, including NO_TOKEN; on main, a reader removed while disconnected comes back after reconnect and keeps receiving data.
  • Add Fixes #112 with the fix commit (this PR targets main, so it will auto-close).
  • Reference #93 as partially addressed: the bucket recovers once the key is repaired, but stalls while the key is the wrong type, and now logs 20 times a second.
  • Reference #73. This PR removes several plausible UB sources: unchecked typed reads, the ThreadPool _go race, terminate on callback exception, and the moved batch. No root cause is claimed.
  • #114 concerns the same RA_Time(string) parser. The rewrite here keeps main's reading of the suffix as nanoseconds, so it neither resolves nor conflicts with #114.
  • Filed while reviewing (pre-existing on main):
    • #115 Unix-socket adapters killed by SIGPIPE
    • #116 watchdog: aborts if destroyed early, never re-registered after it expires
    • #117 a failed reconnect nulls the client
    • #118 reconnect/destructor races → std::terminate
    • #119 / #120 / #121 RedisCache use-after-free, shared-lock race, and writeBuffer keeping the oldest entry of a batch / crashing on an empty one
    • #122 idle readers reconnect every cycle
    • #123 read paths reconnect on rejections
    • #124 mock drift
    • #126 rename() overwrites its destination; generic-key substring match
    • #125 CI/build hygiene

Comment thread RedisAdapter.cpp
Comment thread RedisAdapterTempl.hpp
Comment thread RedisAdapter.cpp Outdated
Comment thread RedisAdapter.cpp Outdated
Comment thread RedisAdapter.cpp Outdated
Comment thread RedisAdapter.hpp Outdated
Comment thread RedisAdapter.hpp Outdated
Comment thread RedisAdapter.cpp Outdated
Comment thread RedisAdapter.cpp Outdated
Comment thread lifecycle_test.cpp Outdated
@derekste

Copy link
Copy Markdown
Member Author

@bigsamich Thanks for the detailed reproductions and for checking both the intended improvements and the failure cases. I'm holding this PR for a revised head and fresh review.

The first fixes need to cover the full callback-lifetime problem: release the running job's captures while unlocked after both catch handlers, and remove callback destruction from the locked legacy-removal paths. The regressions should include normal and throwing callbacks, sibling handles, queued jobs, and one/four-worker configurations.

I'll also address the bounded-read and startup cases: cap each XREAD reply, resolve failed $ tails per key, make the stop-stream wake-up reliable, back off persistent read errors, and prevent readers from restarting during shutdown. Decode failures must not call the legacy cache with an empty batch. Snapshot rejection must remain distinct from transport loss and must not return a replay cursor that appears usable.

The connect/pool timeout change deserves a focused change and its own compatibility tests. I'll separate it from the owned-subscription work, then document the retained legacy behavior changes, cancellation semantics, handle ownership, and parser/decoding rules. The queue-release contract also needs explicit documentation and coverage.

Please share the validated patches and reproducers. I'll use those alongside the existing tests, cover the remaining API/test-hygiene comments, and update the issue links: #84 and #112 only close after the corrected implementation reaches main; #93 and #73 remain qualified follow-ups, with no claim that this establishes the tcache failure's root cause.

The fixes above are in progress. I'll post the changed commits and validation before requesting another code review.

@derekste
derekste requested a review from bigsamich September 30, 2026 22:54
@derekste
derekste changed the base branch from main to dev/reader-connection-policy September 30, 2026 22:54
@derekste

Copy link
Copy Markdown
Member Author

@bigsamich, the reviewed subscription fixes are pushed and ready for another pass. The connection/pool policy is now split into #127, and this PR is stacked above it.

The capture cleanup sits after both catch handlers and before the worker relocks; canceled captures and legacy removals also retire outside their respective locks. The queued-cancellation fence is now tested, and its documented contract allows a callback already past the active check to enter. Reader startup has a shutdown fence, watchdog startup/destruction is joined safely, and reconnect-thread assignment is serialized.

The cursor fixes include stop cursors starting at 0-0, per-key tail retries, an XREAD-only tail fallback on Redis 7.4, a bounded XREAD COUNT, complete-batch sharing, and snapshot rejection distinct from disconnect. Failed snapshots keep $ and the example checks acceptance before subscribing. Strict legacy decoding skips fully rejected callback batches and reports RA_INVALID_PAYLOAD for malformed single items; the cache also guards empty batches. Wire/time helpers compile in mock builds, with byte-count/allocation coverage.

Evidence: the original pool-only deadlock was reproduced before the fix; 70/70 Linux Release/C++20/Redis 7.4.2 tests pass after it. The macOS ASan+UBSan run passed 67 cases, with the two 7.4-only ACL cases exercised on Linux. The lifecycle tests now use GoogleTest and unique keys with cleanup, including the previously missing default-tail, rewind/filter, deferral, remove-all, multi-worker, scoped listen-only retries, capture retirement, shutdown-time legacy removal and NO_TOKEN cases.

Base automatically changed from dev/reader-connection-policy to main October 7, 2026 15:59
@derekste
derekste merged commit 5bdf49e into main Oct 7, 2026
1 check passed
@derekste
derekste deleted the dev/subscription-lifecycle branch October 7, 2026 16:01
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants