11#include <promeki/config.h>
12#if PROMEKI_ENABLE_NETWORK
34PROMEKI_NAMESPACE_BEGIN
38class RtpSeqReorderBuffer;
93class RtpSession :
public ObjectBase {
94 PROMEKI_OBJECT(RtpSession, ObjectBase)
100 RtpSession(ObjectBase *parent =
nullptr);
103 ~RtpSession()
override;
116 Error start(
const SocketAddress &localAddr);
129 Error start(PacketTransport *transport);
164 Error start(PacketTransport *primary, PacketTransport *secondary);
170 bool isRunning()
const {
return _running; }
182 void setRemote(
const SocketAddress &dest) { _remote = dest; }
185 const SocketAddress &remote()
const {
return _remote; }
195 void setRemoteSecondary(
const SocketAddress &dest) { _remoteSecondary = dest; }
198 const SocketAddress &remoteSecondary()
const {
return _remoteSecondary; }
201 bool hasSecondaryLeg()
const {
return _transportSecondary !=
nullptr; }
212 Error sendPacket(
const Buffer &payload, uint32_t timestamp, uint8_t payloadType,
bool marker =
false);
250 Error sendPackets(RtpPacketBatch &batch);
274 struct StreamReceiver {
279 RtpPacket::Queue *outQueue =
nullptr;
283 RtpSeqTracker *seqTracker =
nullptr;
288 RtpSeqReorderBuffer *reorderBuffer =
nullptr;
300 uint32_t clockRateHz = 0;
312 uint8_t payloadType = 0;
361 Error startReceiving(List<StreamReceiver> receivers,
362 const String &threadName =
"rtp-rx");
372 void stopReceiving();
375 bool isReceiving()
const {
return _receiving.value(); }
387 void setReceivePollIntervalMs(
unsigned int timeoutMs) {
388 _receivePollMs = timeoutMs == 0 ? 200 : timeoutMs;
392 unsigned int receivePollIntervalMs()
const {
return _receivePollMs; }
409 Error setPacingRate(uint64_t bytesPerSec);
426 void setScheduler(PacketScheduler::UPtr scheduler);
429 PacketScheduler *scheduler()
const {
return _scheduler.get(); }
444 void setSchedulerSecondary(PacketScheduler::UPtr scheduler);
447 PacketScheduler *schedulerSecondary()
const {
return _schedulerSecondary.get(); }
462 Error configureScheduler(
const PacketScheduler::Spec &spec);
465 uint32_t ssrc()
const {
return _ssrc; }
468 void setSsrc(uint32_t ssrc) { _ssrc = ssrc; }
471 uint16_t sequenceNumber()
const {
return _sequenceNumber; }
474 void setPayloadType(uint8_t pt) { _payloadType = pt; }
477 uint8_t payloadType()
const {
return _payloadType; }
480 void setClockRate(uint32_t hz) { _clockRate = hz; }
483 uint32_t clockRate()
const {
return _clockRate; }
522 void setRtpAnchor(NtpTime captureNtp, uint32_t rtpTs);
555 void setRtpAnchor(
const ClockDomain &domain, uint32_t rtpTs);
565 NtpTime anchorNtp()
const;
571 uint32_t anchorRtpTs()
const;
604 void noteRtpEmission(uint32_t rtpTs);
615 bool hasEmissionRecord()
const;
618 const String &cname()
const {
return _cname; }
621 void setCname(
const String &cname) { _cname = cname; }
644 Error emitRtcpSr(uint32_t senderPacketCount, uint32_t senderOctetCount);
664 Error emitRtcpRr(
const RtcpPacket::ReportBlock &block);
710 ReceivedSr receivedSr()
const;
721 uint32_t srObservedCount()
const;
732 TimeStamp firstSrAt()
const;
749 NtpTime currentSrNtp()
const;
762 PacketTransport *transport()
const {
return _transport; }
772 PROMEKI_SIGNAL(packetReceived, Buffer, uint32_t, uint8_t,
bool);
775 PROMEKI_SIGNAL(ssrcCollision, uint32_t);
791 PROMEKI_SIGNAL(ssrcChange, uint32_t, uint32_t, uint8_t);
804 PROMEKI_SIGNAL(byeReceived, uint32_t);
808 friend class ReceiveThread;
828 void handleRtcp(
const uint8_t *data,
size_t size);
830 void fillHeader(RtpPacket &pkt, uint8_t pt,
bool marker, uint32_t timestamp);
842 void fillTransportHeader(RtpPacket &pkt);
845 PacketTransport *_transport =
nullptr;
846 PacketTransport::UPtr _ownedTransport;
847 PacketScheduler::UPtr _scheduler;
861 PacketTransport *_transportSecondary =
nullptr;
862 PacketTransport::UPtr _ownedTransportSecondary;
863 PacketScheduler::UPtr _schedulerSecondary;
864 SocketAddress _remoteSecondary;
866 bool _running =
false;
867 SocketAddress _remote;
869 uint16_t _sequenceNumber = 0;
870 uint8_t _payloadType = 96;
871 uint32_t _clockRate = 90000;
887 mutable Mutex _rtcpMutex;
889 uint32_t _anchorRtpTs = 0;
890 uint32_t _lastEmissionRtpTs = 0;
891 bool _hasEmission =
false;
897 uint32_t _srObservedCount = 0;
904 TimeStamp _firstSrAt;
914 ReceivedSr _lastReceivedSr;
921 using ReceiveThreadUPtr = UniquePtr<ReceiveThread>;
922 ReceiveThreadUPtr _receiveThread;
923 ReceiveThreadUPtr _receiveThreadSecondary;
924 Mutex _dispatchMutex;
925 List<StreamReceiver> _streamReceivers;
926 Atomic<bool> _receiving;
927 unsigned int _receivePollMs = 200;
933 struct SsrcPinState {
934 uint32_t expectedSsrc = 0;
936 uint32_t mismatchCount = 0;
937 TimeStamp mismatchFirstTime;
939 List<SsrcPinState> _ssrcPinStates;