File: stream_session.cpp

package info (click to toggle)
snapcast 0.31.0-1
  • links: PTS, VCS
  • area: main
  • in suites: forky, sid, trixie
  • size: 4,012 kB
  • sloc: cpp: 37,729; python: 2,543; sh: 455; makefile: 16
file content (122 lines) | stat: -rw-r--r-- 3,442 bytes parent folder | download | duplicates (2)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
/***
    This file is part of snapcast
    Copyright (C) 2014-2025  Johannes Pohl

    This program is free software: you can redistribute it and/or modify
    it under the terms of the GNU General Public License as published by
    the Free Software Foundation, either version 3 of the License, or
    (at your option) any later version.

    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 the
    GNU General Public License for more details.

    You should have received a copy of the GNU General Public License
    along with this program.  If not, see <http://www.gnu.org/licenses/>.
***/

// prototype/interface header file
#include "stream_session.hpp"

// local headers
#include "common/aixlog.hpp"

// 3rd party headers

// standard headers
#include <iostream>


using namespace std;
using namespace streamreader;

static constexpr auto LOG_TAG = "StreamSession";


StreamSession::StreamSession(const boost::asio::any_io_executor& executor, StreamMessageReceiver* receiver)
    : messageReceiver_(receiver), pcmStream_(nullptr), strand_(boost::asio::make_strand(executor))
{
    base_msg_size_ = baseMessage_.getSize();
    buffer_.resize(base_msg_size_);
}


void StreamSession::setPcmStream(PcmStreamPtr pcmStream)
{
    std::lock_guard<std::mutex> lock(mutex_);
    pcmStream_ = std::move(pcmStream);
}


const PcmStreamPtr StreamSession::pcmStream() const
{
    std::lock_guard<std::mutex> lock(mutex_);
    return pcmStream_;
}


void StreamSession::send_next()
{
    auto& buffer = messages_.front();
    buffer.on_air = true;
    boost::asio::post(strand_, [this, self = shared_from_this(), buffer]()
    {
        sendAsync(buffer, [this](boost::system::error_code ec, std::size_t length)
        {
            messages_.pop_front();
            if (ec)
            {
                LOG(ERROR, LOG_TAG) << "StreamSession write error (msg length: " << length << "): " << ec.message() << "\n";
                messageReceiver_->onDisconnect(this);
                return;
            }
            if (!messages_.empty())
                send_next();
        });
    });
}


void StreamSession::send(shared_const_buffer const_buf)
{
    boost::asio::post(strand_, [this, self = shared_from_this(), const_buf]()
    {
        // delete PCM chunks that are older than the overall buffer duration
        messages_.erase(std::remove_if(messages_.begin(), messages_.end(),
                                       [this](const shared_const_buffer& buffer)
        {
            const auto& msg = buffer.message();
            if (!msg.is_pcm_chunk || buffer.on_air)
                return false;
            auto age = chronos::clk::now() - msg.rec_time;
            return (age > std::chrono::milliseconds(bufferMs_) + 100ms);
        }),
                        messages_.end());

        messages_.push_back(const_buf);

        if (messages_.size() > 1)
        {
            LOG(TRACE, LOG_TAG) << "outstanding async_write\n";
            return;
        }
        send_next();
    });
}


void StreamSession::send(msg::message_ptr message)
{
    if (!message)
        return;

    // TODO: better set the timestamp in send_next for more accurate time sync
    send(shared_const_buffer(*message));
}


void StreamSession::setBufferMs(size_t bufferMs)
{
    bufferMs_ = bufferMs;
}