diff --git a/openc3-cosmos-cmd-tlm-api/app/models/logged_streaming_thread.rb b/openc3-cosmos-cmd-tlm-api/app/models/logged_streaming_thread.rb index abfeb9b0ad..f51f3dc4a1 100644 --- a/openc3-cosmos-cmd-tlm-api/app/models/logged_streaming_thread.rb +++ b/openc3-cosmos-cmd-tlm-api/app/models/logged_streaming_thread.rb @@ -206,9 +206,13 @@ def redis_thread_body any_result = false db_shard_groups.each do |db_shard, group| break if @cancel_thread - xread_result = OpenC3::Topic.read_topics(group[:topics], group[:offsets], timeout_per_db_shard, db_shard: db_shard) do |topic, msg_id, msg_hash, _| + xread_result = OpenC3::Topic.read_topics(group[:topics], group[:offsets], timeout_per_db_shard, @max_batch_size, db_shard: db_shard) do |topic, msg_id, msg_hash, _| stored = OpenC3::ConfigParser.handle_true_false(msg_hash["stored"]) - next if stored + if stored + # Advance offsets past skipped messages so a capped read makes progress + advance_offsets(topic, msg_id, item_objects_by_topic, packet_objects_by_topic) + next + end break if @cancel_thread @@ -218,10 +222,7 @@ def redis_thread_body time = msg_hash['time'].to_i if time <= last_time # Skip messages already delivered from TSDB, but advance offsets - objects = item_objects_by_topic[topic] - objects.each { |object| object.offset = msg_id } if objects - objects = packet_objects_by_topic[topic] - objects.each { |object| object.offset = msg_id } if objects + advance_offsets(topic, msg_id, item_objects_by_topic, packet_objects_by_topic) next end # Past the overlap for this topic - clear its filter diff --git a/openc3-cosmos-cmd-tlm-api/app/models/messages_thread.rb b/openc3-cosmos-cmd-tlm-api/app/models/messages_thread.rb index 1477a8154d..61ebd4639f 100644 --- a/openc3-cosmos-cmd-tlm-api/app/models/messages_thread.rb +++ b/openc3-cosmos-cmd-tlm-api/app/models/messages_thread.rb @@ -149,7 +149,7 @@ def file_thread_body def redis_thread_body results = [] - OpenC3::Topic.read_topics(@topics, @offsets) do |topic, msg_id, msg_hash, _redis| + OpenC3::Topic.read_topics(@topics, @offsets, 1000, @max_batch_size) do |topic, msg_id, msg_hash, _redis| @offsets[@offset_index_by_topic[topic]] = msg_id msg_hash[:msg_id] = msg_id result_entry = handle_log_entry(msg_hash) diff --git a/openc3-cosmos-cmd-tlm-api/app/models/streaming_thread.rb b/openc3-cosmos-cmd-tlm-api/app/models/streaming_thread.rb index 01ec7b7d7e..11cd2f3cb2 100644 --- a/openc3-cosmos-cmd-tlm-api/app/models/streaming_thread.rb +++ b/openc3-cosmos-cmd-tlm-api/app/models/streaming_thread.rb @@ -94,9 +94,14 @@ def redis_thread_body any_result = false db_shard_groups.each do |db_shard, group| break if @cancel_thread - xread_result = OpenC3::Topic.read_topics(group[:topics], group[:offsets], timeout_per_db_shard, db_shard: db_shard) do |topic, msg_id, msg_hash, _| + xread_result = OpenC3::Topic.read_topics(group[:topics], group[:offsets], timeout_per_db_shard, @max_batch_size, db_shard: db_shard) do |topic, msg_id, msg_hash, _| stored = OpenC3::ConfigParser.handle_true_false(msg_hash["stored"]) - next if stored # Ignore stored packets while realtime streaming + if stored # Ignore stored packets while realtime streaming + # Still advance offsets past them, otherwise a run of stored packets + # longer than the read count would be re-read forever + advance_offsets(topic, msg_id, item_objects_by_topic, packet_objects_by_topic) + next + end break if @cancel_thread @@ -156,6 +161,13 @@ def redis_thread_body end end + def advance_offsets(topic, msg_id, item_objects_by_topic, packet_objects_by_topic) + objects = item_objects_by_topic[topic] + objects.each { |object| object.offset = msg_id } if objects + objects = packet_objects_by_topic[topic] + objects.each { |object| object.offset = msg_id } if objects + end + def handle_message(msg_hash, objects) first_object = objects[0] time = msg_hash['time'].to_i diff --git a/openc3-cosmos-cmd-tlm-api/app/models/topics_thread.rb b/openc3-cosmos-cmd-tlm-api/app/models/topics_thread.rb index 6857ae8949..c990169821 100644 --- a/openc3-cosmos-cmd-tlm-api/app/models/topics_thread.rb +++ b/openc3-cosmos-cmd-tlm-api/app/models/topics_thread.rb @@ -90,7 +90,7 @@ def thread_setup def thread_body results = [] - OpenC3::Topic.read_topics(@topics, @offsets) do |topic, msg_id, msg_hash, redis| + OpenC3::Topic.read_topics(@topics, @offsets, 1000, @max_batch_size) do |topic, msg_id, msg_hash, redis| @offsets[@offset_index_by_topic[topic]] = msg_id msg_hash[:msg_id] = msg_id if @transmit_msg_id results << msg_hash diff --git a/openc3-cosmos-cmd-tlm-api/spec/models/streaming_thread_spec.rb b/openc3-cosmos-cmd-tlm-api/spec/models/streaming_thread_spec.rb new file mode 100644 index 0000000000..70c6b0b0f9 --- /dev/null +++ b/openc3-cosmos-cmd-tlm-api/spec/models/streaming_thread_spec.rb @@ -0,0 +1,44 @@ +# encoding: ascii-8bit + +# Copyright 2026 OpenC3, Inc. +# All Rights Reserved. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. +# See LICENSE.md for more details. +# +# This file may also be used under the terms of a commercial license +# if purchased from OpenC3, Inc. + +require "rails_helper" + +RSpec.describe StreamingThread, type: :model do + let(:topic) { "DEFAULT__DECOM__{INST}__HEALTH_STATUS" } + let(:streaming_api) { double("StreamingApi", transmit_results: nil) } + let(:object) do + Struct.new(:offset, :db_shard, :topic).new("0-0", 0, topic) + end + let(:collection) do + double("Collection", objects: [object], topics_offsets_and_objects: [[topic], ["0-0"], { topic => [object] }, {}]) + end + + describe "redis_thread_body" do + it "reads with a count bounded by the batch size" do + thread = StreamingThread.new(streaming_api, collection, 25) + expect(OpenC3::Topic).to receive(:read_topics).with([topic], ["0-0"], anything, 25, db_shard: 0).and_return({ topic => [["1-0", {}]] }) + thread.redis_thread_body + end + + it "advances offsets past stored packets so a capped read makes progress" do + thread = StreamingThread.new(streaming_api, collection, 25) + allow(OpenC3::Topic).to receive(:read_topics) do |*_args, **_kwargs, &block| + block.call(topic, "5-0", { "stored" => "true" }, nil) + { topic => [["5-0", {}]] } + end + expect(thread).not_to receive(:handle_message) + thread.redis_thread_body + expect(object.offset).to eql "5-0" + end + end +end diff --git a/openc3/lib/openc3/utilities/store_autoload.rb b/openc3/lib/openc3/utilities/store_autoload.rb index 4ff2dc3f2d..f416969e4d 100644 --- a/openc3/lib/openc3/utilities/store_autoload.rb +++ b/openc3/lib/openc3/utilities/store_autoload.rb @@ -54,6 +54,14 @@ class Store @@db_shard_cache_mutex = Mutex.new DB_SHARD_CACHE_TIMEOUT = 60 # seconds + # Maximum number of entries read_topics returns PER STREAM when the caller + # does not pass a count. XREAD without COUNT returns everything from the + # offset to the tail of the stream, which for a consumer that has fallen + # behind on a high-rate stream can be an enormous reply materialized in + # memory all at once. Callers loop on their offsets so a capped read simply + # continues on the next call. + READ_TOPICS_DEFAULT_COUNT = 1000 + # Mutex used to ensure that only one instance is created @@instance_mutex = Mutex.new @@ -225,8 +233,10 @@ def update_topic_offsets(topics) return offsets end + # @param count [Integer] Maximum entries to return per stream. nil uses READ_TOPICS_DEFAULT_COUNT. def read_topics(topics, offsets = nil, timeout_ms = 1000, count = nil) return {} if topics.empty? + count ||= READ_TOPICS_DEFAULT_COUNT Thread.current[:topic_offsets] ||= {} topic_offsets = Thread.current[:topic_offsets] begin diff --git a/openc3/python/openc3/utilities/store_implementation.py b/openc3/python/openc3/utilities/store_implementation.py index 9b8e7d80e5..8da980fb33 100644 --- a/openc3/python/openc3/utilities/store_implementation.py +++ b/openc3/python/openc3/utilities/store_implementation.py @@ -69,6 +69,7 @@ class StoreMeta(type): "_db_shard_cache", "_db_shard_cache_lock", "DB_SHARD_CACHE_TIMEOUT", + "READ_TOPICS_DEFAULT_COUNT", } ) @@ -98,6 +99,14 @@ class Store(metaclass=StoreMeta): _db_shard_cache_lock = threading.Lock() DB_SHARD_CACHE_TIMEOUT = 60 # seconds + # Maximum number of entries read_topics returns PER STREAM when the caller + # does not pass a count. XREAD without COUNT returns everything from the + # offset to the tail of the stream, which for a consumer that has fallen + # behind on a high-rate stream can be an enormous reply materialized in + # memory all at once. Callers loop on their offsets so a capped read simply + # continues on the next call. + READ_TOPICS_DEFAULT_COUNT = 1000 + # Get the singleton instance for a given db_shard @classmethod def instance(cls, pool_size=100, db_shard=0): @@ -236,8 +245,11 @@ def update_topic_offsets(self, topics): return offsets def read_topics(self, topics, offsets=None, timeout_ms=1000, count=None): + """count is the maximum entries to return per stream. None uses READ_TOPICS_DEFAULT_COUNT.""" if len(topics) == 0: return {} + if count is None: + count = self.READ_TOPICS_DEFAULT_COUNT thread_id = threading.get_native_id() if thread_id not in self.topic_offsets: self.topic_offsets[thread_id] = {} diff --git a/openc3/python/openc3/utilities/store_implementation.pyi b/openc3/python/openc3/utilities/store_implementation.pyi index 3ed921ba62..9f64949ab7 100644 --- a/openc3/python/openc3/utilities/store_implementation.pyi +++ b/openc3/python/openc3/utilities/store_implementation.pyi @@ -43,6 +43,7 @@ class Store(metaclass=StoreMeta): _db_shard_cache: Incomplete _db_shard_cache_lock: Incomplete DB_SHARD_CACHE_TIMEOUT: int + READ_TOPICS_DEFAULT_COUNT: int @classmethod def instance(cls, pool_size: int = 100, db_shard: int = 0): ... diff --git a/openc3/python/test/utilities/test_store_implementation.py b/openc3/python/test/utilities/test_store_implementation.py index 15db74a08e..9f9e6d7510 100644 --- a/openc3/python/test/utilities/test_store_implementation.py +++ b/openc3/python/test/utilities/test_store_implementation.py @@ -16,7 +16,8 @@ from valkey.exceptions import BusyLoadingError, ConnectionError, TimeoutError from valkey.retry import Retry -from openc3.utilities.store_implementation import Store +from openc3.utilities.store_implementation import EphemeralStore, Store +from test.test_helper import mock_redis class TestStoreImplementation(unittest.TestCase): @@ -50,3 +51,25 @@ def test_build_redis_configures_resilience(self): # the final (3rd) retry tops out at the 5s cap (jittered, so 2.5-5s). self.assertEqual(retry._backoff._cap, 5) self.assertLessEqual(retry._backoff.compute(3), 5) + + +class TestStoreReadTopics(unittest.TestCase): + def setUp(self): + mock_redis(self) + patcher = patch.object(Store, "READ_TOPICS_DEFAULT_COUNT", 2) + patcher.start() + self.addCleanup(patcher.stop) + for i in range(5): + EphemeralStore.write_topic("TEST__TOPIC", {"i": str(i)}) + + def test_caps_read_when_no_count_given(self): + ids = [msg_id for _, msg_id, _, _ in EphemeralStore.read_topics(["TEST__TOPIC"], ["0-0"])] + self.assertEqual(len(ids), 2) + # Continuing from the last offset picks up where the capped read stopped + ids += [msg_id for _, msg_id, _, _ in EphemeralStore.read_topics(["TEST__TOPIC"], [ids[-1]])] + self.assertEqual(len(ids), 4) + self.assertEqual(len(set(ids)), 4) + + def test_honors_explicit_count(self): + ids = [msg_id for _, msg_id, _, _ in EphemeralStore.read_topics(["TEST__TOPIC"], ["0-0"], None, 4)] + self.assertEqual(len(ids), 4) diff --git a/openc3/spec/utilities/store_spec.rb b/openc3/spec/utilities/store_spec.rb index b9c246223b..504cb4556f 100644 --- a/openc3/spec/utilities/store_spec.rb +++ b/openc3/spec/utilities/store_spec.rb @@ -51,5 +51,30 @@ module OpenC3 expect(a).not_to eq(b) end end + + describe "read_topics" do + before(:each) do + mock_redis() + stub_const("OpenC3::Store::READ_TOPICS_DEFAULT_COUNT", 2) + end + + it "caps the read at READ_TOPICS_DEFAULT_COUNT when no count is given" do + 5.times { |i| EphemeralStore.write_topic("TEST__TOPIC", { "i" => i.to_s }) } + ids = [] + EphemeralStore.read_topics(["TEST__TOPIC"], ["0-0"]) { |_t, msg_id, _h, _r| ids << msg_id } + expect(ids.length).to eql 2 + # Continuing from the last offset picks up where the capped read stopped + EphemeralStore.read_topics(["TEST__TOPIC"], [ids[-1]]) { |_t, msg_id, _h, _r| ids << msg_id } + expect(ids.length).to eql 4 + expect(ids.uniq.length).to eql 4 + end + + it "honors an explicit count" do + 5.times { |i| EphemeralStore.write_topic("TEST__TOPIC", { "i" => i.to_s }) } + ids = [] + EphemeralStore.read_topics(["TEST__TOPIC"], ["0-0"], nil, 4) { |_t, msg_id, _h, _r| ids << msg_id } + expect(ids.length).to eql 4 + end + end end end