Writing top-level logic - publisher

This commit is contained in:
Pavel Kirienko
2014-03-08 23:01:05 +04:00
parent 35db1858c8
commit 77184fc062
3 changed files with 295 additions and 0 deletions
@@ -0,0 +1,68 @@
/*
* Copyright (C) 2014 Pavel Kirienko <pavel.kirienko@gmail.com>
*/
#pragma once
#include <uavcan/internal/transport/transfer.hpp>
#include <uavcan/internal/transport/transfer_buffer.hpp>
namespace uavcan
{
class IMarshalBuffer : public ITransferBuffer
{
public:
virtual const uint8_t* getDataPtr() const = 0;
virtual unsigned int getDataLength() const = 0;
};
class IMarshalBufferProvider
{
public:
virtual ~IMarshalBufferProvider() { }
virtual IMarshalBuffer* getBuffer(unsigned int size) = 0;
};
template <unsigned int MaxSize_ = MaxTransferPayloadLen>
class MarshalBufferProvider : public IMarshalBufferProvider
{
class Buffer : public IMarshalBuffer
{
StaticTransferBuffer<MaxSize_> buf_;
int read(unsigned int offset, uint8_t* data, unsigned int len) const
{
return buf_.read(offset, data, len);
}
int write(unsigned int offset, const uint8_t* data, unsigned int len)
{
return buf_.write(offset, data, len);
}
const uint8_t* getDataPtr() const { return buf_.getRawPtr(); }
unsigned int getDataLength() const { return buf_.getMaxWritePos(); }
public:
void reset() { buf_.reset(); }
};
Buffer buffer_;
public:
enum { MaxSize = MaxSize_ };
IMarshalBuffer* getBuffer(unsigned int size)
{
if (size > MaxSize)
return NULL;
buffer_.reset();
return &buffer_;
}
};
}
+118
View File
@@ -0,0 +1,118 @@
/*
* Copyright (C) 2014 Pavel Kirienko <pavel.kirienko@gmail.com>
*/
#pragma once
#include <uavcan/scheduler.hpp>
#include <uavcan/data_type.hpp>
#include <uavcan/marshal_buffer.hpp>
#include <uavcan/global_data_type_registry.hpp>
#include <uavcan/internal/debug.hpp>
#include <uavcan/internal/lazy_constructor.hpp>
#include <uavcan/internal/transport/transfer_sender.hpp>
#include <uavcan/internal/marshal/scalar_codec.hpp>
namespace uavcan
{
template <typename DataType_>
class Publisher
{
public:
typedef DataType_ DataType;
private:
enum { MinTxTimeoutUsec = 200 };
const uint64_t max_transfer_interval_; // TODO: memory usage can be reduced
uint64_t tx_timeout_;
Scheduler& scheduler_;
IMarshalBufferProvider& buffer_provider_;
LazyConstructor<TransferSender> sender_;
bool checkInit()
{
if (sender_)
return true;
GlobalDataTypeRegistry::instance().freeze();
const DataTypeDescriptor* const descr =
GlobalDataTypeRegistry::instance().find(DataTypeKindMessage, DataType::getDataTypeFullName());
if (!descr)
{
UAVCAN_TRACE("Publisher", "Type [%s] is not registered", DataType::getDataTypeFullName());
return false;
}
sender_.construct<Dispatcher&, const DataTypeDescriptor&, CanTxQueue::Qos, uint64_t>
(scheduler_.getDispatcher(), *descr, CanTxQueue::Volatile, max_transfer_interval_);
return true;
}
uint64_t getTxDeadline() const { return scheduler_.getMonotonicTimestamp() + tx_timeout_; }
IMarshalBuffer* getBuffer()
{
const int size = (DataType::MaxBitLen + 7) / 8;
return buffer_provider_.getBuffer(size);
}
int genericSend(const DataType& message, TransferType transfer_type, NodeID dst_node_id,
uint64_t monotonic_blocking_deadline)
{
if (!checkInit())
return -1;
IMarshalBuffer* const buf = getBuffer();
if (!buf)
return -1;
BitStream bitstream(*buf);
ScalarCodec codec(bitstream);
const int encode_res = DataType::encode(message, codec);
if (encode_res <= 0)
{
assert(0); // Impossible, internal error
return -1;
}
return sender_->send(buf->getDataPtr(), buf->getDataLength(), getTxDeadline(),
monotonic_blocking_deadline, transfer_type, dst_node_id);
}
public:
Publisher(Scheduler& scheduler, IMarshalBufferProvider& buffer_provider, uint64_t tx_timeout_usec,
uint64_t max_transfer_interval = TransferSender::DefaultMaxTransferInterval)
: max_transfer_interval_(max_transfer_interval)
, tx_timeout_(tx_timeout_usec)
, scheduler_(scheduler)
, buffer_provider_(buffer_provider)
{
setTxTimeout(tx_timeout_usec);
StaticAssert<DataTypeKind(DataType::DataTypeKind) == DataTypeKindMessage>::check();
}
int broadcast(const DataType& message, uint64_t monotonic_blocking_deadline = 0)
{
return genericSend(message, TransferTypeMessageBroadcast, NodeID::Broadcast, monotonic_blocking_deadline);
}
int unicast(const DataType& message, NodeID dst_node_id, uint64_t monotonic_blocking_deadline = 0)
{
if (!dst_node_id.isUnicast())
{
assert(0);
return -1;
}
return genericSend(message, TransferTypeMessageUnicast, dst_node_id, monotonic_blocking_deadline);
}
uint64_t getTxTimeout() const { return tx_timeout_; }
void setTxTimeout(uint64_t usec)
{
tx_timeout_ = std::max(usec, uint64_t(MinTxTimeoutUsec));
}
};
}
+109
View File
@@ -0,0 +1,109 @@
/*
* Copyright (C) 2014 Pavel Kirienko <pavel.kirienko@gmail.com>
*/
#include <gtest/gtest.h>
#include <uavcan/publisher.hpp>
#include <uavcan/mavlink/Message.hpp>
#include "common.hpp"
#include "transport/can/iface_mock.hpp"
TEST(Publisher, Basic)
{
uavcan::PoolAllocator<uavcan::MemPoolBlockSize * 8, uavcan::MemPoolBlockSize> pool;
uavcan::PoolManager<1> poolmgr;
poolmgr.addPool(&pool);
SystemClockMock clock_mock(100);
CanDriverMock can_driver(2, clock_mock);
uavcan::OutgoingTransferRegistry<8> out_trans_reg(poolmgr);
uavcan::Scheduler sch(can_driver, poolmgr, clock_mock, out_trans_reg, uavcan::NodeID(1));
uavcan::MarshalBufferProvider<> buffer_provider;
uavcan::Publisher<uavcan::mavlink::Message> publisher(sch, buffer_provider, 10000);
ASSERT_FALSE(uavcan::GlobalDataTypeRegistry::instance().isFrozen());
/*
* Message layout:
* uint8 seq
* uint8 sysid
* uint8 compid
* uint8 msgid
* uint8[<256] payload
*/
uavcan::mavlink::Message msg;
msg.seq = 0x42;
msg.sysid = 0x72;
msg.compid = 0x08;
msg.msgid = 0xa5;
msg.payload = "Msg";
static const uint8_t expected_transfer_payload[] = {0x42, 0x72, 0x08, 0xa5, 'M', 's', 'g'};
/*
* Broadcast
*/
{
ASSERT_LT(0, publisher.broadcast(msg));
// uint_fast16_t data_type_id, TransferType transfer_type, NodeID src_node_id, NodeID dst_node_id,
// uint_fast8_t frame_index, TransferID transfer_id, bool last_frame = false
uavcan::Frame expected_frame(uavcan::mavlink::Message::DefaultDataTypeID, uavcan::TransferTypeMessageBroadcast,
sch.getDispatcher().getSelfNodeID(), uavcan::NodeID::Broadcast, 0, 0, true);
expected_frame.setPayload(expected_transfer_payload, 7);
uavcan::CanFrame expected_can_frame;
ASSERT_TRUE(expected_frame.compile(expected_can_frame));
ASSERT_TRUE(can_driver.ifaces[0].matchAndPopTx(expected_can_frame, 10000 + 100));
ASSERT_TRUE(can_driver.ifaces[1].matchAndPopTx(expected_can_frame, 10000 + 100));
ASSERT_TRUE(can_driver.ifaces[0].tx.empty());
ASSERT_TRUE(can_driver.ifaces[1].tx.empty());
// Second shot - checking the transfer ID
ASSERT_LT(0, publisher.broadcast(msg));
expected_frame = uavcan::Frame(uavcan::mavlink::Message::DefaultDataTypeID, uavcan::TransferTypeMessageBroadcast,
sch.getDispatcher().getSelfNodeID(), uavcan::NodeID::Broadcast, 0, 1, true);
expected_frame.setPayload(expected_transfer_payload, 7);
ASSERT_TRUE(expected_frame.compile(expected_can_frame));
ASSERT_TRUE(can_driver.ifaces[0].matchAndPopTx(expected_can_frame, 10000 + 100));
ASSERT_TRUE(can_driver.ifaces[1].matchAndPopTx(expected_can_frame, 10000 + 100));
ASSERT_TRUE(can_driver.ifaces[0].tx.empty());
ASSERT_TRUE(can_driver.ifaces[1].tx.empty());
}
clock_mock.advance(1000);
/*
* Unicast
*/
{
ASSERT_LT(0, publisher.unicast(msg, 0x44));
// uint_fast16_t data_type_id, TransferType transfer_type, NodeID src_node_id, NodeID dst_node_id,
// uint_fast8_t frame_index, TransferID transfer_id, bool last_frame = false
uavcan::Frame expected_frame(uavcan::mavlink::Message::DefaultDataTypeID, uavcan::TransferTypeMessageUnicast,
sch.getDispatcher().getSelfNodeID(), uavcan::NodeID(0x44), 0, 0, true);
expected_frame.setPayload(expected_transfer_payload, 7);
uavcan::CanFrame expected_can_frame;
ASSERT_TRUE(expected_frame.compile(expected_can_frame));
ASSERT_TRUE(can_driver.ifaces[0].matchAndPopTx(expected_can_frame, 10000 + 100 + 1000));
ASSERT_TRUE(can_driver.ifaces[1].matchAndPopTx(expected_can_frame, 10000 + 100 + 1000));
ASSERT_TRUE(can_driver.ifaces[0].tx.empty());
ASSERT_TRUE(can_driver.ifaces[1].tx.empty());
}
/*
* Misc
*/
ASSERT_TRUE(uavcan::GlobalDataTypeRegistry::instance().isFrozen());
}