| 1232 | } |
| 1233 | |
| 1234 | void recvThread(broadcast_ctx_t &ctx) { |
| 1235 | std::map<av_session_id_t, message_queue_t> peer_to_video_session; |
| 1236 | std::map<av_session_id_t, message_queue_t> peer_to_audio_session; |
| 1237 | |
| 1238 | auto &video_sock = ctx.video_sock; |
| 1239 | auto &audio_sock = ctx.audio_sock; |
| 1240 | |
| 1241 | auto &message_queue_queue = ctx.message_queue_queue; |
| 1242 | auto broadcast_shutdown_event = mail::man->event<bool>(mail::broadcast_shutdown); |
| 1243 | |
| 1244 | auto &io = ctx.io_context; |
| 1245 | |
| 1246 | udp::endpoint peer; |
| 1247 | |
| 1248 | std::array<char, 2048> buf[2]; |
| 1249 | std::function<void(const boost::system::error_code, size_t)> recv_func[2]; |
| 1250 | |
| 1251 | auto populate_peer_to_session = [&]() { |
| 1252 | while (message_queue_queue->peek()) { |
| 1253 | auto message_queue_opt = message_queue_queue->pop(); |
| 1254 | TUPLE_3D_REF(socket_type, session_id, message_queue, *message_queue_opt); |
| 1255 | |
| 1256 | switch (socket_type) { |
| 1257 | case socket_e::video: |
| 1258 | if (message_queue) { |
| 1259 | peer_to_video_session.emplace(session_id, message_queue); |
| 1260 | } else { |
| 1261 | peer_to_video_session.erase(session_id); |
| 1262 | } |
| 1263 | break; |
| 1264 | case socket_e::audio: |
| 1265 | if (message_queue) { |
| 1266 | peer_to_audio_session.emplace(session_id, message_queue); |
| 1267 | } else { |
| 1268 | peer_to_audio_session.erase(session_id); |
| 1269 | } |
| 1270 | break; |
| 1271 | } |
| 1272 | } |
| 1273 | }; |
| 1274 | |
| 1275 | auto recv_func_init = [&](udp::socket &sock, int buf_elem, std::map<av_session_id_t, message_queue_t> &peer_to_session) { |
| 1276 | recv_func[buf_elem] = [&, buf_elem](const boost::system::error_code &ec, size_t bytes) { |
| 1277 | auto fg = util::fail_guard([&]() { |
| 1278 | sock.async_receive_from(asio::buffer(buf[buf_elem]), peer, 0, recv_func[buf_elem]); |
| 1279 | }); |
| 1280 | |
| 1281 | auto type_str = buf_elem ? "AUDIO"sv : "VIDEO"sv; |
| 1282 | BOOST_LOG(verbose) << "Recv: "sv << peer.address().to_string() << ':' << peer.port() << " :: " << type_str; |
| 1283 | |
| 1284 | populate_peer_to_session(); |
| 1285 | |
| 1286 | // No data, yet no error |
| 1287 | if (ec == boost::system::errc::connection_refused || ec == boost::system::errc::connection_reset) { |
| 1288 | return; |
| 1289 | } |
| 1290 | |
| 1291 | if (ec || !bytes) { |