Skip to content
Open
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
13 changes: 7 additions & 6 deletions openc3-cosmos-cmd-tlm-api/app/models/logged_streaming_thread.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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
Expand Down
2 changes: 1 addition & 1 deletion openc3-cosmos-cmd-tlm-api/app/models/messages_thread.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
16 changes: 14 additions & 2 deletions openc3-cosmos-cmd-tlm-api/app/models/streaming_thread.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion openc3-cosmos-cmd-tlm-api/app/models/topics_thread.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
44 changes: 44 additions & 0 deletions openc3-cosmos-cmd-tlm-api/spec/models/streaming_thread_spec.rb
Original file line number Diff line number Diff line change
@@ -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
10 changes: 10 additions & 0 deletions openc3/lib/openc3/utilities/store_autoload.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down
12 changes: 12 additions & 0 deletions openc3/python/openc3/utilities/store_implementation.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ class StoreMeta(type):
"_db_shard_cache",
"_db_shard_cache_lock",
"DB_SHARD_CACHE_TIMEOUT",
"READ_TOPICS_DEFAULT_COUNT",
}
)

Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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] = {}
Expand Down
1 change: 1 addition & 0 deletions openc3/python/openc3/utilities/store_implementation.pyi

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

25 changes: 24 additions & 1 deletion openc3/python/test/utilities/test_store_implementation.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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)
25 changes: 25 additions & 0 deletions openc3/spec/utilities/store_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Loading