diff options
| author | 2025-01-24 18:14:52 -0500 | |
|---|---|---|
| committer | 2025-01-24 18:14:52 -0500 | |
| commit | f638237a6b4f3d3edf9bdd995c15df537f8d7e7c (patch) | |
| tree | 44efce6d912c43f67548690256879d036e22e742 | |
| parent | 2daa31c0629f2eb4af84d6f4fed8ac89813de056 (diff) | |
| download | camu-f638237a6b4f3d3edf9bdd995c15df537f8d7e7c.tar.gz camu-f638237a6b4f3d3edf9bdd995c15df537f8d7e7c.tar.bz2 camu-f638237a6b4f3d3edf9bdd995c15df537f8d7e7c.zip | |
Fix multiple bugs encountered during stress test
- Basic server resource cleanup
Signed-off-by: Andrew Opalach <andrew@akon.city>
39 files changed, 424 insertions, 177 deletions
diff --git a/env/wsl/bootstrap-configuration.nix b/env/wsl/bootstrap-configuration.nix index 03f4bbf..b66147c 100755 --- a/env/wsl/bootstrap-configuration.nix +++ b/env/wsl/bootstrap-configuration.nix @@ -6,7 +6,7 @@ nix.settings.experimental-features = [ "nix-command" "flakes" ]; - environment.systemPackages = with pkgs; [ git ]; + environment.systemPackages = with pkgs; [ git nfs-utils ]; system.stateVersion = "23.05"; } diff --git a/env/wsl/wsl-dev-host.nix b/env/wsl/wsl-dev-host.nix index fd9a1db..ccd0ac4 100755 --- a/env/wsl/wsl-dev-host.nix +++ b/env/wsl/wsl-dev-host.nix @@ -28,7 +28,7 @@ in { networking.hostName = "wsl-dev"; networking.firewall.enable = false; - hardware.pulseaudio.enable = true; + services.pulseaudio.enable = true; xdg.portal = { enable = true; @@ -21,6 +21,7 @@ libplacebo = super.libplacebo.overrideAttrs (old: { patches = [ ./subprojects/packagefiles/libplacebo/cache_ref_rect.diff + ./subprojects/packagefiles/libplacebo/compiler_warnings.diff ./subprojects/packagefiles/libplacebo/fix_rotate.diff ./subprojects/packagefiles/libplacebo/glClientWaitSync_0_return.diff ./subprojects/packagefiles/libplacebo/info_priv_hash.diff @@ -104,6 +105,7 @@ libblake3 libspng ffmpeg-headless + soxr libcdio libcdio-paranoia libdvdcss @@ -206,6 +208,12 @@ ]; buildInputs = with pkgsCross.ucrt64; [ (zlib.override { shared = false; static = true; }) + (soxr.overrideAttrs (oldAttrs: { + cmakeFlags = [ + "-DBUILD_SHARED_LIBS=OFF" + "-DBUILD_STATIC_LIBS=ON" + ]; + })) (openssl.override { static = true; }) xxHash vulkan-headers diff --git a/scripts/python_env b/scripts/python_env new file mode 100644 index 0000000..2cca301 --- /dev/null +++ b/scripts/python_env @@ -0,0 +1,4 @@ +export CURL_CA_BUNDLE=/etc/ssl/certs/ca-bundle.crt +export PYTHONDONTWRITEBYTECODE=1 +export PYTHONOPTIMIZE=2 +export PYTHONPATH=$HOME/c/camu/src/portal/vendor/vendor diff --git a/scripts/run.sh b/scripts/run.sh index eeebcac..3b47fe4 100755 --- a/scripts/run.sh +++ b/scripts/run.sh @@ -1,6 +1,3 @@ #! /usr/bin/env sh -export CURL_CA_BUNDLE=/etc/ssl/certs/ca-bundle.crt -export PYTHONDONTWRITEBYTECODE=1 -export PYTHONPATH=$HOME/c/camu/src/portal/vendor/vendor -export mesa_glthread=true +source ../scripts/python_env ./src/fruits/cmv/cmv "$@" diff --git a/scripts/run_debug.sh b/scripts/run_debug.sh index bd9a667..a56bf22 100755 --- a/scripts/run_debug.sh +++ b/scripts/run_debug.sh @@ -1,2 +1,3 @@ #! /usr/bin/env sh +source ../scripts/python_env gdb -ex run --args ./src/fruits/cmv/cmv "$@" diff --git a/scripts/run_server.sh b/scripts/run_server.sh index 893f5c1..22945ed 100755 --- a/scripts/run_server.sh +++ b/scripts/run_server.sh @@ -1,8 +1,5 @@ #! /usr/bin/env sh -export CURL_CA_BUNDLE=/etc/ssl/certs/ca-bundle.crt -export PYTHONDONTWRITEBYTECODE=1 -export PYTHONOPTIMIZE=2 -export PYTHONPATH=$HOME/c/camu/src/portal/vendor/vendor +source ../scripts/python_env # screen will use $SHELL. screen -c ../scripts/screenrc #valgrind --leak-check=no --show-error-list=yes --log-file=./server-valgrind.log ./src/fruits/cmsrv/cmsrv $@ diff --git a/scripts/run_valgrind.sh b/scripts/run_valgrind.sh index 033f35a..583327c 100755 --- a/scripts/run_valgrind.sh +++ b/scripts/run_valgrind.sh @@ -1,3 +1,5 @@ #! /usr/bin/env sh -valgrind --leak-check=no --show-error-list=yes ./src/fruits/cmv/cmv "$@" -#valgrind --leak-check=full ./src/fruits/cmv/cmv "$@" +source ../scripts/python_env +valgrind --leak-check=full ./src/fruits/cmv/cmv "$@" +#valgrind --leak-check=full --show-leak-kinds=all ./src/fruits/cmv/cmv "$@" +#valgrind --leak-check=no --show-error-list=yes ./src/fruits/cmv/cmv "$@" diff --git a/src/buffer/audio.c b/src/buffer/audio.c index cf42840..ce027af 100644 --- a/src/buffer/audio.c +++ b/src/buffer/audio.c @@ -11,7 +11,7 @@ #define BUFFER_SIZE 8.0 #define BUFFER_MARK_MIN 3.25 // Must be a most half of the buffer size. -#define BUFFER_MARK_BUFFERED 1.25 +#define BUFFER_MARK_BUFFERED 1.0 #ifdef CAMU_AUDIO_BUFFER_FADE #define FADE_STEP(fmt) (1.75f / (fmt)->sample_rate) @@ -408,8 +408,8 @@ ptrdiff_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, ptrdif if (step != 0.f || buf->fade.volume != 1.f) { if (buf->volume.user > 1.f) { - // Only adjust the step if it's an increase to avoid - // drawn-out fades if the volume is low. + // Only adjust the step if it's an increase to avoid drawn-out + // fades if the volume is low. step *= buf->volume.user; } buf->fade.volume = apply_volume(data, req, fmt, buf->fade.volume, buf->volume.user, step); diff --git a/src/buffer/video.c b/src/buffer/video.c index bb217af..7771939 100644 --- a/src/buffer/video.c +++ b/src/buffer/video.c @@ -17,13 +17,12 @@ bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock) { - buf->clock = clock; buf->latency = 0.0; al_atomic_store(f64)(&buf->pts, -1.0, AL_ATOMIC_RELAXED); buf->last_pts = -1.0; // The least confusing behavior for single_frame is that it can't be - // set unless the buffer is not empty. + // set if the buffer is empty. buf->single_frame = false; buf->avg_frame_duration = 0.0; buf->queue = NULL; diff --git a/src/cache/entry.c b/src/cache/entry.c index c7fb8d2..8930ea0 100644 --- a/src/cache/entry.c +++ b/src/cache/entry.c @@ -6,7 +6,6 @@ bool cch_entry_get_handle(struct cch_entry *entry, struct cch_handle *handle) { handle->entry = entry; handle->pointer = 0; - handle->prev_pointer = 0; nn_cond_init(&handle->wait.cond); nn_mutex_init(&handle->wait.mutex); handle->wait.disabled = false; diff --git a/src/cache/handle.c b/src/cache/handle.c index 98bc2fa..10d1992 100644 --- a/src/cache/handle.c +++ b/src/cache/handle.c @@ -1,6 +1,5 @@ #include <al/lib.h> #include <al/log.h> -#include <al/random.h> #include "../codec/codec.h" @@ -23,9 +22,7 @@ s32 cch_handle_read(struct cch_handle *handle, u8 *buf, s32 size) // filesize < 0: Size is yet to be evaluated and we can wait on it. // filesize = 0: Size is explicitly unknown. off_t filesize = cch_entry_get_size(handle->entry); - if (size < 0) { - filesize = wait_for_size(handle, handler); - } + if (size < 0) filesize = wait_for_size(handle, handler); if (filesize > 0) { if (handle->pointer >= filesize) return CAMU_ERR_EOF; if (handle->pointer + size >= filesize) { @@ -58,13 +55,9 @@ off_t cch_handle_seek(struct cch_handle *handle, off_t offset, s32 whence) { struct cch_handler *handler = handle->entry->handler; off_t filesize = cch_entry_get_size(handle->entry); - if (filesize < 0) { - filesize = wait_for_size(handle, handler); - } + if (filesize < 0) filesize = wait_for_size(handle, handler); if (filesize <= 0) return -1; - if (whence == CAMU_SEEK_SIZE) { - return filesize; - } + if (whence == CAMU_SEEK_SIZE) return filesize; switch (whence) { case SEEK_SET: if (offset >= filesize || offset < 0) { @@ -79,13 +72,11 @@ off_t cch_handle_seek(struct cch_handle *handle, off_t offset, s32 whence) handle->pointer += offset; break; case SEEK_END: - handle->prev_pointer = handle->pointer; handle->pointer = filesize + offset; break; default: return -1; } - handle->prev_pointer = -1; return handle->pointer; } diff --git a/src/cache/handle.h b/src/cache/handle.h index 007e04d..0052c02 100644 --- a/src/cache/handle.h +++ b/src/cache/handle.h @@ -9,8 +9,6 @@ struct cch_entry; struct cch_handle { struct cch_entry *entry; off_t pointer; - // Set after a seek_end and reset after any seek_cur or seek_set. - off_t prev_pointer; struct cch_handler_wait wait; }; diff --git a/src/cache/threaded_waits.c b/src/cache/threaded_waits.c index 5990783..ae43ae9 100644 --- a/src/cache/threaded_waits.c +++ b/src/cache/threaded_waits.c @@ -39,12 +39,12 @@ bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing nn_mutex_unlock(&handler->mutex); return false; } - nn_mutex_unlock(&wait->mutex); nn_mutex_lock(&handler->mutex); if (wait->disabled) { al_array_remove(handler->waits, wait); canceled = true; } + nn_mutex_unlock(&wait->mutex); } nn_mutex_unlock(&handler->mutex); return !canceled; diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c index 5ff99d7..0462090 100644 --- a/src/codec/ffmpeg/decoder.c +++ b/src/codec/ffmpeg/decoder.c @@ -147,14 +147,16 @@ static bool ff_decoder_init(struct camu_decoder *dec, struct camu_renderer *rend av->codec_context = NULL; AVCodecParameters *codecpar = stream->av.stream->codecpar; - void *state = NULL; const AVCodec *codec; +#ifdef CAMU_FF_DECODER_HWACCEL + void *state = NULL; al_array_init(av->supported_hw_codecs); while ((codec = av_codec_iterate(&state))) { if (codec->capabilities & (AV_CODEC_CAP_HARDWARE | AV_CODEC_CAP_HYBRID)) { al_array_push(av->supported_hw_codecs, codec); } } +#endif codec = avcodec_find_decoder(codecpar->codec_id); if (!codec) { @@ -313,15 +315,21 @@ static s32 receive_frames(struct camu_ff_decoder *av) struct camu_codec_frame *frame = al_alloc_object(struct camu_codec_frame); frame->mode = CAMU_FFMPEG_COMPAT; frame->av.frame = av_frame_alloc(); - ret = avcodec_receive_frame(av->codec_context, frame->av.frame); + AVFrame *avframe = frame->av.frame; + ret = avcodec_receive_frame(av->codec_context, avframe); if (ret < 0) { - av_frame_free(&frame->av.frame); + av_frame_free(&avframe); al_free(frame); // Checking for EAGAIN should prevent an infinite loop. if (ret == AVERROR(EAGAIN) || ret == AVERROR_EOF) break; continue; } - al_assert(frame->av.frame->best_effort_timestamp != AV_NOPTS_VALUE); + if (avframe->best_effort_timestamp == AV_NOPTS_VALUE) { + // This trips under what I think are normal conditions, + // more testing is needed. + al_assert(avframe->duration == 0); + avframe->best_effort_timestamp = 0; + } av->callback(av->userdata, frame); } return ret; diff --git a/src/codec/ffmpeg/demuxer.c b/src/codec/ffmpeg/demuxer.c index 1e92186..5c37f1a 100644 --- a/src/codec/ffmpeg/demuxer.c +++ b/src/codec/ffmpeg/demuxer.c @@ -73,30 +73,30 @@ static bool ff_demuxer_init(struct camu_demuxer *demux, struct cch_handle *handl } #define GUESS_STREAM_IS_IMAGE(stream) \ - (stream->duration >= 0L && stream->nb_frames <= 1L && (stream->avg_frame_rate.den == 0 && stream->r_frame_rate.den > 0)) + (stream->duration >= 0 && stream->nb_frames <= 1 && (stream->avg_frame_rate.den == 0 && stream->r_frame_rate.den > 0)) - s64 default_duration = av->format_context->duration < 0L ? 0L : av->format_context->duration; - av->duration = 0Lu; + s64 default_duration = av->format_context->duration < 0 ? 0 : av->format_context->duration; + av->duration = 0; for (u32 i = 0; i < av->format_context->nb_streams; i++) { AVStream *stream = av->format_context->streams[i]; enum AVMediaType type = stream->codecpar->codec_type; - u64 duration = 0Lu; + u64 duration = 0; if (type == AVMEDIA_TYPE_AUDIO || type == AVMEDIA_TYPE_VIDEO || type == AVMEDIA_TYPE_SUBTITLE) { // Try to detect attached images. if (type == AVMEDIA_TYPE_VIDEO && GUESS_STREAM_IS_IMAGE(stream)) { al_log_info("ff_demuxer", "Assuming stream #%u is an image.", i); } else { - if (stream->duration < 0L) { + if (stream->duration < 0) { stream->duration = av_rescale_q(default_duration, AV_TIME_BASE_Q, stream->time_base); - al_assert(stream->duration >= 0L); + al_assert(stream->duration >= 0); al_log_info("ff_demuxer", "Stream #%u has an invalid duration, defaulting to %.3fs.", i, stream->duration * av_q2d(stream->time_base)); } duration = (u64)av_rescale_q(stream->duration, stream->time_base, AV_TIME_BASE_Q); if (duration > av->duration) av->duration = duration; } - if (stream->start_time == AV_NOPTS_VALUE) stream->start_time = 0L; + if (stream->start_time == AV_NOPTS_VALUE) stream->start_time = 0; } else if (type != AVMEDIA_TYPE_ATTACHMENT) { continue; } @@ -112,7 +112,6 @@ static bool ff_demuxer_init(struct camu_demuxer *demux, struct cch_handle *handl //av_dump_format(av->format_context, 0, "", 0); - av->seeked = false; av->eof = false; return true; @@ -123,7 +122,7 @@ err: static bool subscribed_to_index(struct camu_demuxer *demux, s32 index) { - return ((demux->subscribed >> index) & 1); + return (demux->subscribed >> index) & 1; } static s32 ff_demuxer_get_packet(struct camu_demuxer *demux, struct camu_codec_packet *packet) @@ -146,7 +145,6 @@ static s32 ff_demuxer_get_packet(struct camu_demuxer *demux, struct camu_codec_p av_packet_unref(pkt); continue; } - if (pkt->pts == AV_NOPTS_VALUE) pkt->pts = 0L; break; } return ret; @@ -164,14 +162,14 @@ static bool ff_demuxer_seek(struct camu_demuxer *demux, u64 pos) al_assert(pos <= INT64_MAX); if (pos > av->duration) pos = av->duration; if (av->eof) { + avio_flush(av->io_context); + avformat_flush(av->format_context); av->eof = false; } - avformat_flush(av->format_context); s64 ts = (s64)pos; if (avformat_seek_file(av->format_context, -1, INT64_MIN, ts, ts, 0) < 0) { return false; } - av->seeked = true; return true; } diff --git a/src/codec/ffmpeg/demuxer.h b/src/codec/ffmpeg/demuxer.h index 68ad03e..8066513 100644 --- a/src/codec/ffmpeg/demuxer.h +++ b/src/codec/ffmpeg/demuxer.h @@ -12,7 +12,6 @@ struct camu_ff_demuxer { AVIOContext *io_context; AVFormatContext *format_context; u64 duration; - bool seeked; bool eof; }; diff --git a/src/codec/ffmpeg/meson.build b/src/codec/ffmpeg/meson.build index fc0bc91..cba033c 100644 --- a/src/codec/ffmpeg/meson.build +++ b/src/codec/ffmpeg/meson.build @@ -48,14 +48,15 @@ if not (libavutil.found() and libavformat.found() and libavcodec.found() and lib # Android MediaCodec HW decoding. ffmpeg_client_deps += [compiler.find_library('mediandk')] endif - soxr = compiler.find_library('soxr', required: false) - if soxr.found() - ffmpeg_client_deps += [soxr] - ffmpeg_args += ['-DCAMU_HAVE_SOXR'] - endif endif endif +soxr = compiler.find_library('soxr', required: false) +if soxr.found() + ffmpeg_client_deps += [soxr] + ffmpeg_args += ['-DCAMU_HAVE_SOXR'] +endif + if libavutil.found() and libavformat.found() and libavcodec.found() and libswresample.found() and libswscale.found() ffmpeg_server_deps += [libavutil, libavformat, libavcodec] ffmpeg_client_deps += [libavutil, libavformat, libavcodec, libswresample, libswscale] diff --git a/src/fruits/cmc/ui/pane_search.c b/src/fruits/cmc/ui/pane_search.c index 021c42f..c1b47ed 100644 --- a/src/fruits/cmc/ui/pane_search.c +++ b/src/fruits/cmc/ui/pane_search.c @@ -40,7 +40,7 @@ static bool handle_text_input(struct cmc_search_tab *tab, struct ncinput *input) switch (input->id) { case NCKEY_BACKSPACE: if (tab->cursor > 0) { - al_wstr_remove_at(&tab->input_text, tab->cursor); + al_wstr_remove_at(&tab->input_text, tab->cursor - 1); tab->cursor--; tab->input_changed = true; } @@ -79,19 +79,36 @@ static bool handle_text_input(struct cmc_search_tab *tab, struct ncinput *input) tab->cursor++; } break; + case 'U': + for (u32 i = 1; i <= tab->cursor; i++) { + al_wstr_remove_at(&tab->input_text, tab->cursor - i); + } + tab->cursor = 0; + tab->input_changed = true; + break; } break; } + + if (ncinput_alt_p(input) || nckey_synthesized_p(input->id)) { + break; + } + + al_wstr_insert(&tab->input_text, tab->cursor, &al_wstr_w((wchar_t *)input->eff_text, 0, 1)); + tab->cursor++; + tab->input_changed = true; + + /* // https://github.com/dankamongmen/notcurses/blob/c11efe877f2245901f5c9ce110de7a2c83cfa2ed/src/lib/reader.c#L404 for (s32 c = 0; input->eff_text[c] != 0; c++){ uchar egc[5] = { 0 }; if (notcurses_ucs32_to_utf8(&input->eff_text[c], 1, egc, 4) >= 0) { // This probably doesn't work on Windows. al_wstr_insert(&tab->input_text, tab->cursor, &al_wstr_w((wchar_t *)egc, 0, 1)); - tab->cursor++; - tab->input_changed = true; } } + */ + break; } return false; diff --git a/src/liana/client.c b/src/liana/client.c index b6635af..09280cb 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -178,12 +178,13 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) client->callback(client->userdata, LIANA_CLIENT_RECONNECTED, NULL, NULL); } client->reconnect = false; + } else { + al_assert(client->connection_id == 0); } stream->packet_sent_callback = packet_sent_callback; struct nn_packet *packet = nn_packet_create(); - nn_packet_write_u32(packet, client->id); - //nn_packet_write_u32(packet, client->connection_id); - nn_packet_write_u32(packet, 0); + nn_packet_write_u32(packet, client->node_id); + nn_packet_write_u32(packet, client->connection_id); nn_packet_write_s32(packet, client->mask); nn_packet_write_u64(packet, client->pos); if (client->mask == 0) { @@ -220,16 +221,17 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * } void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type, - str *addr, u16 port, u32 id, u64 pos, struct camu_renderer *renderer) + str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer) { client->loop = loop; - client->id = id; + client->node_id = node_id; client->pos = pos; client->mask = 0; client->reconnect = false; lia_vcr_init(&client->vcr, client->loop, &client->data); al_str_clone(&client->addr, addr); client->port = port; + client->connection_id = 0; nn_packet_stream_init(&client->data, connection_callback, connection_closed_callback, client); client->renderer = renderer; #ifdef CAMU_DIRECT_MODE diff --git a/src/liana/client.h b/src/liana/client.h index 55c5b3a..eca47b9 100644 --- a/src/liana/client.h +++ b/src/liana/client.h @@ -8,7 +8,7 @@ struct lia_client { struct nn_event_loop *loop; - u32 id; + u32 node_id; s32 mask; u64 pos; u64 at; @@ -25,7 +25,7 @@ struct lia_client { }; void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type, - str *addr, u16 port, u32 id, u64 pos, struct camu_renderer *renderer); + str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer); void lia_client_seek(struct lia_client *client, u64 pos, u64 at); void lia_client_reseek(struct lia_client *client); void lia_client_disconnect(struct lia_client *client); diff --git a/src/liana/list.c b/src/liana/list.c index fc93fc4..94ff3bf 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -842,18 +842,36 @@ void lia_list_clear(struct lia_list *list) pump_queue(list); } +void lia_list_close(struct lia_list *list) +{ + // TODO: Consider sinks being in use. Wait for list->sinks to be empty? + struct lia_list_entry *entry; + al_array_foreach(list->entries, i, entry) { + list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL); + } +} + void lia_list_free(struct lia_list *list) { + struct lia_list_cmd *cmd; + al_array_foreach(list->queue, i, cmd) { + al_free(cmd); + } + al_array_free(list->queue); + if (list->cmd) al_free(list->cmd); + struct lia_list_entry *entry; al_array_foreach(list->entries, i, entry) { al_wstr_free(&entry->name); al_free(entry); } al_array_free(list->entries); + struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { al_free(sink); } al_array_free(list->sinks); + al_str_free(&list->name); } diff --git a/src/liana/list.h b/src/liana/list.h index 674f194..e61074c 100644 --- a/src/liana/list.h +++ b/src/liana/list.h @@ -135,4 +135,5 @@ void lia_list_sort(struct lia_list *list); void lia_list_shuffle(struct lia_list *list); void lia_list_clear(struct lia_list *list); +void lia_list_close(struct lia_list *list); void lia_list_free(struct lia_list *list); diff --git a/src/liana/server.c b/src/liana/server.c index 87bcf57..2e74d1e 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -79,10 +79,18 @@ static void discard_packet_callback(void *userdata, struct nn_packet_stream *str static void free_connection(struct lia_node_connection *conn) { + struct lia_node *node = conn->node; al_assert(conn->handler); conn->handler->free(&conn->handler); cch_entry_return_handle(conn->node->entry, &conn->handle); nn_packet_pool_free(&conn->pool); + bool removed; + al_array_remove_checked(node->connections, conn, removed); + al_assert(removed); + al_free(conn); + if (node->closed && !node->connections.count) { + cch_entry_free(&node->entry); + } } static void free_connection_stream(struct lia_node_connection *conn) @@ -90,6 +98,20 @@ static void free_connection_stream(struct lia_node_connection *conn) nn_packet_stream_free(conn->stream); al_free(conn->stream); conn->stream = NULL; + conn->ref = false; +} + +static void disable_connection(struct lia_node_connection *conn) +{ + nn_packet_pool_disable(&conn->pool); + cch_handle_disable(&conn->handle); + nn_thread_join(&conn->thread); +} + +static void enable_connection(struct lia_node_connection *conn) +{ + cch_handle_enable(&conn->handle); + nn_packet_pool_enable(&conn->pool); } static void data_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) @@ -97,19 +119,15 @@ static void data_connection_closed_callback(void *userdata, struct nn_packet_str struct lia_node_connection *conn = (struct lia_node_connection *)userdata; (void)stream; - nn_packet_pool_disable(&conn->pool); - cch_handle_disable(&conn->handle); - - // This is joining handler_thread(), we will never be here if init_thread() blocks or fails. - nn_thread_join(&conn->thread); + // We will never be here if init_thread() blocks or fails. + disable_connection(conn); free_connection_stream(conn); if (conn->disconnected) { free_connection(conn); } else { - cch_handle_enable(&conn->handle); - nn_packet_pool_enable(&conn->pool); + enable_connection(conn); } } @@ -117,13 +135,13 @@ static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *s { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; s32 mask = nn_packet_read_s32(packet); + nn_packet_stream_return_packet(stream, packet); conn->handler->subscribe(conn->handler, mask); stream->packet_callback = discard_packet_callback; stream->packet_sent_callback = data_packet_sent_callback; stream->packets_sent_callback = data_packets_sent_callback; stream->connection_closed_callback = data_connection_closed_callback; nn_thread_create(&conn->thread, handler_thread, conn); - nn_packet_stream_return_packet(stream, packet); } static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet) @@ -137,6 +155,7 @@ static void subscribe_connection_closed_callback(void *userdata, struct nn_packe struct lia_node_connection *conn = (struct lia_node_connection *)userdata; (void)stream; free_connection_stream(conn); + free_connection(conn); } static void handle_connection(struct lia_node_connection *conn, struct nn_packet *packet) @@ -146,6 +165,9 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet s32 mask = nn_packet_read_s32(packet); u64 seek_pos = nn_packet_read_u64(packet); + al_assert(!conn->ref); + conn->ref = true; + // Besides being wasteful, seeking to 0 on a new stream can skip data. if (mask != 0 || seek_pos > 0) { conn->handler->seek(conn->handler, seek_pos); @@ -179,13 +201,30 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * al_free(stream); } +static void packet_sent_callback(void *userdata, struct nn_packet *packet) +{ + (void)userdata; + nn_packet_free(packet); +} + +static void demote_and_disconnect_stream(struct lia_server *server, struct nn_packet_stream *stream) +{ + // This connection is now nothing but a packet stream. + stream->userdata = server; + stream->connection_closed_callback = connection_closed_callback; + stream->packet_callback = discard_packet_callback; + stream->packet_sent_callback = packet_sent_callback; + stream->packets_sent_callback = NULL; + nn_packet_stream_disconnect(stream); +} + static void signal_callback(void *userdata) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; struct lia_node *node = conn->node; struct lia_server *server = node->server; - nn_signal_stop(&conn->signal); nn_thread_join(&conn->thread); + nn_signal_stop(&conn->signal); struct nn_packet *packet = conn->packet; conn->packet = NULL; if (!packet || conn->errored) { @@ -194,6 +233,9 @@ static void signal_callback(void *userdata) } if (!packet) { // Connection was closed before init was done. + nn_packet_stream_free(conn->stream); + al_free(conn->stream); + al_free(conn); return; } struct nn_packet_stream *stream = conn->stream; @@ -202,15 +244,14 @@ static void signal_callback(void *userdata) nn_packet_pool_init(&conn->pool, 1024, server->loop, packet_pool_callback, conn); al_array_push(node->connections, conn); handle_connection(conn, packet); - } else { + } + nn_packet_stream_return_packet(stream, packet); + // conn->errored will not have changed but we want to return the packet before disconnecting. + if (conn->errored) { al_free(conn); conn = NULL; - // This connection is now nothing but a packet stream. - stream->userdata = server; - stream->connection_closed_callback = connection_closed_callback; - nn_packet_stream_disconnect(stream); + demote_and_disconnect_stream(server, stream); } - nn_packet_stream_return_packet(stream, packet); } static nn_thread_result NNWT_THREADCALL init_thread(void *userdata) @@ -223,20 +264,22 @@ static nn_thread_result NNWT_THREADCALL init_thread(void *userdata) return 0; } -static struct lia_node_connection *get_connection_from_id(struct lia_node *node, u32 id) +static struct lia_node *get_node_from_id(struct lia_server *server, u32 id) { - struct lia_node_connection *conn = NULL; - al_array_foreach(node->connections, i, conn) { - if (conn->id == id) return conn; + struct lia_node *node; + al_array_foreach(server->nodes, i, node) { + if (node->id == id) return node; } return NULL; } -static struct lia_node *get_node_from_id(struct lia_server *server, u32 id) +static struct lia_node_connection *get_connection_from_id(struct lia_node *node, u32 id) { - struct lia_node *node; - al_array_foreach(server->nodes, i, node) { - if (node->id == id) return node; + struct lia_node_connection *conn; + al_array_foreach(node->connections, i, conn) { + if (conn->id == id) { + return conn; + } } return NULL; } @@ -274,28 +317,35 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str conn->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler(); conn->errored = false; conn->disconnected = false; + conn->ref = false; stream->packet_callback = discard_packet_callback; stream->connection_closed_callback = pre_init_connection_closed_callback; nn_thread_create(&conn->thread, init_thread, conn); } else { // Connections never get removed from node->connection. if ((conn = get_connection_from_id(node, connection_id))) { + if (conn->ref) { + // Cleanup the existing connection's handler and demote it's stream. + // The stream was likely already disconnected client-side but it's still safe + // to disconnect it here to be sure. + disable_connection(conn); + demote_and_disconnect_stream(server, conn->stream); + conn->ref = false; + enable_connection(conn); + } stream->userdata = conn; + al_assert(conn->node == node); conn->stream = stream; handle_connection(conn, packet); - } else { - nn_packet_stream_disconnect(stream); } nn_packet_stream_return_packet(stream, packet); + // Return packet before possibly disconnecting. + if (!conn) { + nn_packet_stream_disconnect(stream); + } } } -static void packet_sent_callback(void *userdata, struct nn_packet *packet) -{ - (void)userdata; - nn_packet_free(packet); -} - static bool connection_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_server *server = (struct lia_server *)userdata; @@ -318,6 +368,7 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en node->id = get_incremental_id(server); node->entry = entry; al_array_init(node->connections); + node->closed = false; node->server = server; al_array_push(server->nodes, node); return node; @@ -355,12 +406,14 @@ void lia_node_get_duration(struct lia_node *node) nn_thread_create(&node->thread, init_duration_thread, node); } -void lia_server_close(struct lia_server *server) +void lia_node_close(struct lia_node *node) { - struct lia_node *node; - al_array_foreach(server->nodes, i, node) { + node->closed = true; + if (!node->connections.count) { + cch_entry_free(&node->entry); + } else { struct lia_node_connection *conn; - al_array_foreach(node->connections, j, conn) { + al_array_foreach_rev(node->connections, i, conn) { if (conn->stream) { conn->disconnected = true; nn_packet_stream_disconnect(conn->stream); @@ -369,7 +422,10 @@ void lia_server_close(struct lia_server *server) } } } +} +void lia_server_close(struct lia_server *server) +{ struct nn_packet_stream *zombie; al_array_foreach_rev(server->zombies, i, zombie) { nn_packet_stream_disconnect(zombie); @@ -380,18 +436,12 @@ void lia_server_free(struct lia_server *server) { struct lia_node *node; al_array_foreach(server->nodes, i, node) { - struct lia_node_connection *conn; - al_array_foreach(node->connections, j, conn) { - al_free(conn); - } + al_assert(!node->connections.count); al_array_free(node->connections); al_free(node); } al_array_free(server->nodes); - struct nn_packet_stream *zombie; - al_array_foreach(server->zombies, i, zombie) { - al_free(zombie); - } + al_assert(!server->zombies.count); al_array_free(server->zombies); } diff --git a/src/liana/server.h b/src/liana/server.h index 849fd38..a9e567a 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -14,6 +14,7 @@ struct lia_node_connection { struct nn_packet_stream *stream; bool errored; bool disconnected; + bool ref; struct lia_server_handler *handler; struct cch_handle handle; struct nn_thread thread; @@ -30,6 +31,7 @@ struct lia_node { u32 id; struct cch_entry *entry; array(struct lia_node_connection *) connections; + bool closed; struct lia_server *server; u64 duration; // Temporary copy from node_connection. We need to @@ -54,5 +56,6 @@ bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop); void lia_server_add_stream(struct lia_server *server, struct nn_packet_stream *stream); struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry); void lia_node_get_duration(struct lia_node *node); +void lia_node_close(struct lia_node *node); void lia_server_close(struct lia_server *server); void lia_server_free(struct lia_server *server); diff --git a/src/liana/vcr.c b/src/liana/vcr.c index 47beed9..33f3365 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -36,12 +36,13 @@ static void reset_metrics(struct lia_vcr *vcr) void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data) { + vcr->data = data; al_array_init(vcr->tracks); al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); vcr->mark.buffered = VCR_BUFFER_BUFFERED; vcr->mark.low = 0; vcr->expand = VCR_EXPAND_UNTOUCHED; - vcr->data = data; + vcr->started = false; #ifndef CAMU_DIRECT_MODE nn_signal_init(&vcr->signal, loop, signal_callback, vcr); #else @@ -168,11 +169,13 @@ void lia_vcr_start(struct lia_vcr *vcr) #endif struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { - if (VCR_TRACK_THREADED(track) && !track->running) { + al_assert(!track->running); + if (VCR_TRACK_THREADED(track)) { nn_thread_create(&track->thread, vcr_track_thread, track); track->running = true; } } + vcr->started = true; } void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track) @@ -327,6 +330,7 @@ void lia_vcr_uncork(struct lia_vcr_track *track) static void vcr_track_close_internal(struct lia_vcr_track *track) { + struct lia_vcr *vcr = track->vcr; // Calling packet_cache_disable() while holding the track mutex can // very possibly deadlock. nn_packet_cache_disable(&track->cache); @@ -336,7 +340,8 @@ static void vcr_track_close_internal(struct lia_vcr_track *track) nn_cond_signal(&track->cond); } nn_mutex_unlock(&track->mutex); - if (track->running) { + if (vcr->started) { + al_assert(track->running); nn_thread_join(&track->thread); track->running = false; } @@ -357,6 +362,7 @@ void lia_vcr_flush(struct lia_vcr *vcr) track->client->flush(track->client); } } + vcr->started = false; #ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); #endif @@ -375,6 +381,7 @@ void lia_vcr_close_all(struct lia_vcr *vcr) vcr_track_close_internal(track); return_entire_cache(track); } + vcr->started = false; #ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); #endif diff --git a/src/liana/vcr.h b/src/liana/vcr.h index bd4f0f9..8330365 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -25,11 +25,12 @@ struct lia_vcr_track { }; struct lia_vcr { + struct nn_packet_stream *data; array(struct lia_vcr_track *) tracks; atomic(u64) count; struct { u64 buffered, low; } mark; u8 expand; - struct nn_packet_stream *data; + bool started; #ifndef CAMU_DIRECT_MODE struct nn_signal signal; #endif diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 0f56c5b..951b230 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -12,6 +12,7 @@ #include "common.h" //#define CAMU_SINK_ONESHOT +//#define CAMU_SINK_TRACE #ifndef CAMU_SINK_NO_VIDEO #include "../render/renderer_libplacebo.h" @@ -191,6 +192,10 @@ static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_e // add_audio/video_if_set_and_buffered() no-ops for ENDED buffers. static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry) { +#ifdef CAMU_SINK_TRACE + al_log_info("sink", "remove_entry_buffers(0x%llx), audio_state: %hhu, video_state: %hhu", + entry, entry->audio.state, entry->video.state); +#endif al_assert(!entry->ended); if (entry->audio.state != BUFFER_ENDED) { remove_entry_audio_buffer(sink, entry); @@ -446,10 +451,14 @@ static void mixer_callback(void *userdata, u8 op) struct camu_sink *sink = (struct camu_sink *)userdata; if (op == CAMU_MIXER_EMPTY) { al_log_info("sink", "Mixer empty."); +#ifndef LIANA_LIST_SCUFFED_LOOP queue_cmd(sink, (struct camu_sink_cmd){ .op = STOP, .value.i = CAMU_SINK_AUDIO }); +#else + (void)sink; +#endif } } @@ -483,6 +492,9 @@ bool camu_sink_init(struct camu_sink *sink, struct nn_event_loop *loop, static void maybe_remove_previous(struct camu_sink *sink) { +#ifdef CAMU_SINK_TRACE + al_log_info("sink", "maybe_remove_previous(), previous_count: %u", sink->previous.count); +#endif struct camu_sink_entry *previous; al_array_foreach(sink->previous, i, previous) { remove_entry_buffers(sink, previous); @@ -495,6 +507,9 @@ static void maybe_remove_previous(struct camu_sink *sink) // An obvious example of this is at the point an entry gets freed. See LIANA_CLIENT_REMOVE_BUFFERS. static void maybe_remove_from_previous(struct camu_sink *sink, struct camu_sink_entry *entry) { +#ifdef CAMU_SINK_TRACE + al_log_info("sink", "maybe_remove_from_previous(0x%llx)", entry); +#endif al_array_remove_all(sink->previous, entry); } @@ -502,6 +517,11 @@ static void maybe_remove_from_previous(struct camu_sink *sink, struct camu_sink_ static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry *previous, struct camu_sink_entry *target) { +#ifdef CAMU_SINK_TRACE + al_log_info("sink", "maybe_add_to_previous(0x%llx, 0x%llx), audio_state: %hhu, video_state: %hhu", + previous, target, previous->audio.state, previous->video.state); +#endif + al_assert(previous != target && !previous->ended); struct camu_sink_entry *entry; @@ -663,6 +683,9 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry) static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) { +#ifdef CAMU_SINK_TRACE + al_log_info("sink", "switch_to(0x%llx)", target); +#endif if (sink->current) { struct camu_sink_entry *current = sink->current; al_assert(current != target); @@ -862,6 +885,9 @@ static void clock_callback(void *userdata, u8 op) struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; struct camu_sink *sink = entry->sink; if (op == CAMU_CLOCK_PAUSED) { +#ifdef CAMU_SINK_TRACE + al_log_info("sink", "clock_callback(), target: 0x%llx", sink->target); +#endif nn_mutex_lock(&sink->mutex); if (entry == sink->current && sink->target) { switch_to(sink, sink->target); @@ -1035,10 +1061,17 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } } - // We need to call this if entry was added to previous then, + // Entry might be in previous here if it was added to previous then, // - it's being cleaned up after ENTRY_MAX_AGE - 1 entries were added but none buffered. // - it was seeked. - maybe_remove_from_previous(sink, entry); + //maybe_remove_from_previous(sink, entry); + struct camu_sink_entry *previous; + al_array_foreach(sink->previous, i, previous) { + if (previous == entry) { + maybe_remove_previous(sink); + break; + } + } bool skip_audio = sink->audio.state == SINK_PAUSED; if (entry->audio.state == BUFFER_ADDED) { @@ -1164,7 +1197,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } bool removed; - al_array_check_remove(sink->entries, entry, removed); + al_array_remove_checked(sink->entries, entry, removed); nn_mutex_unlock(&sink->mutex); @@ -1173,10 +1206,9 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str #ifndef CAMU_SINK_NO_VIDEO camu_video_buffer_free(&entry->video.buf); #endif + al_log_info("sink", "Entry (0x%llx) closed by %s.", entry, removed ? "disconnect" : "cleanup"); al_free(entry); - al_log_info("sink", "Entry closed by %s.", removed ? "disconnect" : "cleanup"); - break; } } @@ -1307,6 +1339,7 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn, camu_clock_resume(&entry->clock, 0); } + al_assert(entry != current); switch_to(sink, entry); #else switch (pause) { diff --git a/src/portal/src/post.c b/src/portal/src/post.c index 4cb0bcc..6ca8f23 100644 --- a/src/portal/src/post.c +++ b/src/portal/src/post.c @@ -27,8 +27,8 @@ void camu_post_clone(struct camu_post *dest, struct camu_post *src) dest->quotes = src->quotes; dest->comments = src->comments; dest->views = src->views; - struct camu_post_media *media; struct camu_post_media media_copy; + struct camu_post_media *media; al_array_foreach_ptr(src->media, i, media) { al_bzero(&media_copy, sizeof(struct camu_post_media)); media_copy.type = media->type; @@ -44,6 +44,30 @@ void camu_post_clone(struct camu_post *dest, struct camu_post *src) al_str_clone(&dest->in_reply_to.unique_id, &src->in_reply_to.unique_id); } +void camu_post_free(struct camu_post *post) +{ + al_str_free(&post->unique_id); + al_str_free(&post->url); + al_str_free(&post->author.unique_id); + al_wstr_free(&post->author.username); + al_wstr_free(&post->author.display_name); + al_str_free(&post->author.profile_picture_url); + al_wstr_free(&post->title); + al_wstr_free(&post->text); + struct camu_post_media *media; + al_array_foreach_ptr(post->media, i, media) { + al_str_free(&media->key); + al_str_free(&media->url); + al_str_free(&media->ext); + al_str_free(&media->thumbnail_url); + al_str_free(&media->thumbnail_ext); + } + al_array_free(post->media); + al_str_free(&post->post.unique_id); + al_str_free(&post->quoted.unique_id); + al_str_free(&post->in_reply_to.unique_id); +} + void camu_post_add_media(struct camu_post *post, u8 type, str *key, str *url, str *ext, str *thumbnail_url, str *thumbnail_ext) { al_array_push(post->media, ((struct camu_post_media){ diff --git a/src/portal/src/post.h b/src/portal/src/post.h index a946968..4bb3785 100644 --- a/src/portal/src/post.h +++ b/src/portal/src/post.h @@ -105,6 +105,7 @@ struct camu_post { void camu_post_reset(struct camu_post *post); void camu_post_clone(struct camu_post *dest, struct camu_post *src); +void camu_post_free(struct camu_post *post); // Python internal. void camu_post_add_date(struct camu_post *post, u8 type, f32 range_start, f32 range_end, u32 meta); diff --git a/src/portal/src/post_cache.c b/src/portal/src/post_cache.c index ad50404..d6b885b 100644 --- a/src/portal/src/post_cache.c +++ b/src/portal/src/post_cache.c @@ -5,13 +5,18 @@ void camu_post_cache_init(struct camu_post_cache *cache) al_array_init(cache->cache); } -void camu_post_cache_push(struct camu_post_cache *cache, struct camu_post *post) +void camu_post_cache_push(struct camu_post_cache *cache, struct camu_post *post, bool clone) { struct camu_post *check = camu_post_cache_get(cache, &post->unique_id); if (!check) { - struct camu_post *npost = al_alloc_object(struct camu_post); - camu_post_clone(npost, post); - al_array_push(cache->cache, npost); + if (clone) { + struct camu_post *cloned = al_alloc_object(struct camu_post); + camu_post_clone(cloned, post); + post = cloned; + } + al_array_push(cache->cache, post); + } else if (!clone) { + camu_post_free(post); } } @@ -28,3 +33,12 @@ struct camu_post *camu_post_cache_get(struct camu_post_cache *cache, str *unique } return NULL; } + +void camu_post_cache_free(struct camu_post_cache *cache) +{ + struct camu_post *post; + al_array_foreach(cache->cache, i, post) { + camu_post_free(post); + } + al_array_free(cache->cache); +} diff --git a/src/portal/src/post_cache.h b/src/portal/src/post_cache.h index a406c72..8d0f975 100644 --- a/src/portal/src/post_cache.h +++ b/src/portal/src/post_cache.h @@ -9,5 +9,6 @@ struct camu_post_cache { }; void camu_post_cache_init(struct camu_post_cache *cache); -void camu_post_cache_push(struct camu_post_cache *cache, struct camu_post *post); +void camu_post_cache_push(struct camu_post_cache *cache, struct camu_post *post, bool clone); struct camu_post *camu_post_cache_get(struct camu_post_cache *cache, str *unique_id); +void camu_post_cache_free(struct camu_post_cache *cache); diff --git a/src/portal/src/search.c b/src/portal/src/search.c index e5d5dca..2584878 100644 --- a/src/portal/src/search.c +++ b/src/portal/src/search.c @@ -59,13 +59,14 @@ static struct camu_search *get_search_by_id(struct camu_portal_bridge *bridge, s static nn_thread_result NNWT_THREADCALL queue_thread(void *userdata) { - struct camu_portal_bridge *bridge = (struct camu_portal_bridge *)userdata; nn_thread_setcanceltype(NNWT_THREAD_CANCEL_ASYNCHRONOUS); + struct camu_portal_bridge *bridge = (struct camu_portal_bridge *)userdata; bool have_python = false; nn_mutex_lock(&bridge->mutex); for (;;) { nn_cond_wait(&bridge->cond, &bridge->mutex); if (bridge->quit) break; + if (!have_python) { // Defer python init. if (!(have_python = camu_python_init())) { @@ -73,6 +74,7 @@ static nn_thread_result NNWT_THREADCALL queue_thread(void *userdata) break; } } + struct camu_portal_cmd *cmd; al_array_foreach_ptr(bridge->queue, i, cmd) { struct camu_portal_result result = { 0 }; @@ -84,15 +86,15 @@ static nn_thread_result NNWT_THREADCALL queue_thread(void *userdata) s32 id = portal_bridge_search(&cmd->module, &cmd->query); if (id >= 0) { struct camu_search *search = al_alloc_object(struct camu_search); - search->page = 0; - al_array_init(search->pages); search->id = id; + result.id = search->id; al_str_clone(&search->module, &cmd->module); al_str_clone(&search->query, &cmd->query); + search->page = 0; + al_array_init(search->pages); search->bridge = bridge; al_array_push(bridge->searches, search); al_log_info("portal", "New search %x (%.*s).", id, al_str_fmt(&cmd->query)); - result.id = search->id; } else { } al_str_free(&cmd->module); @@ -103,22 +105,24 @@ static nn_thread_result NNWT_THREADCALL queue_thread(void *userdata) struct camu_search *search = get_search_by_id(bridge, cmd->id); if (search) { result.id = search->id; + struct camu_result_page *page; al_array_foreach_ptr(search->pages, j, page) { if (page->num == cmd->num) break; } + al_log_info("portal", "Loading page %i (%.*s).", cmd->num, al_str_fmt(&search->query)); if (portal_bridge_get_page(search, search->id, cmd->num) == -1) { break; } - page = &al_array_at(search->pages, cmd->num); + + result.page = &al_array_at(search->pages, cmd->num); if (bridge->cache) { struct camu_post *post; - al_array_foreach_ptr(page->posts, j, post) { - camu_post_cache_push(bridge->cache, post); + al_array_foreach_ptr(result.page->posts, j, post) { + camu_post_cache_push(bridge->cache, post, false); } } - result.page = page; } break; } @@ -205,39 +209,41 @@ void camu_portal_close(struct camu_portal_bridge *bridge) nn_mutex_unlock(&bridge->mutex); nn_thread_join(&bridge->thread); nn_signal_stop(&bridge->results_signal); - camu_queue_free(bridge->results); - al_array_free(bridge->queue); - nn_cond_destroy(&bridge->cond); - nn_mutex_destroy(&bridge->mutex); } -/* -struct camu_search *camu_portal_get_search(struct camu_portal_bridge *bridge, s32 id) +void camu_portal_free(struct camu_portal_bridge *bridge) { + nn_cond_destroy(&bridge->cond); + nn_mutex_destroy(&bridge->mutex); + camu_queue_free(bridge->results); + struct camu_search *search; al_array_foreach(bridge->searches, i, search) { - if (search->id == id) return search; + al_str_free(&search->module); + al_str_free(&search->query); + struct camu_result_page *page; + al_array_foreach_ptr(search->pages, j, page) { + al_array_free(page->posts); + str *unique_id; + al_array_foreach_ptr(page->list, k, unique_id) { + al_str_free(unique_id); + } + al_array_free(page->list); + } + al_array_free(search->pages); } - return NULL; -} -void camu_portal_discard_search(struct camu_portal_bridge *bridge, s32 id) -{ - (void)bridge; - (void)id; -} - -void camu_search_free(struct camu_search *search) -{ - struct camu_result_page *page; - al_array_foreach_ptr(search->pages, i, page) { - // TODO: Free camu_post ? - al_array_free(page->posts); - al_array_free(page->list); + struct camu_portal_cmd *cmd; + al_array_foreach_ptr(bridge->queue, i, cmd) { + switch (cmd->op) { + case CAMU_CLIENT_CREATE_SEARCH: + al_str_free(&cmd->module); + al_str_free(&cmd->query); + break; + } } - al_array_free(search->pages); + al_array_free(bridge->queue); } -*/ static struct camu_result_page *page_at_index(struct camu_search *search, u32 num) { diff --git a/src/portal/src/search.h b/src/portal/src/search.h index f870522..f7e42a2 100644 --- a/src/portal/src/search.h +++ b/src/portal/src/search.h @@ -66,6 +66,7 @@ void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num, void (*callback)(void *, void *, struct camu_portal_result *), void *userdata); void camu_portal_close(struct camu_portal_bridge *bridge); +void camu_portal_free(struct camu_portal_bridge *bridge); /* struct camu_search *camu_portal_get_search(struct camu_portal_bridge *bridge, s32 id); diff --git a/src/render/renderer_libplacebo.c b/src/render/renderer_libplacebo.c index e6cfa0c..cf1ece1 100644 --- a/src/render/renderer_libplacebo.c +++ b/src/render/renderer_libplacebo.c @@ -242,7 +242,7 @@ static void renderer_lp_render(struct camu_renderer *renderer, struct camu_scree while (scr->videos.count > 0) { bool any_eof = false; al_array_foreach_ptr(scr->videos, i, video) { - bool weighted; + bool weighted = false; if (camu_video_buffer_read(video->buf, &mix, &weighted)) { // If mix.frames is NULL, read() returned QUEUE_MORE. if (mix.frames) { diff --git a/src/screen/screen.c b/src/screen/screen.c index 61bbc77..929020e 100644 --- a/src/screen/screen.c +++ b/src/screen/screen.c @@ -2,6 +2,8 @@ #include <nnwt/thread.h> #include <math.h> +#include "../liana/list.h" + #include "view.h" #include "screen.h" diff --git a/src/server/db.c b/src/server/db.c index 4cc449e..e2ad808 100644 --- a/src/server/db.c +++ b/src/server/db.c @@ -65,14 +65,5 @@ bool camu_db_open(struct camu_server *server, str *path) void camu_db_close(struct camu_server *server) { - struct camu_user *user; - al_array_foreach(server->users, i, user) { - al_str_free(&user->name); - } - al_array_free(server->users); - struct lia_list *list; - al_array_foreach(server->lists, i, list) { - lia_list_free(list); - } - al_array_free(server->lists); + (void)server; } diff --git a/src/server/server.c b/src/server/server.c index d75146e..cef05b0 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -292,6 +292,7 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v break; } case LIANA_UNLOAD_ENTRY: + lia_node_close(resource->node); break; case LIANA_LIST_META: switch (*(u8 *)opaque) { @@ -566,6 +567,7 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list #endif } al_assert(resource); + al_array_push(server->data.resources, resource); if (entry) { resource->entry = entry; resource->node = lia_server_create_node(&server->data.server, resource->entry); @@ -750,6 +752,7 @@ void camu_server_init(struct camu_server *server, struct nn_event_loop *loop) nn_rpc_add_command(&server->server, &commands[i]); } + al_array_init(server->data.resources); lia_server_init(&server->data.server, server->loop); #ifdef CAMU_HAVE_PORTAL @@ -778,6 +781,10 @@ void camu_server_close(struct camu_server *server) #ifdef CAMU_HAVE_PORTAL camu_portal_close(&server->bridge); #endif + struct lia_list *list; + al_array_foreach(server->lists, i, list) { + lia_list_close(list); + } lia_server_close(&server->data.server); #ifdef CAMU_DIRECT_MODE nn_multiplex_direct_close(); @@ -788,14 +795,65 @@ void camu_server_close(struct camu_server *server) void camu_server_free(struct camu_server *server) { - // This needs to happen in flight. - // Too many things can block us before we get here. - //struct camu_server_resource *resource; - //al_array_foreach(server->data.resources, i, resource) { - // cch_entry_free(&resource->entry); - //} + al_array_free(server->nodes); + al_array_free(server->clients); + al_array_free(server->sinks); + + struct camu_user *user; + al_array_foreach(server->users, i, user) { + al_str_free(&user->name); + al_free(user); + } + al_array_free(server->users); + + struct camu_resource *resource; + al_array_foreach(server->data.resources, i, resource) { + al_array_free(resource->pending); + switch (resource->type) { + case CAMU_RESOURCE_FILE: { + struct camu_resource_file *file = (struct camu_resource_file *)resource; + al_str_free(&file->path); + al_free(file); + break; + } +#ifdef NAUNET_HAS_CURL + case CAMU_RESOURCE_HTTP: + break; +#endif +#ifdef CACHE_HAVE_CDIO + case CAMU_RESOURCE_CDIO: + break; +#endif +#ifdef CAMU_HAVE_PORTAL + case CAMU_RESOURCE_PORTAL: { + struct camu_resource_portal *portal = (struct camu_resource_portal *)resource; + al_free(portal); + break; + } + case CAMU_RESOURCE_SIMPLE_SEARCH: { + struct camu_resource_portal *portal = (struct camu_resource_portal *)resource; + al_free(portal); + break; + } +#endif + } + } al_array_free(server->data.resources); + + struct lia_list *list; + al_array_foreach(server->lists, i, list) { + lia_list_free(list); + al_free(list); + } + al_array_free(server->lists); + lia_server_free(&server->data.server); + +#ifdef CAMU_HAVE_PORTAL + camu_post_cache_free(&server->cache); + camu_portal_free(&server->bridge); +#endif + nn_rpc_free(&server->server); al_str_free(&server->addr); } diff --git a/src/util/color_palette.c b/src/util/color_palette.c index 73dfc83..9cfdf32 100644 --- a/src/util/color_palette.c +++ b/src/util/color_palette.c @@ -40,31 +40,40 @@ bool camu_color_palette_init(str *path) if (!nn_file_open(&file, path, NNWT_FILE_READONLY)) { return false; } + str s; nn_file_read_as_str(&file, &s); nn_file_close(&file); json_error_t error; json_t *root = json_loadb(s.data, s.length, 0, &error); + al_str_free(&s); if (!root) { - al_log_error("color_palette", "Failed to load %.*s as a json (%s).", - al_str_fmt(path), error.text); + al_log_error("color_palette", "Failed to load %.*s as a json (%s).", al_str_fmt(path), error.text); return false; } + + bool parsed = true; + json_t *special = json_object_get(root, "special"); if (!special) { al_log_error("color_palette", "\"special\" field missing from the color palette."); - return false; + parsed = false; + goto out; } + json_t *colors = json_object_get(root, "colors"); if (!colors) { al_log_error("color_palette", "\"colors\" field missing from the color palette."); - return false; + parsed = false; + goto out; } + json_t *object; const char *color; u32 value; u32 index; bool to_long_error; + for (u32 i = 0; i < ARRAY_SIZE(special_colors); i++) { object = json_object_get(special, special_colors[i]); if (!object) { @@ -80,6 +89,7 @@ bool camu_color_palette_init(str *path) al_log_warn("color_palette", "%s is not a valid hex color.", normal_colors[i]); } } + for (u32 i = 0; i < ARRAY_SIZE(normal_colors); i++) { object = json_object_get(colors, normal_colors[i]); if (!object) { @@ -95,8 +105,13 @@ bool camu_color_palette_init(str *path) al_log_warn("color_palette", "%s is not a valid hex color.", normal_colors[i]); } } + +out: + json_decref(root); + + return parsed; #else (void)path; + return false; #endif - return true; } |