Files
ustreamer/janus/src/plugin.c

1003 lines
29 KiB
C

/*****************************************************************************
# #
# uStreamer - Lightweight and fast MJPEG-HTTP streamer. #
# #
# Copyright (C) 2018-2024 Maxim Devaev <mdevaev@gmail.com> #
# #
# 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 <https://www.gnu.org/licenses/>. #
# #
*****************************************************************************/
#include <stdatomic.h>
#include <stdlib.h>
#include <unistd.h>
#include <fcntl.h>
#include <errno.h>
#include <pthread.h>
#include <jansson.h>
#include <janus/plugins/plugin.h>
#include <janus/rtp.h>
#include <janus/rtcp.h>
#include <alsa/asoundlib.h>
#include <linux/videodev2.h>
#include "uslibs/types.h"
#include "uslibs/const.h"
#include "uslibs/errors.h"
#include "uslibs/tools.h"
#include "uslibs/threading.h"
#include "uslibs/logging.h"
#include "uslibs/list.h"
#include "uslibs/ring.h"
#include "uslibs/memsink.h"
#include "uslibs/chip.h"
#include "uslibs/memsink.h"
#include "const.h"
#include "client.h"
#include "au.h"
#include "acap.h"
#include "rtp.h"
#include "rtpv.h"
#include "rtpa.h"
#include "sdp.h"
#include "config.h"
static const char *const _g_default_ice_url = "stun:stun.l.google.com:19302";
static us_config_s *_g_config = NULL;
static const useconds_t _g_watchers_polling = 100000;
static us_janus_client_s *_g_clients = NULL;
static janus_callbacks *_g_gw = NULL;
static us_ring_s *_g_video_ring = NULL;
static us_ring_s *_g_vplay_ring = NULL;
static us_rtpv_s *_g_rtpv = NULL;
static us_rtpa_s *_g_rtpa = NULL;
static pthread_t _g_video_rtp_tid;
static atomic_bool _g_video_rtp_tid_created = false;
static pthread_t _g_video_sink_tid;
static atomic_bool _g_video_sink_tid_created = false;
static pthread_t _g_vplay_tid;
static atomic_bool _g_vplay_tid_created = false;
static pthread_t _g_acap_tid;
static atomic_bool _g_acap_tid_created = false;
static pthread_t _g_aplay_tid;
static atomic_bool _g_aplay_tid_created = false;
static pthread_mutex_t _g_video_lock = PTHREAD_MUTEX_INITIALIZER;
static pthread_mutex_t _g_vplay_lock = PTHREAD_MUTEX_INITIALIZER;
static pthread_mutex_t _g_acap_lock = PTHREAD_MUTEX_INITIALIZER;
static pthread_mutex_t _g_aplay_lock = PTHREAD_MUTEX_INITIALIZER;
static atomic_bool _g_ready = false;
static atomic_bool _g_stop = false;
static atomic_bool _g_has_watchers = false;
static atomic_bool _g_has_listeners = false;
static atomic_bool _g_has_speakers = false;
static atomic_bool _g_key_required = false;
static struct {
bool active;
us_janus_client_s *client;
uint width;
uint height;
uint fps;
} _g_camera = {0}; // Access should be protected by _g_vplay_lock
#define _LOCK_VIDEO US_MUTEX_LOCK(_g_video_lock)
#define _UNLOCK_VIDEO US_MUTEX_UNLOCK(_g_video_lock)
#define _LOCK_VPLAY US_MUTEX_LOCK(_g_vplay_lock)
#define _UNLOCK_VPLAY US_MUTEX_UNLOCK(_g_vplay_lock)
#define _LOCK_ACAP US_MUTEX_LOCK(_g_acap_lock)
#define _UNLOCK_ACAP US_MUTEX_UNLOCK(_g_acap_lock)
#define _LOCK_APLAY US_MUTEX_LOCK(_g_aplay_lock)
#define _UNLOCK_APLAY US_MUTEX_UNLOCK(_g_aplay_lock)
#define _LOCK_ALL { _LOCK_VIDEO; _LOCK_VPLAY; _LOCK_ACAP; _LOCK_APLAY; }
#define _UNLOCK_ALL { _UNLOCK_APLAY; _UNLOCK_ACAP; _UNLOCK_VPLAY; _UNLOCK_VIDEO; }
#define _READY atomic_load(&_g_ready)
#define _STOP atomic_load(&_g_stop)
#define _HAS_WATCHERS atomic_load(&_g_has_watchers)
#define _HAS_LISTENERS atomic_load(&_g_has_listeners)
#define _HAS_SPEAKERS atomic_load(&_g_has_speakers)
#define _IF_DISABLED(...) { if (!_READY || _STOP) { __VA_ARGS__ } }
janus_plugin *create(void);
static void *_video_rtp_thread(void *arg) {
(void)arg;
US_THREAD_SETTLE("us_p_rtpv");
atomic_store(&_g_video_rtp_tid_created, true);
while (!_STOP) {
const int ri = us_ring_consumer_acquire(_g_video_ring, 0.1);
if (ri >= 0) {
const us_frame_s *const frame = _g_video_ring->items[ri];
_LOCK_VIDEO;
const bool zero_playout_delay = (frame->gop == 0);
us_rtpv_wrap(_g_rtpv, frame, zero_playout_delay);
_UNLOCK_VIDEO;
us_ring_consumer_release(_g_video_ring, ri);
}
}
return NULL;
}
static void *_video_sink_thread(void *arg) {
(void)arg;
US_THREAD_SETTLE("us_p_vcap");
atomic_store(&_g_video_sink_tid_created, true);
us_frame_s *tmp = us_frame_init();
int once = 0;
while (!_STOP) {
if (!_HAS_WATCHERS) {
US_ONCE({ US_LOG_INFO("No active watchers, memsink disconnected"); });
usleep(_g_watchers_polling);
continue;
}
us_memsink_s *sink = us_memsink_init_opened("vcap", _g_config->video_sink_name, false, 0, false, 0, 1);
if (sink == NULL) {
goto close_memsink;
}
once = 0;
US_LOG_INFO("Memsink opened; reading frames ...");
while (!_STOP && _HAS_WATCHERS) {
const us_memsink_wants_s w_put = {.key = atomic_load(&_g_key_required)};
const int got = us_memsink_client_get(sink, tmp, NULL, &w_put);
if (got == 0) {
const int ri = us_ring_producer_acquire(_g_video_ring, 0);
if (ri >= 0) {
us_frame_s *dest = _g_video_ring->items[ri];
us_frame_copy(tmp, dest);
us_ring_producer_release(_g_video_ring, ri);
if (tmp->key) {
atomic_store(&_g_key_required, false);
}
} else {
US_ONCE({ US_LOG_ERROR("Video ring is full"); });
}
} else if (got == US_ERROR_NO_DATA) {
usleep(1000);
} else {
goto close_memsink;
}
}
close_memsink:
US_DELETE(sink, us_memsink_destroy);
US_LOG_INFO("Memsink closed");
sleep(1); // error_delay
}
us_frame_destroy(tmp);
return NULL;
}
static void *_acap_thread(void *arg) {
(void)arg;
US_THREAD_SETTLE("us_p_acap");
atomic_store(&_g_acap_tid_created, true);
US_A(us_str_is_ok(_g_config->acap_dev_name));
US_A(_g_rtpa != NULL);
int once = 0;
while (!_STOP) {
if (!_HAS_WATCHERS || !_HAS_LISTENERS) {
usleep(_g_watchers_polling);
continue;
}
int chip_fd = -1;
us_acap_s *acap = NULL;
if (!us_au_probe(_g_config->acap_dev_name)) {
US_ONCE({ US_LOG_ERROR("No PCM capture device"); });
goto close_acap;
}
int hz = _g_config->acap_hz;
if (hz <= 0) {
US_A(us_str_is_ok(_g_config->tc358743_dev_path));
chip_fd = open(_g_config->tc358743_dev_path, O_RDWR);
if (chip_fd < 0) {
US_ONCE({ US_LOG_PERROR("Can't open TC358743 for the audio request"); });
goto close_acap;
}
hz = us_chip_tc358743_get_audio_hz(chip_fd);
if (hz == US_ERROR_NO_SIGNAL) {
US_ONCE({ US_LOG_INFO("No audio presented from the host"); });
goto close_acap;
} else if (hz <= 0) {
US_ONCE({ US_LOG_PERROR("Can't get audio HZ"); });
goto close_acap;
}
US_ONCE({ US_LOG_INFO("Detected host audio"); });
}
if ((acap = us_acap_init(_g_config->acap_dev_name, hz)) == NULL) {
goto close_acap;
}
once = 0;
while (!_STOP && _HAS_WATCHERS && _HAS_LISTENERS) {
if (chip_fd >= 0 && us_chip_tc358743_get_audio_hz(chip_fd) != (int)acap->pcm_hz) {
goto close_acap;
}
uz size = US_RTP_TOTAL_SIZE - US_RTP_HEADER_SIZE;
u8 data[size];
u64 pts;
const int result = us_acap_get_encoded(acap, data, &size, &pts);
if (result == 0) {
_LOCK_ACAP;
us_rtpa_wrap(_g_rtpa, data, size, pts);
_UNLOCK_ACAP;
} else if (result == -1) {
goto close_acap;
}
}
close_acap:
US_DELETE(acap, us_acap_destroy);
US_CLOSE_FD(chip_fd);
sleep(1); // error_delay
}
return NULL;
}
static void *_aplay_thread(void *arg) {
(void)arg;
US_THREAD_SETTLE("us_p_aplay");
atomic_store(&_g_aplay_tid_created, true);
US_A(us_str_is_ok(_g_config->aplay_dev_name));
int once = 0;
while (!_STOP) {
snd_pcm_t *dev = NULL;
bool skip = true;
while (!_STOP) {
usleep((US_AU_FRAME_MS / 4) * 1000);
us_au_pcm_s mixed = {0};
_LOCK_APLAY;
US_LIST_ITERATE(_g_clients, client, {
us_au_pcm_s last = {0};
do {
const int ri = us_ring_consumer_acquire(client->aplay_pcm_ring, 0);
if (ri >= 0) {
const us_au_pcm_s *pcm = client->aplay_pcm_ring->items[ri];
memcpy(&last, pcm, sizeof(us_au_pcm_s));
us_ring_consumer_release(client->aplay_pcm_ring, ri);
} else {
break;
}
} while (skip && !_STOP);
us_au_pcm_mix(&mixed, &last);
// US_LOG_INFO("++++++ mixed %p", client);
});
_UNLOCK_APLAY;
// US_LOG_INFO("++++++ --------------");
if (skip) {
static uint skipped = 0;
if (skipped < (1000 / (US_AU_FRAME_MS / 4))) {
++skipped;
continue;
} else {
skipped = 0;
}
}
if (!_HAS_WATCHERS || !_HAS_SPEAKERS) {
goto close_aplay;
}
if (dev == NULL) {
if (!us_au_probe(_g_config->aplay_dev_name)) {
US_ONCE({ US_LOG_ERROR("No PCM playback device"); });
goto close_aplay;
}
int err = snd_pcm_open(&dev, _g_config->aplay_dev_name, SND_PCM_STREAM_PLAYBACK, 0);
if (err < 0) {
US_ONCE({ US_LOG_PERROR_ALSA(err, "Can't open PCM playback"); });
goto close_aplay;
}
err = snd_pcm_set_params(dev, SND_PCM_FORMAT_S16_LE, SND_PCM_ACCESS_RW_INTERLEAVED,
US_RTP_OPUS_CH, US_RTP_OPUS_HZ, 1 /* soft resample */, 50000 /* 50000 = 0.05sec */
);
if (err < 0) {
US_ONCE({ US_LOG_PERROR_ALSA(err, "Can't configure PCM playback"); });
goto close_aplay;
}
US_LOG_INFO("Playback opened, playing ...");
once = 0;
}
if (dev != NULL && mixed.frames > 0) {
snd_pcm_sframes_t frames = snd_pcm_writei(dev, mixed.data, mixed.frames);
if (frames < 0) {
frames = snd_pcm_recover(dev, frames, 1);
} else {
if (once != 0) {
US_LOG_INFO("Playing resumed (snd_pcm_writei) ...");
}
once = 0;
skip = false;
}
if (frames < 0) {
US_ONCE({ US_LOG_PERROR_ALSA(frames, "Can't play to PCM playback"); });
if (frames == -ENODEV) {
goto close_aplay;
}
skip = true;
} else {
if (once != 0) {
US_LOG_INFO("Playing resumed (snd_pcm_recover) ...");
}
once = 0;
skip = false;
}
}
}
close_aplay:
if (dev != NULL) {
US_DELETE(dev, snd_pcm_close);
US_LOG_INFO("Playback closed");
}
}
return NULL;
}
static void _push_camera_event(bool requested) {
json_t *const json_event_type = json_string(requested ? "requested" : "released");
json_t *const camera = json_object();
json_object_set_new(camera, "action", json_event_type);
json_t *const result = json_object();
json_object_set_new(result, "ustreamer", json_string("event"));
json_object_set_new(result, "status", json_string("camera"));
json_object_set_new(result, "camera", camera);
json_t *const event = json_object();
json_object_set_new(event, "result", result);
US_LIST_ITERATE(_g_clients, client, {
_g_gw->push_event(client->session, create(), NULL, event, NULL);
});
json_decref(event);
}
static void _camera_set_active(uint width, uint height, uint fps) {
bool again = false;
_LOCK_VPLAY;
if (_g_camera.active) {
_push_camera_event(false);
again = true;
}
_g_camera.active = true;
_g_camera.width = width;
_g_camera.height = height;
_g_camera.fps = fps;
_push_camera_event(true);
_UNLOCK_VPLAY;
US_LOG_INFO("Camera requested%s: %ux%u@%u", (again ? " again" : ""), width, height, fps);
}
static void _camera_set_inactive() {
_LOCK_VPLAY;
if (_g_camera.active) {
_g_camera.active = false;
_push_camera_event(false);
US_LOG_INFO("Camera released");
}
_UNLOCK_VPLAY;
}
static void *_vplay_thread(void *arg) {
(void)arg;
US_THREAD_SETTLE("us_p_vplay");
atomic_store(&_g_vplay_tid_created, true);
US_A(us_str_is_ok(_g_config->vplay_sink_name));
us_frame_s *frame = us_frame_init();
while (!_STOP) {
us_memsink_s* sink = us_memsink_init_opened(
"H264-CAM",
_g_config->vplay_sink_name,
true, // Server
_g_config->vplay_sink_mode,
false, // Don't remote sink on destroy
1, // Client TTL
1); // Timeout
if (sink == NULL) {
goto close_memsink;
}
const ldf memsink_check_interval = 1;
ldf last_memsink_check_time = 0;
int once = 0;
while (!_STOP) {
_LOCK_VPLAY;
const bool c_active = _g_camera.active;
const uint c_width = _g_camera.width;
const uint c_height = _g_camera.height;
_UNLOCK_VPLAY;
if (!c_active) {
frame->used = 0;
}
if (frame->used > 0 && (c_width != frame->width || c_height != frame->height)) {
US_ONCE({ US_LOG_INFO("Got WebRTC frame with wrong resolution"); });
frame->used = 0;
}
const ldf now_ts = us_get_now_monotonic();
if (frame->used > 0 || last_memsink_check_time + memsink_check_interval < now_ts) {
last_memsink_check_time = now_ts;
us_memsink_wants_s w_get = {0};
switch (us_memsink_server_x_lock(sink, &w_get)) {
case US_MSS_SUCCESS:
if (w_get.format == V4L2_PIX_FMT_H264) {
if (w_get.width == c_width && w_get.height == c_height) {
if (us_memsink_server_x_is_consumed(sink)) {
if (!c_active) {
_camera_set_active(w_get.width, w_get.height, w_get.fps);
} else if (frame->used > 0) {
US_ONCE({ US_LOG_INFO("Streaming to the camera ..."); });
us_memsink_server_x_put(sink, frame);
frame->used = 0;
}
} else {
if (!c_active) {
us_memsink_server_x_set_consumed(sink);
}
}
} else { // Notify to changed resolution
_camera_set_active(w_get.width, w_get.height, w_get.fps);
}
} else {
US_ONCE({ US_LOG_ERROR("Got invalid format from the camera: %u", w_get.format); });
}
if (us_memsink_server_x_unlock(sink) < 0) {
goto close_memsink;
}
break;
case US_MSS_NO_CLIENT:
_camera_set_inactive();
break;
case US_MSS_BUSY: // Busy by some client
break;
case US_MSS_ERROR:
default:
goto close_memsink;
}
}
if (frame->used == 0) { // Frame sent (or discarded), get a new one from a client
const int in_ri = us_ring_consumer_acquire(_g_vplay_ring, 0.1);
if (in_ri >= 0) {
us_frame_s *tmp_frame = frame;
frame = _g_vplay_ring->items[in_ri];
_g_vplay_ring->items[in_ri] = tmp_frame;
us_ring_consumer_release(_g_vplay_ring, in_ri);
} else {
US_ONCE({
_LOCK_VPLAY;
if (_g_camera.active) {
US_LOG_INFO("No frames have got from WebRTC");
}
_UNLOCK_VPLAY;
});
}
} else {
usleep(1000); // Don't iterate too frequently
}
}
close_memsink:
US_DELETE(sink, us_memsink_destroy);
US_LOG_INFO("Memsink closed");
sleep(1);
}
us_frame_destroy(frame);
return NULL;
}
static void _relay_rtp_clients(const us_rtp_s *rtp) {
US_LIST_ITERATE(_g_clients, client, {
us_janus_client_send(client, rtp);
});
}
static void _alsa_quiet(const char *file, int line, const char *func, int err, const char *fmt, ...) {
(void)file;
(void)line;
(void)func;
(void)err;
(void)fmt;
}
static int _plugin_init(janus_callbacks *gw, const char *config_dir_path) {
// https://groups.google.com/g/meetecho-janus/c/xoWIQfaoJm8
// sysctl -w net.core.rmem_default=500000
// sysctl -w net.core.wmem_default=500000
// sysctl -w net.core.rmem_max=1000000
// sysctl -w net.core.wmem_max=1000000
US_LOGGING_INIT;
US_LOG_INFO("Initializing PiKVM uStreamer plugin %s ...", US_VERSION);
if (gw == NULL || !us_str_is_ok(config_dir_path)) {
return -1;
}
if ((_g_config = us_config_init(config_dir_path)) == NULL) {
return -1;
}
_g_gw = gw;
snd_lib_error_set_handler(_alsa_quiet);
US_RING_INIT_WITH_ITEMS(_g_video_ring, 64, us_frame_init);
_g_rtpv = us_rtpv_init(_relay_rtp_clients);
if (us_str_is_ok(_g_config->vplay_sink_name)) {
US_RING_INIT_WITH_ITEMS(_g_vplay_ring, 15, us_frame_init);
US_THREAD_CREATE(_g_vplay_tid, _vplay_thread, NULL);
}
if (us_str_is_ok(_g_config->acap_dev_name)) {
_g_rtpa = us_rtpa_init(_relay_rtp_clients);
US_THREAD_CREATE(_g_acap_tid, _acap_thread, NULL);
}
if (us_str_is_ok(_g_config->aplay_dev_name)) {
US_THREAD_CREATE(_g_aplay_tid, _aplay_thread, NULL);
}
US_THREAD_CREATE(_g_video_rtp_tid, _video_rtp_thread, NULL);
US_THREAD_CREATE(_g_video_sink_tid, _video_sink_thread, NULL);
atomic_store(&_g_ready, true);
return 0;
}
static void _plugin_destroy(void) {
US_LOG_INFO("Destroying plugin ...");
atomic_store(&_g_stop, true);
# define JOIN(_tid) { if (atomic_load(&_tid##_created)) { US_THREAD_JOIN(_tid); } }
JOIN(_g_aplay_tid);
JOIN(_g_acap_tid);
JOIN(_g_vplay_tid);
JOIN(_g_video_sink_tid);
JOIN(_g_video_rtp_tid);
# undef JOIN
US_LIST_ITERATE(_g_clients, client, {
if (_g_camera.client == client) {
_g_camera.client = NULL;
}
US_LIST_REMOVE(_g_clients, client);
us_janus_client_destroy(client);
});
US_RING_DELETE_WITH_ITEMS(_g_vplay_ring, us_frame_destroy);
US_RING_DELETE_WITH_ITEMS(_g_video_ring, us_frame_destroy);
US_DELETE(_g_rtpa, us_rtpa_destroy);
US_DELETE(_g_rtpv, us_rtpv_destroy);
US_DELETE(_g_config, us_config_destroy);
US_LOGGING_DESTROY;
}
static void _plugin_create_session(janus_plugin_session *session, int *err) {
_IF_DISABLED({ *err = -1; return; });
_LOCK_ALL;
US_LOG_INFO("Creating session %p ...", session);
us_janus_client_s *const client = us_janus_client_init(_g_gw, session);
US_LIST_APPEND(_g_clients, client);
atomic_store(&_g_has_watchers, true);
_UNLOCK_ALL;
}
static void _plugin_destroy_session(janus_plugin_session* session, int *err) {
_IF_DISABLED({ *err = -1; return; });
_LOCK_ALL;
bool found = false;
bool has_watchers = false;
bool has_listeners = false;
bool has_speakers = false;
US_LIST_ITERATE(_g_clients, client, {
if (client->session == session) {
US_LOG_INFO("Removing session %p ...", session);
if (_g_camera.client == client) {
_g_camera.client = NULL; // Is it required to notify other clients?
}
US_LIST_REMOVE(_g_clients, client);
us_janus_client_destroy(client);
found = true;
} else {
has_watchers = (has_watchers || atomic_load(&client->transmit));
has_listeners = (has_listeners || atomic_load(&client->transmit_acap));
has_speakers = (has_speakers || atomic_load(&client->transmit_aplay));
}
});
if (!found) {
US_LOG_ERROR("No session %p", session);
*err = -2;
}
atomic_store(&_g_has_watchers, has_watchers);
atomic_store(&_g_has_listeners, has_listeners);
atomic_store(&_g_has_speakers, has_speakers);
_UNLOCK_ALL;
}
static json_t *_plugin_query_session(janus_plugin_session *session) {
_IF_DISABLED({ return NULL; });
json_t *info = NULL;
_LOCK_ALL;
US_LIST_ITERATE(_g_clients, client, {
if (client->session == session) {
info = json_string("session_found");
break;
}
});
_UNLOCK_ALL;
return info;
}
static void _set_transmit(janus_plugin_session *session, const char *msg, bool transmit) {
(void)msg;
_IF_DISABLED({ return; });
_LOCK_ALL;
bool found = false;
bool has_watchers = false;
US_LIST_ITERATE(_g_clients, client, {
if (client->session == session) {
atomic_store(&client->transmit, transmit);
// US_LOG_INFO("%s session %p", msg, session);
found = true;
}
has_watchers = (has_watchers || atomic_load(&client->transmit));
});
if (!found) {
US_LOG_ERROR("No session %p", session);
}
atomic_store(&_g_has_watchers, has_watchers);
_UNLOCK_ALL;
}
static void _plugin_setup_media(janus_plugin_session *session) { _set_transmit(session, "Unmuted", true); }
static void _plugin_hangup_media(janus_plugin_session *session) { _set_transmit(session, "Muted", false); }
static bool _is_camera_enabled(void) {
bool enabled = us_str_is_ok(_g_config->vplay_sink_name);
if (enabled && us_str_is_ok(_g_config->vplay_dev_path)) {
if (access(_g_config->vplay_dev_path, F_OK) != 0) {
enabled = false;
}
}
return enabled;
}
static struct janus_plugin_result *_plugin_handle_message(
janus_plugin_session *session, char *transaction, json_t *msg, json_t *jsep) {
janus_plugin_result_type result_type = JANUS_PLUGIN_OK;
char *result_msg = NULL;
if (session == NULL || msg == NULL) {
result_type = JANUS_PLUGIN_ERROR;
result_msg = (msg ? "No session" : "No message");
goto done;
}
# define PUSH_ERROR(x_error, x_reason) { \
/*US_LOG_ERROR("Message error in session %p: %s", session, x_reason);*/ \
json_t *m_event = json_object(); \
json_object_set_new(m_event, "ustreamer", json_string("event")); \
json_object_set_new(m_event, "error_code", json_integer(x_error)); \
json_object_set_new(m_event, "error", json_string(x_reason)); \
_g_gw->push_event(session, create(), NULL, m_event, NULL); \
json_decref(m_event); \
}
json_t *const request = json_object_get(msg, "request");
if (request == NULL) {
PUSH_ERROR(400, "Request missing");
goto done;
}
const char *const request_str = json_string_value(request);
if (request_str == NULL) {
PUSH_ERROR(400, "Request not a string");
goto done;
}
// US_LOG_INFO("Message: %s", request_str);
# define PUSH_STATUS(x_status, x_payload, x_jsep) { \
json_t *const m_event = json_object(); \
json_object_set_new(m_event, "ustreamer", json_string("event")); \
json_t *const m_result = json_object(); \
json_object_set_new(m_result, "status", json_string(x_status)); \
if (x_payload != NULL) { \
json_object_set(m_result, x_status, x_payload); \
} \
json_object_set_new(m_event, "result", m_result); \
_g_gw->push_event(session, create(), NULL, m_event, x_jsep); \
json_decref(m_event); \
}
if (!strcmp(request_str, "start")) {
PUSH_STATUS("started", NULL, NULL);
} else if (!strcmp(request_str, "stop")) {
PUSH_STATUS("stopped", NULL, NULL);
} else if (!strcmp(request_str, "watch")) {
uint video_orient = 0;
bool with_acap = false;
bool with_aplay = false;
bool with_vplay = false;
{
json_t *const params = json_object_get(msg, "params");
if (params != NULL) {
# define READ_BOOL(x_target, x_key, x_cond) { \
json_t *const m_obj = json_object_get(params, x_key); \
if (m_obj != NULL && json_is_boolean(m_obj)) { \
x_target = (json_boolean_value(m_obj) && (x_cond)); \
} \
}
READ_BOOL(with_acap, "audio", us_au_probe(_g_config->acap_dev_name));
READ_BOOL(with_aplay, "mic", us_au_probe(_g_config->aplay_dev_name));
READ_BOOL(with_vplay, "camera", _is_camera_enabled());
# undef READ_BOOL
{
json_t *const obj = json_object_get(params, "orientation");
if (obj != NULL && json_is_integer(obj)) {
video_orient = json_integer_value(obj);
switch (video_orient) {
case 90: case 180: case 270: break;
default: video_orient = 0; break;
}
}
}
}
}
{
_LOCK_ALL;
with_vplay = (with_vplay && _g_camera.active && !_g_camera.client);
bool has_listeners = false;
bool has_speakers = false;
US_LIST_ITERATE(_g_clients, client, {
if (client->session == session) {
char *const sdp = us_sdp_create(
client->video_ssrc,
client->audio_ssrc,
with_acap,
with_aplay,
with_vplay);
json_t *const offer_jsep = json_pack("{ssss}", "type", "offer", "sdp", sdp);
PUSH_STATUS("started", NULL, offer_jsep);
json_decref(offer_jsep);
free(sdp);
atomic_store(&client->transmit_acap, with_acap);
atomic_store(&client->transmit_aplay, with_aplay);
atomic_store(&client->video_orient, video_orient);
if (with_vplay) {
us_janus_client_start_vplay(client, _g_vplay_ring);
_g_camera.client = client;
}
}
has_listeners = (has_listeners || atomic_load(&client->transmit_acap));
has_speakers = (has_speakers || atomic_load(&client->transmit_aplay));
});
atomic_store(&_g_has_listeners, has_listeners);
atomic_store(&_g_has_speakers, has_speakers);
_UNLOCK_ALL;
}
} else if (!strcmp(request_str, "features")) {
const char *const ice_url = getenv("JANUS_USTREAMER_WEB_ICE_URL");
_LOCK_VPLAY
json_t *camera_req = NULL;
if (_g_camera.active && _g_camera.client == NULL) {
json_t *resolution = json_object();
json_object_set_new(resolution, "width", json_integer(_g_camera.width));
json_object_set_new(resolution, "height", json_integer(_g_camera.height));
camera_req = json_object();
json_object_set_new(camera_req, "resolution", resolution);
json_object_set_new(camera_req, "fps", json_integer(_g_camera.fps));
}
_UNLOCK_VPLAY
json_t *const features = json_pack(
"{s:b, s:b, s:{s:b, s:o*}, s:{s:s?}}",
"audio", us_au_probe(_g_config->acap_dev_name),
"mic", us_au_probe(_g_config->aplay_dev_name),
"camera",
"enabled", _is_camera_enabled(),
"request", camera_req,
"ice",
"url", (ice_url != NULL ? ice_url : _g_default_ice_url)
);
PUSH_STATUS("features", features, NULL);
json_decref(features);
} else if (!strcmp(request_str, "key_required")) {
// US_LOG_INFO("Got key_required message");
atomic_store(&_g_key_required, true);
} else {
PUSH_ERROR(405, "Not implemented");
}
done:
US_DELETE(transaction, free);
US_DELETE(msg, json_decref);
US_DELETE(jsep, json_decref);
return janus_plugin_result_new(
result_type, result_msg,
(result_type == JANUS_PLUGIN_OK ? json_pack("{sb}", "ok", 1) : NULL));
# undef PUSH_STATUS
# undef PUSH_ERROR
}
static void _plugin_incoming_rtp(janus_plugin_session *session, janus_plugin_rtp *packet) {
_IF_DISABLED({ return; });
if (session == NULL || packet == NULL) {
return; // Accept only valid packets
}
if (packet->video) {
_LOCK_VPLAY;
if (_g_camera.client != NULL && _g_camera.client->session == session) {
us_janus_client_recv(_g_camera.client, packet);
}
_UNLOCK_VPLAY;
} else {
_LOCK_APLAY;
US_LIST_ITERATE(_g_clients, client, {
if (client->session == session) {
us_janus_client_recv(client, packet);
break;
}
});
_UNLOCK_APLAY;
}
}
static void _plugin_incoming_rtcp(janus_plugin_session *session, janus_plugin_rtcp *packet) {
_IF_DISABLED({ return; });
if (
session == NULL || packet == NULL || !packet->video
|| packet->length < sizeof(janus_rtcp_header)
) {
return;
}
if (
janus_rtcp_has_pli(packet->buffer, packet->length)
|| janus_rtcp_has_fir(packet->buffer, packet->length)
) {
atomic_store(&_g_key_required, true);
}
janus_rtcp_header *header = (janus_rtcp_header *)packet->buffer;
if (header->type == RTCP_SR && janus_rtcp_check_sr(header, packet->length)) {
const janus_rtcp_sr *const sr = (janus_rtcp_sr *)header;
if (sr->si.ntp_ts_msw && sr->si.ntp_ts_lsw) {
_LOCK_VPLAY;
US_LIST_ITERATE(_g_clients, client, {
if (client->session == session) {
const ldf ntp_ts = ntohl(sr->si.ntp_ts_msw) + ntohl(sr->si.ntp_ts_lsw) / 4294967296.L;
us_rtpc_sync_timestamp(client->rtpc, ntp_ts, ntohl(sr->si.rtp_ts));
break;
}
});
_UNLOCK_VPLAY;
}
}
}
// ***** Plugin *****
static int _plugin_get_api_compatibility(void) { return JANUS_PLUGIN_API_VERSION; }
static int _plugin_get_version(void) { return US_VERSION_U; }
static const char *_plugin_get_version_string(void) { return US_VERSION; }
static const char *_plugin_get_description(void) { return "PiKVM uStreamer Janus plugin for H.264 video"; }
static const char *_plugin_get_name(void) { return US_PLUGIN_NAME; }
static const char *_plugin_get_author(void) { return "Maxim Devaev <mdevaev@gmail.com>"; }
static const char *_plugin_get_package(void) { return US_PLUGIN_PACKAGE; }
janus_plugin *create(void) {
# pragma GCC diagnostic push
# pragma GCC diagnostic ignored "-Woverride-init"
static janus_plugin plugin = JANUS_PLUGIN_INIT(
.init = _plugin_init,
.destroy = _plugin_destroy,
.create_session = _plugin_create_session,
.destroy_session = _plugin_destroy_session,
.query_session = _plugin_query_session,
.setup_media = _plugin_setup_media,
.hangup_media = _plugin_hangup_media,
.handle_message = _plugin_handle_message,
.get_api_compatibility = _plugin_get_api_compatibility,
.get_version = _plugin_get_version,
.get_version_string = _plugin_get_version_string,
.get_description = _plugin_get_description,
.get_name = _plugin_get_name,
.get_author = _plugin_get_author,
.get_package = _plugin_get_package,
.incoming_rtp = _plugin_incoming_rtp,
.incoming_rtcp = _plugin_incoming_rtcp,
);
# pragma GCC diagnostic pop
return &plugin;
}