Conversation
XREAD without COUNT returns everything from the consumer's offset to the stream tail, so a consumer that falls behind on high-rate telemetry pulls an unbounded reply into memory at once. Default to 1000 entries per stream and have the cmd-tlm-api streaming threads read at most one batch per stream. Streaming threads now advance offsets past skipped stored packets so a capped read cannot stall re-reading them. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
|
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #3950 +/- ##
==========================================
+ Coverage 80.08% 80.14% +0.06%
==========================================
Files 901 901
Lines 68356 68367 +11
Branches 2645 2698 +53
==========================================
+ Hits 54743 54793 +50
+ Misses 12946 12908 -38
+ Partials 667 666 -1
Flags with carried forward coverage won't be shown. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
| patcher.start() | ||
| self.addCleanup(patcher.stop) | ||
| for i in range(5): | ||
| EphemeralStore.write_topic("TEST__TOPIC", {"i": str(i)}) |
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.


What
Store#read_topics(Ruby) andStoreImplementation.read_topics(Python) now defaultcounttoREAD_TOPICS_DEFAULT_COUNT(1000 entries per stream) when the caller passes none. The positional signature is unchanged, so existing plugin callers keep working.StreamingThread,LoggedStreamingThread,TopicsThreadandMessagesThreadpass an explicit count equal to their batch size, so each XREAD reads at most one batch per stream.StreamingThread/LoggedStreamingThreadnow advance object offsets paststoredpackets they skip. Before, offsets only moved on realtime packets. With a capped read, a run of stored packets longer than the count would have been re-read forever.Why
XREADwith noCOUNTreturns every entry from the consumer's offset to the end of the stream, and the client holds the whole reply in memory before the caller's block runs. With high-rate telemetry, a consumer that falls behind can pull hundreds of thousands of entries in one call, which spikes the process's memory (RSS). In Ruby that memory is often never returned to the OS. A bounded read turns a burst like that into a series of normal-sized reads.Caller audit
I checked every
read_topicscaller in core Ruby/Python, cmd-tlm-api, script-runner-api and Enterprise. All of them either loop and continue from the offsets tracked per thread, or pass and update explicit offsets. A capped read therefore just continues on the next iteration.get_packets,ConfigTopic.readandLimitsEventTopic.readalready passed their own count. The one caller that could stall was the realtime streaming path skipping stored packets without advancing offsets; that is fixed here.Testing
openc3/spec/utilities/store_spec.rb(default cap and explicit count) andopenc3-cosmos-cmd-tlm-api/spec/models/streaming_thread_spec.rb(count passed through; offsets advance past stored packets).test/utilities/test_store_implementation.py::TestStoreReadTopics.openc3rspec on utilities, topics, api, models and microservices: 1893 examples. The 2 failures also fail onmain(local AWS credential issue).spec/models: 210 examples, 0 failures.ruff checkis clean.🤖 Generated with Claude Code