diff --git a/docs.openc3.com/docs/configuration/interfaces.md b/docs.openc3.com/docs/configuration/interfaces.md index db580d4cd2..c95e99b7ba 100644 --- a/docs.openc3.com/docs/configuration/interfaces.md +++ b/docs.openc3.com/docs/configuration/interfaces.md @@ -187,6 +187,16 @@ Maximum number of bytes buffered on the interface read queue. The stream and UDP Because the read thread keeps draining the socket, a sender is no longer slowed by TCP flow control until the queue is full. If the interface can't keep up, telemetry can fall behind by up to READ_QUEUE_MAX_SIZE bytes before the sender is throttled. Lower READ_QUEUE_MAX_SIZE if you'd rather the sender block sooner. ::: +Packets read through the queue are timestamped with the time their data was read from the socket, not the time they were processed, so the received time stays accurate while the queue is backed up. + +:::note Sizing memory +READ_QUEUE_MAX_SIZE applies to each queue, so the memory an interface can hold is READ_QUEUE_MAX_SIZE times the number of queues. The TCP/IP server interface has one queue per connected client, so 10 clients at the default 20MB can buffer up to 200MB. All the interfaces in a microservice container share its memory, so lower READ_QUEUE_MAX_SIZE when running many interfaces or TCP/IP server interfaces that accept many clients. +::: + +:::note Disconnecting +When the interface disconnects it closes the stream, which unblocks the read thread, and then waits up to 2 seconds for the thread to exit. The built in interfaces wake up their pending read when they disconnect. A custom stream whose `disconnect` doesn't unblock a pending `read` delays each disconnect by those 2 seconds. In Ruby the thread is then killed and a warning is logged. In Python the thread can't be killed, so it is left blocked until its read returns and then exits without using the data. Make sure `disconnect` wakes up any blocked `read` (for example by closing the socket or writing to a wake up pipe). +::: + :::warning Custom interfaces which read the stream directly Custom interfaces which subclass the stream, TCP/IP, serial, MQTT stream or UDP interfaces and read the stream or socket themselves now compete with the read thread for data, which silently splits and loses bytes. This includes overriding `read_interface` without calling `super`, reading a handshake response in `connect`, reading an acknowledgement in `write_interface`, and protocols which read `interface.stream` directly. Either update the interface to use the data passed through `read_interface` and the protocols, or set `OPTION READ_QUEUE_MAX_SIZE 0` to restore the previous inline reads. ::: diff --git a/openc3/lib/openc3/interfaces/interface.rb b/openc3/lib/openc3/interfaces/interface.rb index 191008e132..3b8da2323a 100644 --- a/openc3/lib/openc3/interfaces/interface.rb +++ b/openc3/lib/openc3/interfaces/interface.rb @@ -329,6 +329,12 @@ def read # Return packet @read_count += 1 Logger.warn("#{@name}: Interface unexpectedly requested disconnect") unless packet + # Timestamp buffered reads with when the data arrived rather than when + # it was processed (otherwise received_time is set when it is handled) + if packet and !packet.received_time + data_time = read_queue_data_time + packet.received_time = data_time if data_time + end return packet end rescue Exception => e @@ -452,6 +458,13 @@ def read_queue_bytes 0 end + # @return [Time, nil] When the data most recently returned by + # read_interface was read from the underlying source, or nil if unknown. + # Interfaces which buffer raw reads override this (see ReadQueue). + def read_queue_data_time + nil + end + # @return [Boolean] Whether reading is allowed def read_allowed? @read_allowed @@ -573,7 +586,7 @@ def convert_packet_to_data(packet) # @return [String] Raw packet data def read_interface_base(data, _extra = nil) if @save_raw_data - @read_raw_data_time = Time.now + @read_raw_data_time = read_queue_data_time || Time.now @read_raw_data = data.clone end @bytes_read += data.length diff --git a/openc3/lib/openc3/utilities/read_queue.rb b/openc3/lib/openc3/utilities/read_queue.rb index a1d1d258a1..0d6fad7012 100644 --- a/openc3/lib/openc3/utilities/read_queue.rb +++ b/openc3/lib/openc3/utilities/read_queue.rb @@ -13,6 +13,7 @@ require 'openc3/top_level' require 'openc3/core_ext/exception' +require 'openc3/core_ext/time' require 'openc3/utilities/logger' module OpenC3 @@ -48,16 +49,24 @@ module ReadQueue # @return [Integer] Maximum number of bytes buffered on the queue attr_accessor :read_queue_max_size + # @return [Time, nil] When the read thread read the data most recently + # returned by read_queue_pop (nil when the read queue is disabled) + attr_reader :read_queue_data_time + # Initialize the read queue attributes. Must be called from the including # class initialize method. def initialize_read_queue(max_size = DEFAULT_READ_QUEUE_MAX_SIZE) @raw_read_queue = nil @raw_read_thread = nil - @raw_read_cancel = false @raw_read_bytes = 0 @raw_read_budget = 0 - # Guards @raw_read_bytes / @raw_read_budget and wakes the read thread when - # the bytes it queued are consumed and there is room to queue more + @read_queue_data_time = nil + # Guards @raw_read_queue / @raw_read_bytes / @raw_read_budget and wakes the + # read thread when the bytes it queued are consumed and there is room to + # queue more. Each read thread is cancelled by closing its own queue under + # this mutex, so a thread which outlives stop_read_queue_thread (and a new + # thread starting) keeps seeing its old queue closed and can never touch + # the byte counters or read again. @raw_read_mutex = Mutex.new @raw_read_condition = ConditionVariable.new @read_queue_max_size = max_size @@ -110,11 +119,10 @@ def start_read_queue_thread queue = Queue.new @raw_read_mutex.synchronize do - @raw_read_cancel = false @raw_read_bytes = 0 @raw_read_budget = 0 + @raw_read_queue = queue end - @raw_read_queue = queue @raw_read_thread = Thread.new do read_queue_thread_body(queue) rescue Exception => e @@ -125,17 +133,18 @@ def start_read_queue_thread # Stop the read thread and discard anything left on the queue. The read # source should be disconnected first so the thread isn't blocked reading. def stop_read_queue_thread - queue = @raw_read_queue - @raw_read_queue = nil @raw_read_mutex.synchronize do - @raw_read_cancel = true + # Closing the queue under the mutex cancels the read thread (it either + # already charged its bytes, which are reset below, or never will) and + # unblocks read_queue_pop + @raw_read_queue.close if @raw_read_queue + @raw_read_queue = nil @raw_read_bytes = 0 @raw_read_budget = 0 + @read_queue_data_time = nil # Unblock the read thread if it is waiting for room on the queue @raw_read_condition.broadcast end - # Closing the queue unblocks read_queue_pop - queue.close if queue thread = @raw_read_thread @raw_read_thread = nil # The read source should already be disconnected which unblocks the read @@ -161,12 +170,17 @@ def read_queue_pop end data = queue.pop - if data.kind_of?(String) + if data.kind_of?(Array) + data, time = data @raw_read_mutex.synchronize do - @raw_read_bytes -= data.length - @raw_read_budget -= data.length + READ_QUEUE_ENTRY_OVERHEAD - # Tell the read thread there is room for more data - @raw_read_condition.broadcast + # Counters were reset if the queue was stopped / replaced + if queue.equal?(@raw_read_queue) + @raw_read_bytes -= data.length + @raw_read_budget -= data.length + READ_QUEUE_ENTRY_OVERHEAD + @read_queue_data_time = time + # Tell the read thread there is room for more data + @raw_read_condition.broadcast + end end end # Exceptions raised by the read thread are re-raised here so they are @@ -179,7 +193,7 @@ def read_queue_pop # protected def read_queue_thread_body(queue) - loop do + until queue.closed? begin data = read_queue_data() rescue Exception => e @@ -194,9 +208,11 @@ def read_queue_thread_body(queue) # Charge the bytes against the budget before the push so # read_queue_bytes never goes negative if the data is dequeued # before we get back here - break unless reserve_read_queue_bytes(data.length) + break unless reserve_read_queue_bytes(queue, data.length) - queue.push(data) + # Queue the time of the read so packets are timestamped when the data + # arrived rather than when it was processed + queue.push([data, Time.now.sys]) end rescue ClosedQueueError # Interface disconnected while we were pushing @@ -211,15 +227,15 @@ def read_queue_thread_body(queue) # than never fitting, so the budget can be exceeded by at most one read. # # @return [Boolean] Whether the bytes were reserved (false if disconnected) - def reserve_read_queue_bytes(length) + def reserve_read_queue_bytes(queue, length) cost = length + READ_QUEUE_ENTRY_OVERHEAD @raw_read_mutex.synchronize do while @raw_read_budget > 0 and (@raw_read_budget + cost) > @read_queue_max_size - return false if @raw_read_cancel + return false if queue.closed? @raw_read_condition.wait(@raw_read_mutex, QUEUE_POLL_TIMEOUT) end - return false if @raw_read_cancel + return false if queue.closed? @raw_read_bytes += length @raw_read_budget += cost diff --git a/openc3/python/openc3/interfaces/interface.py b/openc3/python/openc3/interfaces/interface.py index f54f174860..6ac82ae8f2 100644 --- a/openc3/python/openc3/interfaces/interface.py +++ b/openc3/python/openc3/interfaces/interface.py @@ -69,6 +69,10 @@ def __init__(self): self.read_raw_data = "" self.written_raw_data = "" self.read_raw_data_time = None + # When the data most recently returned by read_interface was read from the + # underlying source, or None if unknown. Set by interfaces which buffer + # raw reads (see ReadQueue). + self.read_queue_data_time = None self.written_raw_data_time = None self.config_params = [] self.interfaces = [] @@ -221,6 +225,10 @@ def read(self): self.read_count += 1 if not packet: Logger.warn(f"{self.name}: Interface unexpectedly requested disconnect") + # Timestamp buffered reads with when the data arrived rather than + # when it was processed (otherwise received_time is set when handled) + if packet and packet.received_time is None and self.read_queue_data_time is not None: + packet.received_time = self.read_queue_data_time return packet except Exception as error: Logger.error(f"{self.name}: Error reading from interface") @@ -425,7 +433,7 @@ def convert_packet_to_data(self, packet): # self.return [String] Raw packet data def read_interface_base(self, data, extra=None): if self.save_raw_data: - self.read_raw_data_time = datetime.now(timezone.utc) + self.read_raw_data_time = self.read_queue_data_time or datetime.now(timezone.utc) self.read_raw_data = data self.bytes_read += len(data) if self.stream_log_pair: diff --git a/openc3/python/openc3/interfaces/stream_interface.py b/openc3/python/openc3/interfaces/stream_interface.py index 479d1598df..5cd640694f 100644 --- a/openc3/python/openc3/interfaces/stream_interface.py +++ b/openc3/python/openc3/interfaces/stream_interface.py @@ -22,7 +22,6 @@ def __init__(self, protocol_type=None, protocol_args=None): if protocol_args is None: protocol_args = [] super().__init__() - self.initialize_read_queue() self._stream = None self.protocol_type = ConfigParser.handle_none(protocol_type) self.protocol_args = protocol_args diff --git a/openc3/python/openc3/interfaces/udp_interface.py b/openc3/python/openc3/interfaces/udp_interface.py index 0d243e04c7..261b4523fa 100644 --- a/openc3/python/openc3/interfaces/udp_interface.py +++ b/openc3/python/openc3/interfaces/udp_interface.py @@ -46,7 +46,6 @@ def __init__( bind_address="0.0.0.0", ): super().__init__() - self.initialize_read_queue() self.hostname = ConfigParser.handle_none(hostname) if self.hostname is not None: self.hostname = str(hostname) diff --git a/openc3/python/openc3/utilities/read_queue.py b/openc3/python/openc3/utilities/read_queue.py index c2555702ec..1eb75fb79e 100644 --- a/openc3/python/openc3/utilities/read_queue.py +++ b/openc3/python/openc3/utilities/read_queue.py @@ -12,6 +12,7 @@ import queue import threading import traceback +from datetime import datetime, timezone from openc3.interfaces.interface import Interface from openc3.utilities.logger import Logger @@ -42,16 +43,23 @@ # # Subclasses must implement read_queue_data which performs # a single blocking read and returns the data read or None to indicate the read -# source is done (which disconnects the interface). +# source is done (which disconnects the interface). The read queue is set up by +# __init__ so subclasses only need to call super().__init__(). # # Setting READ_QUEUE_MAX_SIZE to 0 disables the read thread and reads inline # from read_queue_pop, which is how interfaces read before the queue existed. class ReadQueue(Interface): - # Initialize the read queue attributes. Must be called from the subclass - # __init__ method. + def __init__(self): + super().__init__() + self.initialize_read_queue() + + # Initialize (or reset) the read queue attributes. Called by __init__. def initialize_read_queue(self, max_size=DEFAULT_READ_QUEUE_MAX_SIZE): self._read_queue = None self.read_queue_thread = None + # When the read thread read the data most recently returned by + # read_queue_pop (None when the read queue is disabled) + self.read_queue_data_time = None # Each read thread gets its own cancel event. A thread which is still # blocked in a read when it is replaced (Python can't kill it) keeps # seeing its own event set, so it can never touch the byte counters or @@ -131,6 +139,7 @@ def stop_read_queue_thread(self): cancel_event.set() self._read_queue_bytes = 0 self._read_queue_budget = 0 + self.read_queue_data_time = None # Unblock the read thread if it is waiting for room on the queue self._read_queue_condition.notify_all() if read_queue is not None: @@ -177,12 +186,14 @@ def read_queue_pop(self): if read_queue is not self._read_queue or thread is None or not thread.is_alive(): return None continue - if isinstance(data, (bytes, bytearray)): + if isinstance(data, tuple): + data, data_time = data with self._read_queue_condition: # Counters were reset if the queue was stopped / replaced if read_queue is self._read_queue: self._read_queue_bytes -= len(data) self._read_queue_budget -= len(data) + READ_QUEUE_ENTRY_OVERHEAD + self.read_queue_data_time = data_time # Tell the read thread there is room for more data self._read_queue_condition.notify_all() # Exceptions raised by the read thread are re-raised here so they are @@ -228,6 +239,8 @@ def _read_queue_thread_body(self, read_queue, cancel_event): # before we get back here if not self._reserve_read_queue_bytes(len(data), cancel_event): break - read_queue.put(data) + # Queue the time of the read so packets are timestamped when the + # data arrived rather than when it was processed + read_queue.put((data, datetime.now(timezone.utc))) except Exception: Logger.error(f"{self.name}: Read queue thread unexpectedly died: {traceback.format_exc()}") diff --git a/openc3/python/test/interfaces/protocols/test_fixed_protocol.py b/openc3/python/test/interfaces/protocols/test_fixed_protocol.py index 1165118181..a5ecc14ae4 100644 --- a/openc3/python/test/interfaces/protocols/test_fixed_protocol.py +++ b/openc3/python/test/interfaces/protocols/test_fixed_protocol.py @@ -85,7 +85,9 @@ def test_returns_unknown_packets(self): # The read thread buffers ahead so flush it to pick up the new data self.interface.stop_read_queue_thread() packet = self.interface.read() - self.assertIsNone(packet.received_time) + # Unknown packets are timestamped with when their data was read + self.assertIsNotNone(packet.received_time) + self.assertEqual(packet.received_time, self.interface.read_queue_data_time) self.assertIsNone(packet.target_name) self.assertIsNone(packet.packet_name) self.assertEqual(packet.buffer, b"\x00") @@ -110,7 +112,9 @@ def test_handles_targets_with_no_defined_telemetry(self): self.interface.tlm_target_names = ["EMPTY"] TestFixedProtocol.index = 1 packet = self.interface.read() - self.assertIsNone(packet.received_time) + # Unknown packets are timestamped with when their data was read + self.assertIsNotNone(packet.received_time) + self.assertEqual(packet.received_time, self.interface.read_queue_data_time) self.assertIsNone(packet.target_name) self.assertIsNone(packet.packet_name) self.assertEqual(packet.buffer, b"\x01") diff --git a/openc3/python/test/interfaces/test_stream_interface.py b/openc3/python/test/interfaces/test_stream_interface.py index 5a44475da6..858be2357a 100644 --- a/openc3/python/test/interfaces/test_stream_interface.py +++ b/openc3/python/test/interfaces/test_stream_interface.py @@ -13,6 +13,7 @@ import threading import time import unittest +from datetime import datetime, timezone from unittest.mock import Mock, patch from openc3.interfaces.protocols.burst_protocol import BurstProtocol @@ -245,6 +246,26 @@ def test_reads_inline_without_a_read_thread_when_the_max_size_is_0(self): stream.disconnect() self.assertEqual(self.interface.read_interface(), (None, None)) + def test_timestamps_packets_with_when_the_data_was_read_rather_than_processed(self): + self.interface.stream = QueueStream(b"\x01\x02") + self.interface.connect() + self.wait_for_queue_size(1) + queued = datetime.now(timezone.utc) + time.sleep(0.05) + + packet = self.interface.read() + self.assertEqual(packet.buffer, b"\x01\x02") + self.assertLessEqual(packet.received_time, queued) + self.assertEqual(self.interface.read_queue_data_time, packet.received_time) + + def test_doesnt_timestamp_packets_when_reading_inline(self): + self.interface.set_option("READ_QUEUE_MAX_SIZE", ["0"]) + self.interface.stream = QueueStream(b"\x01\x02") + self.interface.connect() + packet = self.interface.read() + self.assertEqual(packet.buffer, b"\x01\x02") + self.assertIsNone(packet.received_time) + def test_disconnect_stops_the_read_thread_and_clears_the_queue(self): self.interface.stream = QueueStream(b"\x01", b"\x02") self.interface.connect() diff --git a/openc3/spec/interfaces/protocols/fixed_protocol_spec.rb b/openc3/spec/interfaces/protocols/fixed_protocol_spec.rb index dcd907e323..ff0a65d60a 100644 --- a/openc3/spec/interfaces/protocols/fixed_protocol_spec.rb +++ b/openc3/spec/interfaces/protocols/fixed_protocol_spec.rb @@ -94,7 +94,9 @@ def read # The read thread buffers ahead so flush it to pick up the new data @interface.stop_read_queue_thread packet = @interface.read - expect(packet.received_time.to_f).to eql 0.0 + # Unknown packets are timestamped with when their data was read + expect(packet.received_time).to_not be_nil + expect(packet.received_time).to eql @interface.read_queue_data_time expect(packet.target_name).to eql nil expect(packet.packet_name).to eql nil expect(packet.buffer).to eql "\x00" @@ -117,7 +119,9 @@ def read @interface.tlm_target_names = ['EMPTY'] $index = 1 packet = @interface.read - expect(packet.received_time.to_f).to eql 0.0 + # Unknown packets are timestamped with when their data was read + expect(packet.received_time).to_not be_nil + expect(packet.received_time).to eql @interface.read_queue_data_time expect(packet.target_name).to eql nil expect(packet.packet_name).to eql nil expect(packet.buffer).to eql "\x01" diff --git a/openc3/spec/interfaces/stream_interface_spec.rb b/openc3/spec/interfaces/stream_interface_spec.rb index c847f648e9..9a8413a1e6 100644 --- a/openc3/spec/interfaces/stream_interface_spec.rb +++ b/openc3/spec/interfaces/stream_interface_spec.rb @@ -54,6 +54,37 @@ def write(_data) end end + # Stream whose read blocks until released, even across a disconnect, like a + # stream whose disconnect doesn't wake up a pending read + class StuckStream < Stream + attr_reader :reads, :reading + + def initialize + @reading = Queue.new + @release = Queue.new + @reads = 0 + end + + def release + @release << true + end + + def connect; end + + def connected?; true; end + + def disconnect; end + + def read + @reads += 1 + @reading << true + @release.pop(timeout: 5) + "\x01\x02\x03" + end + + def write(_data); end + end + let(:interface) { StreamInterface.new } after(:each) do @@ -208,6 +239,32 @@ def write(_data) end end + describe "read" do + it "timestamps packets with when the data was read rather than processed" do + stream = QueueStream.new("\x01\x02") + interface.stream = stream + interface.connect + start = Time.now + sleep(0.001) while interface.read_queue_size < 1 and (Time.now - start) < 2 + queued = Time.now.sys + sleep(0.05) + + packet = interface.read + expect(packet.buffer).to eql "\x01\x02" + expect(packet.received_time).to be <= queued + expect(interface.read_queue_data_time).to eql packet.received_time + end + + it "doesn't timestamp packets when reading inline" do + interface.set_option('READ_QUEUE_MAX_SIZE', ['0']) + interface.stream = QueueStream.new("\x01\x02") + interface.connect + packet = interface.read + expect(packet.buffer).to eql "\x01\x02" + expect(packet.received_time).to be_nil + end + end + describe "disconnect" do it "stops the read thread and clears the queue" do stream = QueueStream.new("\x01", "\x02") @@ -233,6 +290,32 @@ def write(_data) expect(interface.read_queue_size).to eql 0 expect(interface.read_interface()[0]).to eql "\x02" end + + it "doesn't let a read thread which survived being stopped queue or read again" do + stub_const("OpenC3::ReadQueue::THREAD_JOIN_TIMEOUT", 0.05) + # Simulate the thread surviving being killed + allow(OpenC3).to receive(:kill_thread) + stuck = StuckStream.new + interface.stream = stuck + interface.connect + expect(stuck.reading.pop(timeout: 2)).to be true + old_thread = interface.instance_variable_get(:@raw_read_thread) + + interface.stream = QueueStream.new("\x0a") + interface.connect + start = Time.now + sleep(0.001) while interface.read_queue_size < 1 and (Time.now - start) < 2 + expect(old_thread.alive?).to be true + + # Once its read finally returns the old thread must exit without + # charging the bytes against the new queue or reading the new stream + stuck.release + expect(old_thread.join(2)).to_not be_nil + expect(stuck.reads).to eql 1 + expect(interface.read_queue_bytes).to eql 1 + expect(interface.read_interface()[0]).to eql "\x0a" + expect(interface.read_queue_bytes).to eql 0 + end end end end