F´ Flight Software - C/C++ Documentation
A framework for building embedded system applications to NASA flight quality standards.
ComAggregator.cpp
Go to the documentation of this file.
1 // ======================================================================
2 // \title ComAggregator.cpp
3 // \author lestarch
4 // \brief cpp file for ComAggregator component implementation class
5 // ======================================================================
6 
8 #include <cstring>
11 
12 namespace Svc {
13 
14 // Definition for ODR-use of static constexpr member (required until C++17)
15 constexpr U16 ComAggregator::FHP_UNSET;
17 
18 // ----------------------------------------------------------------------
19 // Component construction and destruction
20 // ----------------------------------------------------------------------
21 
22 ComAggregator ::ComAggregator(const char* const compName)
23  : ComAggregatorComponentBase(compName),
24  m_allocator(nullptr),
25  m_allocationId(static_cast<FwEnumStoreType>(-1)),
26  m_allocation(nullptr),
27  m_bufferState(Fw::Buffer::OwnershipState::OWNED),
28  m_frameBuffer(),
29  m_frameSerializer(),
30  m_allow_timeout(false),
31  m_spanning(false),
32  m_aggregationSize(0),
33  m_heldOffset(0),
34  m_fhp(FHP_UNSET),
35  m_pendingIdleCount(0),
36  m_leadingIdleCount(0),
37  m_lastFrameLost(false) {}
38 
40 
42  bool spanningEnabled,
43  FwEnumStoreType allocationId,
44  Fw::MemAllocator& allocator) {
45  // Storage is configured once, before any data is aggregated; cleanup() must run before reconfiguring
46  FW_ASSERT(this->m_allocation == nullptr);
47  // An aggregate must have room for data beyond the minimum idle packet
48  FW_ASSERT(aggregationSize > Ccsds::Utils::IdlePacket::MIN_SIZE, static_cast<FwAssertArgType>(aggregationSize));
49  if (spanningEnabled) {
50  // Every packet header offset in an aggregate must be representable as an 11-bit First Header Pointer
51  // and distinct from the reserved values (CCSDS 132.0-B-3 4.1.2.7.6)
52  const FwSizeType fhpRange = static_cast<FwSizeType>(Ccsds::TMSubfields::FHP_IDLE_DATA_ONLY);
53  FW_ASSERT(aggregationSize <= fhpRange, static_cast<FwAssertArgType>(aggregationSize));
54  } else {
55  // Without spanning, a full com buffer or file buffer Space Packet must fit next to a minimum idle packet
56  FW_ASSERT(aggregationSize >= MIN_NON_SPANNING_AGGREGATION_SIZE, static_cast<FwAssertArgType>(aggregationSize));
57  // The residual idle packet, up to aggregationSize - 1 bytes, must have a representable SPP length field
58  FW_ASSERT((aggregationSize - 1) <= Ccsds::Utils::IdlePacket::MAX_SIZE,
59  static_cast<FwAssertArgType>(aggregationSize));
60  }
61  this->m_spanning = spanningEnabled;
62  this->m_aggregationSize = aggregationSize;
63 
64  this->m_allocator = &allocator;
65  this->m_allocationId = allocationId;
66  FwSizeType allocatedSize = aggregationSize;
67  this->m_allocation = allocator.checkedAllocate(allocationId, allocatedSize);
68  this->m_frameBuffer.set(static_cast<U8*>(this->m_allocation), aggregationSize);
69  this->m_frameSerializer.setExtBuffer(static_cast<U8*>(this->m_allocation), aggregationSize);
70 }
71 
73  if ((this->m_allocator != nullptr) && (this->m_allocation != nullptr)) {
74  // The aggregate must not be held downstream when its storage is released
75  FW_ASSERT(this->m_bufferState == Fw::Buffer::OwnershipState::OWNED,
76  static_cast<FwAssertArgType>(this->m_bufferState.load()));
77  // Drop the per-aggregate state so a later configure() starts clean
78  this->m_held = Svc::ComDataContextPair();
79  this->m_heldOffset = 0;
80  this->m_fhp = FHP_UNSET;
81  this->m_pendingIdleCount = 0;
82  this->m_leadingIdleCount = 0;
83  this->m_lastFrameLost = false;
84  this->m_frameSerializer.setExtBuffer(nullptr, 0);
85  this->m_frameBuffer.set(nullptr, 0);
86  this->m_aggregationSize = 0;
87  this->m_allocator->deallocate(this->m_allocationId, this->m_allocation);
88  this->m_allocation = nullptr;
89  this->m_allocator = nullptr;
90  }
91 }
92 
95  this->comStatusOut_out(0, good);
96 }
97 
98 // ----------------------------------------------------------------------
99 // Handler implementations for typed input ports
100 // ----------------------------------------------------------------------
101 
102 void ComAggregator ::comStatusIn_handler(FwIndexType portNum, Fw::Success& condition) {
103  this->aggregationMachine_sendSignal_status(condition);
104 }
105 
106 void ComAggregator ::dataIn_handler(FwIndexType portNum, Fw::Buffer& data, const ComCfg::FrameContext& context) {
107  FW_ASSERT(this->m_allocation != nullptr);
108  // Without spanning, any packet must fit in an empty aggregate next to a minimum idle packet
109  FW_ASSERT(this->m_spanning || data.getSize() <= (this->m_aggregationSize - Ccsds::Utils::IdlePacket::MIN_SIZE),
110  static_cast<FwAssertArgType>(data.getSize()));
111  Svc::ComDataContextPair pair(data, context);
113 }
114 
115 void ComAggregator ::dataReturnIn_handler(FwIndexType portNum, Fw::Buffer& data, const ComCfg::FrameContext& context) {
116  // This handler runs on the returning caller's thread: take ownership atomically
117  const Fw::Buffer::OwnershipState previousState = this->m_bufferState.exchange(Fw::Buffer::OwnershipState::OWNED);
118  FW_ASSERT(previousState == Fw::Buffer::OwnershipState::NOT_OWNED, static_cast<FwAssertArgType>(previousState));
119 }
120 
121 void ComAggregator ::timeout_handler(FwIndexType portNum, U32 context) {
122  // Timeout is ignored in WAIT_STATUS state. However, the queue may not process timeout messages until the wait
123  // status is returned because the port chain may be synchronous and downstream components (radio, retry, etc) may
124  // take a long time to complete the transmission of data. This can cause the queue to overflow with messages that
125  // will soon be discarded.
126  //
127  // Therefore, to fix the risk of queue overflow we only queue timeout messages when they would be processed by the
128  // state machine (i.e. in the FILL state). Otherwise, these messages are not queued.
129  //
130  // Behaviorally, this solution will work exactly like the naive implementation with an infinite queue depth, but
131  // prevents queue overflow when using finite queues.
132  if (this->m_allow_timeout) {
134  }
135 }
136 
137 // ----------------------------------------------------------------------
138 // Implementations for internal state machine actions
139 // ----------------------------------------------------------------------
140 
141 void ComAggregator ::Svc_AggregationMachine_action_doClear(SmId smId, Svc_AggregationMachine::Signal signal) {
142  this->m_allow_timeout = true; // Allow timeout messages in FILL state
143  this->m_frameSerializer.resetSer();
144  this->m_frameBuffer.setSize(this->m_aggregationSize);
145  this->m_lastContext = ComCfg::FrameContext();
146  this->m_fhp = FHP_UNSET;
147  this->dropLostFrameState();
148  this->m_leadingIdleCount = this->m_pendingIdleCount;
149  // Write out any idle packet bytes spanning over from the previous aggregate
150  if (this->m_pendingIdleCount > 0) {
151  Fw::SerializeStatus status = this->m_frameSerializer.serializeFrom(
152  this->m_pendingIdle, this->m_pendingIdleCount, Fw::Serialization::OMIT_LENGTH);
154  this->m_pendingIdleCount = 0;
155  }
156  // Fill from the held buffer (whole packet, or remainder of a spanned packet)
157  this->fillFromHeld();
158 }
159 
160 void ComAggregator ::Svc_AggregationMachine_action_doFill(SmId smId,
162  const Svc::ComDataContextPair& value) {
163  this->markFirstHeaderIfUnset();
164  Fw::SerializeStatus status = this->m_frameSerializer.serializeFrom(
167  this->m_lastContext = value.get_context();
168  this->returnAndSignalReady(value);
169 }
170 
171 void ComAggregator ::Svc_AggregationMachine_action_doSend(SmId smId, Svc_AggregationMachine::Signal signal) {
172  // Send only when the buffer will be valid
173  if (this->m_frameSerializer.getSize() > 0) {
174  this->fillResidualWithIdle();
175  FW_ASSERT(this->m_frameSerializer.getSize() == this->m_aggregationSize,
176  static_cast<FwAssertArgType>(this->m_frameSerializer.getSize()));
177  if (this->m_spanning) {
178  this->m_lastContext.set_firstHeaderPointer(
179  (this->m_fhp == FHP_UNSET) ? static_cast<U16>(Ccsds::TMSubfields::FHP_NO_PACKET_START) : this->m_fhp);
180  } else {
181  // Packets are never split: the first packet header is always at the start of the aggregate
182  this->m_lastContext.set_firstHeaderPointer(0);
183  }
184  const Fw::Buffer::OwnershipState previousState =
185  this->m_bufferState.exchange(Fw::Buffer::OwnershipState::NOT_OWNED);
186  FW_ASSERT(previousState == Fw::Buffer::OwnershipState::OWNED, static_cast<FwAssertArgType>(previousState));
187  this->m_allow_timeout = false; // Timeout messages should be discarded in WAIT_STATUS state
188  this->dataOut_out(0, this->m_frameBuffer, this->m_lastContext);
189  }
190 }
191 
192 void ComAggregator ::Svc_AggregationMachine_action_doHold(SmId smId,
194  const Svc::ComDataContextPair& value) {
195  FW_ASSERT(not this->m_held.get_data().isValid());
196  this->m_held = value;
197  this->m_heldOffset = 0;
198 }
199 
200 void ComAggregator ::Svc_AggregationMachine_action_doSplitHold(SmId smId,
202  const Svc::ComDataContextPair& value) {
203  this->Svc_AggregationMachine_action_doHold(smId, signal, value);
204  if (this->m_spanning) {
205  // Split the leading bytes of the held packet into the remaining aggregation space
206  const FwSizeType remaining = this->remainingCapacity();
207  if (remaining > 0) {
208  // The held packet's header starts at the current fill offset of this aggregate
209  this->markFirstHeaderIfUnset();
210  Fw::SerializeStatus status = this->m_frameSerializer.serializeFrom(value.get_data().getData(), remaining,
213  this->m_heldOffset = remaining;
214  this->m_lastContext = value.get_context();
215  }
216  }
217 }
218 
219 void ComAggregator ::Svc_AggregationMachine_action_doNoteFailure(SmId smId, Svc_AggregationMachine::Signal signal) {
220  this->m_lastFrameLost = true;
221 }
222 
223 void ComAggregator ::Svc_AggregationMachine_action_assertNoStatus(SmId smId, Svc_AggregationMachine::Signal signal) {
224  // Status is not possible in this state, confirm by assertion
225  FW_ASSERT(false);
226 }
227 
228 // ----------------------------------------------------------------------
229 // Implementations for internal state machine guards
230 // ----------------------------------------------------------------------
231 
232 bool ComAggregator ::Svc_AggregationMachine_guard_isFull(SmId smId,
234  const Svc::ComDataContextPair& value) const {
235  return not this->accepts(value.get_data().getSize());
236 }
237 
238 bool ComAggregator ::Svc_AggregationMachine_guard_willFill(SmId smId,
240  const Svc::ComDataContextPair& value) const {
241  return (this->remainingCapacity() == value.get_data().getSize());
242 }
243 
244 bool ComAggregator ::Svc_AggregationMachine_guard_isNotEmpty(SmId smId, Svc_AggregationMachine::Signal signal) const {
245  // Carried-over idle bytes are not payload: an aggregate holding only those is empty for timeout purposes
246  return this->m_frameSerializer.getSize() > this->m_leadingIdleCount;
247 }
248 
249 bool ComAggregator ::Svc_AggregationMachine_guard_isGood(SmId smId,
251  const Fw::Success& value) const {
252  return value == Fw::Success::SUCCESS;
253 }
254 
255 bool ComAggregator ::Svc_AggregationMachine_guard_isSpanFull(SmId smId, Svc_AggregationMachine::Signal signal) const {
256  return this->m_spanning && (this->m_frameSerializer.getSize() == this->m_aggregationSize);
257 }
258 
259 // ----------------------------------------------------------------------
260 // Helper functions
261 // ----------------------------------------------------------------------
262 
263 FwSizeType ComAggregator ::remainingCapacity() const {
264  FW_ASSERT(this->m_frameSerializer.getSize() <= this->m_aggregationSize,
265  static_cast<FwAssertArgType>(this->m_frameSerializer.getSize()));
266  return this->m_aggregationSize - this->m_frameSerializer.getSize();
267 }
268 
269 bool ComAggregator ::accepts(FwSizeType size) const {
270  const FwSizeType remaining = this->remainingCapacity();
271  if (size > remaining) {
272  return false;
273  }
274  // Without spanning, an idle packet cannot continue into the next aggregate: the packet must either complete
275  // the aggregate or leave room for a whole minimum idle packet
276  const FwSizeType residual = remaining - size;
277  return this->m_spanning || (residual == 0) || (residual >= Ccsds::Utils::IdlePacket::MIN_SIZE);
278 }
279 
280 void ComAggregator ::markFirstHeaderIfUnset() {
281  if (this->m_spanning && this->m_fhp == FHP_UNSET) {
282  this->m_fhp = static_cast<U16>(this->m_frameSerializer.getSize());
283  }
284 }
285 
286 void ComAggregator ::fillFromHeld() {
287  if (this->m_held.get_data().isValid()) {
288  const Fw::Buffer& held = this->m_held.get_data();
289  const FwSizeType heldRemaining = held.getSize() - this->m_heldOffset;
290  const FwSizeType fillSize = FW_MIN(this->remainingCapacity(), heldRemaining);
291  if (this->m_heldOffset == 0) {
292  // The held packet's header starts at the current fill offset of this aggregate
293  this->markFirstHeaderIfUnset();
294  }
295  Fw::SerializeStatus status = this->m_frameSerializer.serializeFrom(held.getData() + this->m_heldOffset,
298  this->m_lastContext = this->m_held.get_context();
299  this->m_heldOffset += fillSize;
300  if (this->m_heldOffset == held.getSize()) {
301  // Held buffer fully consumed: return it and request more data
302  this->returnAndSignalReady(this->m_held);
303  this->m_held = Svc::ComDataContextPair();
304  this->m_heldOffset = 0;
305  }
306  }
307 }
308 
309 void ComAggregator ::fillResidualWithIdle() {
310  const FwSizeType residual = this->remainingCapacity();
311  if (residual == 0) {
312  return;
313  }
314  // The idle packet's header starts at the current fill offset of this aggregate
315  this->markFirstHeaderIfUnset();
316  // Idle packet size: fill the residual space exactly, spanning a minimum-size idle packet
317  // into the next aggregate when the residual space is too small (CCSDS 132.0-B-3 4.1.4)
319  if (residual >= Ccsds::Utils::IdlePacket::MIN_SIZE) {
320  // Idle packet fits entirely within this aggregate
321  status = Ccsds::Utils::IdlePacket::serialize(this->m_frameSerializer, residual);
323  } else {
324  // Without spanning, accepts() keeps the residual at zero or at least a minimum idle packet
325  FW_ASSERT(this->m_spanning);
326  // Stage a minimum-size idle packet, emit the leading bytes now and span the rest
328  Fw::ExternalSerializeBuffer stager(staging, sizeof(staging));
331  status = this->m_frameSerializer.serializeFrom(staging, residual, Fw::Serialization::OMIT_LENGTH);
333  this->m_pendingIdleCount = Ccsds::Utils::IdlePacket::MIN_SIZE - residual;
334  (void)memcpy(this->m_pendingIdle, &staging[residual], this->m_pendingIdleCount);
335  }
336 }
337 
338 void ComAggregator ::dropLostFrameState() {
339  if (this->m_lastFrameLost) {
340  this->m_pendingIdleCount = 0;
341  if (this->m_held.get_data().isValid() && this->m_heldOffset > 0) {
342  this->returnAndSignalReady(this->m_held);
343  this->m_held = Svc::ComDataContextPair();
344  this->m_heldOffset = 0;
345  }
346  this->m_lastFrameLost = false;
347  }
348 }
349 
350 void ComAggregator ::returnAndSignalReady(const Svc::ComDataContextPair& pair) {
351  // Return port does not alter data and thus const-cast is safe
352  this->dataReturnOut_out(0, const_cast<Fw::Buffer&>(pair.get_data()), pair.get_context());
354  this->comStatusOut_out(0, good);
355 }
356 
357 } // namespace Svc
Serialization/Deserialization operation was successful.
void set(U8 *data, FwSizeType size, U32 context=NO_CONTEXT)
Definition: Buffer.cpp:144
Representing success.
PlatformSizeType FwSizeType
I32 FwEnumStoreType
void setSize(FwSizeType size)
Definition: Buffer.cpp:131
Serializable::SizeType getSize() const override
Get current buffer size.
Fw::SerializeStatus serialize(Fw::SerialBufferBase &serializer, FwSizeType size)
Serialize an idle packet of exactly size bytes (header and idle data) into serializer ...
Definition: IdlePacket.cpp:16
void dataReturnOut_out(FwIndexType portNum, Fw::Buffer &data, const ComCfg::FrameContext &context) const
Invoke output port dataReturnOut.
void aggregationMachine_sendSignal_status(const Fw::Success &value)
Send signal status to state machine aggregationMachine.
U8 * getData() const
Definition: Buffer.cpp:82
static constexpr FwSizeType MIN_NON_SPANNING_AGGREGATION_SIZE
Smallest aggregate that holds a full com buffer and a full file buffer Space Packet without spanning...
SerializeStatus serializeFrom(U8 val, Endianness mode=Endianness::BIG) override
Serialize an 8-bit unsigned integer value.
The buffer is currently not owned.
SerializeStatus
forward declaration for string
~ComAggregator()
Destroy ComAggregator object.
void preamble() override
A function that will be called before the event loop is entered.
#define FW_MIN(a, b)
MIN macro (deprecated in C++, use std::min)
Definition: BasicTypes.h:99
Omit length from serialization.
bool isValid() const
Definition: Buffer.cpp:78
External serialize buffer with no copy semantics.
void * checkedAllocate(const FwEnumStoreType identifier, FwSizeType &size, bool &recoverable, FwSizeType alignment=alignof(std::max_align_t))
void resetSer() override
Reset serialization pointer to beginning of buffer.
ComCfg::FrameContext & get_context()
Get member context.
constexpr FwSizeType MIN_SIZE
Minimum idle packet size: header plus one byte of idle data.
Definition: IdlePacket.hpp:24
uint8_t U8
8-bit unsigned integer
Definition: BasicTypes.h:54
FwSizeType getSize() const
Definition: Buffer.cpp:90
Auto-generated base for ComAggregator component.
void aggregationMachine_sendSignal_timeout()
Send signal timeout to state machine aggregationMachine.
Fw::Buffer & get_data()
Get member data.
void setExtBuffer(U8 *buffPtr, Serializable::SizeType size)
Set the external buffer.
Memory Allocation base class.
ComAggregator(const char *const compName)
Construct ComAggregator object.
PlatformIndexType FwIndexType
void comStatusOut_out(FwIndexType portNum, Fw::Success &condition) const
Invoke output port comStatusOut.
OwnershipState
Definition: Buffer.hpp:60
Type used to pass context info between components during framing/deframing.
RateGroupDivider component implementation.
constexpr FwSizeType MAX_SIZE
Maximum idle packet size: header plus the largest data length the SPP length field can express...
Definition: IdlePacket.hpp:27
virtual void deallocate(const FwEnumStoreType identifier, void *ptr)=0
The buffer is currently owned.
Implementation of malloc based allocator.
void set_firstHeaderPointer(U16 firstHeaderPointer)
Set member firstHeaderPointer.
void dataOut_out(FwIndexType portNum, Fw::Buffer &data, const ComCfg::FrameContext &context) const
Invoke output port dataOut.
void configure(FwSizeType aggregationSize, bool spanningEnabled, FwEnumStoreType allocationId, Fw::MemAllocator &allocator)
FpySequencer_SequencerStateMachineStateMachineBase::Signal Signal
#define FW_ASSERT(...)
Definition: Assert.hpp:14
void aggregationMachine_sendSignal_fill(const Svc::ComDataContextPair &value)
Send signal fill to state machine aggregationMachine.
Success/Failure.
PlatformAssertArgType FwAssertArgType
The type of arguments to assert functions.