diff options
| author | 2024-12-24 14:52:51 -0500 | |
|---|---|---|
| committer | 2024-12-24 14:52:51 -0500 | |
| commit | 5206f05fdf77bb65c125ddb133cf46a59608c671 (patch) | |
| tree | 69707685b0245f1db508a6bce98d8ab5399a932c /src | |
| parent | 3d55d2722a3129449ab1418e73abd97caa7fd2ae (diff) | |
| download | camu-5206f05fdf77bb65c125ddb133cf46a59608c671.tar.gz camu-5206f05fdf77bb65c125ddb133cf46a59608c671.tar.bz2 camu-5206f05fdf77bb65c125ddb133cf46a59608c671.zip | |
Bring in line with dependency changes
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src')
94 files changed, 1491 insertions, 1393 deletions
diff --git a/src/buffer/audio.c b/src/buffer/audio.c index 2cedcc8..5a956ec 100644 --- a/src/buffer/audio.c +++ b/src/buffer/audio.c @@ -11,7 +11,7 @@ #define BUFFER_SIZE 9.0 #define BUFFER_MARK_MIN 4.3 // Must be a most half of the buffer size. -#define BUFFER_MARK_BUFFERED 2.5 +#define BUFFER_MARK_BUFFERED 2.0 #ifdef CAMU_AUDIO_BUFFER_FADE #define FADE_STEP(fmt) (1.75f / (fmt)->sample_rate) @@ -316,7 +316,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re pts = -pts; ret = camu_audio_format_sec_to_bytes(&buf->fmt.req, pts); ret = MIN(ret, req); - al_log_info("audio_buffer", "Delaying audio by %fs (%zu bytes).", pts, ret); + al_log_info("audio_buffer", "Delaying audio by %fs.", pts); al_memset(data, 0, ret); data += ret; req -= ret; diff --git a/src/buffer/clock.c b/src/buffer/clock.c index 624c534..cbf63d5 100644 --- a/src/buffer/clock.c +++ b/src/buffer/clock.c @@ -1,5 +1,5 @@ #include <al/lib.h> -#include <aki/thread.h> +#include <nnwt/thread.h> #include "clock.h" @@ -31,13 +31,13 @@ void camu_clock_seek(struct camu_clock *clock, f64 base, u64 target) { clock->base = base; clock->offset = 0.0; - f64 tick = aki_get_tick(); + f64 tick = nn_get_tick(); if (clock->paused_at == -1.0) { // We are safe to directly edit the tick here. if (target == 0) { al_atomic_store(f64)(&clock->tick, -1.0, AL_ATOMIC_RELAXED); } else { - tick = calc_tick_offset(tick, aki_get_timestamp(), target); + tick = calc_tick_offset(tick, nn_get_timestamp(), target); al_atomic_store(f64)(&clock->tick, tick, AL_ATOMIC_RELAXED); } } else { @@ -54,9 +54,9 @@ void camu_clock_offset(struct camu_clock *clock, f64 amount) void camu_clock_pause(struct camu_clock *clock, u64 target) { al_assert(clock->paused_at == -1.0); - f64 tick = aki_get_tick(); + f64 tick = nn_get_tick(); if (target > 0) { - tick = calc_tick_offset(tick, aki_get_timestamp(), target); + tick = calc_tick_offset(tick, nn_get_timestamp(), target); al_atomic_store(f64)(&clock->pause, tick, AL_ATOMIC_RELAXED); } else { al_atomic_store(f64)(&clock->pause, PAUSED, AL_ATOMIC_RELAXED); @@ -67,12 +67,12 @@ void camu_clock_pause(struct camu_clock *clock, u64 target) void camu_clock_resume(struct camu_clock *clock, u64 target) { al_assert(clock->paused_at != -1.0); - f64 tick = aki_get_tick(); + f64 tick = nn_get_tick(); //if (pause != PAUSED && tick < pause) { // tick = pause; //} else if (target > 0) { if (target > 0) { - tick = calc_tick_offset(tick, aki_get_timestamp(), target); + tick = calc_tick_offset(tick, nn_get_timestamp(), target); } if (clock->paused_at == 0.0) { if (target == 0) { @@ -107,7 +107,7 @@ f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset) { f64 pause = al_atomic_load(f64)(&clock->pause, AL_ATOMIC_ACQUIRE); if (pause == PAUSED) return -1.0; - f64 current = aki_get_tick(); + f64 current = nn_get_tick(); f64 tick = al_atomic_load(f64)(&clock->tick, AL_ATOMIC_RELAXED); if (tick == -1.0) { tick = al_atomic_compare_and_swap(f64)(&clock->tick, -1.0, current); diff --git a/src/buffer/peak_buffer.c b/src/buffer/peak_buffer.c index 005e77f..eec031a 100644 --- a/src/buffer/peak_buffer.c +++ b/src/buffer/peak_buffer.c @@ -4,27 +4,27 @@ void camu_peak_buffer_init(struct camu_peak_buffer *buf, size_t size) { - aki_buffer_init(&buf->buf); - aki_buffer_ensure_space(&buf->buf, size); + nn_buffer_init(&buf->buf); + nn_buffer_ensure_space(&buf->buf, size); } void camu_peak_buffer_push(struct camu_peak_buffer *buf, u8 *data, size_t size) { - aki_buffer_append(&buf->buf, data, size); + nn_buffer_append(&buf->buf, data, size); } size_t camu_peak_buffer_size(struct camu_peak_buffer *buf) { - return aki_buffer_get_size(&buf->buf); + return nn_buffer_get_size(&buf->buf); } u8 *camu_peak_buffer_flush(struct camu_peak_buffer *buf) { - aki_buffer_set_size(&buf->buf, 0); - return aki_buffer_get_ptr(&buf->buf, 0); + nn_buffer_set_size(&buf->buf, 0); + return nn_buffer_get_ptr(&buf->buf, 0); } void camu_peak_buffer_free(struct camu_peak_buffer *buf) { - aki_buffer_free(&buf->buf); + nn_buffer_free(&buf->buf); } diff --git a/src/buffer/peak_buffer.h b/src/buffer/peak_buffer.h index c518ade..ebb9939 100644 --- a/src/buffer/peak_buffer.h +++ b/src/buffer/peak_buffer.h @@ -1,10 +1,10 @@ #pragma once #include <al/types.h> -#include <aki/common.h> +#include <nnwt/common.h> struct camu_peak_buffer { - struct aki_buffer buf; + struct nn_buffer buf; }; void camu_peak_buffer_init(struct camu_peak_buffer *buf, size_t size); diff --git a/src/cache/backings/file.c b/src/cache/backings/file.c index a09d535..bf62190 100644 --- a/src/cache/backings/file.c +++ b/src/cache/backings/file.c @@ -6,64 +6,64 @@ static void file_backing_lock(struct cch_backing *backing) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_lock(&file->mutex); + nn_mutex_lock(&file->mutex); } static void file_backing_write(struct cch_backing *backing, u8 *buf, off_t index, size_t *size) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_lock(&file->mutex); + nn_mutex_lock(&file->mutex); if (file->u.pointer != index) { - file->u.pointer = aki_file_seek(&file->file, index, SEEK_SET); + file->u.pointer = nn_file_seek(&file->file, index, SEEK_SET); al_assert(file->u.pointer == index); } if (file->size >= 0 && index + ((off_t)*size) >= file->size) { *size = MAX(file->size - index, 0L); } - aki_file_write(&file->file, buf, *size); + nn_file_write(&file->file, buf, *size); file->u.pointer += *size; cch_backing_fill_range(&file->backing, index, *size); - aki_mutex_unlock(&file->mutex); + nn_mutex_unlock(&file->mutex); } static void file_backing_read(struct cch_backing *backing, u8 *buf, off_t index, size_t *size) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_lock(&file->mutex); + nn_mutex_lock(&file->mutex); if (file->u.pointer != index) { - file->u.pointer = aki_file_seek(&file->file, index, SEEK_SET); + file->u.pointer = nn_file_seek(&file->file, index, SEEK_SET); al_assert(file->u.pointer == index); } if (file->size >= 0 && index + ((off_t)*size) >= file->size) { *size = MAX(file->size - index, 0L); } - aki_file_read(&file->file, buf, *size); + nn_file_read(&file->file, buf, *size); file->u.pointer += *size; - aki_mutex_unlock(&file->mutex); + nn_mutex_unlock(&file->mutex); } static void file_backing_unlock(struct cch_backing *backing) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_unlock(&file->mutex); + nn_mutex_unlock(&file->mutex); } static void file_backing_resize(struct cch_backing *backing, size_t size) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_lock(&file->mutex); + nn_mutex_lock(&file->mutex); file->size = size; if ((off_t)size > file->filesize) { - aki_file_truncate(&file->file, size); + nn_file_truncate(&file->file, size); file->filesize = size; } - aki_mutex_unlock(&file->mutex); + nn_mutex_unlock(&file->mutex); } static void file_backing_free(struct cch_backing **backing) { struct cch_backing_file *file = (struct cch_backing_file *)*backing; - aki_file_close(&file->file); + nn_file_close(&file->file); al_array_free(file->backing.available); al_free(file); *backing = NULL; @@ -87,6 +87,6 @@ struct cch_backing *cch_backing_file_create(str *path, size_t size) return NULL; } file->u.pointer = 0; - aki_mutex_init(&file->mutex); + nn_mutex_init(&file->mutex); return (struct cch_backing *)file; } diff --git a/src/cache/backings/file.h b/src/cache/backings/file.h index 3f8437b..8a2bf7a 100644 --- a/src/cache/backings/file.h +++ b/src/cache/backings/file.h @@ -1,22 +1,22 @@ #pragma once #include <al/str.h> -#include <aki/thread.h> -#include <aki/common.h> -#include <aki/file.h> +#include <nnwt/thread.h> +#include <nnwt/common.h> +#include <nnwt/file.h> #include "../backing.h" struct cch_backing_file { struct cch_backing backing; - struct aki_file file; + struct nn_file file; off_t size; off_t filesize; union { void *map; // CACHE_BACKING_MAPPED off_t pointer; // CACHE_BACKING_READ } u; - struct aki_mutex mutex; + struct nn_mutex mutex; }; struct cch_backing *cch_backing_file_create(str *path, size_t filesize); diff --git a/src/cache/backings/file_common.c b/src/cache/backings/file_common.c index 35322fd..3914f95 100644 --- a/src/cache/backings/file_common.c +++ b/src/cache/backings/file_common.c @@ -1,24 +1,24 @@ static off_t file_backing_get_size_estimate(struct cch_backing *backing) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_lock(&file->mutex); + nn_mutex_lock(&file->mutex); off_t size = file->size; - aki_mutex_unlock(&file->mutex); + nn_mutex_unlock(&file->mutex); return size; } static bool file_open_internal(struct cch_backing_file *file, str *path, size_t size) { - s32 flags = size != 0 ? AKI_FILE_CREATE : 0; - if (!aki_file_open(&file->file, path, flags)) return false; - file->filesize = aki_file_get_filesize(&file->file); + s32 flags = size != 0 ? NNWT_FILE_CREATE : 0; + if (!nn_file_open(&file->file, path, flags)) return false; + file->filesize = nn_file_get_filesize(&file->file); if (!size) { file->size = file->filesize; cch_backing_fill_range(&file->backing, 0, file->size); } else { file->size = -1; if (file->filesize < (off_t)size) { - aki_file_truncate(&file->file, size); + nn_file_truncate(&file->file, size); } } return true; diff --git a/src/cache/backings/file_mapped.c b/src/cache/backings/file_mapped.c index 22f9ae5..35ab267 100644 --- a/src/cache/backings/file_mapped.c +++ b/src/cache/backings/file_mapped.c @@ -6,25 +6,25 @@ static void file_backing_lock(struct cch_backing *backing) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_lock(&file->mutex); + nn_mutex_lock(&file->mutex); } static void file_backing_write(struct cch_backing *backing, u8 *buf, off_t index, size_t *size) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_lock(&file->mutex); + nn_mutex_lock(&file->mutex); if (file->size >= 0 && index + (off_t)*size >= file->size) { *size = MAX(file->size - index, 0L); } al_memcpy(file->u.map + index, buf, *size); cch_backing_fill_range(&file->backing, index, *size); - aki_mutex_unlock(&file->mutex); + nn_mutex_unlock(&file->mutex); } static u8 *file_backing_get_ptr(struct cch_backing *backing, off_t index, size_t *size) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_lock(&file->mutex); + nn_mutex_lock(&file->mutex); if (file->size >= 0 && index + (off_t)*size >= file->size) { *size = MAX(file->size - index, 0L); } @@ -34,28 +34,28 @@ static u8 *file_backing_get_ptr(struct cch_backing *backing, off_t index, size_t static void file_backing_unlock(struct cch_backing *backing) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_unlock(&file->mutex); + nn_mutex_unlock(&file->mutex); } static void file_backing_resize(struct cch_backing *backing, size_t size) { struct cch_backing_file *file = (struct cch_backing_file *)backing; - aki_mutex_lock(&file->mutex); + nn_mutex_lock(&file->mutex); file->size = size; if ((off_t)size > file->filesize) { - aki_file_munmap(&file->file, file->u.map); - aki_file_truncate(&file->file, size); - file->u.map = aki_file_mmap(&file->file); + nn_file_munmap(&file->file, file->u.map); + nn_file_truncate(&file->file, size); + file->u.map = nn_file_mmap(&file->file); file->filesize = size; } - aki_mutex_unlock(&file->mutex); + nn_mutex_unlock(&file->mutex); } static void file_backing_free(struct cch_backing **backing) { struct cch_backing_file *file = (struct cch_backing_file *)*backing; - aki_file_munmap(&file->file, file->u.map); - aki_file_close(&file->file); + nn_file_munmap(&file->file, file->u.map); + nn_file_close(&file->file); al_array_free(file->backing.available); al_free(file); *backing = NULL; @@ -74,10 +74,10 @@ struct cch_backing *cch_backing_file_create(str *path, size_t size) file->backing.resize = file_backing_resize; file->backing.free = file_backing_free; al_array_init(file->backing.available); - if (!file_open_internal(file, path, size) || !(file->u.map = aki_file_mmap(&file->file))) { + if (!file_open_internal(file, path, size) || !(file->u.map = nn_file_mmap(&file->file))) { al_free(file); return NULL; } - aki_mutex_init(&file->mutex); + nn_mutex_init(&file->mutex); return (struct cch_backing *)file; } diff --git a/src/cache/backings/memory.c b/src/cache/backings/memory.c index 5db2063..aa5c1b4 100644 --- a/src/cache/backings/memory.c +++ b/src/cache/backings/memory.c @@ -13,28 +13,28 @@ static void ensure_alloced(struct cch_backing_memory *mem, size_t size) static void memory_backing_lock(struct cch_backing *backing) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; - aki_mutex_lock(&mem->mutex); + nn_mutex_lock(&mem->mutex); } static void memory_backing_write(struct cch_backing *backing, u8 *buf, off_t index, size_t *size) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; - aki_mutex_lock(&mem->mutex); + nn_mutex_lock(&mem->mutex); if (mem->size >= 0 && index + (off_t)*size >= mem->size) { - *size = MAX(mem->size - index, 0L); + *size = MAX(mem->size - index, (off_t)0L); } ensure_alloced(mem, index + *size); al_memcpy(mem->data + index, buf, *size); cch_backing_fill_range(&mem->backing, index, *size); - aki_mutex_unlock(&mem->mutex); + nn_mutex_unlock(&mem->mutex); } static u8 *memory_backing_get_ptr(struct cch_backing *backing, off_t index, size_t *size) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; - aki_mutex_lock(&mem->mutex); + nn_mutex_lock(&mem->mutex); if (mem->size >= 0 && index + (off_t)*size >= mem->size) { - *size = MAX(mem->size - index, 0L); + *size = MAX(mem->size - index, (off_t)0L); } return (u8 *)(mem->data + index); } @@ -42,32 +42,32 @@ static u8 *memory_backing_get_ptr(struct cch_backing *backing, off_t index, size static void memory_backing_unlock(struct cch_backing *backing) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; - aki_mutex_unlock(&mem->mutex); + nn_mutex_unlock(&mem->mutex); } static off_t memory_backing_get_size_estimate(struct cch_backing *backing) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; - aki_mutex_lock(&mem->mutex); + nn_mutex_lock(&mem->mutex); off_t size = mem->size; - aki_mutex_unlock(&mem->mutex); + nn_mutex_unlock(&mem->mutex); return size; } static void memory_backing_resize(struct cch_backing *backing, size_t size) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; - aki_mutex_lock(&mem->mutex); + nn_mutex_lock(&mem->mutex); ensure_alloced(mem, size); mem->size = size; - aki_mutex_unlock(&mem->mutex); + nn_mutex_unlock(&mem->mutex); } static void memory_backing_free(struct cch_backing **backing) { struct cch_backing_memory *mem = (struct cch_backing_memory *)*backing; al_free(mem->data); - aki_mutex_destroy(&mem->mutex); + nn_mutex_destroy(&mem->mutex); al_array_free(mem->backing.available); al_free(mem); *backing = NULL; @@ -87,7 +87,7 @@ struct cch_backing *cch_backing_memory_create(size_t size) mem->data = al_malloc(size); mem->alloc = size; mem->size = -1; - aki_mutex_init(&mem->mutex); + nn_mutex_init(&mem->mutex); al_array_init(mem->backing.available); return (struct cch_backing *)mem; } diff --git a/src/cache/backings/memory.h b/src/cache/backings/memory.h index 107ae4b..e290b48 100644 --- a/src/cache/backings/memory.h +++ b/src/cache/backings/memory.h @@ -1,7 +1,7 @@ #pragma once #include <al/lib.h> -#include <aki/thread.h> +#include <nnwt/thread.h> #include "../backing.h" @@ -10,7 +10,7 @@ struct cch_backing_memory { u8 *data; size_t alloc; off_t size; - struct aki_mutex mutex; + struct nn_mutex mutex; }; struct cch_backing *cch_backing_memory_create(size_t size); diff --git a/src/cache/entry.c b/src/cache/entry.c index 8068397..c7fb8d2 100644 --- a/src/cache/entry.c +++ b/src/cache/entry.c @@ -7,12 +7,12 @@ bool cch_entry_get_handle(struct cch_entry *entry, struct cch_handle *handle) handle->entry = entry; handle->pointer = 0; handle->prev_pointer = 0; - aki_cond_init(&handle->wait.cond); - aki_mutex_init(&handle->wait.mutex); + nn_cond_init(&handle->wait.cond); + nn_mutex_init(&handle->wait.mutex); handle->wait.disabled = false; - aki_mutex_lock(&entry->mutex); + nn_mutex_lock(&entry->mutex); entry->ref_count++; - aki_mutex_unlock(&entry->mutex); + nn_mutex_unlock(&entry->mutex); return true; } @@ -43,15 +43,15 @@ void cch_entry_return_handle(struct cch_entry *entry, struct cch_handle *handle) { handle->entry = NULL; entry->ref_count--; - aki_mutex_destroy(&handle->wait.mutex); - aki_cond_destroy(&handle->wait.cond); + nn_mutex_destroy(&handle->wait.mutex); + nn_cond_destroy(&handle->wait.cond); } void cch_entry_free(struct cch_entry **entry) { (*entry)->handler->free(&(*entry)->handler); (*entry)->backing->free(&(*entry)->backing); - aki_mutex_destroy(&(*entry)->mutex); + nn_mutex_destroy(&(*entry)->mutex); al_free(*entry); *entry = NULL; } diff --git a/src/cache/entry.h b/src/cache/entry.h index 3196809..4a20058 100644 --- a/src/cache/entry.h +++ b/src/cache/entry.h @@ -19,7 +19,7 @@ struct cch_entry { struct cch_chapter *chapter; struct cch_backing *backing; struct cch_handler *handler; - struct aki_mutex mutex; + struct nn_mutex mutex; }; bool cch_entry_get_handle(struct cch_entry *entry, struct cch_handle *handle); diff --git a/src/cache/handle.h b/src/cache/handle.h index ae936f3..fc84109 100644 --- a/src/cache/handle.h +++ b/src/cache/handle.h @@ -1,7 +1,7 @@ #pragma once #include <al/types.h> -#include <aki/thread.h> +#include <nnwt/thread.h> #include "wait.h" diff --git a/src/cache/handler.h b/src/cache/handler.h index 9dd1896..b4c6975 100644 --- a/src/cache/handler.h +++ b/src/cache/handler.h @@ -3,7 +3,7 @@ #include <al/types.h> #include <al/array.h> #include <al/str.h> -#include <aki/thread.h> +#include <nnwt/thread.h> #include "handle.h" @@ -15,6 +15,6 @@ struct cch_handler { struct cch_entry *entry; str liana; bool disabled; - struct aki_mutex mutex; + struct nn_mutex mutex; array(struct cch_handler_wait *) waits; }; diff --git a/src/cache/handlers/cdio.c b/src/cache/handlers/cdio.c index 23f87b6..b10bbf1 100644 --- a/src/cache/handlers/cdio.c +++ b/src/cache/handlers/cdio.c @@ -31,7 +31,7 @@ static bool handler_cdio_can_seek(struct cch_handler *handler) return true; } -static aki_thread_result AKI_THREADCALL cd_read_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL cd_read_thread(void *userdata) { struct cch_handler_cdio *cdio = (struct cch_handler_cdio *)userdata; while (1) { @@ -41,7 +41,7 @@ static aki_thread_result AKI_THREADCALL cd_read_thread(void *userdata) } // cdio_read() here can block for a very long time so we buffer each read // before writing it to the backing. - long ret = cdio_cddap_read(cdio->drive, aki_buffer_get_ptr(&cdio->buffer, 0), cdio->sector, SECTORS_PER_STEP); + long ret = cdio_cddap_read(cdio->drive, nn_buffer_get_ptr(&cdio->buffer, 0), cdio->sector, SECTORS_PER_STEP); cdio_log_messages(cdio); al_assert(ret <= SECTORS_PER_STEP); off_t pointer = cdio->sector * CDIO_CD_FRAMESIZE_RAW; @@ -53,7 +53,7 @@ static aki_thread_result AKI_THREADCALL cd_read_thread(void *userdata) break; } if (ret == SECTORS_PER_STEP) { - al_memcpy(ptr, aki_buffer_get_ptr(&cdio->buffer, 0), BYTES_PER_STEP); + al_memcpy(ptr, nn_buffer_get_ptr(&cdio->buffer, 0), BYTES_PER_STEP); } else if (ret == -10) { al_log_error("cdio", "I/O error, skipping sector."); al_memset(ptr, 0, BYTES_PER_STEP); @@ -86,12 +86,12 @@ static void handler_cdio_maybe_spawn_worker(struct cch_handler *handler, size_t if (cdio->start == (lsn_t)index) return; if (al_atomic_load(s32)(&cdio->running, AL_ATOMIC_ACQUIRE)) { al_atomic_store(s32)(&cdio->running, 0, AL_ATOMIC_RELEASE); - aki_thread_join(&cdio->thread); + nn_thread_join(&cdio->thread); } cdio->start = index; cdio->sector = cdio->start; al_atomic_store(s32)(&cdio->running, 1, AL_ATOMIC_RELEASE); - aki_thread_create(&cdio->thread, cd_read_thread, cdio); + nn_thread_create(&cdio->thread, cd_read_thread, cdio); } static bool handler_cdio_wait_for_range(struct cch_handler *handler, struct cch_handler_wait *wait) @@ -106,7 +106,7 @@ static void handler_cdio_free(struct cch_handler **handler) cch_threaded_waits_disable_all(&cdio->handler); if (al_atomic_load(s32)(&cdio->running, AL_ATOMIC_RELAXED)) { al_atomic_store(s32)(&cdio->running, 0, AL_ATOMIC_RELAXED); - aki_thread_join(&cdio->thread); + nn_thread_join(&cdio->thread); } if (cdio->drive) { cdio_log_messages(cdio); @@ -189,7 +189,7 @@ struct cch_entry *cch_handler_cdio_create(void) entry->ref_count = 0; entry->unknown_size = false; entry->hash = 0; - aki_mutex_init(&entry->mutex); + nn_mutex_init(&entry->mutex); entry->handler->can_seek = handler_cdio_can_seek; entry->handler->maybe_spawn_worker = handler_cdio_maybe_spawn_worker; entry->handler->wait_for_range = handler_cdio_wait_for_range; @@ -198,9 +198,9 @@ struct cch_entry *cch_handler_cdio_create(void) al_str_from(&entry->handler->liana, "cdio"); if (!open_cd_drive(cdio)) return NULL; al_atomic_store(s32)(&cdio->running, 0, AL_ATOMIC_RELAXED); - aki_buffer_init(&cdio->buffer); - aki_buffer_ensure_space(&cdio->buffer, BYTES_PER_STEP); - aki_buffer_set_size(&cdio->buffer, BYTES_PER_STEP); + nn_buffer_init(&cdio->buffer); + nn_buffer_ensure_space(&cdio->buffer, BYTES_PER_STEP); + nn_buffer_set_size(&cdio->buffer, BYTES_PER_STEP); cch_threaded_waits_init(&cdio->handler); return entry; } diff --git a/src/cache/handlers/cdio.h b/src/cache/handlers/cdio.h index 953c988..3a3966a 100644 --- a/src/cache/handlers/cdio.h +++ b/src/cache/handlers/cdio.h @@ -1,5 +1,5 @@ #include <al/atomic.h> -#include <aki/common.h> +#include <nnwt/common.h> #include <cdio/paranoia/cdda.h> #include <cdio/cd_types.h> @@ -13,9 +13,9 @@ struct cch_handler_cdio { lsn_t start; lsn_t end; lsn_t sector; + struct nn_buffer buffer; atomic(s32) running; - struct aki_thread thread; - struct aki_buffer buffer; + struct nn_thread thread; }; struct cch_entry *cch_handler_cdio_create(void); diff --git a/src/cache/handlers/file.c b/src/cache/handlers/file.c index 03e2814..c09aac0 100644 --- a/src/cache/handlers/file.c +++ b/src/cache/handlers/file.c @@ -42,7 +42,7 @@ struct cch_entry *cch_handler_file_create(str *path) entry->ref_count = 0; entry->unknown_size = false; entry->hash = 0; - aki_mutex_init(&entry->mutex); + nn_mutex_init(&entry->mutex); entry->handler->can_seek = handler_file_can_seek; entry->handler->maybe_spawn_worker = handler_file_maybe_spawn_worker; entry->handler->wait_for_range = handler_file_wait_for_range; diff --git a/src/cache/handlers/http.c b/src/cache/handlers/http.c index 71ae186..308633e 100644 --- a/src/cache/handlers/http.c +++ b/src/cache/handlers/http.c @@ -21,39 +21,37 @@ static size_t http_callback(void *userdata, u8 op, u8 *buf, s64 int0) struct cch_handler_http *http = (struct cch_handler_http *)userdata; size_t ret = (size_t)int0; switch (op) { - case AKI_HTTP_FINISHED: { - al_log_debug("cache_handler_http", "Transfer finished."); - break; - } - case AKI_HTTP_ERROR: { - cch_threaded_waits_disable_all(&http->handler); + case NNWT_HTTP_READ: + al_assert_and_return(0); + case NNWT_HTTP_WRITE: { + // If we got data before a Content-Length, assume the size is unknown. + if (cch_entry_get_size(http->handler.entry) < 0) { + cch_entry_set_size(http->handler.entry, 0); + } + http->backing->write(http->backing, buf, http->pointer, &ret); + http->pointer += ret; + cch_threaded_waits_evaluate(&http->handler, http->backing); break; } - case AKI_HTTP_RESPONSE_CODE: { + case NNWT_HTTP_RESPONSE_CODE: al_log_debug("cache_handler_http", "HTTP %ld.", int0); break; - } - case AKI_HTTP_REDIRECT: { - al_log_debug("cache_handler_http", "Redirect."); - break; - } - case AKI_HTTP_CONTENT_LENGTH: { + case NNWT_HTTP_CONTENT_LENGTH: { off_t content_length = (off_t)int0; al_log_debug("cache_handler_http", "Content-Length: %lld.", content_length); cch_entry_set_size(http->handler.entry, content_length); cch_threaded_waits_signal_any(&http->handler); break; } - case AKI_HTTP_WRITE: { - // If we got data before a Content-Length, assume the size is unknown. - if (cch_entry_get_size(http->handler.entry) < 0) { - cch_entry_set_size(http->handler.entry, 0); - } - http->backing->write(http->backing, buf, http->pointer, &ret); - http->pointer += ret; - cch_threaded_waits_evaluate(&http->handler, http->backing); + case NNWT_HTTP_REDIRECT: + al_log_debug("cache_handler_http", "Redirect %ld.", int0); + break; + case NNWT_HTTP_FINISHED: + al_log_debug("cache_handler_http", "Transfer finished."); + break; + case NNWT_HTTP_ERROR: + cch_threaded_waits_disable_all(&http->handler); break; - } default: break; } @@ -64,12 +62,12 @@ static void handler_http_maybe_spawn_worker(struct cch_handler *handler, size_t { struct cch_handler_http *http = (struct cch_handler_http *)handler; (void)index; - al_array_push(http->requests, (struct aki_http){ 0 }); - struct aki_http *request = &al_array_last(http->requests); - aki_http_init(request); - aki_http_set_url(request, &http->url); - aki_http_set_user_agent(request, USER_AGENT); - aki_http_request_stream(request, AKI_HTTP_GET, http->loop, http_callback, http); + al_array_push(http->requests, (struct nn_http){ 0 }); + struct nn_http *request = &al_array_last(http->requests); + nn_http_init(request); + nn_http_set_url(request, &http->url); + nn_http_set_user_agent(request, USER_AGENT); + nn_http_request_stream(request, NNWT_HTTP_GET, http->loop, http_callback, http); } static bool handler_http_wait_for_range(struct cch_handler *handler, struct cch_handler_wait *wait) @@ -82,9 +80,9 @@ static void handler_http_free(struct cch_handler **handler) { struct cch_handler_http *http = (struct cch_handler_http *)*handler; cch_threaded_waits_disable_all(&http->handler); - struct aki_http *request; + struct nn_http *request; al_array_foreach_ptr(http->requests, i, request) { - aki_http_close(request); + nn_http_close(request); } cch_threaded_waits_close(&http->handler); al_array_free(http->requests); @@ -94,7 +92,7 @@ static void handler_http_free(struct cch_handler **handler) *handler = NULL; } -struct cch_entry *cch_handler_http_create(str *url, struct aki_event_loop *loop) +struct cch_entry *cch_handler_http_create(str *url, struct nn_event_loop *loop) { struct cch_backing *backing = cch_backing_memory_create(1024 * 128); if (!backing) return NULL; @@ -106,7 +104,7 @@ struct cch_entry *cch_handler_http_create(str *url, struct aki_event_loop *loop) entry->ref_count = 0; entry->unknown_size = false; entry->hash = 0; - aki_mutex_init(&entry->mutex); + nn_mutex_init(&entry->mutex); entry->handler->can_seek = handler_http_can_seek; entry->handler->maybe_spawn_worker = handler_http_maybe_spawn_worker; entry->handler->wait_for_range = handler_http_wait_for_range; diff --git a/src/cache/handlers/http.h b/src/cache/handlers/http.h index 5f9df86..ef47173 100644 --- a/src/cache/handlers/http.h +++ b/src/cache/handlers/http.h @@ -1,8 +1,8 @@ #pragma once -#include <aki/http.h> -#include <aki/thread.h> -#include <aki/signal.h> +#include <nnwt/http.h> +#include <nnwt/thread.h> +#include <nnwt/signal.h> #include "../handler.h" #include "../entry.h" @@ -12,8 +12,8 @@ struct cch_handler_http { struct cch_backing *backing; str url; off_t pointer; - array(struct aki_http) requests; - struct aki_event_loop *loop; + array(struct nn_http) requests; + struct nn_event_loop *loop; }; -struct cch_entry *cch_handler_http_create(str *url, struct aki_event_loop *loop); +struct cch_entry *cch_handler_http_create(str *url, struct nn_event_loop *loop); diff --git a/src/cache/meson.build b/src/cache/meson.build index 5fb47ae..5b70c3c 100644 --- a/src/cache/meson.build +++ b/src/cache/meson.build @@ -6,6 +6,7 @@ cache_src = [ 'backings/memory.c' ] cache_deps = [] +cache_args = [] cache_have_cdio = false libcdio_paranoia = dependency('libcdio_paranoia', required: false, allow_fallback: true) @@ -13,6 +14,7 @@ libcdio_cdda = dependency('libcdio_cdda', required: false, allow_fallback: true) if libcdio_paranoia.found() and libcdio_cdda.found() cache_src += ['handlers/cdio.c'] cache_deps += [libcdio_paranoia, libcdio_cdda] + cache_args += ['-DCACHE_HAVE_CDIO'] cache_have_cdio = true endif @@ -22,14 +24,15 @@ endif #if libdvdcss.found() and libdvdread.found() and libdvdnav.found() #endif -if akiyo_has_mmap +if naunet_has_mmap cache_src += ['backings/file_mapped.c'] else cache_src += ['backings/file.c'] endif -if akiyo_has_curl +if naunet_has_curl cache_src += ['handlers/http.c'] endif -cache = declare_dependency(sources: cache_src, dependencies: cache_deps) +cache = declare_dependency(sources: cache_src, dependencies: cache_deps, + compile_args: cache_args) diff --git a/src/cache/threaded_waits.c b/src/cache/threaded_waits.c index 10f08f0..3bf67eb 100644 --- a/src/cache/threaded_waits.c +++ b/src/cache/threaded_waits.c @@ -4,7 +4,7 @@ void cch_threaded_waits_init(struct cch_handler *handler) { handler->disabled = false; - aki_mutex_init(&handler->mutex); + nn_mutex_init(&handler->mutex); al_array_init(handler->waits); } @@ -22,104 +22,98 @@ static bool wait_range_satisfied(struct cch_backing *backing, struct cch_handler bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing, struct cch_handler_wait *wait) { - aki_mutex_lock(&handler->mutex); + nn_mutex_lock(&handler->mutex); bool canceled = handler->disabled; off_t size = cch_entry_get_size(handler->entry); backing->lock(backing); bool satisfied = size >= 0 && wait_range_satisfied(backing, wait); backing->unlock(backing); if (!canceled && !satisfied) { - aki_mutex_lock(&wait->mutex); + nn_mutex_lock(&wait->mutex); if (!wait->disabled) { al_array_push(handler->waits, wait); - aki_mutex_unlock(&handler->mutex); - aki_cond_wait(&wait->cond, &wait->mutex); + nn_mutex_unlock(&handler->mutex); + nn_cond_wait(&wait->cond, &wait->mutex); } else { - aki_mutex_unlock(&wait->mutex); - aki_mutex_unlock(&handler->mutex); + nn_mutex_unlock(&wait->mutex); + nn_mutex_unlock(&handler->mutex); return false; } - aki_mutex_unlock(&wait->mutex); - aki_mutex_lock(&handler->mutex); + nn_mutex_unlock(&wait->mutex); + nn_mutex_lock(&handler->mutex); if (wait->disabled) { - struct cch_handler_wait *rwait; - al_array_foreach(handler->waits, i, rwait) { - if (rwait == wait) { - al_array_remove_at(handler->waits, i); - break; - } - } + al_array_remove(handler->waits, wait); canceled = true; } } - aki_mutex_unlock(&handler->mutex); + nn_mutex_unlock(&handler->mutex); return !canceled; } void cch_threaded_waits_signal_any(struct cch_handler *handler) { - aki_mutex_lock(&handler->mutex); + nn_mutex_lock(&handler->mutex); struct cch_handler_wait *wait; al_array_foreach_rev(handler->waits, i, wait) { if (wait->end < 0) { al_array_remove_at(handler->waits, i); - aki_mutex_lock(&wait->mutex); - aki_cond_signal(&wait->cond); - aki_mutex_unlock(&wait->mutex); + nn_mutex_lock(&wait->mutex); + nn_cond_signal(&wait->cond); + nn_mutex_unlock(&wait->mutex); } } - aki_mutex_unlock(&handler->mutex); + nn_mutex_unlock(&handler->mutex); } void cch_threaded_wait_disable(struct cch_handler_wait *wait) { - aki_mutex_lock(&wait->mutex); + nn_mutex_lock(&wait->mutex); wait->disabled = true; // cond_is_waiting() could be false if the handler is disabled. - if (aki_cond_is_waiting(&wait->cond)) { - aki_cond_signal(&wait->cond); + if (nn_cond_is_waiting(&wait->cond)) { + nn_cond_signal(&wait->cond); } - aki_mutex_unlock(&wait->mutex); + nn_mutex_unlock(&wait->mutex); } void cch_threaded_waits_disable_all(struct cch_handler *handler) { - aki_mutex_lock(&handler->mutex); + nn_mutex_lock(&handler->mutex); if (handler->disabled) { - aki_mutex_unlock(&handler->mutex); + nn_mutex_unlock(&handler->mutex); return; } // Disallow any further waits. handler->disabled = true; struct cch_handler_wait *wait; al_array_foreach(handler->waits, i, wait) { - aki_mutex_lock(&wait->mutex); + nn_mutex_lock(&wait->mutex); wait->disabled = true; - aki_cond_signal(&wait->cond); - aki_mutex_unlock(&wait->mutex); + nn_cond_signal(&wait->cond); + nn_mutex_unlock(&wait->mutex); } - aki_mutex_unlock(&handler->mutex); + nn_mutex_unlock(&handler->mutex); } void cch_threaded_waits_evaluate(struct cch_handler *handler, struct cch_backing *backing) { - aki_mutex_lock(&handler->mutex); + nn_mutex_lock(&handler->mutex); backing->lock(backing); struct cch_handler_wait *wait; al_array_foreach_rev(handler->waits, i, wait) { - aki_mutex_lock(&wait->mutex); + nn_mutex_lock(&wait->mutex); if (wait_range_satisfied(backing, wait)) { - aki_cond_signal(&wait->cond); + nn_cond_signal(&wait->cond); al_array_remove_at(handler->waits, i); } - aki_mutex_unlock(&wait->mutex); + nn_mutex_unlock(&wait->mutex); } backing->unlock(backing); - aki_mutex_unlock(&handler->mutex); + nn_mutex_unlock(&handler->mutex); } void cch_threaded_waits_close(struct cch_handler *handler) { - aki_mutex_destroy(&handler->mutex); + nn_mutex_destroy(&handler->mutex); al_array_free(handler->waits); } diff --git a/src/cache/wait.h b/src/cache/wait.h index 8654362..4c9b228 100644 --- a/src/cache/wait.h +++ b/src/cache/wait.h @@ -1,10 +1,10 @@ #pragma once -#include <aki/thread.h> +#include <nnwt/thread.h> struct cch_handler_wait { off_t start, end; - struct aki_cond cond; - struct aki_mutex mutex; + struct nn_cond cond; + struct nn_mutex mutex; bool disabled; }; diff --git a/src/codec/codec.h b/src/codec/codec.h index a9b188a..697f456 100644 --- a/src/codec/codec.h +++ b/src/codec/codec.h @@ -1,7 +1,7 @@ #pragma once #include <al/array.h> -#include <aki/common.h> +#include <nnwt/common.h> #ifdef CAMU_HAVE_FFMPEG #include <libavutil/pixdesc.h> #include <libavformat/avformat.h> @@ -141,7 +141,7 @@ struct camu_codec_stream { struct camu_codec_packet { u8 mode; - struct aki_buffer *buffer; + struct nn_buffer *buffer; #ifdef CAMU_HAVE_FFMPEG struct { AVPacket *pkt; } av; #endif diff --git a/src/codec/ffmpeg/common.c b/src/codec/ffmpeg/common.c index c147d81..16bbeab 100644 --- a/src/codec/ffmpeg/common.c +++ b/src/codec/ffmpeg/common.c @@ -22,7 +22,7 @@ static void av_log_callback(void *userdata, int level, const char *fmt, va_list if (level < AV_LOG_DEBUG) { pos += al_vsnprintf(&buf[pos], CHUNK_SIZE, fmt, args); if (buf[pos - 1] == '\n' || pos >= CHUNK_SIZE) { - al_log_info("ff_log", buf); + al_log_info("ff", buf); pos = 0; } } diff --git a/src/codec/ffmpeg/encoder.c b/src/codec/ffmpeg/encoder.c index 6e914c8..7c7b010 100644 --- a/src/codec/ffmpeg/encoder.c +++ b/src/codec/ffmpeg/encoder.c @@ -99,7 +99,7 @@ void camu_ff_encoder_push(struct camu_ff_encoder *enc, u8 **data, s32 sample_cou return; } AVPacket pkt = { 0 }; - do { + for (;;) { ret = avcodec_receive_packet(enc->codec_context, &pkt); if (ret < 0) { if (!(ret = AVERROR_EOF || ret == AVERROR(EAGAIN))) { @@ -109,7 +109,7 @@ void camu_ff_encoder_push(struct camu_ff_encoder *enc, u8 **data, s32 sample_cou } enc->callback(enc->userdata, &pkt); enc->frame->pts += SAMPLES_PER_FRAME; - } while (1); + } } void camu_ff_encoder_reset(struct camu_ff_encoder *enc) diff --git a/src/codec/ffmpeg/packet_ext.c b/src/codec/ffmpeg/packet_ext.c index 16f5748..759719b 100644 --- a/src/codec/ffmpeg/packet_ext.c +++ b/src/codec/ffmpeg/packet_ext.c @@ -1,201 +1,201 @@ #include "packet_ext.h" -void aki_packet_write_av_codec_parameters(struct aki_packet *packet, AVCodecParameters *codecpar) +void nn_packet_write_av_codec_parameters(struct nn_packet *packet, AVCodecParameters *codecpar) { - AKI_PACKET_WRITE_TYPE(packet, enum AVMediaType, codecpar->codec_type); - AKI_PACKET_WRITE_TYPE(packet, enum AVCodecID, codecpar->codec_id); - AKI_PACKET_WRITE_TYPE(packet, u32, codecpar->codec_tag); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->extradata_size); - AKI_PACKET_WRITE_DATA(packet, codecpar->extradata, codecpar->extradata_size); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->format); - AKI_PACKET_WRITE_TYPE(packet, s64, codecpar->bit_rate); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->bits_per_coded_sample); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->bits_per_raw_sample); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->profile); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->level); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->width); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->height); - AKI_PACKET_WRITE_TYPE(packet, AVRational, codecpar->sample_aspect_ratio); - AKI_PACKET_WRITE_TYPE(packet, AVRational, codecpar->framerate); - AKI_PACKET_WRITE_TYPE(packet, enum AVFieldOrder, codecpar->field_order); - AKI_PACKET_WRITE_TYPE(packet, enum AVColorRange, codecpar->color_range); - AKI_PACKET_WRITE_TYPE(packet, enum AVColorPrimaries, codecpar->color_primaries); - AKI_PACKET_WRITE_TYPE(packet, enum AVColorTransferCharacteristic, codecpar->color_trc); - AKI_PACKET_WRITE_TYPE(packet, enum AVColorSpace, codecpar->color_space); - AKI_PACKET_WRITE_TYPE(packet, enum AVChromaLocation, codecpar->chroma_location); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->video_delay); - AKI_PACKET_WRITE_TYPE(packet, AVChannelLayout, codecpar->ch_layout); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->sample_rate); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->block_align); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->frame_size); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->initial_padding); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->trailing_padding); - AKI_PACKET_WRITE_TYPE(packet, s32, codecpar->seek_preroll); + NNWT_PACKET_WRITE_TYPE(packet, enum AVMediaType, codecpar->codec_type); + NNWT_PACKET_WRITE_TYPE(packet, enum AVCodecID, codecpar->codec_id); + NNWT_PACKET_WRITE_TYPE(packet, u32, codecpar->codec_tag); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->extradata_size); + NNWT_PACKET_WRITE_DATA(packet, codecpar->extradata, codecpar->extradata_size); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->format); + NNWT_PACKET_WRITE_TYPE(packet, s64, codecpar->bit_rate); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->bits_per_coded_sample); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->bits_per_raw_sample); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->profile); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->level); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->width); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->height); + NNWT_PACKET_WRITE_TYPE(packet, AVRational, codecpar->sample_aspect_ratio); + NNWT_PACKET_WRITE_TYPE(packet, AVRational, codecpar->framerate); + NNWT_PACKET_WRITE_TYPE(packet, enum AVFieldOrder, codecpar->field_order); + NNWT_PACKET_WRITE_TYPE(packet, enum AVColorRange, codecpar->color_range); + NNWT_PACKET_WRITE_TYPE(packet, enum AVColorPrimaries, codecpar->color_primaries); + NNWT_PACKET_WRITE_TYPE(packet, enum AVColorTransferCharacteristic, codecpar->color_trc); + NNWT_PACKET_WRITE_TYPE(packet, enum AVColorSpace, codecpar->color_space); + NNWT_PACKET_WRITE_TYPE(packet, enum AVChromaLocation, codecpar->chroma_location); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->video_delay); + NNWT_PACKET_WRITE_TYPE(packet, AVChannelLayout, codecpar->ch_layout); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->sample_rate); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->block_align); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->frame_size); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->initial_padding); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->trailing_padding); + NNWT_PACKET_WRITE_TYPE(packet, s32, codecpar->seek_preroll); } -void aki_packet_write_av_codec_id(struct aki_packet *packet, enum AVCodecID codec_id) +void nn_packet_write_av_codec_id(struct nn_packet *packet, enum AVCodecID codec_id) { - AKI_PACKET_WRITE_TYPE(packet, enum AVCodecID, codec_id); + NNWT_PACKET_WRITE_TYPE(packet, enum AVCodecID, codec_id); } -void aki_packet_write_av_dictionary(struct aki_packet *packet, AVDictionary *dict) +void nn_packet_write_av_dictionary(struct nn_packet *packet, AVDictionary *dict) { s32 count = av_dict_count(dict); - AKI_PACKET_WRITE_TYPE(packet, s32, count); + NNWT_PACKET_WRITE_TYPE(packet, s32, count); const AVDictionaryEntry *entry = NULL; while ((entry = av_dict_iterate(dict, entry))) { size_t len = al_strlen(entry->key); al_assert(len <= UINT32_MAX); - AKI_PACKET_WRITE_TYPE(packet, u32, len); - AKI_PACKET_WRITE_DATA(packet, entry->key, len); + NNWT_PACKET_WRITE_TYPE(packet, u32, len); + NNWT_PACKET_WRITE_DATA(packet, entry->key, len); len = al_strlen(entry->value); al_assert(len <= UINT32_MAX); - AKI_PACKET_WRITE_TYPE(packet, u32, len); - AKI_PACKET_WRITE_DATA(packet, entry->value, len); + NNWT_PACKET_WRITE_TYPE(packet, u32, len); + NNWT_PACKET_WRITE_DATA(packet, entry->value, len); } } -void aki_packet_write_av_stream(struct aki_packet *packet, AVStream *stream) +void nn_packet_write_av_stream(struct nn_packet *packet, AVStream *stream) { - AKI_PACKET_WRITE_TYPE(packet, s32, stream->index); - aki_packet_write_av_codec_parameters(packet, stream->codecpar); - AKI_PACKET_WRITE_TYPE(packet, AVRational, stream->time_base); - AKI_PACKET_WRITE_TYPE(packet, s64, stream->duration); - AKI_PACKET_WRITE_TYPE(packet, s64, stream->start_time); - AKI_PACKET_WRITE_TYPE(packet, s64, stream->nb_frames); - aki_packet_write_av_dictionary(packet, stream->metadata); + NNWT_PACKET_WRITE_TYPE(packet, s32, stream->index); + nn_packet_write_av_codec_parameters(packet, stream->codecpar); + NNWT_PACKET_WRITE_TYPE(packet, AVRational, stream->time_base); + NNWT_PACKET_WRITE_TYPE(packet, s64, stream->duration); + NNWT_PACKET_WRITE_TYPE(packet, s64, stream->start_time); + NNWT_PACKET_WRITE_TYPE(packet, s64, stream->nb_frames); + nn_packet_write_av_dictionary(packet, stream->metadata); if (stream->avg_frame_rate.den == 0) { - // r_frame_rate.den == 0 handled on the client. - AKI_PACKET_WRITE_TYPE(packet, AVRational, stream->r_frame_rate); + // r_frame_rate.den = 0 handled on the client. + NNWT_PACKET_WRITE_TYPE(packet, AVRational, stream->r_frame_rate); } else { - AKI_PACKET_WRITE_TYPE(packet, AVRational, stream->avg_frame_rate); + NNWT_PACKET_WRITE_TYPE(packet, AVRational, stream->avg_frame_rate); } } -void aki_packet_write_av_packet(struct aki_packet *packet, AVPacket *pkt) +void nn_packet_write_av_packet(struct nn_packet *packet, AVPacket *pkt) { - AKI_PACKET_WRITE_TYPE(packet, s64, pkt->pts); - AKI_PACKET_WRITE_TYPE(packet, s64, pkt->dts); - AKI_PACKET_WRITE_TYPE(packet, s32, pkt->size); - AKI_PACKET_WRITE_DATA(packet, pkt->data, pkt->size); - AKI_PACKET_WRITE_TYPE(packet, s32, pkt->stream_index); - AKI_PACKET_WRITE_TYPE(packet, s32, pkt->flags); - AKI_PACKET_WRITE_TYPE(packet, s32, pkt->side_data_elems); + NNWT_PACKET_WRITE_TYPE(packet, s64, pkt->pts); + NNWT_PACKET_WRITE_TYPE(packet, s64, pkt->dts); + NNWT_PACKET_WRITE_TYPE(packet, s32, pkt->size); + NNWT_PACKET_WRITE_DATA(packet, pkt->data, pkt->size); + NNWT_PACKET_WRITE_TYPE(packet, s32, pkt->stream_index); + NNWT_PACKET_WRITE_TYPE(packet, s32, pkt->flags); + NNWT_PACKET_WRITE_TYPE(packet, s32, pkt->side_data_elems); for (s32 i = 0; i < pkt->side_data_elems; i++) { // Possibly upcast size_t. s64 size = (s64)pkt->side_data[i].size; - AKI_PACKET_WRITE_TYPE(packet, s64, size); - AKI_PACKET_WRITE_DATA(packet, pkt->side_data[i].data, pkt->side_data[i].size); - AKI_PACKET_WRITE_TYPE(packet, enum AVPacketSideDataType, pkt->side_data[i].type); + NNWT_PACKET_WRITE_TYPE(packet, s64, size); + NNWT_PACKET_WRITE_DATA(packet, pkt->side_data[i].data, pkt->side_data[i].size); + NNWT_PACKET_WRITE_TYPE(packet, enum AVPacketSideDataType, pkt->side_data[i].type); } - AKI_PACKET_WRITE_TYPE(packet, s64, pkt->duration); - AKI_PACKET_WRITE_TYPE(packet, s64, pkt->pos); - AKI_PACKET_WRITE_TYPE(packet, AVRational, pkt->time_base); + NNWT_PACKET_WRITE_TYPE(packet, s64, pkt->duration); + NNWT_PACKET_WRITE_TYPE(packet, s64, pkt->pos); + NNWT_PACKET_WRITE_TYPE(packet, AVRational, pkt->time_base); } -void aki_packet_read_av_codec_parameters(struct aki_packet *packet, AVCodecParameters *codecpar) +void nn_packet_read_av_codec_parameters(struct nn_packet *packet, AVCodecParameters *codecpar) { - AKI_PACKET_READ_TYPE(packet, enum AVMediaType, codecpar->codec_type); - AKI_PACKET_READ_TYPE(packet, enum AVCodecID, codecpar->codec_id); - AKI_PACKET_READ_TYPE(packet, u32, codecpar->codec_tag); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->extradata_size); + NNWT_PACKET_READ_TYPE(packet, enum AVMediaType, codecpar->codec_type); + NNWT_PACKET_READ_TYPE(packet, enum AVCodecID, codecpar->codec_id); + NNWT_PACKET_READ_TYPE(packet, u32, codecpar->codec_tag); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->extradata_size); codecpar->extradata = av_malloc(codecpar->extradata_size); u8 *extradata; - AKI_PACKET_READ_DATA(packet, codecpar->extradata_size, extradata); + NNWT_PACKET_READ_DATA(packet, codecpar->extradata_size, extradata); al_memcpy(codecpar->extradata, extradata, codecpar->extradata_size); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->format); - AKI_PACKET_READ_TYPE(packet, s64, codecpar->bit_rate); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->bits_per_coded_sample); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->bits_per_raw_sample); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->profile); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->level); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->width); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->height); - AKI_PACKET_READ_TYPE(packet, AVRational, codecpar->sample_aspect_ratio); - AKI_PACKET_READ_TYPE(packet, AVRational, codecpar->framerate); - AKI_PACKET_READ_TYPE(packet, enum AVFieldOrder, codecpar->field_order); - AKI_PACKET_READ_TYPE(packet, enum AVColorRange, codecpar->color_range); - AKI_PACKET_READ_TYPE(packet, enum AVColorPrimaries, codecpar->color_primaries); - AKI_PACKET_READ_TYPE(packet, enum AVColorTransferCharacteristic, codecpar->color_trc); - AKI_PACKET_READ_TYPE(packet, enum AVColorSpace, codecpar->color_space); - AKI_PACKET_READ_TYPE(packet, enum AVChromaLocation, codecpar->chroma_location); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->video_delay); - AKI_PACKET_READ_TYPE(packet, AVChannelLayout, codecpar->ch_layout); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->sample_rate); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->block_align); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->frame_size); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->initial_padding); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->trailing_padding); - AKI_PACKET_READ_TYPE(packet, s32, codecpar->seek_preroll); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->format); + NNWT_PACKET_READ_TYPE(packet, s64, codecpar->bit_rate); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->bits_per_coded_sample); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->bits_per_raw_sample); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->profile); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->level); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->width); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->height); + NNWT_PACKET_READ_TYPE(packet, AVRational, codecpar->sample_aspect_ratio); + NNWT_PACKET_READ_TYPE(packet, AVRational, codecpar->framerate); + NNWT_PACKET_READ_TYPE(packet, enum AVFieldOrder, codecpar->field_order); + NNWT_PACKET_READ_TYPE(packet, enum AVColorRange, codecpar->color_range); + NNWT_PACKET_READ_TYPE(packet, enum AVColorPrimaries, codecpar->color_primaries); + NNWT_PACKET_READ_TYPE(packet, enum AVColorTransferCharacteristic, codecpar->color_trc); + NNWT_PACKET_READ_TYPE(packet, enum AVColorSpace, codecpar->color_space); + NNWT_PACKET_READ_TYPE(packet, enum AVChromaLocation, codecpar->chroma_location); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->video_delay); + NNWT_PACKET_READ_TYPE(packet, AVChannelLayout, codecpar->ch_layout); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->sample_rate); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->block_align); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->frame_size); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->initial_padding); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->trailing_padding); + NNWT_PACKET_READ_TYPE(packet, s32, codecpar->seek_preroll); } -enum AVCodecID aki_packet_read_av_codec_id(struct aki_packet *packet) +enum AVCodecID nn_packet_read_av_codec_id(struct nn_packet *packet) { enum AVCodecID id; - AKI_PACKET_READ_TYPE(packet, enum AVCodecID, id); + NNWT_PACKET_READ_TYPE(packet, enum AVCodecID, id); return id; } -void aki_packet_read_av_dictionary(struct aki_packet *packet, AVDictionary **dict) +void nn_packet_read_av_dictionary(struct nn_packet *packet, AVDictionary **dict) { s32 count; - AKI_PACKET_READ_TYPE(packet, s32, count); + NNWT_PACKET_READ_TYPE(packet, s32, count); for (s32 i = 0; i < count; i++) { char *key, *value; u32 len; - AKI_PACKET_READ_TYPE(packet, u32, len); - AKI_PACKET_READ_DATA(packet, len, key); + NNWT_PACKET_READ_TYPE(packet, u32, len); + NNWT_PACKET_READ_DATA(packet, len, key); // Key and value MUST be allocated with av_malloc functions // just like all FFmpeg structures. key = av_strndup(key, len); - AKI_PACKET_READ_TYPE(packet, u32, len); - AKI_PACKET_READ_DATA(packet, len, value); + NNWT_PACKET_READ_TYPE(packet, u32, len); + NNWT_PACKET_READ_DATA(packet, len, value); value = av_strndup(value, len); av_dict_set(dict, key, value, AV_DICT_DONT_STRDUP_KEY | AV_DICT_DONT_STRDUP_VAL); } } -AVStream *aki_packet_read_av_stream(AVFormatContext *format_context, const AVCodec *codec, struct aki_packet *packet) +AVStream *nn_packet_read_av_stream(AVFormatContext *format_context, const AVCodec *codec, struct nn_packet *packet) { AVStream *stream = avformat_new_stream(format_context, codec); - AKI_PACKET_READ_TYPE(packet, s32, stream->index); - aki_packet_read_av_codec_parameters(packet, stream->codecpar); - AKI_PACKET_READ_TYPE(packet, AVRational, stream->time_base); - AKI_PACKET_READ_TYPE(packet, s64, stream->duration); - AKI_PACKET_READ_TYPE(packet, s64, stream->start_time); - AKI_PACKET_READ_TYPE(packet, s64, stream->nb_frames); + NNWT_PACKET_READ_TYPE(packet, s32, stream->index); + nn_packet_read_av_codec_parameters(packet, stream->codecpar); + NNWT_PACKET_READ_TYPE(packet, AVRational, stream->time_base); + NNWT_PACKET_READ_TYPE(packet, s64, stream->duration); + NNWT_PACKET_READ_TYPE(packet, s64, stream->start_time); + NNWT_PACKET_READ_TYPE(packet, s64, stream->nb_frames); stream->metadata = NULL; - aki_packet_read_av_dictionary(packet, &stream->metadata); - AKI_PACKET_READ_TYPE(packet, AVRational, stream->avg_frame_rate); + nn_packet_read_av_dictionary(packet, &stream->metadata); + NNWT_PACKET_READ_TYPE(packet, AVRational, stream->avg_frame_rate); return stream; } -AVPacket *aki_packet_read_av_packet(struct aki_packet *packet) +AVPacket *nn_packet_read_av_packet(struct nn_packet *packet) { AVPacket *pkt = av_packet_alloc(); - AKI_PACKET_READ_TYPE(packet, s64, pkt->pts); - AKI_PACKET_READ_TYPE(packet, s64, pkt->dts); - AKI_PACKET_READ_TYPE(packet, s32, pkt->size); - AKI_PACKET_READ_DATA(packet, pkt->size, pkt->data); - AKI_PACKET_READ_TYPE(packet, s32, pkt->stream_index); - AKI_PACKET_READ_TYPE(packet, s32, pkt->flags); + NNWT_PACKET_READ_TYPE(packet, s64, pkt->pts); + NNWT_PACKET_READ_TYPE(packet, s64, pkt->dts); + NNWT_PACKET_READ_TYPE(packet, s32, pkt->size); + NNWT_PACKET_READ_DATA(packet, pkt->size, pkt->data); + NNWT_PACKET_READ_TYPE(packet, s32, pkt->stream_index); + NNWT_PACKET_READ_TYPE(packet, s32, pkt->flags); s32 side_data_elems; - AKI_PACKET_READ_TYPE(packet, s32, side_data_elems); + NNWT_PACKET_READ_TYPE(packet, s32, side_data_elems); pkt->side_data_elems = 0; for (s32 i = 0; i < side_data_elems; i++) { u8 *data; enum AVPacketSideDataType type; s64 storage_for_size; - AKI_PACKET_READ_TYPE(packet, s64, storage_for_size); + NNWT_PACKET_READ_TYPE(packet, s64, storage_for_size); al_assert((sizeof(size_t) > 4 || storage_for_size <= INT32_MAX)); size_t size = (size_t)storage_for_size; - AKI_PACKET_READ_DATA(packet, size, data); - AKI_PACKET_READ_TYPE(packet, enum AVPacketSideDataType, type); + NNWT_PACKET_READ_DATA(packet, size, data); + NNWT_PACKET_READ_TYPE(packet, enum AVPacketSideDataType, type); u8 *side_data = av_packet_new_side_data(pkt, type, size); al_memcpy(side_data, data, size); } - AKI_PACKET_READ_TYPE(packet, s64, pkt->duration); - AKI_PACKET_READ_TYPE(packet, s64, pkt->pos); - AKI_PACKET_READ_TYPE(packet, AVRational, pkt->time_base); + NNWT_PACKET_READ_TYPE(packet, s64, pkt->duration); + NNWT_PACKET_READ_TYPE(packet, s64, pkt->pos); + NNWT_PACKET_READ_TYPE(packet, AVRational, pkt->time_base); return pkt; } diff --git a/src/codec/ffmpeg/packet_ext.h b/src/codec/ffmpeg/packet_ext.h index 26a1392..91254c2 100644 --- a/src/codec/ffmpeg/packet_ext.h +++ b/src/codec/ffmpeg/packet_ext.h @@ -1,19 +1,17 @@ #pragma once -#include <aki/common.h> -#include <aki/event_loop.h> - +#include <nnwt/packet.h> #include <libavcodec/packet.h> #include <libavcodec/codec_par.h> #include <libavformat/avformat.h> #include <libavutil/dict.h> -void aki_packet_write_av_codec_parameters(struct aki_packet *packet, AVCodecParameters *codecpar); -void aki_packet_write_av_codec_id(struct aki_packet *packet, enum AVCodecID codec_id); -void aki_packet_write_av_stream(struct aki_packet *packet, AVStream *stream); -void aki_packet_write_av_packet(struct aki_packet *packet, AVPacket *pkt); +void nn_packet_write_av_codec_parameters(struct nn_packet *packet, AVCodecParameters *codecpar); +void nn_packet_write_av_codec_id(struct nn_packet *packet, enum AVCodecID codec_id); +void nn_packet_write_av_stream(struct nn_packet *packet, AVStream *stream); +void nn_packet_write_av_packet(struct nn_packet *packet, AVPacket *pkt); -void aki_packet_read_av_codec_parameters(struct aki_packet *packet, AVCodecParameters *codecpar); -enum AVCodecID aki_packet_read_av_codec_id(struct aki_packet *packet); -AVStream *aki_packet_read_av_stream(AVFormatContext *format_context, const AVCodec *codec, struct aki_packet *packet); -AVPacket *aki_packet_read_av_packet(struct aki_packet *packet); +void nn_packet_read_av_codec_parameters(struct nn_packet *packet, AVCodecParameters *codecpar); +enum AVCodecID nn_packet_read_av_codec_id(struct nn_packet *packet); +AVStream *nn_packet_read_av_stream(AVFormatContext *format_context, const AVCodec *codec, struct nn_packet *packet); +AVPacket *nn_packet_read_av_packet(struct nn_packet *packet); diff --git a/src/codec/spng/demuxer.h b/src/codec/spng/demuxer.h index 1947bb7..bd89ef1 100644 --- a/src/codec/spng/demuxer.h +++ b/src/codec/spng/demuxer.h @@ -5,7 +5,7 @@ struct camu_spng_demuxer { struct camu_demuxer demux; struct cch_handle *handle; - struct aki_buffer buffer; + struct nn_buffer buffer; bool eof; }; diff --git a/src/codec/spng/impl.c b/src/codec/spng/impl.c index 198d514..58b68f6 100644 --- a/src/codec/spng/impl.c +++ b/src/codec/spng/impl.c @@ -1,5 +1,5 @@ #include <al/lib.h> -#include <aki/common.h> +#include <nnwt/common.h> #include <spng.h> #include "../../cache/entry.h" @@ -14,17 +14,17 @@ static bool spng_demuxer_init(struct camu_demuxer *demux, struct cch_handle *han struct camu_spng_demuxer *spng = (struct camu_spng_demuxer *)demux; spng->handle = handle; - aki_buffer_init(&spng->buffer); + nn_buffer_init(&spng->buffer); spng->eof = false; cch_handle_seek(spng->handle, 0, SEEK_SET); off_t size = cch_entry_get_size(spng->handle->entry); - aki_buffer_ensure_space(&spng->buffer, size); - cch_handle_read(spng->handle, aki_buffer_get_ptr(&spng->buffer, 0), size); + nn_buffer_ensure_space(&spng->buffer, size); + cch_handle_read(spng->handle, nn_buffer_get_ptr(&spng->buffer, 0), size); spng->buffer.size = size; spng_ctx *ctx = spng_ctx_new(SPNG_CTX_IGNORE_ADLER32); - spng_set_png_buffer(ctx, aki_buffer_get_ptr(&spng->buffer, 0), spng->buffer.size); + spng_set_png_buffer(ctx, nn_buffer_get_ptr(&spng->buffer, 0), spng->buffer.size); spng_set_option(ctx, SPNG_CHUNK_COUNT_LIMIT, 6250); spng_set_crc_action(ctx, SPNG_CRC_DISCARD, SPNG_CRC_DISCARD); @@ -73,7 +73,7 @@ static bool spng_demuxer_seek(struct camu_demuxer *demux, u64 pos) static void spng_demuxer_free(struct camu_demuxer **demux) { struct camu_spng_demuxer *spng = (struct camu_spng_demuxer *)*demux; - aki_buffer_free(&spng->buffer); + nn_buffer_free(&spng->buffer); al_array_free(spng->demux.streams); al_free(spng); *demux = NULL; @@ -96,7 +96,7 @@ static s32 spng_decoder_push(struct camu_decoder *dec, struct camu_codec_packet if (!packet) return CAMU_OK; spng_ctx *ctx = spng_ctx_new(SPNG_CTX_IGNORE_ADLER32); - spng_set_png_buffer(ctx, aki_buffer_get_ptr(packet->buffer, 0), packet->buffer->size); + spng_set_png_buffer(ctx, nn_buffer_get_ptr(packet->buffer, 0), packet->buffer->size); spng_set_option(ctx, SPNG_CHUNK_COUNT_LIMIT, 6250); spng_set_crc_action(ctx, SPNG_CRC_DISCARD, SPNG_CRC_DISCARD); diff --git a/src/codec/stb_image/demuxer.h b/src/codec/stb_image/demuxer.h index 52f9256..d64e311 100644 --- a/src/codec/stb_image/demuxer.h +++ b/src/codec/stb_image/demuxer.h @@ -5,7 +5,7 @@ struct camu_stbi_demuxer { struct camu_demuxer demux; struct cch_handle *handle; - struct aki_buffer buffer; + struct nn_buffer buffer; bool eof; }; diff --git a/src/codec/stb_image/impl.c b/src/codec/stb_image/impl.c index ac1cbc5..bb38844 100644 --- a/src/codec/stb_image/impl.c +++ b/src/codec/stb_image/impl.c @@ -1,5 +1,5 @@ #include <al/lib.h> -#include <aki/common.h> +#include <nnwt/common.h> #include "../common.h" #define STB_IMAGE_IMPLEMENTATION #define STB_IMAGE_STATIC @@ -22,17 +22,17 @@ static bool stbi_demuxer_init(struct camu_demuxer *demux, struct cch_handle *han struct camu_stbi_demuxer *stb = (struct camu_stbi_demuxer *)demux; stb->handle = handle; - aki_buffer_init(&stb->buffer); + nn_buffer_init(&stb->buffer); stb->eof = false; cch_handle_seek(stb->handle, 0, SEEK_SET); off_t size = cch_entry_get_size(stb->handle->entry); - aki_buffer_ensure_space(&stb->buffer, size); - cch_handle_read(stb->handle, aki_buffer_get_ptr(&stb->buffer, 0), size); + nn_buffer_ensure_space(&stb->buffer, size); + cch_handle_read(stb->handle, nn_buffer_get_ptr(&stb->buffer, 0), size); stb->buffer.size = size; s32 w, h, channels; - u8 *data = aki_buffer_get_ptr(&stb->buffer, 0); + u8 *data = nn_buffer_get_ptr(&stb->buffer, 0); if (!stbi_info_from_memory(data, stb->buffer.size, &w, &h, &channels)) { goto err; } @@ -80,7 +80,7 @@ static bool stbi_demuxer_seek(struct camu_demuxer *demux, u64 pos) static void stbi_demuxer_free(struct camu_demuxer **demux) { struct camu_stbi_demuxer *stb = (struct camu_stbi_demuxer *)*demux; - aki_buffer_free(&stb->buffer); + nn_buffer_free(&stb->buffer); al_array_free(stb->demux.streams); al_free(stb); *demux = NULL; @@ -104,7 +104,7 @@ static s32 stbi_decoder_push(struct camu_decoder *dec, struct camu_codec_packet if (!packet) return CAMU_OK; s32 w, h, channels; - u8 *data = aki_buffer_get_ptr(packet->buffer, 0); + u8 *data = nn_buffer_get_ptr(packet->buffer, 0); u8 *pixels = stbi_load_from_memory(data, packet->buffer->size, &w, &h, &channels, 0); if (!pixels) return CAMU_ERR_ERROR; struct camu_video_format *fmt = &stb->stream->video.fmt; diff --git a/src/codec/wuffs/demuxer.h b/src/codec/wuffs/demuxer.h index 6e9e391..c2cb8a5 100644 --- a/src/codec/wuffs/demuxer.h +++ b/src/codec/wuffs/demuxer.h @@ -5,7 +5,7 @@ struct camu_wuffs_demuxer { struct camu_demuxer demux; struct cch_handle *handle; - struct aki_buffer buffer; + struct nn_buffer buffer; bool eof; }; diff --git a/src/codec/wuffs/impl.c b/src/codec/wuffs/impl.c index 4d9ea7e..378ac00 100644 --- a/src/codec/wuffs/impl.c +++ b/src/codec/wuffs/impl.c @@ -1,4 +1,4 @@ -#include <aki/common.h> +#include <nnwt/common.h> #define WUFFS_IMPLEMENTATION #define WUFFS_CONFIG__STATIC_FUNCTIONS #define WUFFS_CONFIG__MODULES @@ -22,13 +22,13 @@ static bool wuffs_demuxer_init(struct camu_demuxer *demux, struct cch_handle *ha struct camu_wuffs_demuxer *wfs = (struct camu_wuffs_demuxer *)demux; wfs->handle = handle; - aki_buffer_init(&wfs->buffer); + nn_buffer_init(&wfs->buffer); wfs->eof = false; cch_handle_seek(wfs->handle, 0, SEEK_SET); off_t size = cch_entry_get_size(wfs->handle->entry); - aki_buffer_ensure_space(&wfs->buffer, size); - cch_handle_read(wfs->handle, aki_buffer_get_ptr(&wfs->buffer, 0), size); + nn_buffer_ensure_space(&wfs->buffer, size); + cch_handle_read(wfs->handle, nn_buffer_get_ptr(&wfs->buffer, 0), size); wfs->buffer.size = size; wuffs_png__decoder wf_dec = { 0 }; @@ -38,7 +38,7 @@ static bool wuffs_demuxer_init(struct camu_demuxer *demux, struct cch_handle *ha wuffs_png__decoder__set_quirk(&wf_dec, WUFFS_BASE__QUIRK_IGNORE_CHECKSUM, true); wuffs_base__io_buffer src = { 0 }; - src.data.ptr = aki_buffer_get_ptr(&wfs->buffer, 0); + src.data.ptr = nn_buffer_get_ptr(&wfs->buffer, 0); src.data.len = wfs->buffer.size; src.meta.wi = wfs->buffer.size; src.meta.closed = true; @@ -90,7 +90,7 @@ static bool wuffs_demuxer_seek(struct camu_demuxer *demux, u64 pos) static void wuffs_demuxer_free(struct camu_demuxer **demux) { struct camu_wuffs_demuxer *wfs = (struct camu_wuffs_demuxer *)*demux; - aki_buffer_free(&wfs->buffer); + nn_buffer_free(&wfs->buffer); al_array_free(wfs->demux.streams); al_free(wfs); *demux = NULL; @@ -127,7 +127,7 @@ static s32 wuffs_decoder_push(struct camu_decoder *dec, struct camu_codec_packet wuffs_png__decoder__set_quirk(&wf_dec, WUFFS_BASE__QUIRK_IGNORE_CHECKSUM, true); wuffs_base__io_buffer src = { 0 }; - src.data.ptr = aki_buffer_get_ptr(packet->buffer, 0); + src.data.ptr = nn_buffer_get_ptr(packet->buffer, 0); src.data.len = packet->buffer->size; src.meta.wi = packet->buffer->size; src.meta.closed = true; diff --git a/src/fruits/cmc/cmc.c b/src/fruits/cmc/cmc.c index 3c62ba8..35f2d9b 100644 --- a/src/fruits/cmc/cmc.c +++ b/src/fruits/cmc/cmc.c @@ -1,8 +1,11 @@ -#include <aki/common.h> -#include <aki/event_loop.h> +#include <nnwt/common.h> +#include <nnwt/event_loop.h> +#include <nnwt/packet.h> #include "../../portal/src/packet_ext.h" -#include "../../server/common.c" +#include "../../server/common.h" + +#include "../common.h" #include "cmc.h" @@ -23,16 +26,16 @@ static struct cmc_search *get_search_by_id(struct cmc *c, s32 id) return NULL; } -static void parse_user_state(struct cmc *c, struct aki_packet *packet) +static void parse_user_state(struct cmc *c, struct nn_packet *packet) { - u32 size = aki_packet_read_u32(packet); + u32 size = nn_packet_read_u32(packet); for (u32 i = 0; i < size; i++) { struct cmc_search *search = al_alloc_object(struct cmc_search); - search->id = aki_packet_read_s32(packet); + search->id = nn_packet_read_s32(packet); str s; - aki_packet_read_str(packet, &s); + nn_packet_read_str(packet, &s); al_str_clone(&search->module, &s); - aki_packet_read_str(packet, &s); + nn_packet_read_str(packet, &s); al_str_clone(&search->query, &s); al_array_init(search->pages); al_array_push(c->searches, search); @@ -44,12 +47,12 @@ static void client_callback(void *userdata, u8 op, void *opaque) struct cmc *c = (struct cmc *)userdata; switch (op) { case CAMU_CLIENT_LOGIN: { - struct aki_packet *packet = (struct aki_packet *)opaque; + struct nn_packet *packet = (struct nn_packet *)opaque; switch (c->command) { case CLI_OPEN_UI: parse_user_state(c, packet); if (!cmc_ui_init(&c->ui, &c->loop, c)) { - aki_event_loop_break_one(&c->loop); + nn_event_loop_break_one(&c->loop); } break; case CLI_ADD: { @@ -82,22 +85,22 @@ static void client_callback(void *userdata, u8 op, void *opaque) break; } case CAMU_CLIENT_PAGE_RESULTS: { - struct aki_packet *packet = (struct aki_packet *)opaque; - s32 id = aki_packet_read_s32(packet); + struct nn_packet *packet = (struct nn_packet *)opaque; + s32 id = nn_packet_read_s32(packet); struct cmc_search *search = get_search_by_id(c, id); if (!search) return; struct cmc_search_page *page = al_alloc_object(struct cmc_search_page); - page->num = aki_packet_read_u32(packet); - u32 posts = aki_packet_read_u32(packet); + page->num = nn_packet_read_u32(packet); + u32 posts = nn_packet_read_u32(packet); for (u32 i = 0; i < posts; i++) { struct camu_post post; - aki_packet_read_post(packet, &post); + nn_packet_read_post(packet, &post); camu_post_cache_push(&c->cache, &post); } - u32 ids = aki_packet_read_u32(packet); + u32 ids = nn_packet_read_u32(packet); for (u32 i = 0; i < ids; i++) { str unique_id; - aki_packet_read_str(packet, &unique_id); + nn_packet_read_str(packet, &unique_id); str s; al_str_clone(&s, &unique_id); al_array_push(page->list, s); @@ -154,9 +157,9 @@ s32 main(s32 argc, char *argv[]) s32 wmain(s32 argc, wchar_t **argv) #endif { - if (!aki_common_init()) return EXIT_FAILURE; + if (!nn_common_init()) return EXIT_FAILURE; - aki_event_loop_init(&c.loop); + nn_event_loop_init(&c.loop); al_array_init(c.args); if (!parse_cmd(argc, argv)) return EXIT_FAILURE; @@ -166,11 +169,11 @@ s32 wmain(s32 argc, wchar_t **argv) c.client.callback = client_callback; c.client.userdata = &c; - if (!camu_client_login(&c.client, &c.loop, CAMU_LOCAL_TYPE, CAMU_LOCAL_ADDR, CAMU_PORT, al_str_c("andrew"))) { + if (!camu_client_login(&c.client, &c.loop, CAMU_TEST_TYPE, CAMU_TEST_ADDR, CAMU_PORT, al_str_c("andrew"))) { return EXIT_FAILURE; } - aki_event_loop_run(&c.loop); + nn_event_loop_run(&c.loop); cmc_ui_close(&c.ui); diff --git a/src/fruits/cmc/cmc.h b/src/fruits/cmc/cmc.h index 353bd87..96de27b 100644 --- a/src/fruits/cmc/cmc.h +++ b/src/fruits/cmc/cmc.h @@ -22,7 +22,7 @@ struct cmc_search { }; struct cmc { - struct aki_event_loop loop; + struct nn_event_loop loop; u8 command; array(str) args; struct camu_client client; diff --git a/src/fruits/cmc/ui/ui.c b/src/fruits/cmc/ui/ui.c index 58d630f..b4c196d 100644 --- a/src/fruits/cmc/ui/ui.c +++ b/src/fruits/cmc/ui/ui.c @@ -27,7 +27,7 @@ static void input_poll_callback(void *userdata, s32 revents) if (input.evtype == NCTYPE_PRESS || input.evtype == NCTYPE_UNKNOWN) { switch (input.id) { case 'q': - aki_event_loop_break_one(ui->loop); + nn_event_loop_break_one(ui->loop); break; } } @@ -44,7 +44,7 @@ static void input_poll_callback(void *userdata, s32 revents) } while (1); } -bool cmc_ui_init(struct cmc_ui *ui, struct aki_event_loop *loop, struct cmc *c) +bool cmc_ui_init(struct cmc_ui *ui, struct nn_event_loop *loop, struct cmc *c) { al_memset(ui, 0, sizeof(struct cmc_ui)); ui->c = c; @@ -57,9 +57,9 @@ bool cmc_ui_init(struct cmc_ui *ui, struct aki_event_loop *loop, struct cmc *c) struct ncplane *stdplane = notcurses_stdplane(ui->nc); ncplane_set_resizecb(stdplane, resize_cb); ncplane_set_userptr(stdplane, ui); - aki_poll_init(&ui->input_poll, input_poll_callback, ui); - aki_poll_set(&ui->input_poll, notcurses_inputready_fd(ui->nc), AKI_POLL_READ); - aki_poll_start(&ui->input_poll, ui->loop); + nn_poll_init(&ui->input_poll, input_poll_callback, ui); + nn_poll_set(&ui->input_poll, notcurses_inputready_fd(ui->nc), NNWT_POLL_READ); + nn_poll_start(&ui->input_poll, ui->loop); ui->pane = CMC_PANE_SEARCH; return true; } diff --git a/src/fruits/cmc/ui/ui.h b/src/fruits/cmc/ui/ui.h index 1a49985..1c3b9d9 100644 --- a/src/fruits/cmc/ui/ui.h +++ b/src/fruits/cmc/ui/ui.h @@ -1,7 +1,7 @@ #pragma once #include <al/types.h> -#include <aki/event_loop.h> +#include <nnwt/event_loop.h> #include <notcurses/notcurses.h> enum { @@ -11,11 +11,11 @@ enum { struct cmc; struct cmc_ui { struct notcurses *nc; - struct aki_event_loop *loop; + struct nn_event_loop *loop; u32 term_cols; u32 term_rows; bool pending_layout; - struct aki_poll input_poll; + struct nn_poll input_poll; u8 pane; struct { struct ncplane *n; @@ -24,7 +24,7 @@ struct cmc_ui { struct cmc *c; }; -bool cmc_ui_init(struct cmc_ui *ui, struct aki_event_loop *loop, struct cmc *c); +bool cmc_ui_init(struct cmc_ui *ui, struct nn_event_loop *loop, struct cmc *c); void cmc_ui_close(struct cmc_ui *ui); void cmc_sp_init(struct cmc_ui *ui); diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c index 5eb9fee..53555d3 100644 --- a/src/fruits/cmsrv/cmsrv.c +++ b/src/fruits/cmsrv/cmsrv.c @@ -1,32 +1,34 @@ #define CAMU_LOCAL_SOCKET #include <al/log.h> -#include <aki/common.h> +#include <nnwt/common.h> #ifdef CAMU_LOCAL_SOCKET -#include <aki/line_processor.h> +#include <nnwt/line_processor.h> #endif -#include <aki/timer.h> +#include <nnwt/timer.h> #include "../../server/server.h" #include "../../server/common.h" #include "../../cache/handlers/cdio.h" #include "../../codec/ffmpeg/common.h" +#include "../common.h" + #include "ui.h" struct cmsrv { - struct aki_event_loop loop; - struct aki_signal quit_signal; + struct nn_event_loop loop; + struct nn_signal quit_signal; struct camu_server server; #ifdef CAMU_LOCAL_SOCKET struct { - struct aki_socket sock; - struct aki_line_processor cli; + struct nn_socket sock; + struct nn_line_processor cli; } local; #endif struct cmsrv_ui ui; - struct aki_poll input_poll; - struct aki_timer render_timer; + struct nn_poll input_poll; + struct nn_timer render_timer; }; #ifdef CAMU_LOCAL_SOCKET @@ -49,25 +51,34 @@ static u8 server_line_callback(void *userdata, str *line) } else if (al_str_eq(line, al_str_c(";CLEAR"))) { lia_list_clear(list); } else { - struct aki_packet *packet = aki_packet_create(); + struct nn_packet *packet = nn_packet_create(); +#ifdef CAMU_HAVE_PORTAL if (al_str_at(line, 0) == ';' || camu_is_url(line, 0)) { - aki_packet_write_u8(packet, CAMU_RESOURCE_SIMPLE_SEARCH); + nn_packet_write_u8(packet, CAMU_RESOURCE_SIMPLE_SEARCH); +#elif defined NAUNET_HAS_CURL + if (camu_is_url(line, 0)) { + nn_packet_write_u8(packet, CAMU_RESOURCE_HTTP); +#endif +#if CACHE_HAVE_CDIO + } else if (al_str_cmp(line, al_str_c("cdda://"), 0, 7) == 0) { + nn_packet_write_u8(packet, CAMU_RESOURCE_CDIO); +#endif } else { - aki_packet_write_u8(packet, CAMU_RESOURCE_FILE); + nn_packet_write_u8(packet, CAMU_RESOURCE_FILE); } - aki_packet_write_str(packet, line); + nn_packet_write_str(packet, line); camu_server_local_add(&s->server, packet); } - return AKI_LINE_PROCESSOR_CONTINUE; + return NNWT_LINE_PROCESSOR_CONTINUE; } #endif -static void render_timer_callback(void *userdata, struct aki_timer *timer) +static void render_timer_callback(void *userdata, struct nn_timer *timer) { struct cmsrv *s = (struct cmsrv *)userdata; (void)timer; cmsrv_ui_render(&s->ui); - aki_timer_again(&s->render_timer); + nn_timer_again(&s->render_timer); } static void input_poll_callback(void *userdata, s32 revents) @@ -80,7 +91,7 @@ static void input_poll_callback(void *userdata, s32 revents) if (input.evtype == NCTYPE_PRESS || input.evtype == NCTYPE_UNKNOWN) { switch (input.id) { case 'q': - aki_event_loop_break_one(&s->loop); + nn_event_loop_break_one(&s->loop); break; case '1': cmsrv_ui_set_pane(&s->ui, CMSRV_UI_LISTS); @@ -99,13 +110,13 @@ static void input_poll_callback(void *userdata, s32 revents) static void quit_signal_callback(void *userdata) { struct cmsrv *s = (struct cmsrv *)userdata; - aki_event_loop_break_one(&s->loop); + nn_event_loop_break_one(&s->loop); } -static s32 log_callback(void *userdata, char *message, const char *color) +static s32 log_callback(void *userdata, u8 level, char *message) { struct cmsrv *s = (struct cmsrv *)userdata; - (void)color; + (void)level; cmsrv_ui_push_message(&s->ui, message); return al_strlen(message); } @@ -115,7 +126,7 @@ static struct cmsrv s = { 0 }; static void sigint_handler(s32 signum) { (void)signum; - aki_signal_send(&s.quit_signal); + nn_signal_send(&s.quit_signal); } #ifndef _WIN32 @@ -126,7 +137,7 @@ s32 wmain(s32 argc, wchar_t **argv) { (void)argc; (void)argv; - if (!aki_common_init() || !cmsrv_ui_init(&s.ui, &s.server)) return EXIT_FAILURE; + if (!nn_common_init() || !cmsrv_ui_init(&s.ui, &s.server)) return EXIT_FAILURE; signal(SIGINT, sigint_handler); @@ -136,36 +147,35 @@ s32 wmain(s32 argc, wchar_t **argv) camu_ff_set_default_log_callback(); #endif - aki_event_loop_init(&s.loop); + nn_event_loop_init(&s.loop); - aki_signal_init(&s.quit_signal, &s.loop, quit_signal_callback, &s); - aki_signal_start(&s.quit_signal); + nn_signal_init(&s.quit_signal, &s.loop, quit_signal_callback, &s); + nn_signal_start(&s.quit_signal); - if (!camu_server_init(&s.server, CAMU_LOCAL_TYPE, &s.loop)) return EXIT_FAILURE; - camu_server_listen(&s.server, CAMU_LOCAL_ADDR, CAMU_PORT); + if (!camu_server_init(&s.server, CAMU_TEST_TYPE, &s.loop)) return EXIT_FAILURE; + camu_server_listen(&s.server, CAMU_TEST_ADDR, CAMU_PORT); #ifdef CAMU_LOCAL_SOCKET - s.local.sock.type = AKI_SOCKET_UNIX; - aki_socket_init(&s.local.sock); - aki_socket_set_blocking(&s.local.sock, false); + s.local.sock.type = NNWT_SOCKET_UNIX; + nn_socket_init(&s.local.sock, NNWT_SOCKET_NONBLOCKING); s.local.cli.callback = server_line_callback; s.local.cli.userdata = &s; - aki_line_processor_init(&s.local.cli, al_str_c("\n")); - aki_line_processor_open_socket(&s.local.cli, &s.local.sock); - if (aki_socket_bind(&s.local.sock, CAMU_UNIX_LOCAL, 0) && aki_socket_listen(&s.local.sock)) { - aki_line_processor_run(&s.local.cli, &s.loop); + nn_line_processor_init(&s.local.cli, al_str_c("\n")); + nn_line_processor_open_socket(&s.local.cli, &s.local.sock); + if (nn_socket_bind(&s.local.sock, CAMU_TEST_CONTROL_PATH, 0) && nn_socket_listen(&s.local.sock)) { + nn_line_processor_run(&s.local.cli, &s.loop); } #endif - aki_poll_init(&s.input_poll, input_poll_callback, &s); - aki_poll_set(&s.input_poll, cmsrv_ui_get_input_fd(&s.ui), AKI_POLL_READ); - aki_poll_start(&s.input_poll, &s.loop); + nn_poll_init(&s.input_poll, input_poll_callback, &s); + nn_poll_set(&s.input_poll, cmsrv_ui_get_input_fd(&s.ui), NNWT_POLL_READ); + nn_poll_start(&s.input_poll, &s.loop); - aki_timer_init(&s.render_timer, &s.loop, render_timer_callback, &s); - aki_timer_set_repeat(&s.render_timer, AKI_TS_FROM_USEC(100000)); - aki_timer_again(&s.render_timer); + nn_timer_init(&s.render_timer, &s.loop, render_timer_callback, &s); + nn_timer_set_repeat(&s.render_timer, NNWT_TS_FROM_USEC(100000)); + nn_timer_again(&s.render_timer); - aki_event_loop_run(&s.loop); + nn_event_loop_run(&s.loop); cmsrv_ui_close(&s.ui); @@ -173,7 +183,7 @@ s32 wmain(s32 argc, wchar_t **argv) camu_ff_free_default_log_callback(); #endif - aki_common_close(); + nn_common_close(); return EXIT_SUCCESS; } diff --git a/src/fruits/cmsrv/ui.c b/src/fruits/cmsrv/ui.c index 0c26881..5d35380 100644 --- a/src/fruits/cmsrv/ui.c +++ b/src/fruits/cmsrv/ui.c @@ -1,6 +1,8 @@ #include <al/lib.h> #include <al/log.h> +#include "../../server/common.h" + #include "ui.h" #define LOG_RATIO 1.5 @@ -145,7 +147,15 @@ static void render_lists(struct cmsrv_ui *ui) ncplane_putchar_yx(n, y, 1, '>'); x += 2; } - putnwstr_maxwidth_yx(n, y, x, max_width - x, &entry->name); + wstr *name = &entry->name; +#ifdef CAMU_HAVE_PORTAL + struct camu_resource *resource = (struct camu_resource *)entry->opaque; + if (resource->type == CAMU_RESOURCE_PORTAL) { + struct camu_resource_portal *portal = (struct camu_resource_portal *)resource; + name = &portal->post->title; + } +#endif + putnwstr_maxwidth_yx(n, y, x, max_width - x, name); if (++current_line >= max_height) break; } } diff --git a/src/fruits/cmv/cmv.c b/src/fruits/cmv/cmv.c index 0fe33ad..8509485 100644 --- a/src/fruits/cmv/cmv.c +++ b/src/fruits/cmv/cmv.c @@ -1,18 +1,21 @@ #include <al/log.h> -#include <aki/event_loop.h> -#include <aki/thread.h> +#include <nnwt/event_loop.h> +#include <nnwt/thread.h> #include "../../sink/desktop.h" #include "../../server/common.h" #include "../../codec/ffmpeg/common.h" #ifndef CAMU_SINK_ONLY #include "../../server/server.h" -#else -#include "../../server/common.c" #endif +#include "../../server/common.h" + +#include "../common.h" + +static str *CMV_UNIX_PATH = al_str_c("/tmp/cmv_sock"); struct cmv { - struct aki_event_loop loop; + struct nn_event_loop loop; struct camu_desktop desktop; #ifndef CAMU_SINK_ONLY struct camu_server server; @@ -28,10 +31,10 @@ static void exit_callback(void *userdata, struct camu_desktop *desktop) } #endif -static aki_thread_result AKI_THREADCALL event_loop_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL event_loop_thread(void *userdata) { struct cmv *c = (struct cmv *)userdata; - aki_event_loop_run(&c->loop); + nn_event_loop_run(&c->loop); return 0; } @@ -53,7 +56,7 @@ s32 main(s32 argc, char *argv[]) s32 wmain(s32 argc, wchar_t **argv) #endif { - if (!aki_common_init() || !stl_global_init(false)) { + if (!nn_common_init() || !stl_global_init(false)) { return EXIT_FAILURE; } @@ -65,22 +68,26 @@ s32 wmain(s32 argc, wchar_t **argv) camu_ff_set_default_log_callback(); #endif - aki_event_loop_init(&c.loop); + nn_event_loop_init(&c.loop); u8 type; str *addr; #ifndef CAMU_SINK_ONLY bool local = argc > 1; if (local) { - type = AKI_SOCKET_UNIX; addr = CAMU_UNIX_LOCAL; +#ifndef _WIN32 + type = NNWT_SOCKET_UNIX; addr = CMV_UNIX_PATH; +#else + type = NNWT_SOCKET_TCP; addr = CAMU_LOCALHOST; +#endif if (!camu_server_init(&c.server, type, &c.loop)) failure(); if (!camu_server_listen(&c.server, addr, CAMU_PORT)) failure(); } else { - type = CAMU_LOCAL_TYPE; - addr = CAMU_LOCAL_ADDR; + type = CAMU_TEST_TYPE; + addr = CAMU_TEST_ADDR; } #else - type = AKI_SOCKET_TCP; addr = CAMU_SERVER_IP; + type = NNWT_SOCKET_TCP; addr = CAMU_SERVER_IP; #endif #ifndef CAMU_SINK_ONLY @@ -90,18 +97,29 @@ s32 wmain(s32 argc, wchar_t **argv) str arg = *al_str_cr(argv[i]); #else str arg; - al_wstr_to_str(al_wstr_cr(argv[i]), &arg); + if (!al_wstr_to_str(al_wstr_cr(argv[i]), &arg)) { + al_log_error("cmv", "Failed to parse argument #%i.", i); + continue; + } +#endif + struct nn_packet *packet = nn_packet_create(); +#ifdef CAMU_HAVE_PORTAL + if (al_str_at(&arg, 0) == ';' || camu_is_url(&arg, 0)) { + nn_packet_write_u8(packet, CAMU_RESOURCE_SIMPLE_SEARCH); +#elif defined NAUNET_HAS_CURL + if (camu_is_url(&arg, 0)) { + nn_packet_write_u8(packet, CAMU_RESOURCE_HTTP); +#endif +#if CACHE_HAVE_CDIO + } else if (al_str_cmp(&arg, al_str_c("cdda://"), 0, 7) == 0) { + nn_packet_write_u8(packet, CAMU_RESOURCE_CDIO); #endif - struct aki_packet *packet = aki_packet_create(); - if (al_str_at(&arg, 0) == ';') { - aki_packet_write_u8(packet, CAMU_RESOURCE_PORTAL); - aki_packet_write_str(packet, al_str_substr(&arg, 1, arg.len)); } else { - aki_packet_write_u8(packet, CAMU_RESOURCE_FILE); - aki_packet_write_str(packet, &arg); + nn_packet_write_u8(packet, CAMU_RESOURCE_FILE); } + nn_packet_write_str(packet, &arg); camu_server_local_add(&c.server, packet); - aki_packet_free(packet); + nn_packet_free(packet); } } #else @@ -118,13 +136,13 @@ s32 wmain(s32 argc, wchar_t **argv) #endif if (!camu_desktop_connect(&c.desktop, type, &c.loop, addr, CAMU_PORT)) failure(); - struct aki_thread thread0; - aki_thread_create(&thread0, event_loop_thread, &c); + struct nn_thread thread0; + nn_thread_create(&thread0, event_loop_thread, &c); while (camu_desktop_tick(&c.desktop)) {} camu_desktop_stop(&c.desktop); - aki_thread_join(&thread0); + nn_thread_join(&thread0); camu_desktop_free(&c.desktop); #ifndef CAMU_SINK_ONLY camu_server_free(&c.server); @@ -132,14 +150,14 @@ s32 wmain(s32 argc, wchar_t **argv) ret = EXIT_SUCCESS; out: - aki_event_loop_destroy(&c.loop); + nn_event_loop_destroy(&c.loop); #ifdef CAMU_HAVE_FFMPEG camu_ff_free_default_log_callback(); #endif stl_global_close(); - aki_common_close(); + nn_common_close(); return ret; } diff --git a/src/fruits/common.h b/src/fruits/common.h new file mode 100644 index 0000000..c7b432a --- /dev/null +++ b/src/fruits/common.h @@ -0,0 +1,17 @@ +#include <al/str.h> +#include <nnwt/socket.h> + +AL_UNUSED_VARIABLE_PUSH + +static str *CAMU_DB_PATH = al_str_c("/home/andrew/c/camu/data/camu_db_test"); + +static str *CAMU_TEST_IP = al_str_c("108.52.160.112"); +static str *CAMU_LOCALHOST = al_str_c("127.0.0.1"); + +static str *CAMU_TEST_PATH = al_str_c("/tmp/camu_sock"); +static str *CAMU_TEST_CONTROL_PATH = al_str_c("/tmp/camu_control_sock"); + +AL_UNUSED_VARIABLE_POP + +#define CAMU_TEST_TYPE NNWT_SOCKET_TCP +#define CAMU_TEST_ADDR CAMU_TEST_IP diff --git a/src/fruits/ctv/ctv.c b/src/fruits/ctv/ctv.c index 5f4a819..c7bd250 100644 --- a/src/fruits/ctv/ctv.c +++ b/src/fruits/ctv/ctv.c @@ -11,12 +11,12 @@ #include "../../server/common.c" struct ctv { - struct aki_event_loop loop; + struct nn_event_loop loop; struct camu_screen scr; struct camu_renderer *renderer; struct camu_mixer mixer; struct camu_sink sink; - struct aki_thread thread; + struct nn_thread thread; bool created; }; @@ -37,10 +37,10 @@ static void screen_callback(void *userdata, u8 op, void *opaque) (void)opaque; } -static aki_thread_result AKI_THREADCALL event_loop_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL event_loop_thread(void *userdata) { struct ctv *c = (struct ctv *)userdata; - aki_event_loop_run(&c->loop); + nn_event_loop_run(&c->loop); return 0; } @@ -66,14 +66,14 @@ static void onSurfaceCreated(GLFMDisplay *display, s32 width, s32 height) c->mixer.audio->configure_stream(c->mixer.audio, NULL); camu_mixer_pick_format(&c->mixer, &c->mixer.fmt); - aki_event_loop_init(&c->loop); + nn_event_loop_init(&c->loop); camu_sink_init(&c->sink, &c->loop, &c->mixer, c->renderer); c->sink.callback = sink_callback; c->sink.userdata = c; - camu_sink_connect(&c->sink, AKI_SOCKET_TCP, CAMU_SERVER_IP, CAMU_PORT, al_str_c("ctv")); + camu_sink_connect(&c->sink, NNWT_SOCKET_TCP, CAMU_SERVER_IP, CAMU_PORT, al_str_c("ctv")); - aki_thread_create(&c->thread, event_loop_thread, c); + nn_thread_create(&c->thread, event_loop_thread, c); c->created = true; } @@ -92,9 +92,10 @@ static void onDraw(GLFMDisplay *display) } } -static s32 log_callback(void *userdata, char *message) +static s32 log_callback(void *userdata, u8 level, char *message) { (void)userdata; + (void)level; __android_log_print(ANDROID_LOG_DEBUG, "CTV", "%s", message); s32 ret = al_strlen(message); al_free(message); diff --git a/src/liana/client.c b/src/liana/client.c index ad365b3..1ba09cc 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -9,24 +9,24 @@ #include "handlers.h" #include "list.h" -static void data_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +static void data_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) { struct lia_client *client = (struct lia_client *)userdata; (void)stream; lia_vcr_push_packet(&client->vcr, packet); } -static void parse_info_packet(struct lia_client *client, struct aki_packet *packet) +static void parse_info_packet(struct lia_client *client, struct nn_packet *packet) { str liana; - aki_packet_read_str(packet, &liana); - client->duration = aki_packet_read_u64(packet); - u32 count = aki_packet_read_u32(packet); + nn_packet_read_str(packet, &liana); + client->duration = nn_packet_read_u64(packet); + u32 count = nn_packet_read_u32(packet); for (u32 i = 0; i < count; i++) { - u8 mode = aki_packet_read_u8(packet); - u8 type = aki_packet_read_u8(packet); - u64 duration = aki_packet_read_u64(packet); - s32 index = aki_packet_read_s32(packet); + u8 mode = nn_packet_read_u8(packet); + u8 type = nn_packet_read_u8(packet); + u64 duration = nn_packet_read_u64(packet); + s32 index = nn_packet_read_s32(packet); al_assert(index < 32); struct lia_vcr_track *track = NULL; switch (mode) { @@ -35,17 +35,17 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack track = al_alloc_object(struct lia_vcr_track); if (type == CAMU_STREAM_AUDIO) { struct camu_audio_format *fmt = &track->stream.audio.fmt; - fmt->format = aki_packet_read_s32(packet); - fmt->sample_rate = aki_packet_read_s32(packet); - fmt->channel_count = aki_packet_read_s32(packet); + fmt->format = nn_packet_read_s32(packet); + fmt->sample_rate = nn_packet_read_s32(packet); + fmt->channel_count = nn_packet_read_s32(packet); #ifdef CAMU_HAVE_FFMPEG av_channel_layout_default(&fmt->channel_layout, fmt->channel_count); #endif } else if (type == CAMU_STREAM_VIDEO) { struct camu_video_format *fmt = &track->stream.video.fmt; - fmt->width = aki_packet_read_s32(packet); - fmt->height = aki_packet_read_s32(packet); - fmt->format = aki_packet_read_s32(packet); + fmt->width = nn_packet_read_s32(packet); + fmt->height = nn_packet_read_s32(packet); + fmt->format = nn_packet_read_s32(packet); } break; } @@ -67,10 +67,10 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack default: continue; } - enum AVCodecID codec_id = aki_packet_read_av_codec_id(packet); + enum AVCodecID codec_id = nn_packet_read_av_codec_id(packet); const AVCodec *codec = avcodec_find_decoder(codec_id); AVFormatContext *format_context = avformat_alloc_context(); - AVStream *stream = aki_packet_read_av_stream(format_context, codec, packet); + AVStream *stream = nn_packet_read_av_stream(format_context, codec, packet); if (type == CAMU_STREAM_SUBTITLE && codec_id != AV_CODEC_ID_ASS) { client->mask &= ~(1 << index); continue; @@ -120,30 +120,30 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack } } -static void info_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +static void info_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) { struct lia_client *client = (struct lia_client *)userdata; - client->connection_id = aki_packet_read_u32(packet); + client->connection_id = nn_packet_read_u32(packet); parse_info_packet(client, packet); - aki_packet_free(packet); + nn_packet_free(packet); if (client->mask == 0 || lia_vcr_is_empty(&client->vcr)) { client->reconnect = false; - aki_packet_stream_disconnect(&client->data); + nn_packet_stream_disconnect(&client->data); return; } stream->packet_callback = data_packet_callback; - struct aki_packet *rpacket = aki_packet_create(); - aki_packet_write_s32(rpacket, client->mask); - aki_packet_stream_send_packet(stream, rpacket); + struct nn_packet *rpacket = nn_packet_create(); + nn_packet_write_s32(rpacket, client->mask); + nn_packet_stream_send_packet(stream, rpacket); } -static void packet_sent_callback(void *userdata, struct aki_packet *packet) +static void packet_sent_callback(void *userdata, struct nn_packet *packet) { (void)userdata; - aki_packet_free(packet); + nn_packet_free(packet); } -static void connection_callback(void *userdata, struct aki_packet_stream *stream) +static bool connection_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_client *client = (struct lia_client *)userdata; if (client->reconnect) { @@ -164,20 +164,21 @@ static void connection_callback(void *userdata, struct aki_packet_stream *stream } lia_vcr_start(&client->vcr); stream->packet_sent_callback = packet_sent_callback; - struct aki_packet *packet = aki_packet_create(); - aki_packet_write_u32(packet, client->id); - aki_packet_write_u32(packet, 0); - aki_packet_write_s32(packet, client->mask); - aki_packet_write_u64(packet, client->pos); + struct nn_packet *packet = nn_packet_create(); + nn_packet_write_u32(packet, client->id); + nn_packet_write_u32(packet, 0); + nn_packet_write_s32(packet, client->mask); + nn_packet_write_u64(packet, client->pos); if (client->mask == 0) { stream->packet_callback = info_packet_callback; } else { stream->packet_callback = data_packet_callback; } - aki_packet_stream_send_packet(stream, packet); + nn_packet_stream_send_packet(stream, packet); + return true; } -static void connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_client *client = (struct lia_client *)userdata; if (client->reconnect) { @@ -190,13 +191,13 @@ static void connection_closed_callback(void *userdata, struct aki_packet_stream // any unexpected behavior. client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->reconnect); if (client->reconnect) { - aki_packet_stream_reconnect(stream, &client->addr, client->port); + nn_packet_stream_reconnect(stream, &client->addr, client->port); } else { client->callback(client->userdata, LIANA_CLIENT_CLOSED, NULL, NULL); } } -void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop, u8 type, +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) { client->loop = loop; @@ -207,12 +208,12 @@ void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop, lia_vcr_init(&client->vcr, client->loop, &client->data); al_str_clone(&client->addr, addr); client->port = port; - if (!aki_packet_stream_init(&client->data, type, connection_callback, connection_closed_callback, client)) { + if (!nn_packet_stream_init(&client->data, type, connection_callback, connection_closed_callback, client)) { connection_closed_callback(client, &client->data); } - aki_packet_stream_set_multiplex(&client->data, CAMU_MULTIPLEX_LIANA); + nn_packet_stream_set_multiplex(&client->data, CAMU_MULTIPLEX_LIANA); client->renderer = renderer; - aki_packet_stream_connect(&client->data, client->loop, &client->addr, client->port); + nn_packet_stream_connect(&client->data, client->loop, &client->addr, client->port); } void lia_client_set_renderer(struct lia_client *client, struct camu_renderer *renderer) @@ -226,7 +227,7 @@ void lia_client_seek(struct lia_client *client, u64 pos, u64 at) client->at = at; if (!client->reconnect) { client->reconnect = true; - aki_packet_stream_disconnect(&client->data); + nn_packet_stream_disconnect(&client->data); } } @@ -238,12 +239,12 @@ void lia_client_reseek(struct lia_client *client) void lia_client_disconnect(struct lia_client *client) { client->reconnect = false; - aki_packet_stream_disconnect(&client->data); + nn_packet_stream_disconnect(&client->data); } void lia_client_free(struct lia_client *client) { lia_vcr_free(&client->vcr); - aki_packet_stream_free(&client->data); + nn_packet_stream_free(&client->data); al_str_free(&client->addr); } diff --git a/src/liana/client.h b/src/liana/client.h index cd1297e..55c5b3a 100644 --- a/src/liana/client.h +++ b/src/liana/client.h @@ -1,13 +1,13 @@ #pragma once -#include <aki/packet_stream.h> +#include <nnwt/packet_stream.h> #include "../codec/codec.h" #include "vcr.h" struct lia_client { - struct aki_event_loop *loop; + struct nn_event_loop *loop; u32 id; s32 mask; u64 pos; @@ -16,7 +16,7 @@ struct lia_client { str addr; u16 port; u32 connection_id; - struct aki_packet_stream data; + struct nn_packet_stream data; u64 duration; struct lia_vcr vcr; struct camu_renderer *renderer; @@ -24,7 +24,7 @@ struct lia_client { void *userdata; }; -void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop, u8 type, +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); void lia_client_seek(struct lia_client *client, u64 pos, u64 at); void lia_client_reseek(struct lia_client *client); diff --git a/src/liana/handler.h b/src/liana/handler.h index 21437c4..9f93e5c 100644 --- a/src/liana/handler.h +++ b/src/liana/handler.h @@ -1,18 +1,18 @@ #pragma once -#include <aki/packet_pool.h> +#include <nnwt/packet_pool.h> #include "../codec/codec.h" #include "../cache/handle.h" struct lia_server_handler { bool (*init)(struct lia_server_handler *, struct cch_handle *); - void (*write_info)(struct lia_server_handler *, struct aki_packet *); + void (*write_info)(struct lia_server_handler *, struct nn_packet *); void (*subscribe)(struct lia_server_handler *, s32); u64 (*get_duration)(struct lia_server_handler *); bool (*seek)(struct lia_server_handler *, u64); void (*step)(struct lia_server_handler *); - void (*write_packet)(struct lia_server_handler *, struct aki_packet *); + void (*write_packet)(struct lia_server_handler *, struct nn_packet *); void (*free)(struct lia_server_handler **); s32 status; }; @@ -49,7 +49,7 @@ struct lia_seek_req { struct lia_client_handler { bool (*init)(struct lia_client_handler *, struct camu_renderer *, struct camu_codec_stream *); - bool (*handle_packet)(struct lia_client_handler *, struct aki_packet *); + bool (*handle_packet)(struct lia_client_handler *, struct nn_packet *); void (*flush)(struct lia_client_handler *); void (*free)(struct lia_client_handler **); struct camu_codec_stream *stream; diff --git a/src/liana/handlers.c b/src/liana/handlers.c index 936f801..1a28024 100644 --- a/src/liana/handlers.c +++ b/src/liana/handlers.c @@ -13,7 +13,7 @@ struct lia_handler_entry liana_handlers[] = { .create_client_handler = lia_codec_client_create #endif }, -#ifdef LIANA_HAVE_CDIO +#ifdef CACHE_HAVE_CDIO { .name = al_str_c("cdio"), #ifdef LIANA_SERVER diff --git a/src/liana/handlers/cdio.h b/src/liana/handlers/cdio.h index 2727616..4c40986 100644 --- a/src/liana/handlers/cdio.h +++ b/src/liana/handlers/cdio.h @@ -6,8 +6,8 @@ struct lia_cdio_server { struct lia_server_handler handler; struct cch_handle *handle; struct camu_audio_format fmt; + struct nn_buffer buffer; f64 pts; - struct aki_buffer buffer; }; struct lia_cdio_client { diff --git a/src/liana/handlers/cdio_client.c b/src/liana/handlers/cdio_client.c index 93b1a27..665858f 100644 --- a/src/liana/handlers/cdio_client.c +++ b/src/liana/handlers/cdio_client.c @@ -9,7 +9,7 @@ static bool cdio_client_init(struct lia_client_handler *handler, struct camu_ren return true; } -static bool cdio_client_handle_packet(struct lia_client_handler *handler, struct aki_packet *packet) +static bool cdio_client_handle_packet(struct lia_client_handler *handler, struct nn_packet *packet) { struct lia_cdio_client *cdio = (struct lia_cdio_client *)handler; if (!packet) { @@ -18,12 +18,12 @@ static bool cdio_client_handle_packet(struct lia_client_handler *handler, struct } struct camu_codec_frame *frame = al_alloc_object(struct camu_codec_frame); frame->mode = CAMU_NORMAL; - frame->pts = aki_packet_read_f64(packet); - frame->audio.sample_count = aki_packet_read_s32(packet); - struct aki_buffer buffer; - aki_packet_read_buffer(packet, &buffer); + frame->pts = nn_packet_read_f64(packet); + frame->audio.sample_count = nn_packet_read_s32(packet); + struct nn_buffer buffer; + nn_packet_read_buffer(packet, &buffer); frame->data = al_malloc(buffer.size); - al_memcpy(frame->data, aki_buffer_get_ptr(&buffer, 0), buffer.size); + al_memcpy(frame->data, nn_buffer_get_ptr(&buffer, 0), buffer.size); cdio->handler.callback(cdio->handler.userdata, LIANA_CLIENT_DATA, cdio->handler.stream, frame); return true; } diff --git a/src/liana/handlers/cdio_server.c b/src/liana/handlers/cdio_server.c index 641ed92..01a9265 100644 --- a/src/liana/handlers/cdio_server.c +++ b/src/liana/handlers/cdio_server.c @@ -15,8 +15,8 @@ static bool cdio_server_init(struct lia_server_handler *handler, struct cch_hand cdio->fmt.format = CAMU_SAMPLE_FORMAT_S16; cdio->fmt.sample_rate = 44100; cdio->fmt.channel_count = 2; + nn_buffer_init(&cdio->buffer); cdio->pts = 0.0; - aki_buffer_init(&cdio->buffer); if (handle->entry->chapter) { cch_handle_seek(handle, handle->entry->chapter->start * CDIO_CD_FRAMESIZE_RAW, SEEK_SET); @@ -25,18 +25,18 @@ static bool cdio_server_init(struct lia_server_handler *handler, struct cch_hand return true; } -static void cdio_server_write_info(struct lia_server_handler *handler, struct aki_packet *packet) +static void cdio_server_write_info(struct lia_server_handler *handler, struct nn_packet *packet) { struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; - aki_packet_write_u64(packet, 0); - aki_packet_write_u32(packet, 1); - aki_packet_write_u8(packet, CAMU_NORMAL); - aki_packet_write_u8(packet, CAMU_STREAM_AUDIO); - aki_packet_write_u64(packet, cdio->handler.get_duration(&cdio->handler)); - aki_packet_write_s32(packet, 0); - aki_packet_write_s32(packet, cdio->fmt.format); - aki_packet_write_s32(packet, cdio->fmt.sample_rate); - aki_packet_write_s32(packet, cdio->fmt.channel_count); + nn_packet_write_u64(packet, 0); + nn_packet_write_u32(packet, 1); + nn_packet_write_u8(packet, CAMU_NORMAL); + nn_packet_write_u8(packet, CAMU_STREAM_AUDIO); + nn_packet_write_u64(packet, cdio->handler.get_duration(&cdio->handler)); + nn_packet_write_s32(packet, 0); + nn_packet_write_s32(packet, cdio->fmt.format); + nn_packet_write_s32(packet, cdio->fmt.sample_rate); + nn_packet_write_s32(packet, cdio->fmt.channel_count); } static void cdio_server_subscribe(struct lia_server_handler *handler, s32 mask) @@ -69,26 +69,26 @@ static void cdio_server_step(struct lia_server_handler *handler) (void)cdio; } -static void cdio_server_write_packet(struct lia_server_handler *handler, struct aki_packet *packet) +static void cdio_server_write_packet(struct lia_server_handler *handler, struct nn_packet *packet) { struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; off_t size = CDIO_CD_FRAMESIZE_RAW * SECTORS_PER_PACKET; - aki_buffer_ensure_space(&cdio->buffer, size); - s32 ret = cch_handle_read(cdio->handle, aki_buffer_get_ptr(&cdio->buffer, 0), size); + nn_buffer_ensure_space(&cdio->buffer, size); + s32 ret = cch_handle_read(cdio->handle, nn_buffer_get_ptr(&cdio->buffer, 0), size); if (ret == CAMU_ERR_EOF) { cdio->handler.status = CAMU_ERR_EOF; - aki_packet_write_u8(packet, LIANA_PACKET_EOF); + nn_packet_write_u8(packet, LIANA_PACKET_EOF); return; } cdio->handler.status = CAMU_OK; - aki_packet_write_u8(packet, LIANA_PACKET_DATA); - aki_packet_write_s32(packet, 0); - aki_packet_write_f64(packet, cdio->pts); + nn_packet_write_u8(packet, LIANA_PACKET_DATA); + nn_packet_write_s32(packet, 0); + nn_packet_write_f64(packet, cdio->pts); s32 sample_count = ret / (camu_audio_format_bytes_per_sample(&cdio->fmt) * cdio->fmt.channel_count); + nn_packet_write_s32(packet, sample_count); + nn_buffer_set_size(&cdio->buffer, ret); + nn_packet_write_buffer(packet, &cdio->buffer); cdio->pts += camu_audio_format_samples_to_sec(&cdio->fmt, sample_count); - aki_packet_write_s32(packet, sample_count); - aki_buffer_set_size(&cdio->buffer, ret); - aki_packet_write_buffer(packet, &cdio->buffer); } static void cdio_server_free(struct lia_server_handler **handler) diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c index dc53129..3722197 100644 --- a/src/liana/handlers/codec_client.c +++ b/src/liana/handlers/codec_client.c @@ -44,7 +44,7 @@ static bool push_av_packet(struct lia_codec_client *codec, AVPacket *pkt) } #endif -static bool push_packet(struct lia_codec_client *codec, struct aki_buffer *buffer) +static bool push_packet(struct lia_codec_client *codec, struct nn_buffer *buffer) { struct camu_codec_packet packet; packet.buffer = buffer; @@ -52,7 +52,7 @@ static bool push_packet(struct lia_codec_client *codec, struct aki_buffer *buffe return ret == CAMU_OK; } -static bool codec_client_handle_packet(struct lia_client_handler *handler, struct aki_packet *packet) +static bool codec_client_handle_packet(struct lia_client_handler *handler, struct nn_packet *packet) { struct lia_codec_client *codec = (struct lia_codec_client *)handler; @@ -72,12 +72,12 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc // Push packet. bool success; - u8 type = aki_packet_read_u8(packet); + u8 type = nn_packet_read_u8(packet); switch (type) { case CAMU_NORMAL: { if (codec->dec) { - struct aki_buffer buffer; - aki_packet_read_buffer(packet, &buffer); + struct nn_buffer buffer; + nn_packet_read_buffer(packet, &buffer); success = push_packet(codec, &buffer); } else { success = true; @@ -86,7 +86,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc } #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: { - AVPacket *pkt = aki_packet_read_av_packet(packet); + AVPacket *pkt = nn_packet_read_av_packet(packet); if (codec->dec) { success = push_av_packet(codec, pkt); } else { diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c index 94b18df..7167f97 100644 --- a/src/liana/handlers/codec_server.c +++ b/src/liana/handlers/codec_server.c @@ -33,30 +33,30 @@ static bool codec_server_init(struct lia_server_handler *handler, struct cch_han return true; } -static void codec_server_write_info(struct lia_server_handler *handler, struct aki_packet *packet) +static void codec_server_write_info(struct lia_server_handler *handler, struct nn_packet *packet) { struct lia_codec_server *codec = (struct lia_codec_server *)handler; u64 duration = codec->demux->get_duration(codec->demux); - aki_packet_write_u64(packet, duration); - aki_packet_write_u32(packet, codec->demux->streams.size); + nn_packet_write_u64(packet, duration); + nn_packet_write_u32(packet, codec->demux->streams.size); struct camu_codec_stream *stream; al_array_foreach_ptr(codec->demux->streams, i, stream) { - aki_packet_write_u8(packet, stream->mode); - aki_packet_write_u8(packet, stream->type); - aki_packet_write_u64(packet, stream->duration); - aki_packet_write_s32(packet, i); + nn_packet_write_u8(packet, stream->mode); + nn_packet_write_u8(packet, stream->type); + nn_packet_write_u64(packet, stream->duration); + nn_packet_write_s32(packet, i); switch (stream->mode) { case CAMU_NORMAL: { struct camu_video_format *fmt = &stream->video.fmt; - aki_packet_write_s32(packet, fmt->width); - aki_packet_write_s32(packet, fmt->height); - aki_packet_write_s32(packet, fmt->format); + nn_packet_write_s32(packet, fmt->width); + nn_packet_write_s32(packet, fmt->height); + nn_packet_write_s32(packet, fmt->format); break; } #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: - aki_packet_write_av_codec_id(packet, stream->av.stream->codecpar->codec_id); - aki_packet_write_av_stream(packet, stream->av.stream); + nn_packet_write_av_codec_id(packet, stream->av.stream->codecpar->codec_id); + nn_packet_write_av_stream(packet, stream->av.stream); break; #endif } @@ -87,33 +87,33 @@ static void codec_server_step(struct lia_server_handler *handler) codec->handler.status = codec->demux->get_packet(codec->demux, &codec->packet); } -static void codec_server_write_packet(struct lia_server_handler *handler, struct aki_packet *packet) +static void codec_server_write_packet(struct lia_server_handler *handler, struct nn_packet *packet) { struct lia_codec_server *codec = (struct lia_codec_server *)handler; if (codec->handler.status == CAMU_OK) { - aki_packet_write_u8(packet, LIANA_PACKET_DATA); + nn_packet_write_u8(packet, LIANA_PACKET_DATA); switch (codec->packet.mode) { case CAMU_NORMAL: { - aki_packet_write_s32(packet, 0); - aki_packet_write_u8(packet, codec->packet.mode); - aki_packet_write_buffer(packet, codec->packet.buffer); + nn_packet_write_s32(packet, 0); + nn_packet_write_u8(packet, codec->packet.mode); + nn_packet_write_buffer(packet, codec->packet.buffer); break; } #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: { AVPacket *pkt = codec->packet.av.pkt; - aki_packet_write_s32(packet, pkt->stream_index); - aki_packet_write_u8(packet, codec->packet.mode); - aki_packet_write_av_packet(packet, pkt); + nn_packet_write_s32(packet, pkt->stream_index); + nn_packet_write_u8(packet, codec->packet.mode); + nn_packet_write_av_packet(packet, pkt); av_packet_unref(pkt); break; } #endif } } else if (codec->handler.status == CAMU_ERR_EOF) { - aki_packet_write_u8(packet, LIANA_PACKET_EOF); + nn_packet_write_u8(packet, LIANA_PACKET_EOF); } else { - aki_packet_write_u8(packet, LIANA_PACKET_ERROR); + nn_packet_write_u8(packet, LIANA_PACKET_ERROR); } } diff --git a/src/liana/list.c b/src/liana/list.c index 073cdf0..7c4979b 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -1,4 +1,4 @@ -#include <aki/thread.h> +#include <nnwt/thread.h> #include <al/random.h> #include <al/lib.h> #include <al/log.h> @@ -112,7 +112,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) if (error) pump_queue(list); return false; } - u64 now = aki_get_timestamp(); + u64 now = nn_get_timestamp(); u8 pause; u64 at = LIANA_TIMESTAMP_INVALID; u64 seek_pos = current->offset; @@ -168,7 +168,7 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) } list->current++; list->idle = false; - entry->start = aki_get_timestamp() + LIANA_BASE_DELAY; + entry->start = nn_get_timestamp() + LIANA_BASE_DELAY; struct lia_timing time = { .at = entry->start, .seek_pos = entry->offset, @@ -299,7 +299,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) return false; } - u64 now = aki_get_timestamp(); + u64 now = nn_get_timestamp(); u64 at = now + LIANA_BASE_DELAY; u8 pause; @@ -405,7 +405,7 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) struct lia_list_entry *entry = get_entry_from_sequence(list, sequence); al_assert(entry && !entry->held); - u64 now = aki_get_timestamp(); + u64 now = nn_get_timestamp(); u8 pause = entry->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_PAUSE : LIANA_PAUSE_RESUME; u64 at; bool ended = assume_ended(entry, now); @@ -461,7 +461,7 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent } entry->ended = false; - u64 now = aki_get_timestamp(); + u64 now = nn_get_timestamp(); u64 pos = (u64)(entry->duration * percent); u64 at = now + LIANA_BASE_DELAY; u8 pause = entry->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE; @@ -527,10 +527,12 @@ static bool handle_end(struct lia_list *list, s32 id) return false; } else { list->idle = true; + /* struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { sink->set = -1; } + */ } } diff --git a/src/liana/meson.build b/src/liana/meson.build index 32e6a00..b68ba8b 100644 --- a/src/liana/meson.build +++ b/src/liana/meson.build @@ -16,7 +16,6 @@ liana_args = [] if cache_have_cdio liana_server_src += ['handlers/cdio_server.c'] liana_client_src += ['handlers/cdio_client.c'] - liana_args += ['-DLIANA_HAVE_CDIO'] endif liana_server = declare_dependency(sources: liana_server_src, diff --git a/src/liana/server.c b/src/liana/server.c index 44b3f12..27c253b 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -5,7 +5,7 @@ #include "handlers.h" #include "list.h" -bool lia_server_init(struct lia_server *server, struct aki_event_loop *loop) +bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop) { server->loop = loop; server->increment = 1; @@ -14,52 +14,46 @@ bool lia_server_init(struct lia_server *server, struct aki_event_loop *loop) return true; } -static void remove_zombie(struct lia_server *server, struct aki_packet_stream *stream) +static void remove_zombie(struct lia_server *server, struct nn_packet_stream *stream) { - struct aki_packet_stream *zombie; - al_array_foreach(server->zombies, i, zombie) { - if (zombie == stream) { - al_array_remove_at(server->zombies, i); - break; - } - } + al_array_remove(server->zombies, stream); } -static void data_packet_sent_callback(void *userdata, struct aki_packet *packet) +static void data_packet_sent_callback(void *userdata, struct nn_packet *packet) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; - aki_packet_pool_return(&conn->pool, packet); + nn_packet_pool_return(&conn->pool, packet); } -static u8 packet_pool_callback(void *userdata, struct aki_packet *packet) +static u8 packet_pool_callback(void *userdata, struct nn_packet *packet) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; - if (!aki_packet_stream_send_packet(conn->stream, packet)) { - return AKI_PACKET_POOL_RETURN; + if (!nn_packet_stream_send_packet(conn->stream, packet)) { + return NNWT_PACKET_POOL_RETURN; } - return AKI_PACKET_POOL_KEEP; + return NNWT_PACKET_POOL_KEEP; } -static aki_thread_result AKI_THREADCALL handler_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; do { - struct aki_packet *packet = aki_packet_pool_get(&conn->pool); + struct nn_packet *packet = nn_packet_pool_get(&conn->pool); if (!packet) break; conn->handler->step(conn->handler); conn->handler->write_packet(conn->handler, packet); - aki_packet_pool_submit(&conn->pool, packet); + nn_packet_pool_submit(&conn->pool, packet); if (conn->handler->status != CAMU_OK) break; } while (1); - aki_packet_pool_flush(&conn->pool); + nn_packet_pool_flush(&conn->pool); return 0; } -static void discard_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +static void discard_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) { (void)userdata; (void)stream; - aki_packet_free(packet); + nn_packet_free(packet); // We should never be here. Although, we also shouldn't assert because // any erroneous connection can bring us here. al_assert(false); @@ -70,83 +64,83 @@ static void close_connection_internal(struct lia_node_connection *conn) al_assert(conn->handler); conn->handler->free(&conn->handler); cch_entry_return_handle(conn->node->entry, &conn->handle); - aki_packet_stream_free(conn->stream); + nn_packet_stream_free(conn->stream); al_free(conn->stream); conn->stream = NULL; } -static void data_connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +static void data_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; (void)stream; - aki_packet_pool_disable(&conn->pool); + 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. - aki_thread_join(&conn->thread); + nn_thread_join(&conn->thread); close_connection_internal(conn); - aki_packet_pool_free(&conn->pool); + nn_packet_pool_free(&conn->pool); } -static void subscribe_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; - s32 mask = aki_packet_read_s32(packet); + s32 mask = nn_packet_read_s32(packet); conn->handler->subscribe(conn->handler, mask); stream->packet_callback = discard_packet_callback; stream->packet_sent_callback = data_packet_sent_callback; stream->connection_closed_callback = data_connection_closed_callback; - aki_thread_create(&conn->thread, handler_thread, conn); - aki_packet_free(packet); + nn_thread_create(&conn->thread, handler_thread, conn); + nn_packet_free(packet); } -static void subscribe_packet_sent_callback(void *userdata, struct aki_packet *packet) +static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet) { (void)userdata; - aki_packet_free(packet); + nn_packet_free(packet); } -static void subscribe_connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +static void subscribe_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; (void)stream; close_connection_internal(conn); - aki_packet_pool_free(&conn->pool); + nn_packet_pool_free(&conn->pool); } -static void handle_connection(struct lia_node_connection *conn, struct aki_packet *packet) +static void handle_connection(struct lia_node_connection *conn, struct nn_packet *packet) { - struct aki_packet_stream *stream = conn->stream; + struct nn_packet_stream *stream = conn->stream; - s32 mask = aki_packet_read_s32(packet); - u64 seek_pos = aki_packet_read_u64(packet); + s32 mask = nn_packet_read_s32(packet); + u64 seek_pos = nn_packet_read_u64(packet); // Besides being wasteful, seeking to 0 on a new stream can skip data. if (seek_pos > 0) conn->handler->seek(conn->handler, seek_pos); if (mask == 0) { - struct aki_packet *rpacket = aki_packet_create(); - aki_packet_write_u32(rpacket, conn->id); - aki_packet_write_str(rpacket, cch_entry_get_liana(conn->node->entry)); + struct nn_packet *rpacket = nn_packet_create(); + nn_packet_write_u32(rpacket, conn->id); + nn_packet_write_str(rpacket, cch_entry_get_liana(conn->node->entry)); conn->handler->write_info(conn->handler, rpacket); stream->packet_callback = subscribe_packet_callback; stream->packet_sent_callback = subscribe_packet_sent_callback; stream->connection_closed_callback = subscribe_connection_closed_callback; - aki_packet_stream_send_packet(stream, rpacket); + nn_packet_stream_send_packet(stream, rpacket); } else { conn->handler->subscribe(conn->handler, mask); stream->packet_callback = discard_packet_callback; stream->packet_sent_callback = data_packet_sent_callback; stream->connection_closed_callback = data_connection_closed_callback; - aki_thread_create(&conn->thread, handler_thread, conn); + nn_thread_create(&conn->thread, handler_thread, conn); } } -static void connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_server *server = (struct lia_server *)userdata; // This connection might no longer be a zombie, but that's fine. remove_zombie(server, stream); - aki_packet_stream_free(stream); + nn_packet_stream_free(stream); al_free(stream); } @@ -155,9 +149,9 @@ 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; - aki_signal_stop(&conn->signal); - aki_thread_join(&conn->thread); - struct aki_packet *packet = conn->packet; + nn_signal_stop(&conn->signal); + nn_thread_join(&conn->thread); + struct nn_packet *packet = conn->packet; if (!packet) { // Connection was closed before init was done. close_connection_internal(conn); @@ -166,30 +160,30 @@ static void signal_callback(void *userdata) if (!conn->errored) { conn->id = server->increment; server->increment = al_u32_inc_wrap(server->increment); - aki_packet_pool_init(&conn->pool, 96, server->loop, packet_pool_callback, conn); + nn_packet_pool_init(&conn->pool, 96, server->loop, packet_pool_callback, conn); al_array_push(node->connections, conn); handle_connection(conn, packet); } else { conn->handler->free(&conn->handler); cch_entry_return_handle(node->entry, &conn->handle); - struct aki_packet_stream *stream = conn->stream; + struct nn_packet_stream *stream = conn->stream; al_free(conn); conn = NULL; // This connection is now nothing but a packet stream. stream->userdata = server; stream->connection_closed_callback = connection_closed_callback; - aki_packet_stream_disconnect(stream); + nn_packet_stream_disconnect(stream); } - aki_packet_free(packet); + nn_packet_free(packet); } -static aki_thread_result AKI_THREADCALL init_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL init_thread(void *userdata) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; if (!conn->handler->init(conn->handler, &conn->handle)) { conn->errored = true; } - aki_signal_send(&conn->signal); + nn_signal_send(&conn->signal); return 0; } @@ -211,24 +205,24 @@ static struct lia_node *get_node_from_id(struct lia_server *server, u32 id) return NULL; } -static void pre_init_connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +static void pre_init_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; (void)stream; - aki_packet_free(conn->packet); + nn_packet_free(conn->packet); // Checked in signal_callback and will signal to cleanup the connection. conn->packet = NULL; } -static void packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) { struct lia_server *server = (struct lia_server *)userdata; // We got a packet, this connection is no longer a zombie. remove_zombie(server, stream); - u32 node_id = aki_packet_read_u32(packet); - u32 connection_id = aki_packet_read_u32(packet); + u32 node_id = nn_packet_read_u32(packet); + u32 connection_id = nn_packet_read_u32(packet); struct lia_node *node = get_node_from_id(server, node_id); struct lia_node_connection *conn = NULL; @@ -238,46 +232,47 @@ static void packet_callback(void *userdata, struct aki_packet_stream *stream, st conn->node = node; conn->stream = stream; conn->packet = packet; - aki_signal_init(&conn->signal, server->loop, signal_callback, conn); - aki_signal_start(&conn->signal); + nn_signal_init(&conn->signal, server->loop, signal_callback, conn); + nn_signal_start(&conn->signal); cch_entry_get_handle(node->entry, &conn->handle); conn->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler(); conn->errored = false; stream->packet_callback = discard_packet_callback; stream->connection_closed_callback = pre_init_connection_closed_callback; - aki_thread_create(&conn->thread, init_thread, conn); + nn_thread_create(&conn->thread, init_thread, conn); } else { // This is completely unused and connections never get removed from node->connection. if ((conn = get_connection_from_id(node, connection_id))) { handle_connection(conn, packet); } else { - aki_packet_stream_disconnect(stream); + nn_packet_stream_disconnect(stream); } - aki_packet_free(packet); + nn_packet_free(packet); } } -static void packet_sent_callback(void *userdata, struct aki_packet *packet) +static void packet_sent_callback(void *userdata, struct nn_packet *packet) { (void)userdata; - aki_packet_free(packet); + nn_packet_free(packet); } -static void connection_callback(void *userdata, struct aki_packet_stream *stream) +static bool connection_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_server *server = (struct lia_server *)userdata; stream->packet_callback = packet_callback; stream->packet_sent_callback = packet_sent_callback; al_array_push(server->zombies, stream); + return true; } -void lia_server_add_socket(struct lia_server *server, struct aki_socket *sock) +void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock) { - struct aki_packet_stream *stream = al_alloc_object(struct aki_packet_stream); + struct nn_packet_stream *stream = al_alloc_object(struct nn_packet_stream); stream->connection_callback = connection_callback; stream->connection_closed_callback = connection_closed_callback; stream->userdata = server; - aki_packet_stream_from_socket(stream, server->loop, sock); + nn_packet_stream_from_socket(stream, server->loop, sock); } struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry) @@ -292,7 +287,7 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en return node; } -static aki_thread_result AKI_THREADCALL init_duration_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL init_duration_thread(void *userdata) { struct lia_node *node = (struct lia_node *)userdata; if (!node->handler->init(node->handler, &node->handle)) { @@ -301,15 +296,15 @@ static aki_thread_result AKI_THREADCALL init_duration_thread(void *userdata) } else { node->duration = node->handler->get_duration(node->handler); } - aki_signal_send(&node->signal); + nn_signal_send(&node->signal); return 0; } static void duration_signal_callback(void *userdata) { struct lia_node *node = (struct lia_node *)userdata; - aki_signal_stop(&node->signal); - aki_thread_join(&node->thread); + nn_signal_stop(&node->signal); + nn_thread_join(&node->thread); node->handler->free(&node->handler); cch_entry_return_handle(node->entry, &node->handle); node->callback(node->userdata, LIANA_NODE_DURATION, node->duration); @@ -317,11 +312,11 @@ static void duration_signal_callback(void *userdata) void lia_node_get_duration(struct lia_node *node) { - aki_signal_init(&node->signal, node->server->loop, duration_signal_callback, node); - aki_signal_start(&node->signal); + nn_signal_init(&node->signal, node->server->loop, duration_signal_callback, node); + nn_signal_start(&node->signal); cch_entry_get_handle(node->entry, &node->handle); node->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler(); - aki_thread_create(&node->thread, init_duration_thread, node); + nn_thread_create(&node->thread, init_duration_thread, node); } void lia_server_close(struct lia_server *server) @@ -331,13 +326,13 @@ void lia_server_close(struct lia_server *server) struct lia_node_connection *connection; al_array_foreach(node->connections, j, connection) { if (connection->stream) { - aki_packet_stream_disconnect(connection->stream); + nn_packet_stream_disconnect(connection->stream); } } } - struct aki_packet_stream *zombie; + struct nn_packet_stream *zombie; al_array_foreach_rev(server->zombies, i, zombie) { - aki_packet_stream_disconnect(zombie); + nn_packet_stream_disconnect(zombie); } } @@ -353,7 +348,7 @@ void lia_server_free(struct lia_server *server) al_free(node); } al_array_free(server->nodes); - struct aki_packet_stream *zombie; + struct nn_packet_stream *zombie; al_array_foreach(server->zombies, i, zombie) { al_free(zombie); } diff --git a/src/liana/server.h b/src/liana/server.h index c62eb29..b4ada58 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -1,23 +1,23 @@ #pragma once -#include <aki/event_loop.h> -#include <aki/packet_stream.h> -#include <aki/packet_pool.h> -#include <aki/socket.h> -#include <aki/signal.h> +#include <nnwt/event_loop.h> +#include <nnwt/packet_stream.h> +#include <nnwt/packet_pool.h> +#include <nnwt/socket.h> +#include <nnwt/signal.h> #include "../cache/entry.h" struct lia_node_connection { u32 id; - struct aki_packet *packet; - struct aki_packet_stream *stream; + struct nn_packet *packet; + struct nn_packet_stream *stream; struct lia_server_handler *handler; bool errored; struct cch_handle handle; - struct aki_thread thread; - struct aki_signal signal; - struct aki_packet_pool pool; + struct nn_thread thread; + struct nn_signal signal; + struct nn_packet_pool pool; struct lia_node *node; }; @@ -36,21 +36,21 @@ struct lia_node { struct lia_server_handler *handler; bool errored; struct cch_handle handle; - struct aki_thread thread; - struct aki_signal signal; + struct nn_thread thread; + struct nn_signal signal; void (*callback)(void *, u8, u64); void *userdata; }; struct lia_server { - struct aki_event_loop *loop; + struct nn_event_loop *loop; u32 increment; array(struct lia_node *) nodes; - array(struct aki_packet_stream *) zombies; + array(struct nn_packet_stream *) zombies; }; -bool lia_server_init(struct lia_server *server, struct aki_event_loop *loop); -void lia_server_add_socket(struct lia_server *server, struct aki_socket *sock); +bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop); +void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock); 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_server_close(struct lia_server *server); diff --git a/src/liana/vcr.c b/src/liana/vcr.c index c4bf583..4131c80 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -19,7 +19,7 @@ enum { static void signal_callback(void *userdata) { struct lia_vcr *vcr = (struct lia_vcr *)userdata; - aki_packet_stream_cork(vcr->data, false); + nn_packet_stream_cork(vcr->data, false); } static void reset_metrics(struct lia_vcr *vcr) @@ -28,7 +28,7 @@ static void reset_metrics(struct lia_vcr *vcr) vcr->metric.last_report_ts = 0Lu; } -void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_packet_stream *data) +void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data) { al_array_init(vcr->tracks); al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); @@ -36,25 +36,25 @@ void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_p vcr->mark.low = 0; vcr->expand = VCR_EXPAND_UNTOUCHED; vcr->data = data; - aki_signal_init(&vcr->signal, loop, signal_callback, vcr); + nn_signal_init(&vcr->signal, loop, signal_callback, vcr); reset_metrics(vcr); } void lia_vcr_start(struct lia_vcr *vcr) { - aki_signal_start(&vcr->signal); + nn_signal_start(&vcr->signal); } -static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) { struct lia_vcr_track *track = (struct lia_vcr_track *)userdata; struct lia_vcr *vcr = track->vcr; u32 packets; - struct aki_packet *packet = NULL; - while (aki_packet_cache_wait(&track->cache, &packets)) { + struct nn_packet *packet = NULL; + while (nn_packet_cache_wait(&track->cache, &packets)) { s32 state = 0; for (u32 i = 0; i < packets; i++) { - packet = aki_packet_cache_pop(&track->cache); + packet = nn_packet_cache_pop(&track->cache); state = 0; if (packet) { state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED); @@ -62,11 +62,11 @@ static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata) // We were signaled to close, exit thread. goto out; } - u32 size = aki_packet_get_size(packet); + u32 size = nn_packet_get_size(packet); u8 buffered = al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED); u64 buffer = al_atomic_sub(u64)(&vcr->count, size, AL_ATOMIC_RELAXED); if (buffered && buffer <= vcr->mark.low) { - aki_signal_send(&vcr->signal); + nn_signal_send(&vcr->signal); } } if (!track->client->handle_packet(track->client, packet)) { @@ -74,14 +74,14 @@ static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata) goto out; } if (packet) { - aki_packet_free(packet); + nn_packet_free(packet); packet = NULL; if (state == LIANA_STREAM_STOPPED) { // Unlock here to accumulate packets while waiting. - aki_packet_cache_unlock(&track->cache); - aki_mutex_lock(&track->mutex); - aki_cond_wait(&track->cond, &track->mutex); - aki_mutex_unlock(&track->mutex); + nn_packet_cache_unlock(&track->cache); + nn_mutex_lock(&track->mutex); + nn_cond_wait(&track->cond, &track->mutex); + nn_mutex_unlock(&track->mutex); break; } } else { @@ -91,25 +91,25 @@ static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata) } if (state != LIANA_STREAM_STOPPED) { // If state = STOPPED, we already unlocked. - aki_packet_cache_unlock(&track->cache); + nn_packet_cache_unlock(&track->cache); } } // packet_cache_wait() returned false. return 0; out: // packet_cache_wait() returned true and we are jumping out of the loop. - if (packet) aki_packet_free(packet); - aki_packet_cache_unlock(&track->cache); + if (packet) nn_packet_free(packet); + nn_packet_cache_unlock(&track->cache); return 0; } void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track) { track->vcr = vcr; - aki_cond_init(&track->cond); - aki_mutex_init(&track->mutex); + nn_cond_init(&track->cond); + nn_mutex_init(&track->mutex); al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED); - aki_packet_cache_init(&track->cache, 256); + nn_packet_cache_init(&track->cache, 256); al_array_push(vcr->tracks, track); al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); } @@ -145,9 +145,9 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) vcr->mark.low = vcr->mark.buffered - MB(2); vcr->expand = VCR_EXPAND_COMPLETE; } - aki_packet_stream_cork(vcr->data, true); + nn_packet_stream_cork(vcr->data, true); al_array_foreach(vcr->tracks, i, track) { - aki_packet_cache_flush(&track->cache); + nn_packet_cache_flush(&track->cache); } } } @@ -155,7 +155,7 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) static void update_metrics(struct lia_vcr *vcr, u32 size) { vcr->metric.current_frame += size; - u64 now = aki_get_timestamp(); + u64 now = nn_get_timestamp(); if (!vcr->metric.last_report_ts) { vcr->metric.last_report_ts = now; return; @@ -166,6 +166,7 @@ static void update_metrics(struct lia_vcr *vcr, u32 size) u64 frame = vcr->metric.current_frame; vcr->metric.current_frame = 0Lu; if (diff > 2500000Lu) { + // We are buffering fast enough for it to not matter. al_log_debug("vcr", "Ignoring %lu bytes in metrics.", frame); return; } @@ -174,25 +175,25 @@ static void update_metrics(struct lia_vcr *vcr, u32 size) } } -void lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet) +void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) { struct lia_vcr_track *track; - u8 op = aki_packet_read_u8(packet); + u8 op = nn_packet_read_u8(packet); switch (op) { case LIANA_PACKET_DATA: - track = get_track_from_index(vcr, aki_packet_read_s32(packet)); + track = get_track_from_index(vcr, nn_packet_read_s32(packet)); if (!track) { al_log_warn("liana", "Received data from errored or unknown track."); - aki_packet_free(packet); + nn_packet_free(packet); return; } if (!track->running) { - aki_thread_create(&track->thread, vcr_track_thread, track); + nn_thread_create(&track->thread, vcr_track_thread, track); track->running = true; } - u32 size = aki_packet_get_size(packet); - if (!aki_packet_cache_send_packet(&track->cache, packet)) { - aki_packet_free(packet); + u32 size = nn_packet_get_size(packet); + if (!nn_packet_cache_send_packet(&track->cache, packet)) { + nn_packet_free(packet); return; } u64 buffer; @@ -203,14 +204,14 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet) break; case LIANA_PACKET_EOF: al_array_foreach(vcr->tracks, i, track) { - aki_packet_cache_send_packet(&track->cache, NULL); + nn_packet_cache_send_packet(&track->cache, NULL); } - aki_packet_free(packet); - aki_signal_stop(&vcr->signal); + nn_packet_free(packet); + nn_signal_stop(&vcr->signal); break; case LIANA_PACKET_ERROR: al_log_warn("liana", "Unhandled error packet."); - aki_packet_free(packet); + nn_packet_free(packet); break; default: al_assert(false); @@ -230,29 +231,29 @@ void lia_vcr_uncork(struct lia_vcr_track *track) // sense is up for consideration. } al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); - aki_mutex_lock(&track->mutex); - if (aki_cond_is_waiting(&track->cond)) { - aki_cond_signal(&track->cond); + nn_mutex_lock(&track->mutex); + if (nn_cond_is_waiting(&track->cond)) { + nn_cond_signal(&track->cond); } - aki_mutex_unlock(&track->mutex); + nn_mutex_unlock(&track->mutex); } static void vcr_track_close_internal(struct lia_vcr_track *track) { al_atomic_store(s32)(&track->state, LIANA_STREAM_CLOSED, AL_ATOMIC_RELAXED); - aki_packet_cache_disable(&track->cache); - aki_mutex_lock(&track->mutex); - if (aki_cond_is_waiting(&track->cond)) { - aki_cond_signal(&track->cond); + nn_packet_cache_disable(&track->cache); + nn_mutex_lock(&track->mutex); + if (nn_cond_is_waiting(&track->cond)) { + nn_cond_signal(&track->cond); } - aki_mutex_unlock(&track->mutex); + nn_mutex_unlock(&track->mutex); if (track->running) { - aki_thread_join(&track->thread); + nn_thread_join(&track->thread); track->running = false; } - struct aki_packet *packet; - while ((packet = aki_packet_cache_pop(&track->cache))) { - aki_packet_free(packet); + struct nn_packet *packet; + while ((packet = nn_packet_cache_pop(&track->cache))) { + nn_packet_free(packet); } } @@ -262,11 +263,11 @@ void lia_vcr_flush(struct lia_vcr *vcr) al_array_foreach(vcr->tracks, i, track) { vcr_track_close_internal(track); track->client->flush(track->client); - aki_packet_cache_enable(&track->cache); + nn_packet_cache_enable(&track->cache); al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED); al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); } - aki_signal_stop(&vcr->signal); + nn_signal_stop(&vcr->signal); al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); reset_metrics(vcr); if (vcr->expand == VCR_EXPAND_COMPLETE) { @@ -281,14 +282,14 @@ void lia_vcr_close_all(struct lia_vcr *vcr) al_array_foreach(vcr->tracks, i, track) { vcr_track_close_internal(track); } - aki_signal_stop(&vcr->signal); + nn_signal_stop(&vcr->signal); } void lia_vcr_free(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { - aki_packet_cache_free(&track->cache); + nn_packet_cache_free(&track->cache); track->client->free(&track->client); #ifdef CAMU_HAVE_FFMPEG if (track->stream.mode == CAMU_FFMPEG_COMPAT) { diff --git a/src/liana/vcr.h b/src/liana/vcr.h index 8041d96..ce47104 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -1,9 +1,9 @@ #pragma once #include <al/atomic.h> -#include <aki/packet_cache.h> -#include <aki/packet_stream.h> -#include <aki/signal.h> +#include <nnwt/packet_cache.h> +#include <nnwt/packet_stream.h> +#include <nnwt/signal.h> #include "../codec/codec.h" @@ -19,11 +19,11 @@ struct lia_vcr_track { struct lia_client_handler *client; atomic(s32) state; atomic(u8) buffered; - struct aki_packet_cache cache; - struct aki_cond cond; - struct aki_mutex mutex; + struct nn_packet_cache cache; + struct nn_cond cond; + struct nn_mutex mutex; bool running; - struct aki_thread thread; + struct nn_thread thread; struct lia_vcr *vcr; }; @@ -32,19 +32,19 @@ struct lia_vcr { atomic(u64) count; struct { u64 buffered, low; } mark; u8 expand; - struct aki_packet_stream *data; - struct aki_signal signal; + struct nn_packet_stream *data; + struct nn_signal signal; struct { u64 current_frame; u64 last_report_ts; } metric; }; -void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_packet_stream *data); +void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data); void lia_vcr_start(struct lia_vcr *vcr); void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track); bool lia_vcr_is_empty(struct lia_vcr *vcr); -void lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet); +void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet); void lia_vcr_cork(struct lia_vcr_track *track); void lia_vcr_uncork(struct lia_vcr_track *track); void lia_vcr_flush(struct lia_vcr *vcr); diff --git a/src/libclient/client.c b/src/libclient/client.c index 4b9f46a..b4e44da 100644 --- a/src/libclient/client.c +++ b/src/libclient/client.c @@ -3,17 +3,17 @@ #include "../server/common.h" -static bool results_callback(void *userdata, struct aki_rpc_connection *conn, - struct aki_packet *packet, struct aki_packet *rpacket) +static bool results_callback(void *userdata, struct nn_rpc_connection *conn, + struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_client *client = (struct camu_client *)userdata; (void)conn; (void)rpacket; - u8 op = aki_packet_read_u8(packet); + u8 op = nn_packet_read_u8(packet); switch (op) { case CAMU_CLIENT_CREATE_SEARCH: { - s32 id = aki_packet_read_s32(packet); + s32 id = nn_packet_read_s32(packet); client->callback(client->userdata, CAMU_CLIENT_SEARCH_CREATED, &id); break; } @@ -22,115 +22,115 @@ static bool results_callback(void *userdata, struct aki_rpc_connection *conn, break; } - aki_packet_free(packet); + nn_packet_free(packet); return false; } -static struct aki_rpc_command commands[] = { +static struct nn_rpc_command commands[] = { { .op = CAMU_CLIENT_RESULTS, .callback = results_callback, .userdata = NULL } }; -static void idd_callback(void *userdata, struct aki_packet *packet) +static void idd_callback(void *userdata, struct nn_packet *packet) { struct camu_client *client = (struct camu_client *)userdata; client->callback(client->userdata, CAMU_CLIENT_LOGIN, packet); } -static void connection_callback(void *userdata, struct aki_rpc_connection *conn) +static void connection_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_client *client = (struct camu_client *)userdata; client->conn = conn; - struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_IDENTIFY); - aki_packet_write_u8(packet, CAMU_CLIENT); - aki_packet_write_str(packet, &client->username); - aki_rpc_connection_command(client->conn, packet, idd_callback, client); + struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_IDENTIFY); + nn_packet_write_u8(packet, CAMU_CLIENT); + nn_packet_write_str(packet, &client->username); + nn_rpc_connection_command(client->conn, packet, idd_callback, client); } -static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) +static void connection_closed_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_client *client = (struct camu_client *)userdata; al_assert(client->conn == NULL || client->conn == conn); client->conn = NULL; } -bool camu_client_login(struct camu_client *client, struct aki_event_loop *loop, +bool camu_client_login(struct camu_client *client, struct nn_event_loop *loop, u8 type, str *addr, u16 port, str *username) { client->loop = loop; al_str_clone(&client->username, username); client->conn = NULL; - aki_rpc_init(&client->client, client->loop, connection_callback, connection_closed_callback, client); + nn_rpc_init(&client->client, client->loop, connection_callback, connection_closed_callback, client); for (u32 i = 0; i < ARRAY_SIZE(commands); i++) { commands[i].userdata = client; - aki_rpc_add_command(&client->client, &commands[i]); + nn_rpc_add_command(&client->client, &commands[i]); } - if (!aki_rpc_prepare_client(&client->client, type, CAMU_MULTIPLEX_RPC)) { + if (!nn_rpc_prepare_client(&client->client, type, CAMU_MULTIPLEX_RPC)) { return false; } - aki_rpc_connect(&client->client, addr, port); + nn_rpc_connect(&client->client, addr, port); return true; } void camu_client_create_list(struct camu_client *client, str *name, - void (*callback)(void *, struct aki_packet *), void *userdata) + void (*callback)(void *, struct nn_packet *), void *userdata) { - struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); - aki_packet_write_u8(packet, CAMU_CLIENT_CREATE_LIST); - aki_packet_write_str(packet, name); - aki_rpc_connection_command(client->conn, packet, callback, userdata); + struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); + nn_packet_write_u8(packet, CAMU_CLIENT_CREATE_LIST); + nn_packet_write_str(packet, name); + nn_rpc_connection_command(client->conn, packet, callback, userdata); } void camu_client_toggle_sink(struct camu_client *client, str *sink, str *list, bool enable, - void (*callback)(void *, struct aki_packet *), void *userdata) + void (*callback)(void *, struct nn_packet *), void *userdata) { - struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); - aki_packet_write_u8(packet, CAMU_CLIENT_TOGGLE_SINK); - aki_packet_write_str(packet, sink); - aki_packet_write_str(packet, list); - aki_packet_write_bool(packet, enable); - aki_rpc_connection_command(client->conn, packet, callback, userdata); + struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); + nn_packet_write_u8(packet, CAMU_CLIENT_TOGGLE_SINK); + nn_packet_write_str(packet, sink); + nn_packet_write_str(packet, list); + nn_packet_write_bool(packet, enable); + nn_rpc_connection_command(client->conn, packet, callback, userdata); } void camu_client_create_search(struct camu_client *client, str *module, str *query) { - struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); - aki_packet_write_u8(packet, CAMU_CLIENT_CREATE_SEARCH); - aki_packet_write_str(packet, module); - aki_packet_write_str(packet, query); - aki_rpc_connection_command(client->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); + nn_packet_write_u8(packet, CAMU_CLIENT_CREATE_SEARCH); + nn_packet_write_str(packet, module); + nn_packet_write_str(packet, query); + nn_rpc_connection_command(client->conn, packet, NULL, NULL); } void camu_client_get_page(struct camu_client *client, s32 id, u32 num) { - struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); - aki_packet_write_u8(packet, CAMU_CLIENT_GET_PAGE); - aki_packet_write_s32(packet, id); - aki_packet_write_u32(packet, num); - aki_rpc_connection_command(client->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); + nn_packet_write_u8(packet, CAMU_CLIENT_GET_PAGE); + nn_packet_write_s32(packet, id); + nn_packet_write_u32(packet, num); + nn_rpc_connection_command(client->conn, packet, NULL, NULL); } void camu_client_add_from_path(struct camu_client *client, str *list, str *path) { - struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, list); - aki_packet_write_u8(packet, CAMU_LIST_ADD); - aki_packet_write_u8(packet, CAMU_RESOURCE_FILE); - aki_packet_write_str(packet, path); - aki_rpc_connection_command(client->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_LIST_ACTION); + nn_packet_write_str(packet, list); + nn_packet_write_u8(packet, CAMU_LIST_ADD); + nn_packet_write_u8(packet, CAMU_RESOURCE_FILE); + nn_packet_write_str(packet, path); + nn_rpc_connection_command(client->conn, packet, NULL, NULL); } void camu_client_add_from_post(struct camu_client *client, str *list, str *unique_id, u32 index) { - struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, list); - aki_packet_write_u8(packet, CAMU_LIST_ADD); - aki_packet_write_u8(packet, CAMU_RESOURCE_PORTAL); - aki_packet_write_str(packet, unique_id); - aki_packet_write_u32(packet, index); - aki_rpc_connection_command(client->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_LIST_ACTION); + nn_packet_write_str(packet, list); + nn_packet_write_u8(packet, CAMU_LIST_ADD); + nn_packet_write_u8(packet, CAMU_RESOURCE_PORTAL); + nn_packet_write_str(packet, unique_id); + nn_packet_write_u32(packet, index); + nn_rpc_connection_command(client->conn, packet, NULL, NULL); } void camu_client_disconnect(struct camu_client *client) { - aki_rpc_conn_flush(client->conn); + nn_rpc_conn_flush(client->conn); } diff --git a/src/libclient/client.h b/src/libclient/client.h index c073c9e..63d877d 100644 --- a/src/libclient/client.h +++ b/src/libclient/client.h @@ -1,6 +1,6 @@ #pragma once -#include <aki/rpc2.h> +#include <nnwt/rpc2.h> enum { CAMU_CLIENT_LOGIN = 0, @@ -9,21 +9,21 @@ enum { }; struct camu_client { - struct aki_event_loop *loop; + struct nn_event_loop *loop; str username; - struct aki_rpc client; - struct aki_rpc_connection *conn; + struct nn_rpc client; + struct nn_rpc_connection *conn; void (*callback)(void *, u8, void *); void *userdata; }; -bool camu_client_login(struct camu_client *client, struct aki_event_loop *loop, +bool camu_client_login(struct camu_client *client, struct nn_event_loop *loop, u8 type, str *addr, u16 port, str *username); void camu_client_create_list(struct camu_client *client, str *name, - void (*callback)(void *, struct aki_packet *), void *userdata); + void (*callback)(void *, struct nn_packet *), void *userdata); void camu_client_toggle_sink(struct camu_client *client, str *sink, str *list, bool enable, - void (*callback)(void *, struct aki_packet *), void *userdata); + void (*callback)(void *, struct nn_packet *), void *userdata); void camu_client_create_search(struct camu_client *client, str *module, str *query); void camu_client_get_page(struct camu_client *client, s32 id, u32 num); diff --git a/src/libsink/sink.c b/src/libsink/sink.c index bd33f46..984a41c 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -50,7 +50,7 @@ enum { }; // Number of entries to keep buffered at one time. -#define ENTRY_MAX_AGE 6 +#define ENTRY_MAX_AGE 7 // If a buffer is still INIT or QUEUED after an entry is configured, it's "empty". #define BUFFER_EMPTY(buf) ((buf)->state == BUFFER_INIT || (buf)->state == BUFFER_QUEUED) @@ -82,9 +82,9 @@ enum { #endif #if defined CAMU_SCREEN_THREADED && defined CAMU_MIXER_THREADED -#define BLOCKING_SLEEP(delay) aki_thread_sleep(delay) +#define BLOCKING_SLEEP(delay) nn_thread_sleep(delay) #else -#define BLOCKING_SLEEP(delay) aki_event_loop_sleep(sink->loop, delay) +#define BLOCKING_SLEEP(delay) nn_event_loop_sleep(sink->loop, delay) #endif static inline bool entry_audio_buffer_held(struct camu_sink_entry *entry) @@ -223,7 +223,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) { switch (cmd->op) { case START: { - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); switch (cmd->value.i) { case CAMU_SINK_AUDIO: if (sink->audio.state == SINK_PAUSED) { @@ -240,11 +240,11 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; #endif } - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; } case STOP: { - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); switch (cmd->value.i) { case CAMU_SINK_AUDIO: if (sink->audio.state == SINK_PLAYING) { @@ -261,25 +261,25 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; #endif } - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; } case SKIP: { if (!sink->conn) return; - struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, &sink->default_list); - aki_packet_write_u8(packet, CAMU_LIST_SKIP); - aki_packet_write_s32(packet, get_sequence_for_command(sink)); - aki_packet_write_s32(packet, (s32)cmd->value.i); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); + nn_packet_write_str(packet, &sink->default_list); + nn_packet_write_u8(packet, CAMU_LIST_SKIP); + nn_packet_write_s32(packet, get_sequence_for_command(sink)); + nn_packet_write_s32(packet, (s32)cmd->value.i); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case SHUFFLE: { if (!sink->conn) return; - struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, &sink->default_list); - aki_packet_write_u8(packet, CAMU_LIST_SHUFFLE); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); + nn_packet_write_str(packet, &sink->default_list); + nn_packet_write_u8(packet, CAMU_LIST_SHUFFLE); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case TOGGLE_PAUSE: { @@ -287,30 +287,30 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) sink_local_pause(sink, (struct camu_sink_entry *)cmd->opaque); #else if (!sink->conn) return; - struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, &sink->default_list); - aki_packet_write_u8(packet, CAMU_LIST_TOGGLE_PAUSE); - aki_packet_write_s32(packet, get_sequence_for_command(sink)); - aki_packet_write_f64(packet, cmd->value.f); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); + nn_packet_write_str(packet, &sink->default_list); + nn_packet_write_u8(packet, CAMU_LIST_TOGGLE_PAUSE); + nn_packet_write_s32(packet, get_sequence_for_command(sink)); + nn_packet_write_f64(packet, cmd->value.f); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); #endif break; } case SEEK: { if (!sink->conn) return; - struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, &sink->default_list); - aki_packet_write_u8(packet, CAMU_LIST_SEEK); + struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); + nn_packet_write_str(packet, &sink->default_list); + nn_packet_write_u8(packet, CAMU_LIST_SEEK); struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; if (entry) { - aki_packet_write_s32(packet, entry->sequence); - aki_packet_write_u32(packet, entry->id); + nn_packet_write_s32(packet, entry->sequence); + nn_packet_write_u32(packet, entry->id); } else { - aki_packet_write_s32(packet, LIANA_SEQUENCE_ANY); - aki_packet_write_u32(packet, 0); + nn_packet_write_s32(packet, LIANA_SEQUENCE_ANY); + nn_packet_write_u32(packet, 0); } - aki_packet_write_f64(packet, cmd->value.f); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + nn_packet_write_f64(packet, cmd->value.f); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case RESEEK: { @@ -320,23 +320,23 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) } case UNSET: { if (!sink->conn) return; - struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, &sink->default_list); - aki_packet_write_u8(packet, CAMU_LIST_UNSET); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); + nn_packet_write_str(packet, &sink->default_list); + nn_packet_write_u8(packet, CAMU_LIST_UNSET); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case END: { if (!sink->conn) return; - struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, &sink->default_list); - aki_packet_write_u8(packet, CAMU_LIST_END); - aki_packet_write_u32(packet, (u32)cmd->value.u); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); + nn_packet_write_str(packet, &sink->default_list); + nn_packet_write_u8(packet, CAMU_LIST_END); + nn_packet_write_u32(packet, (u32)cmd->value.u); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case CLOSE: { - aki_signal_stop(&sink->queue_signal); + nn_signal_stop(&sink->queue_signal); sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL); return; } @@ -358,7 +358,7 @@ static void queue_signal_callback(void *userdata) static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd) { camu_queue_push(sink->queue, cmd); - aki_signal_send(&sink->queue_signal); + nn_signal_send(&sink->queue_signal); } static void mixer_callback(void *userdata, u8 op) @@ -373,7 +373,7 @@ static void mixer_callback(void *userdata, u8 op) } } -bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, +bool camu_sink_init(struct camu_sink *sink, struct nn_event_loop *loop, struct camu_mixer *mixer #ifndef CAMU_SINK_NO_VIDEO , struct camu_renderer *renderer @@ -381,9 +381,9 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, ) { sink->loop = loop; - aki_mutex_init(&sink->mutex); - aki_signal_init(&sink->queue_signal, sink->loop, queue_signal_callback, sink); - aki_signal_start(&sink->queue_signal); + nn_mutex_init(&sink->mutex); + nn_signal_init(&sink->queue_signal, sink->loop, queue_signal_callback, sink); + nn_signal_start(&sink->queue_signal); camu_queue_init(sink->queue); sink->queued = NULL; sink->current = NULL; @@ -415,13 +415,7 @@ 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) { - struct camu_sink_entry *rentry; - // al_array_remove_all() - al_array_foreach_rev(sink->previous, i, rentry) { - if (rentry == entry) { - al_array_remove_at(sink->previous, i); - } - } + al_array_remove_all(sink->previous, entry); } // Every call to maybe_add_to_previous() must map to a remove_entry_buffers(). @@ -430,10 +424,10 @@ static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry { al_assert(previous != target && !previous->ended); - struct camu_sink_entry *rentry; - al_array_foreach_rev(sink->previous, i, rentry) { + struct camu_sink_entry *entry; + al_array_foreach_rev(sink->previous, i, entry) { // If the entry we are about to set is in previous, run the queue. - if (rentry == target) { + if (entry == target) { maybe_remove_previous(sink); break; } @@ -441,7 +435,12 @@ static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry // If neither the entries audio or video buffer is ADDED, we don't care about adding it // to previous (waiting for the next added entry to remove it). - if (!(previous->audio.state == BUFFER_ADDED || previous->video.state == BUFFER_ADDED)) { +#ifndef CAMU_SINK_NO_VIDEO + if (!(previous->audio.state == BUFFER_ADDED || previous->video.state == BUFFER_ADDED)) +#else + if (!(previous->audio.state == BUFFER_ADDED)) +#endif + { remove_entry_buffers(sink, previous); return; } @@ -500,7 +499,7 @@ static void maybe_cleanup_old_entries(struct camu_sink *sink) void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) { al_assert(!entry->ended); - al_assert(entry->audio.state != BUFFER_ADDED && entry->audio.state != BUFFER_QUEUED); + al_assert(entry->audio.state != BUFFER_ADDED); if (entry->audio.state == BUFFER_ENDED) { al_log_warn("sink", "Tried to add an ended audio buffer."); return; @@ -548,7 +547,7 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) #ifndef CAMU_SINK_NO_VIDEO void add_video_if_set_and_buffered(struct camu_sink_entry *entry) { - al_assert(entry->video.state != BUFFER_ADDED && entry->video.state != BUFFER_QUEUED); + al_assert(entry->video.state != BUFFER_ADDED); if (entry->video.state == BUFFER_ENDED) { al_log_warn("sink", "Tried to add an ended video buffer."); return; @@ -585,6 +584,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) if (sink->current) { struct camu_sink_entry *current = sink->current; al_assert(current != target); + if (!current->ended) { if (target->ended) { remove_entry_buffers(sink, current); @@ -599,8 +599,9 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) #endif al_assert(AUDIO_REMOVED_OR_EMPTY(current) && VIDEO_REMOVED_OR_EMPTY(current)); } + if (sink->reconnecting) { - al_assert(sink->reconnecting == sink->current); + al_assert(sink->reconnecting == current); sink->reconnecting = NULL; } } @@ -656,9 +657,9 @@ static void audio_buffer_callback(void *userdata, u8 op) struct camu_sink *sink = entry->sink; switch (op) { case CAMU_BUFFER_BUFFERED: - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); add_audio_if_set_and_buffered(entry); - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; case CAMU_BUFFER_CORK: lia_vcr_cork(entry->audio.track); @@ -667,7 +668,7 @@ static void audio_buffer_callback(void *userdata, u8 op) lia_vcr_uncork(entry->audio.track); break; case CAMU_BUFFER_PAUSED: - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); al_log_info("sink", "Audio buffer paused."); if (!entry->audio.ignore_paused) { queue_cmd(sink, (struct camu_sink_cmd){ @@ -675,11 +676,11 @@ static void audio_buffer_callback(void *userdata, u8 op) .value.i = CAMU_SINK_AUDIO }); } - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; case CAMU_BUFFER_EOF: { lia_vcr_cork(entry->audio.track); - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); al_log_info("sink", "Audio EOF."); if (entry->audio.state == BUFFER_ADDED) { remove_entry_audio_buffer(sink, entry); @@ -694,7 +695,7 @@ static void audio_buffer_callback(void *userdata, u8 op) if (run_queue) { end_entry_and_advance_queue(sink, entry); } - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; } } @@ -707,9 +708,9 @@ static void video_buffer_callback(void *userdata, u8 op) struct camu_sink *sink = entry->sink; switch (op) { case CAMU_BUFFER_BUFFERED: { - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); add_video_if_set_and_buffered(entry); - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; } case CAMU_BUFFER_CORK: @@ -721,7 +722,7 @@ static void video_buffer_callback(void *userdata, u8 op) case CAMU_BUFFER_EOF: { lia_vcr_cork(entry->video.track); bool swapped = false; - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); al_log_info("sink", "Video EOF."); if (!ENTRY_IS_SINGLE_FRAME(entry)) { if (entry->video.state == BUFFER_ADDED) { @@ -733,7 +734,7 @@ static void video_buffer_callback(void *userdata, u8 op) swapped = end_entry_and_advance_queue(sink, entry); } } - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); if (!swapped) { queue_cmd(sink, (struct camu_sink_cmd){ .op = STOP, @@ -762,14 +763,12 @@ 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) { - aki_mutex_lock(&sink->mutex); - if (entry == sink->current) { - if (sink->target) { - switch_to(sink, sink->target); - sink->target = NULL; - } + nn_mutex_lock(&sink->mutex); + if (entry == sink->current && sink->target) { + switch_to(sink, sink->target); + sink->target = NULL; } - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); } } @@ -817,7 +816,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str lia_client_disconnect(&entry->client); return; } - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); if (entry->audio.state == BUFFER_QUEUED) { entry->audio.state = BUFFER_SET_OR_BUFFERED; } else if (entry->audio.state == BUFFER_INIT) { @@ -826,7 +825,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str assert(false); } if (VIDEO_ADDED_OR_EMPTY(entry)) evaluate_latency(sink, entry); - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; #ifndef CAMU_SINK_NO_VIDEO case CAMU_STREAM_VIDEO: @@ -836,7 +835,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str lia_client_disconnect(&entry->client); return; } - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); if (entry->video.state == BUFFER_QUEUED) { entry->video.state = BUFFER_SET_OR_BUFFERED; } else if (entry->video.state == BUFFER_INIT) { @@ -845,7 +844,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str assert(false); } if (AUDIO_ADDED_OR_EMPTY(entry)) evaluate_latency(sink, entry); - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; case CAMU_STREAM_SUBTITLE: #ifdef CAMU_HAVE_FFMPEG @@ -907,7 +906,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str case LIANA_CLIENT_REMOVE_BUFFERS: { bool reconnect = *(bool *)opaque; - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); // We need to call this if entry was added to previous then, // - it's being cleaned up after (ENTRY_MAX_AGE - 1) entries were added but none buffered. @@ -917,11 +916,10 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str if (reconnect) { entry->ended = false; if (entry == sink->current) { + sink->reconnecting = entry; if (sink->target) { switch_to(sink, sink->target); sink->target = NULL; - } else { - sink->reconnecting = sink->current; } } } @@ -949,7 +947,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } #endif - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); while ( // Block until buffers are removed. #ifndef CAMU_SINK_NO_VIDEO @@ -957,7 +955,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str #else (!skip_audio && entry_audio_buffer_held(entry)) #endif - ) { BLOCKING_SLEEP(AKI_TS_FROM_USEC(2500)); } + ) { BLOCKING_SLEEP(NNWT_TS_FROM_USEC(2500)); } if (reconnect) { if (BUFFER_NOT_EMPTY(&entry->audio)) { @@ -974,13 +972,13 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } case LIANA_CLIENT_RESUME_AT: { struct lia_timing *time = (struct lia_timing *)opaque; - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); camu_clock_seek(&entry->clock, time->seek_pos / 1000000.0, time->at); - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; } case LIANA_CLIENT_RECONNECTED: { - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); if (entry == sink->reconnecting) { if (entry->audio.state > BUFFER_QUEUED) { add_audio_if_set_and_buffered(entry); @@ -991,7 +989,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } #endif } - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); break; } case LIANA_CLIENT_EOF: { @@ -1015,11 +1013,16 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str break; } case LIANA_CLIENT_CLOSED: { - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); + if (entry == sink->target) { maybe_unset_current(sink); sink->target = NULL; } else if (entry == sink->current) { + if (sink->reconnecting) { + al_assert(sink->reconnecting == sink->current); + sink->reconnecting = NULL; + } // current's buffers will already be removed. sink->current = NULL; if (sink->target) { @@ -1031,26 +1034,17 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str }/* else if (entry == sink->queued) { sink->queued = NULL; }*/ + lia_client_free(&entry->client); camu_audio_buffer_free(&entry->audio.buf); #ifndef CAMU_SINK_NO_VIDEO camu_video_buffer_free(&entry->video.buf); #endif - bool removed = false; - struct camu_sink_entry *rentry; - al_array_foreach(sink->entries, i, rentry) { - if (rentry == entry) { - al_array_remove_at(sink->entries, i); - removed = true; - al_log_info("sink", "Entry closed by disconnect."); - break; - } - } + bool removed; + al_array_check_remove(sink->entries, entry, removed); + nn_mutex_unlock(&sink->mutex); al_free(entry); - if (!removed) { - al_log_info("sink", "Entry closed by cleanup."); - } - aki_mutex_unlock(&sink->mutex); + al_log_info("sink", "Entry closed by %s.", removed ? "disconnect" : "cleanup"); break; } } @@ -1097,16 +1091,16 @@ static struct camu_sink_entry *get_entry_from_id(struct camu_sink *sink, u32 id) return NULL; } -static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn, - struct aki_packet *packet, struct aki_packet *rpacket) +static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn, + struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_sink *sink = (struct camu_sink *)userdata; (void)conn; (void)rpacket; - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); - u8 op = aki_packet_read_u8(packet); + u8 op = nn_packet_read_u8(packet); if (op == LIANA_SINK_UNSET) { maybe_unset_current(sink); goto out; @@ -1114,18 +1108,18 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn // Liana node info. str addr; - aki_packet_read_str(packet, &addr); - u16 port = aki_packet_read_u16(packet); - u32 node_id = aki_packet_read_u32(packet); + nn_packet_read_str(packet, &addr); + u16 port = nn_packet_read_u16(packet); + u32 node_id = nn_packet_read_u32(packet); // List entry info. - u32 id = aki_packet_read_u32(packet); - s32 sequence = aki_packet_read_s32(packet); - u64 at = aki_packet_read_u64(packet); - u64 seek_pos = aki_packet_read_u64(packet); - u8 pause = aki_packet_read_u8(packet); - bool previous_ended = aki_packet_read_bool(packet); - bool ended = aki_packet_read_bool(packet); + u32 id = nn_packet_read_u32(packet); + s32 sequence = nn_packet_read_s32(packet); + u64 at = nn_packet_read_u64(packet); + u64 seek_pos = nn_packet_read_u64(packet); + u8 pause = nn_packet_read_u8(packet); + bool previous_ended = nn_packet_read_bool(packet); + bool ended = nn_packet_read_bool(packet); (void)previous_ended; bool created = false; @@ -1177,7 +1171,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn if (sink->target) { sink->target = NULL; } else { - al_assert(sink->current != entry); + al_assert(entry != sink->current); } if (entry != sink->current) { switch_to(sink, entry); @@ -1188,7 +1182,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn if (sink->target) { sink->target = NULL; } else { - al_assert(sink->current != entry); + al_assert(entry != sink->current); } if (entry != sink->current) { switch_to(sink, entry); @@ -1240,32 +1234,32 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn #endif out: - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); if (op != LIANA_SINK_BUFFER) { maybe_cleanup_old_entries(sink); } - aki_packet_free(packet); + nn_packet_free(packet); return false; } -static bool pause_command_callback(void *userdata, struct aki_rpc_connection *conn, - struct aki_packet *packet, struct aki_packet *rpacket) +static bool pause_command_callback(void *userdata, struct nn_rpc_connection *conn, + struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_sink *sink = (struct camu_sink *)userdata; (void)conn; (void)rpacket; - u32 id = aki_packet_read_u32(packet); - s32 sequence = aki_packet_read_s32(packet); - u64 at = aki_packet_read_u64(packet); - u8 pause = aki_packet_read_u8(packet); + u32 id = nn_packet_read_u32(packet); + s32 sequence = nn_packet_read_s32(packet); + u64 at = nn_packet_read_u64(packet); + u8 pause = nn_packet_read_u8(packet); struct camu_sink_entry *entry = get_entry_from_id(sink, id); if (!entry) goto out; al_assert(entry->sequence == sequence); - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); #ifdef CAMU_SINK_LOCAL (void)at; (void)pause; @@ -1287,25 +1281,25 @@ static bool pause_command_callback(void *userdata, struct aki_rpc_connection *co break; } #endif - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); out: - aki_packet_free(packet); + nn_packet_free(packet); return false; } -static bool seek_command_callback(void *userdata, struct aki_rpc_connection *conn, - struct aki_packet *packet, struct aki_packet *rpacket) +static bool seek_command_callback(void *userdata, struct nn_rpc_connection *conn, + struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_sink *sink = (struct camu_sink *)userdata; (void)conn; (void)rpacket; - u32 id = aki_packet_read_u32(packet); - s32 sequence = aki_packet_read_s32(packet); + u32 id = nn_packet_read_u32(packet); + s32 sequence = nn_packet_read_s32(packet); (void)sequence; - u64 at = aki_packet_read_u64(packet); - u64 pos = aki_packet_read_u64(packet); + u64 at = nn_packet_read_u64(packet); + u64 pos = nn_packet_read_u64(packet); struct camu_sink_entry *entry = get_entry_from_id(sink, id); if (!entry) goto out; @@ -1317,81 +1311,81 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con lia_client_seek(&entry->client, pos, at); out: - aki_packet_free(packet); + nn_packet_free(packet); return false; } -static struct aki_rpc_command commands[] = { +static struct nn_rpc_command commands[] = { { .op = CAMU_SINK_SET, .callback = set_command_callback, .userdata = NULL }, { .op = CAMU_SINK_PAUSE, .callback = pause_command_callback, .userdata = NULL }, { .op = CAMU_SINK_SEEK, .callback = seek_command_callback, .userdata = NULL } }; -static void idd_callback(void *userdata, struct aki_packet *packet) +static void idd_callback(void *userdata, struct nn_packet *packet) { struct camu_sink *sink = (struct camu_sink *)userdata; (void)sink; - aki_packet_free(packet); + nn_packet_free(packet); } -static void connection_callback(void *userdata, struct aki_rpc_connection *conn) +static void connection_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_sink *sink = (struct camu_sink *)userdata; sink->conn = conn; - aki_timer_stop(&sink->reconnect_timer); - struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_IDENTIFY); - aki_packet_write_u8(packet, CAMU_SINK); - aki_packet_write_str(packet, &sink->name); - aki_rpc_connection_command(sink->conn, packet, idd_callback, sink); + nn_timer_stop(&sink->reconnect_timer); + struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_IDENTIFY); + nn_packet_write_u8(packet, CAMU_SINK); + nn_packet_write_str(packet, &sink->name); + nn_rpc_connection_command(sink->conn, packet, idd_callback, sink); } -static void reconnect_timer_callback(void *userdata, struct aki_timer *timer) +static void reconnect_timer_callback(void *userdata, struct nn_timer *timer) { struct camu_sink *sink = (struct camu_sink *)userdata; (void)timer; - aki_rpc_reconnect(&sink->client, &sink->addr, sink->port); + nn_rpc_reconnect(&sink->client, &sink->addr, sink->port); } -static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) +static void connection_closed_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_sink *sink = (struct camu_sink *)userdata; if (sink->conn) { al_assert(sink->conn == conn); sink->conn = NULL; } - aki_timer_again(&sink->reconnect_timer); + nn_timer_again(&sink->reconnect_timer); } bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str *name) { al_str_clone(&sink->name, name); sink->type = type; - aki_timer_init(&sink->reconnect_timer, sink->loop, reconnect_timer_callback, sink); - aki_timer_set_repeat(&sink->reconnect_timer, AKI_TS_FROM_USEC(1000000)); - aki_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink); + nn_timer_init(&sink->reconnect_timer, sink->loop, reconnect_timer_callback, sink); + nn_timer_set_repeat(&sink->reconnect_timer, NNWT_TS_FROM_USEC(1000000)); + nn_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink); for (u32 i = 0; i < ARRAY_SIZE(commands); i++) { commands[i].userdata = sink; al_assert(sink->callback); - aki_rpc_add_command(&sink->client, &commands[i]); + nn_rpc_add_command(&sink->client, &commands[i]); } - if (!aki_rpc_prepare_client(&sink->client, sink->type, CAMU_MULTIPLEX_RPC)) { + if (!nn_rpc_prepare_client(&sink->client, sink->type, CAMU_MULTIPLEX_RPC)) { return false; } al_str_clone(&sink->addr, addr); sink->port = port; - aki_rpc_connect(&sink->client, &sink->addr, sink->port); + nn_rpc_connect(&sink->client, &sink->addr, sink->port); return true; } struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink) { - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); return sink->current; } void camu_sink_return_current(struct camu_sink *sink) { - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); } void camu_sink_skip(struct camu_sink *sink, s32 n) @@ -1411,9 +1405,9 @@ void camu_sink_shuffle(struct camu_sink *sink) void camu_sink_toggle_pause(struct camu_sink *sink) { - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); struct camu_sink_entry *current = sink->current; - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); if (current) { queue_cmd(sink, (struct camu_sink_cmd){ .op = TOGGLE_PAUSE, @@ -1425,9 +1419,9 @@ void camu_sink_toggle_pause(struct camu_sink *sink) void camu_sink_seek(struct camu_sink *sink, f64 precent) { - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); struct camu_sink_entry *current = sink->current; - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); queue_cmd(sink, (struct camu_sink_cmd){ .op = SEEK, .value.f = precent, @@ -1437,9 +1431,9 @@ void camu_sink_seek(struct camu_sink *sink, f64 precent) void camu_sink_reseek(struct camu_sink *sink) { - aki_mutex_lock(&sink->mutex); + nn_mutex_lock(&sink->mutex); struct camu_sink_entry *current = sink->current; - aki_mutex_unlock(&sink->mutex); + nn_mutex_unlock(&sink->mutex); queue_cmd(sink, (struct camu_sink_cmd){ .op = RESEEK, .opaque = current @@ -1471,9 +1465,9 @@ void camu_sink_stop(struct camu_sink *sink) void camu_sink_close(struct camu_sink *sink) { - aki_timer_stop(&sink->reconnect_timer); - aki_timer_disable(&sink->reconnect_timer); - if (sink->conn) aki_rpc_conn_disconnect(sink->conn); + nn_timer_stop(&sink->reconnect_timer); + nn_timer_disable(&sink->reconnect_timer); + if (sink->conn) nn_rpc_conn_disconnect(sink->conn); struct camu_sink_entry *entry; al_array_foreach_rev(sink->entries, i, entry) { al_array_remove_at(sink->entries, i); @@ -1488,9 +1482,9 @@ void camu_sink_free(struct camu_sink *sink) al_free(entry); } al_array_free(sink->entries); - aki_rpc_free(&sink->client); + nn_rpc_free(&sink->client); camu_queue_free(sink->queue); - aki_mutex_destroy(&sink->mutex); + nn_mutex_destroy(&sink->mutex); al_str_free(&sink->addr); al_str_free(&sink->name); } diff --git a/src/libsink/sink.h b/src/libsink/sink.h index 48b2b87..4f719d2 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -2,9 +2,9 @@ #include <al/str.h> #include <al/array.h> -#include <aki/rpc2.h> -#include <aki/signal.h> -#include <aki/timer.h> +#include <nnwt/rpc2.h> +#include <nnwt/signal.h> +#include <nnwt/timer.h> #include "../util/queue.h" @@ -69,16 +69,16 @@ struct camu_sink_cmd { }; struct camu_sink { - struct aki_event_loop *loop; + struct nn_event_loop *loop; str name; u8 type; str addr; u16 port; - struct aki_rpc client; - struct aki_rpc_connection *conn; - struct aki_mutex mutex; - struct aki_timer reconnect_timer; - struct aki_signal queue_signal; + struct nn_rpc client; + struct nn_rpc_connection *conn; + struct nn_mutex mutex; + struct nn_timer reconnect_timer; + struct nn_signal queue_signal; queue(struct camu_sink_cmd) queue; str default_list; struct camu_sink_entry *current; @@ -102,7 +102,7 @@ struct camu_sink { void *userdata; }; -bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, +bool camu_sink_init(struct camu_sink *sink, struct nn_event_loop *loop, struct camu_mixer *mixer #ifndef CAMU_SINK_NO_VIDEO , struct camu_renderer *renderer diff --git a/src/mixer/audio_miniaudio.c b/src/mixer/audio_miniaudio.c index 526e464..f0afeee 100644 --- a/src/mixer/audio_miniaudio.c +++ b/src/mixer/audio_miniaudio.c @@ -130,7 +130,7 @@ static s32 camu_sample_format_from_miniaudio(s32 fmt) return CAMU_SAMPLE_FORMAT_FLT; case ma_format_unknown: default: - al_assert_and_return(false); + al_assert_and_return(CAMU_SAMPLE_FORMAT_NONE); } } diff --git a/src/mixer/audio_miniaudio.h b/src/mixer/audio_miniaudio.h index 5bae185..c067024 100644 --- a/src/mixer/audio_miniaudio.h +++ b/src/mixer/audio_miniaudio.h @@ -1,9 +1,9 @@ #pragma once -// Include <aki/thread.h> here to ensure OS related headers +// Include <nnwt/thread.h> here to ensure OS related headers // get included the way we want. // Specifically, on Windows we want to include <winsock2.h> before <windows.h>. -#include <aki/thread.h> +#include <nnwt/thread.h> #include <al/macros.h> #include <miniaudio.h> diff --git a/src/mixer/mixer.c b/src/mixer/mixer.c index 086498e..18d6b8a 100644 --- a/src/mixer/mixer.c +++ b/src/mixer/mixer.c @@ -76,7 +76,7 @@ bool camu_mixer_init(struct camu_mixer *mixer, struct camu_audio *audio) al_array_init(mixer->add_queue); al_array_init(mixer->rem_queue); al_atomic_store(u8)(&mixer->queued, 0, AL_ATOMIC_RELAXED); - aki_mutex_init(&mixer->mutex); + nn_mutex_init(&mixer->mutex); #endif return true; } @@ -89,7 +89,7 @@ void camu_mixer_pick_format(struct camu_mixer *mixer, struct camu_resampler_form void camu_mixer_set_volume(struct camu_mixer *mixer, f32 volume) { #ifdef CAMU_MIXER_THREADED - aki_mutex_lock(&mixer->mutex); + nn_mutex_lock(&mixer->mutex); #endif mixer->volume = volume; struct camu_audio_buffer *buf; @@ -97,14 +97,14 @@ void camu_mixer_set_volume(struct camu_mixer *mixer, f32 volume) camu_audio_buffer_set_volume(buf, mixer->volume); } #ifdef CAMU_MIXER_THREADED - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); #endif } void camu_mixer_offset_volume(struct camu_mixer *mixer, f32 amount) { #ifdef CAMU_MIXER_THREADED - aki_mutex_lock(&mixer->mutex); + nn_mutex_lock(&mixer->mutex); #endif if (FLT_MAX - mixer->volume < amount) { mixer->volume = FLT_MAX; @@ -116,7 +116,7 @@ void camu_mixer_offset_volume(struct camu_mixer *mixer, f32 amount) camu_audio_buffer_set_volume(buf, mixer->volume); } #ifdef CAMU_MIXER_THREADED - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); #endif } @@ -152,15 +152,19 @@ static void remove_buffer_internal(struct camu_mixer *mixer, struct camu_audio_b static void run_queue_internal(struct camu_mixer *mixer) { struct camu_audio_buffer *buf; - al_array_foreach(mixer->rem_queue, i, buf) { - remove_buffer_internal(mixer, buf); + if (mixer->rem_queue.size > 0) { + al_array_foreach(mixer->rem_queue, i, buf) { + remove_buffer_internal(mixer, buf); + } + mixer->rem_queue.size = 0; } - mixer->rem_queue.size = 0; - al_array_reserve(mixer->buffers, mixer->buffers.size + mixer->add_queue.size); - al_array_foreach(mixer->add_queue, i, buf) { - add_buffer_internal(mixer, buf); + if (mixer->add_queue.size > 0) { + al_array_reserve(mixer->buffers, mixer->buffers.size + mixer->add_queue.size); + al_array_foreach(mixer->add_queue, i, buf) { + add_buffer_internal(mixer, buf); + } + mixer->add_queue.size = 0; } - mixer->add_queue.size = 0; al_atomic_store(u8)(&mixer->queued, 0, AL_ATOMIC_RELAXED); if (al_array_is_empty(mixer->buffers)) { mixer->empty_after = MIXER_TRAILING_SILENCE; @@ -173,11 +177,11 @@ static void run_queue_internal(struct camu_mixer *mixer) void camu_mixer_add_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *buf) { #ifdef CAMU_MIXER_THREADED - aki_mutex_lock(&mixer->mutex); + nn_mutex_lock(&mixer->mutex); struct camu_audio_buffer *rbuf; al_array_foreach(mixer->add_queue, i, rbuf) { if (rbuf == buf) { - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); return; } } @@ -188,13 +192,13 @@ void camu_mixer_add_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *b if (queue_empty) { al_atomic_store(u8)(&mixer->queued, 0, AL_ATOMIC_RELAXED); } - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); return; } } al_array_push(mixer->add_queue, buf); al_atomic_store(u8)(&mixer->queued, 1, AL_ATOMIC_RELAXED); - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); #else add_buffer_internal(mixer, buf); mixer->empty_after = 0; @@ -204,11 +208,11 @@ void camu_mixer_add_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *b void camu_mixer_remove_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *buf) { #ifdef CAMU_MIXER_THREADED - aki_mutex_lock(&mixer->mutex); + nn_mutex_lock(&mixer->mutex); struct camu_audio_buffer *rbuf; al_array_foreach(mixer->rem_queue, i, rbuf) { if (rbuf == buf) { - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); return; } } @@ -219,7 +223,7 @@ void camu_mixer_remove_buffer(struct camu_mixer *mixer, struct camu_audio_buffer if (queue_empty) { al_atomic_store(u8)(&mixer->queued, 0, AL_ATOMIC_RELAXED); } - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); return; } } @@ -228,7 +232,7 @@ void camu_mixer_remove_buffer(struct camu_mixer *mixer, struct camu_audio_buffer if (mixer->paused) { run_queue_internal(mixer); } - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); #else remove_buffer_internal(mixer, buf); if (al_array_is_empty(mixer->buffers)) { @@ -240,9 +244,9 @@ void camu_mixer_remove_buffer(struct camu_mixer *mixer, struct camu_audio_buffer #ifdef CAMU_MIXER_THREADED void camu_mixer_run_queue(struct camu_mixer *mixer) { - aki_mutex_lock(&mixer->mutex); + nn_mutex_lock(&mixer->mutex); run_queue_internal(mixer); - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); } void camu_mixer_clear(struct camu_mixer *mixer) @@ -257,7 +261,7 @@ void camu_mixer_clear(struct camu_mixer *mixer) void camu_mixer_pause(struct camu_mixer *mixer) { #ifdef CAMU_MIXER_THREADED - aki_mutex_lock(&mixer->mutex); + nn_mutex_lock(&mixer->mutex); #endif if (!mixer->paused) { mixer->audio->stop(mixer->audio); @@ -265,14 +269,14 @@ void camu_mixer_pause(struct camu_mixer *mixer) } #ifdef CAMU_MIXER_THREADED run_queue_internal(mixer); - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); #endif } void camu_mixer_resume(struct camu_mixer *mixer) { #ifdef CAMU_MIXER_THREADED - aki_mutex_lock(&mixer->mutex); + nn_mutex_lock(&mixer->mutex); #endif if (mixer->paused) { // start() can internally call data_callback once before returning. @@ -283,7 +287,7 @@ void camu_mixer_resume(struct camu_mixer *mixer) mixer->paused = false; } #ifdef CAMU_MIXER_THREADED - aki_mutex_unlock(&mixer->mutex); + nn_mutex_unlock(&mixer->mutex); #endif } @@ -293,7 +297,7 @@ void camu_mixer_close(struct camu_mixer *mixer) #ifdef CAMU_MIXER_THREADED al_array_free(mixer->add_queue); al_array_free(mixer->rem_queue); - aki_mutex_destroy(&mixer->mutex); + nn_mutex_destroy(&mixer->mutex); #endif mixer->audio->free(&mixer->audio); } diff --git a/src/mixer/mixer.h b/src/mixer/mixer.h index a925e0c..16ec924 100644 --- a/src/mixer/mixer.h +++ b/src/mixer/mixer.h @@ -4,7 +4,7 @@ #include <al/array.h> #ifdef CAMU_MIXER_THREADED #include <al/atomic.h> -#include <aki/thread.h> +#include <nnwt/thread.h> #endif #include "../codec/codec.h" @@ -24,7 +24,7 @@ struct camu_mixer { array(struct camu_audio_buffer *) add_queue; array(struct camu_audio_buffer *) rem_queue; atomic(u8) queued; - struct aki_mutex mutex; + struct nn_mutex mutex; #endif void (*callback)(void *, u8); void *userdata; diff --git a/src/portal/meson.build b/src/portal/meson.build index d3904ea..e87b588 100644 --- a/src/portal/meson.build +++ b/src/portal/meson.build @@ -10,7 +10,7 @@ portal_args = ['-DCAMU_HAVE_PORTAL'] cpy_dir = join_paths(meson.current_source_dir(), 'cpy') run_command('sh', join_paths(cpy_dir, 'build.sh'), cpy_dir, check: false) -python3_embed = import('python').find_installation('python3.12').dependency(embed: true) +python3_embed = import('python').find_installation('python3').dependency(embed: true) portal_deps += [python3_embed] portal = declare_dependency(sources: portal_src, dependencies: portal_deps, compile_args: portal_args) diff --git a/src/portal/py/modules/__init__.py b/src/portal/py/modules/__init__.py index c44e571..cc34b82 100644 --- a/src/portal/py/modules/__init__.py +++ b/src/portal/py/modules/__init__.py @@ -14,11 +14,11 @@ ALL_MODULES = { # config.TWITTER_ACCESS_TOKEN, config.TWITTER_ACCESS_TOKEN_SECRET, # config.TWITTER_CONSUMER_TOKEN, config.TWITTER_CONSUMER_TOKEN_SECRET), []), # 'twitter': (TwitterScrapeModule(config.TWITTER_COOKIES_PATH), []), - 'pixiv_web': (PixivWebModule(config.PIXIV_SESSID, config.PIXIV_USERID), []), +# 'pixiv_web': (PixivWebModule(config.PIXIV_SESSID, config.PIXIV_USERID), []), # 'pixiv_app': (PixivAppModule(config.PIXIV_REFRESH_TOKEN), []), 'fanbox': (FanboxModule(config.FANBOX_SESSID, config.FANBOX_CF_CLEARANCE), []), - 'instagram': (InstagramModule( - config.INSTAGRAM_USER_AGENT, config.INSTAGRAM_SETTINGS_PATH, - config.INSTAGRAM_SESSION_ID), []), +# 'instagram': (InstagramModule( +# config.INSTAGRAM_USER_AGENT, config.INSTAGRAM_SETTINGS_PATH, +# config.INSTAGRAM_SESSION_ID), []), # 'patreon': (PatreonModule(config.PATREON_SESSION_ID, config.PATREON_UUID, config.PATREON_USER_ID), []) } diff --git a/src/portal/py/modules/youtube.py b/src/portal/py/modules/youtube.py index 5125f85..582a16c 100644 --- a/src/portal/py/modules/youtube.py +++ b/src/portal/py/modules/youtube.py @@ -26,7 +26,8 @@ ydl_opts = { 'quiet': False, 'logger': YDLLogger(), 'cachedir': False, - 'socket_timeout': 10 + 'socket_timeout': 10, + 'extractor_args': {'youtube': {'skip': ['hls'], 'player_client': ['android_vr']}} # 'cookiefile': '' } @@ -94,7 +95,7 @@ class YoutubeBase(Search): if info.get('direct'): url = info['url'] else: - url = get_playback_url(info, video=True) + url = get_playback_url(info, video=False) if not url: return False unique_id = f'youtube:v:{info['id']}' diff --git a/src/portal/py/tests/archive_query.py b/src/portal/py/tests/archive_query.py index cff5e5d..1ed9db4 100644 --- a/src/portal/py/tests/archive_query.py +++ b/src/portal/py/tests/archive_query.py @@ -9,10 +9,12 @@ sys.path.append('../') from base import Method from post import PostEncoder, PostType, DateType, User from modules import ALL_MODULES -from logger import Logger +from tests.logger import Logger +#mode = 'pixiv_web' mode = 'fanbox' cmd = 'user' +#cmd = 'bookmarks' module = ALL_MODULES[mode][0] if not module.init(): sys.exit(1) @@ -170,8 +172,8 @@ class QueryDownloadThread(threading.Thread): #RUNTIME_PATH = "./run" RUNTIME_PATH = "/mnt/store/files/tmp/run" -LOG_FILE = f'{RUNTIME_PATH}/archive3.log' -ARG_FILE = f'{RUNTIME_PATH}/completed_args3.log' +LOG_FILE = f'{RUNTIME_PATH}/archive4.log' +ARG_FILE = f'{RUNTIME_PATH}/completed_args4.log' def read_completed_args(path): args = [] @@ -204,7 +206,23 @@ def download_args(): os.mkdir(output_dir) #args = module.search(f'following:{22781328}') - args = ['anoh223'] + #args = ['22781328'] + #args = ['anoh223'] + #args = ['532891', '1324074',] + #args = ['6600030', '163436'] + #args = ['85117546', '29506', '3298353', '21674944', '2034628', '3085731'] + #args = ['6827013', '11742', '4671', '36323560', '5128785', '30970895'] + #args = ['mugi49'] + #args = ['71598064', '9204502', '689283', '19956185', '42987864', '103682732', '4675776', '7780', '160755', '250142', '3444594', '1995476', '357647'] + #args = ['10473280'] + #args = ['178217', '106414', '882896', '3040749'] + #args = ['21245269'] + #args = ['1969907', '292644', '746841'] + #args = ['1319954', '43718238', '2721319', '53459337'] + # TODO: + args = ['yzr'] + #args = ['11187954'] + # twitter test: butcha_u, sep__rina threads = [] start = True diff --git a/src/portal/requirements.txt b/src/portal/requirements.txt index f4982dc..c2e2915 100644 --- a/src/portal/requirements.txt +++ b/src/portal/requirements.txt @@ -1,4 +1,3 @@ httpx[http2,brotli,zstd] blake3 pillow -notcurses diff --git a/src/portal/scripts/create_vendor.sh b/src/portal/scripts/create_vendor.sh index 414a4c2..833d143 100755 --- a/src/portal/scripts/create_vendor.sh +++ b/src/portal/scripts/create_vendor.sh @@ -24,13 +24,13 @@ cd instagrapi/ pip3 install . --upgrade --target=../vendor cd ../ -cd CyberDropDownloader/ -pip3 install . --upgrade --target=../vendor -cd ../ - -cd gallery-dl/ -pip3 install . --upgrade --target=../vendor -cd ../ +#cd CyberDropDownloader/ +#pip3 install . --upgrade --target=../vendor +#cd ../ +# +#cd gallery-dl/ +#pip3 install . --upgrade --target=../vendor +#cd ../ cd ../ rm -r venv diff --git a/src/portal/scripts/run_py.sh b/src/portal/scripts/run_py.sh index 4414f93..e745a72 100755 --- a/src/portal/scripts/run_py.sh +++ b/src/portal/scripts/run_py.sh @@ -1,5 +1,6 @@ #! /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 $@ diff --git a/src/portal/src/packet_ext.c b/src/portal/src/packet_ext.c index 833e7d9..9b6d6be 100644 --- a/src/portal/src/packet_ext.c +++ b/src/portal/src/packet_ext.c @@ -1,102 +1,102 @@ #include "packet_ext.h" -static void aki_packet_write_optional_int(struct aki_packet *packet, optional_int *o) +static void nn_packet_write_optional_int(struct nn_packet *packet, optional_int *o) { - AKI_PACKET_WRITE_TYPE(packet, s64, o->i); - AKI_PACKET_WRITE_TYPE(packet, bool, o->set); + NNWT_PACKET_WRITE_TYPE(packet, s64, o->i); + NNWT_PACKET_WRITE_TYPE(packet, bool, o->set); } -void aki_packet_write_post(struct aki_packet *packet, struct camu_post *post) +void nn_packet_write_post(struct nn_packet *packet, struct camu_post *post) { - AKI_PACKET_WRITE_TYPE(packet, u16, post->version); - AKI_PACKET_WRITE_TYPE(packet, u8, post->type); - aki_packet_write_str(packet, &post->unique_id); - aki_packet_write_str(packet, &post->url); - AKI_PACKET_WRITE_TYPE(packet, u32, post->dates.size); + NNWT_PACKET_WRITE_TYPE(packet, u16, post->version); + NNWT_PACKET_WRITE_TYPE(packet, u8, post->type); + nn_packet_write_str(packet, &post->unique_id); + nn_packet_write_str(packet, &post->url); + NNWT_PACKET_WRITE_TYPE(packet, u32, post->dates.size); struct camu_post_date *date; al_array_foreach_ptr(post->dates, i, date) { - AKI_PACKET_WRITE_TYPE(packet, u8, date->type); - AKI_PACKET_WRITE_TYPE(packet, f32, date->range_start); - AKI_PACKET_WRITE_TYPE(packet, f32, date->range_end); - AKI_PACKET_WRITE_TYPE(packet, u32, date->meta); + NNWT_PACKET_WRITE_TYPE(packet, u8, date->type); + NNWT_PACKET_WRITE_TYPE(packet, f32, date->range_start); + NNWT_PACKET_WRITE_TYPE(packet, f32, date->range_end); + NNWT_PACKET_WRITE_TYPE(packet, u32, date->meta); } - aki_packet_write_str(packet, &post->author.unique_id); - aki_packet_write_wstr(packet, &post->author.username); - aki_packet_write_wstr(packet, &post->author.display_name); - aki_packet_write_str(packet, &post->author.profile_picture_url); - aki_packet_write_wstr(packet, &post->title); - aki_packet_write_wstr(packet, &post->text); - aki_packet_write_optional_int(packet, &post->likes); - aki_packet_write_optional_int(packet, &post->bookmarks); - aki_packet_write_optional_int(packet, &post->reposts); - aki_packet_write_optional_int(packet, &post->quotes); - aki_packet_write_optional_int(packet, &post->comments); - aki_packet_write_optional_int(packet, &post->views); - AKI_PACKET_WRITE_TYPE(packet, u32, post->media.size); + nn_packet_write_str(packet, &post->author.unique_id); + nn_packet_write_wstr(packet, &post->author.username); + nn_packet_write_wstr(packet, &post->author.display_name); + nn_packet_write_str(packet, &post->author.profile_picture_url); + nn_packet_write_wstr(packet, &post->title); + nn_packet_write_wstr(packet, &post->text); + nn_packet_write_optional_int(packet, &post->likes); + nn_packet_write_optional_int(packet, &post->bookmarks); + nn_packet_write_optional_int(packet, &post->reposts); + nn_packet_write_optional_int(packet, &post->quotes); + nn_packet_write_optional_int(packet, &post->comments); + nn_packet_write_optional_int(packet, &post->views); + NNWT_PACKET_WRITE_TYPE(packet, u32, post->media.size); struct camu_post_media *media; al_array_foreach_ptr(post->media, i, media) { - AKI_PACKET_WRITE_TYPE(packet, u8, media->type); - aki_packet_write_str(packet, &media->key); - aki_packet_write_str(packet, &media->url); - aki_packet_write_str(packet, &media->ext); - aki_packet_write_str(packet, &media->thumbnail_url); - aki_packet_write_str(packet, &media->thumbnail_ext); + NNWT_PACKET_WRITE_TYPE(packet, u8, media->type); + nn_packet_write_str(packet, &media->key); + nn_packet_write_str(packet, &media->url); + nn_packet_write_str(packet, &media->ext); + nn_packet_write_str(packet, &media->thumbnail_url); + nn_packet_write_str(packet, &media->thumbnail_ext); } - aki_packet_write_str(packet, &post->post.unique_id); - aki_packet_write_str(packet, &post->quoted.unique_id); - aki_packet_write_str(packet, &post->in_reply_to.unique_id); + nn_packet_write_str(packet, &post->post.unique_id); + nn_packet_write_str(packet, &post->quoted.unique_id); + nn_packet_write_str(packet, &post->in_reply_to.unique_id); } -static void aki_packet_read_optional_int(struct aki_packet *packet, optional_int *o) +static void nn_packet_read_optional_int(struct nn_packet *packet, optional_int *o) { - AKI_PACKET_READ_TYPE(packet, s64, o->i); - AKI_PACKET_READ_TYPE(packet, bool, o->set); + NNWT_PACKET_READ_TYPE(packet, s64, o->i); + NNWT_PACKET_READ_TYPE(packet, bool, o->set); } -void aki_packet_read_post(struct aki_packet *packet, struct camu_post *post) +void nn_packet_read_post(struct nn_packet *packet, struct camu_post *post) { camu_post_reset(post); - AKI_PACKET_READ_TYPE(packet, u16, post->version); - AKI_PACKET_READ_TYPE(packet, u8, post->type); - aki_packet_read_str(packet, &post->unique_id); - aki_packet_read_str(packet, &post->url); + NNWT_PACKET_READ_TYPE(packet, u16, post->version); + NNWT_PACKET_READ_TYPE(packet, u8, post->type); + nn_packet_read_str(packet, &post->unique_id); + nn_packet_read_str(packet, &post->url); u32 dates_size; - AKI_PACKET_READ_TYPE(packet, u32, dates_size); + NNWT_PACKET_READ_TYPE(packet, u32, dates_size); al_array_reserve(post->dates, dates_size); for (u32 i = 0; i < dates_size; i++) { struct camu_post_date date; - AKI_PACKET_READ_TYPE(packet, u8, date.type); - AKI_PACKET_READ_TYPE(packet, f32, date.range_start); - AKI_PACKET_READ_TYPE(packet, f32, date.range_end); - AKI_PACKET_READ_TYPE(packet, u32, date.meta); + NNWT_PACKET_READ_TYPE(packet, u8, date.type); + NNWT_PACKET_READ_TYPE(packet, f32, date.range_start); + NNWT_PACKET_READ_TYPE(packet, f32, date.range_end); + NNWT_PACKET_READ_TYPE(packet, u32, date.meta); al_array_push(post->dates, date); } - aki_packet_read_str(packet, &post->author.unique_id); - aki_packet_read_wstr(packet, &post->author.username); - aki_packet_read_wstr(packet, &post->author.display_name); - aki_packet_read_str(packet, &post->author.profile_picture_url); - aki_packet_read_wstr(packet, &post->title); - aki_packet_read_wstr(packet, &post->text); - aki_packet_read_optional_int(packet, &post->likes); - aki_packet_read_optional_int(packet, &post->bookmarks); - aki_packet_read_optional_int(packet, &post->reposts); - aki_packet_read_optional_int(packet, &post->quotes); - aki_packet_read_optional_int(packet, &post->comments); - aki_packet_read_optional_int(packet, &post->views); + nn_packet_read_str(packet, &post->author.unique_id); + nn_packet_read_wstr(packet, &post->author.username); + nn_packet_read_wstr(packet, &post->author.display_name); + nn_packet_read_str(packet, &post->author.profile_picture_url); + nn_packet_read_wstr(packet, &post->title); + nn_packet_read_wstr(packet, &post->text); + nn_packet_read_optional_int(packet, &post->likes); + nn_packet_read_optional_int(packet, &post->bookmarks); + nn_packet_read_optional_int(packet, &post->reposts); + nn_packet_read_optional_int(packet, &post->quotes); + nn_packet_read_optional_int(packet, &post->comments); + nn_packet_read_optional_int(packet, &post->views); u32 media_size; - AKI_PACKET_READ_TYPE(packet, u32, media_size); + NNWT_PACKET_READ_TYPE(packet, u32, media_size); al_array_reserve(post->media, media_size); for (u32 i = 0; i < media_size; i++) { struct camu_post_media media; - AKI_PACKET_READ_TYPE(packet, u8, media.type); - aki_packet_read_str(packet, &media.key); - aki_packet_read_str(packet, &media.url); - aki_packet_read_str(packet, &media.ext); - aki_packet_read_str(packet, &media.thumbnail_url); - aki_packet_read_str(packet, &media.thumbnail_ext); + NNWT_PACKET_READ_TYPE(packet, u8, media.type); + nn_packet_read_str(packet, &media.key); + nn_packet_read_str(packet, &media.url); + nn_packet_read_str(packet, &media.ext); + nn_packet_read_str(packet, &media.thumbnail_url); + nn_packet_read_str(packet, &media.thumbnail_ext); al_array_push(post->media, media); } - aki_packet_read_str(packet, &post->post.unique_id); - aki_packet_read_str(packet, &post->quoted.unique_id); - aki_packet_read_str(packet, &post->in_reply_to.unique_id); + nn_packet_read_str(packet, &post->post.unique_id); + nn_packet_read_str(packet, &post->quoted.unique_id); + nn_packet_read_str(packet, &post->in_reply_to.unique_id); } diff --git a/src/portal/src/packet_ext.h b/src/portal/src/packet_ext.h index eae1e07..6da61fc 100644 --- a/src/portal/src/packet_ext.h +++ b/src/portal/src/packet_ext.h @@ -1,10 +1,9 @@ #pragma once #include <al/types.h> -#include <aki/common.h> -#include <aki/event_loop.h> +#include <nnwt/packet.h> #include "post.h" -void aki_packet_write_post(struct aki_packet *packet, struct camu_post *post); -void aki_packet_read_post(struct aki_packet *packet, struct camu_post *post); +void nn_packet_write_post(struct nn_packet *packet, struct camu_post *post); +void nn_packet_read_post(struct nn_packet *packet, struct camu_post *post); diff --git a/src/portal/src/search.c b/src/portal/src/search.c index e8f30c0..88cf1f1 100644 --- a/src/portal/src/search.c +++ b/src/portal/src/search.c @@ -8,26 +8,35 @@ bool camu_python_init(void) { - PyPreConfig config; - PyPreConfig_InitPythonConfig(&config); - config.dev_mode = 0; - config.utf8_mode = 1; - PyStatus status = Py_PreInitialize(&config); + PyPreConfig pre_config; + PyPreConfig_InitPythonConfig(&pre_config); + + pre_config.dev_mode = 0; + pre_config.utf8_mode = 1; + pre_config.isolated = 0; + pre_config.parse_argv = 0; + pre_config.use_environment = 1; + + PyStatus status = Py_PreInitialize(&pre_config); if (PyStatus_Exception(status)) { al_log_error("portal", "Preinitialization failed."); return false; } + if (PyImport_AppendInittab("portal", PyInit_portal) == -1) { al_log_error("portal", "Could not extend in-built modules table."); return false; } + Py_Initialize(); + PyObject *module = PyImport_ImportModule("portal"); if (!module) { PyErr_Print(); al_log_error("portal", "Could not import module."); goto err; } + return true; err: Py_Finalize(); @@ -48,14 +57,14 @@ static struct camu_search *get_search_by_id(struct camu_portal_bridge *bridge, s return NULL; } -static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL queue_thread(void *userdata) { struct camu_portal_bridge *bridge = (struct camu_portal_bridge *)userdata; - aki_thread_setcanceltype(AKI_THREAD_CANCEL_ASYNCHRONOUS); + nn_thread_setcanceltype(NNWT_THREAD_CANCEL_ASYNCHRONOUS); bool have_python = false; - aki_mutex_lock(&bridge->mutex); + nn_mutex_lock(&bridge->mutex); do { - aki_cond_wait(&bridge->cond, &bridge->mutex); + nn_cond_wait(&bridge->cond, &bridge->mutex); if (bridge->quit) break; if (!have_python) { // Defer python init. @@ -117,9 +126,9 @@ static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) camu_queue_push(bridge->results, result); } bridge->queue.size = 0; - aki_signal_send(&bridge->results_signal); + nn_signal_send(&bridge->results_signal); } while (1); - aki_mutex_unlock(&bridge->mutex); + nn_mutex_unlock(&bridge->mutex); if (have_python) camu_python_close(); return 0; } @@ -137,19 +146,19 @@ static void results_signal_callback(void *userdata) } void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache, - struct aki_event_loop *loop, void *userdata) + struct nn_event_loop *loop, void *userdata) { al_array_init(bridge->searches); bridge->cache = cache; bridge->quit = 0; - aki_mutex_init(&bridge->mutex); - aki_cond_init(&bridge->cond); + nn_mutex_init(&bridge->mutex); + nn_cond_init(&bridge->cond); al_array_init(bridge->queue); camu_queue_init(bridge->results); - aki_signal_init(&bridge->results_signal, loop, results_signal_callback, bridge); - aki_signal_start(&bridge->results_signal); + nn_signal_init(&bridge->results_signal, loop, results_signal_callback, bridge); + nn_signal_start(&bridge->results_signal); bridge->userdata = userdata; - aki_thread_create(&bridge->thread, queue_thread, bridge); + nn_thread_create(&bridge->thread, queue_thread, bridge); } void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query, @@ -161,12 +170,12 @@ void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, s al_str_clone(&cmd.query, query); cmd.callback = callback; cmd.userdata = userdata; - aki_mutex_lock(&bridge->mutex); + nn_mutex_lock(&bridge->mutex); al_array_push(bridge->queue, cmd); - if (aki_cond_is_waiting(&bridge->cond)) { - aki_cond_signal(&bridge->cond); + if (nn_cond_is_waiting(&bridge->cond)) { + nn_cond_signal(&bridge->cond); } - aki_mutex_unlock(&bridge->mutex); + nn_mutex_unlock(&bridge->mutex); } void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num, @@ -178,28 +187,28 @@ void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num, cmd.num = num; cmd.callback = callback; cmd.userdata = userdata; - aki_mutex_lock(&bridge->mutex); + nn_mutex_lock(&bridge->mutex); al_array_push(bridge->queue, cmd); - if (aki_cond_is_waiting(&bridge->cond)) { - aki_cond_signal(&bridge->cond); + if (nn_cond_is_waiting(&bridge->cond)) { + nn_cond_signal(&bridge->cond); } - aki_mutex_unlock(&bridge->mutex); + nn_mutex_unlock(&bridge->mutex); } void camu_portal_close(struct camu_portal_bridge *bridge) { - aki_mutex_lock(&bridge->mutex); + nn_mutex_lock(&bridge->mutex); bridge->quit = 1; - if (aki_cond_is_waiting(&bridge->cond)) { - aki_cond_signal(&bridge->cond); + if (nn_cond_is_waiting(&bridge->cond)) { + nn_cond_signal(&bridge->cond); } - aki_mutex_unlock(&bridge->mutex); - aki_thread_join(&bridge->thread); - aki_signal_stop(&bridge->results_signal); + 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); - aki_cond_destroy(&bridge->cond); - aki_mutex_destroy(&bridge->mutex); + nn_cond_destroy(&bridge->cond); + nn_mutex_destroy(&bridge->mutex); } /* diff --git a/src/portal/src/search.h b/src/portal/src/search.h index 199f491..f870522 100644 --- a/src/portal/src/search.h +++ b/src/portal/src/search.h @@ -1,7 +1,7 @@ #pragma once -#include <aki/thread.h> -#include <aki/signal.h> +#include <nnwt/thread.h> +#include <nnwt/signal.h> #include "post.h" #include "post_cache.h" @@ -45,12 +45,12 @@ struct camu_portal_bridge { array(struct camu_search *) searches; struct camu_post_cache *cache; u8 quit; - struct aki_thread thread; - struct aki_mutex mutex; - struct aki_cond cond; + struct nn_thread thread; + struct nn_mutex mutex; + struct nn_cond cond; array(struct camu_portal_cmd) queue; queue(struct camu_portal_result) results; - struct aki_signal results_signal; + struct nn_signal results_signal; void *userdata; }; @@ -58,7 +58,7 @@ bool camu_python_init(void); void camu_python_close(void); void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache, - struct aki_event_loop *loop, void *userdata); + struct nn_event_loop *loop, void *userdata); void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query, void (*callback)(void *, void *, struct camu_portal_result *), void *userdata); diff --git a/src/render/queue_libplacebo.c b/src/render/queue_libplacebo.c index 8ecf9ce..07ce379 100644 --- a/src/render/queue_libplacebo.c +++ b/src/render/queue_libplacebo.c @@ -234,13 +234,7 @@ static void unmap_av_frame(pl_gpu gpu, struct pl_frame *frame, const struct pl_s } al_free((struct pl_overlay_part *)current->parts); } - struct camu_overlay_lp *roverlay; - al_array_foreach(lq->overlays, i, roverlay) { - if (roverlay == overlay) { - al_array_remove_at(lq->overlays, i); - break; - } - } + al_array_remove(lq->overlays, overlay); al_free(overlay); } } diff --git a/src/render/renderer_libplacebo.c b/src/render/renderer_libplacebo.c index 49addb0..9fedc5f 100644 --- a/src/render/renderer_libplacebo.c +++ b/src/render/renderer_libplacebo.c @@ -1,6 +1,6 @@ #include <al/lib.h> #include <al/log.h> -#include <aki/file.h> +#include <nnwt/file.h> #ifdef CAMU_HAVE_FFMPEG #include <libplacebo/utils/libav.h> #endif @@ -129,9 +129,9 @@ static bool renderer_lp_create_renderer(struct camu_renderer *renderer, s32 *wid pl_swapchain_resize(lr->swapchain, width, height); al_memset(&lr->params, 0, sizeof(struct pl_render_params)); - lr->params = pl_render_fast_params; + //lr->params = pl_render_fast_params; //lr->params.upscaler = &pl_filter_nearest; - //lr->params = pl_render_default_params; + lr->params = pl_render_default_params; //lr->params = pl_render_high_quality_params; lr->params.deband_params = NULL; @@ -139,14 +139,14 @@ static bool renderer_lp_create_renderer(struct camu_renderer *renderer, s32 *wid lr->params.skip_caching_single_frame = true; // Clear manually so we can draw multiple images per frame. lr->params.border = PL_CLEAR_SKIP; - // Prioritize a more consistant image. + // Prioritize a more consistent image. lr->params.correct_subpixel_offsets = true; #if 0 - struct aki_file file; - if (aki_file_open(&file, al_str_c(""), 0)) { + struct nn_file file; + if (nn_file_open(&file, al_str_c(""), 0)) { char *str; - size_t len = aki_file_read_as_c_str(&file, &str); + size_t len = nn_file_read_as_c_str(&file, &str); const struct pl_hook *hook = pl_mpv_user_shader_parse(lr->gpu, (const char *)str, len); const struct pl_hook **hooks = al_malloc(sizeof(struct pl_hook *)); hooks[0] = al_malloc(sizeof(struct pl_hook)); diff --git a/src/screen/screen.c b/src/screen/screen.c index 5b5d175..17f4e3b 100644 --- a/src/screen/screen.c +++ b/src/screen/screen.c @@ -1,5 +1,5 @@ #include <al/log.h> -#include <aki/thread.h> +#include <nnwt/thread.h> #include "view.h" #include "screen.h" @@ -79,11 +79,11 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button) switch (state) { case STELA_BUTTON_PRESSED: scr->flags |= CAMU_SCREEN_DRAGGING; - scr->last_click_ts = aki_get_timestamp(); + scr->last_click_ts = nn_get_timestamp(); break; case STELA_BUTTON_RELEASED: scr->flags &= ~CAMU_SCREEN_DRAGGING; - if (scr->last_click_ts && aki_get_timestamp() - scr->last_click_ts <= 300000) { + if (scr->last_click_ts && nn_get_timestamp() - scr->last_click_ts <= 300000) { if (scr->flags & CAMU_SCREEN_MODIFIER) { f64 percent = scr->last_mouse_x / scr->width; percent = CLAMP(percent, 0.0, 100.0); @@ -285,7 +285,7 @@ bool camu_screen_init(struct camu_screen *scr, void *context) al_array_init(scr->add_queue); al_array_init(scr->rem_queue); al_atomic_store(u8)(&scr->queued, 0, AL_ATOMIC_RELAXED); - aki_mutex_init(&scr->mutex); + nn_mutex_init(&scr->mutex); #endif scr->last_mouse_x = 0.0; scr->last_mouse_y = 0.0; @@ -294,7 +294,7 @@ bool camu_screen_init(struct camu_screen *scr, void *context) } #ifdef STELA_EVENT_BUFFER -static aki_thread_result AKI_THREADCALL event_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL event_thread(void *userdata) { struct camu_screen *scr = (struct camu_screen *)userdata; while (scr->window->poll(scr->window, true) && al_atomic_load(s32)(&scr->state, AL_ATOMIC_RELAXED) != CAMU_SCREEN_STOPPED) { @@ -313,7 +313,7 @@ bool camu_screen_create_window(struct camu_screen *scr, const char *name) scr->width = scr->window->width; scr->height = scr->window->height; #ifdef STELA_EVENT_BUFFER - aki_thread_create(&scr->thread, event_thread, scr); + nn_thread_create(&scr->thread, event_thread, scr); #endif return true; } @@ -407,11 +407,11 @@ static void run_queue_internal(struct camu_screen *scr) void camu_screen_add_buffer(struct camu_screen *scr, struct camu_video_buffer *buf) { #ifdef CAMU_SCREEN_THREADED - aki_mutex_lock(&scr->mutex); + nn_mutex_lock(&scr->mutex); struct camu_video_buffer *rbuf; al_array_foreach(scr->add_queue, i, rbuf) { if (rbuf == buf) { - aki_mutex_unlock(&scr->mutex); + nn_mutex_unlock(&scr->mutex); return; } } @@ -422,13 +422,13 @@ void camu_screen_add_buffer(struct camu_screen *scr, struct camu_video_buffer *b if (queue_empty) { al_atomic_store(u8)(&scr->queued, 0, AL_ATOMIC_RELAXED); } - aki_mutex_unlock(&scr->mutex); + nn_mutex_unlock(&scr->mutex); return; } } al_array_push(scr->add_queue, buf); al_atomic_store(u8)(&scr->queued, 1, AL_ATOMIC_RELAXED); - aki_mutex_unlock(&scr->mutex); + nn_mutex_unlock(&scr->mutex); #else add_buffer_internal(scr, buf); #endif @@ -438,11 +438,11 @@ void camu_screen_add_buffer(struct camu_screen *scr, struct camu_video_buffer *b void camu_screen_remove_buffer(struct camu_screen *scr, struct camu_video_buffer *buf) { #ifdef CAMU_SCREEN_THREADED - aki_mutex_lock(&scr->mutex); + nn_mutex_lock(&scr->mutex); struct camu_video_buffer *rbuf; al_array_foreach(scr->rem_queue, i, rbuf) { if (rbuf == buf) { - aki_mutex_unlock(&scr->mutex); + nn_mutex_unlock(&scr->mutex); return; } } @@ -453,13 +453,13 @@ void camu_screen_remove_buffer(struct camu_screen *scr, struct camu_video_buffer if (queue_empty) { al_atomic_store(u8)(&scr->queued, 0, AL_ATOMIC_RELAXED); } - aki_mutex_unlock(&scr->mutex); + nn_mutex_unlock(&scr->mutex); return; } } al_array_push(scr->rem_queue, buf); al_atomic_store(u8)(&scr->queued, 1, AL_ATOMIC_RELAXED); - aki_mutex_unlock(&scr->mutex); + nn_mutex_unlock(&scr->mutex); #else remove_buffer_internal(scr, buf); #endif @@ -469,21 +469,21 @@ void camu_screen_remove_buffer(struct camu_screen *scr, struct camu_video_buffer #ifdef CAMU_SCREEN_THREADED void camu_screen_run_queue(struct camu_screen *scr) { - aki_mutex_lock(&scr->mutex); + nn_mutex_lock(&scr->mutex); run_queue_internal(scr); al_atomic_store(u8)(&scr->queued, 0, AL_ATOMIC_RELAXED); - aki_mutex_unlock(&scr->mutex); + nn_mutex_unlock(&scr->mutex); } void camu_screen_clear(struct camu_screen *scr) { - aki_mutex_lock(&scr->mutex); + nn_mutex_lock(&scr->mutex); struct camu_screen_video *video; al_array_foreach_ptr(scr->videos, i, video) { struct camu_video_buffer *buf = video->buf; al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED); } - aki_mutex_unlock(&scr->mutex); + nn_mutex_unlock(&scr->mutex); } #endif @@ -498,7 +498,7 @@ bool camu_screen_poll(struct camu_screen *scr, bool block) { #ifdef STELA_EVENT_BUFFER if (block) { - aki_thread_sleep(AKI_TS_FROM_USEC(1000000/CAMU_SCREEN_POLL_HZ)); + nn_thread_sleep(NNWT_TS_FROM_USEC(1000000/CAMU_SCREEN_POLL_HZ)); } stl_window_read_events(scr->window); #else @@ -534,12 +534,12 @@ void camu_screen_close(struct camu_screen *scr) #ifdef CAMU_SCREEN_THREADED al_array_free(scr->add_queue); al_array_free(scr->rem_queue); - aki_mutex_destroy(&scr->mutex); + nn_mutex_destroy(&scr->mutex); #endif al_array_free(scr->videos); al_atomic_store(s32)(&scr->state, CAMU_SCREEN_STOPPED, AL_ATOMIC_RELAXED); #ifdef STELA_EVENT_BUFFER - aki_thread_join(&scr->thread); + nn_thread_join(&scr->thread); #endif scr->window->free(&scr->window); } diff --git a/src/screen/screen.h b/src/screen/screen.h index 349d41f..80f888c 100644 --- a/src/screen/screen.h +++ b/src/screen/screen.h @@ -3,7 +3,7 @@ #include <al/array.h> #include <al/atomic.h> #ifdef CAMU_SCREEN_THREADED -#include <aki/thread.h> +#include <nnwt/thread.h> #endif #include <stl/window.h> @@ -53,7 +53,7 @@ struct camu_screen { u32 flags; struct stl_window *window; #ifdef STELA_EVENT_BUFFER - struct aki_thread thread; + struct nn_thread thread; #endif struct camu_renderer *renderer; s32 width; @@ -66,7 +66,7 @@ struct camu_screen { array(struct camu_video_buffer *) add_queue; array(struct camu_video_buffer *) rem_queue; atomic(u8) queued; - struct aki_mutex mutex; + struct nn_mutex mutex; #endif void (*callback)(void *, u8, void *); void *userdata; diff --git a/src/server/common.c b/src/server/common.c deleted file mode 100644 index 765418c..0000000 --- a/src/server/common.c +++ /dev/null @@ -1,10 +0,0 @@ -#include "common.h" - -//const str *CAMU_DB_PATH = al_str_c("/mnt/store/files/camu_db"); -str *CAMU_DB_PATH = al_str_c("/home/andrew/c/camu/data/camu_db_test"); - -str *CAMU_SERVER_IP = al_str_c("108.52.160.112"); -//str *CAMU_SERVER_IP = al_str_c("127.0.0.1"); - -str *CAMU_UNIX_PATH = al_str_c("/tmp/camu_sock"); -str *CAMU_UNIX_LOCAL = al_str_c("/tmp/cmv_sock"); diff --git a/src/server/common.h b/src/server/common.h index b6ee7a1..b4c6344 100644 --- a/src/server/common.h +++ b/src/server/common.h @@ -2,19 +2,10 @@ #include <al/str.h> -extern str *CAMU_DB_PATH; - #define CAMU_PORT 14356 #define CAMU_MULTIPLEX_RPC 0x53 #define CAMU_MULTIPLEX_LIANA 0x85 -extern str *CAMU_SERVER_IP; -extern str *CAMU_UNIX_PATH; -extern str *CAMU_UNIX_LOCAL; - -#define CAMU_LOCAL_TYPE AKI_SOCKET_TCP -#define CAMU_LOCAL_ADDR CAMU_SERVER_IP - enum { CAMU_NODE = 0, CAMU_CLIENT, @@ -47,9 +38,16 @@ enum { enum { CAMU_RESOURCE_FILE = 0, +#ifdef NAUNET_HAS_CURL + CAMU_RESOURCE_HTTP, +#endif +#ifdef CACHE_HAVE_CDIO CAMU_RESOURCE_CDIO, +#endif +#ifdef CAMU_HAVE_PORTAL CAMU_RESOURCE_PORTAL, CAMU_RESOURCE_SIMPLE_SEARCH +#endif }; AL_UNUSED_FUNCTION_PUSH @@ -57,8 +55,7 @@ AL_UNUSED_FUNCTION_PUSH static bool camu_is_url(str *s, u32 i) { return al_str_cmp(s, al_str_c("https://"), i, 8) == 0 - || al_str_cmp(s, al_str_c("http://"), i, 7) == 0 - || al_str_cmp(s, al_str_c("cdda://"), i, 7) == 0; + || al_str_cmp(s, al_str_c("http://"), i, 7) == 0; } AL_UNUSED_FUNCTION_POP diff --git a/src/server/db.c b/src/server/db.c index 6f76864..ce4c99e 100644 --- a/src/server/db.c +++ b/src/server/db.c @@ -1,17 +1,17 @@ #include <al/log.h> -#include <aki/file.h> +#include <nnwt/file.h> #include <jansson.h> #include "server.h" -static bool open_user(struct camu_server *server, struct aki_dir_entry *dir) +static bool open_user(struct camu_server *server, struct nn_dir_entry *dir) { - struct aki_file file; - if (!aki_file_open(&file, &dir->path, 0)) { + struct nn_file file; + if (!nn_file_open(&file, &dir->path, 0)) { return false; } str s; - aki_file_read_as_str(&file, &s); + nn_file_read_as_str(&file, &s); json_error_t error; json_t *root = json_loadb(s.data, s.len, 0, &error); if (!root) { @@ -30,28 +30,28 @@ static bool open_user(struct camu_server *server, struct aki_dir_entry *dir) bool camu_db_open(struct camu_server *server, str *path) { - struct aki_dir camu_db; - if (!aki_dir_open(&camu_db, path)) { + struct nn_dir camu_db; + if (!nn_dir_open(&camu_db, path)) { return false; } - struct aki_dir_entry entry; - while (aki_dir_read(&camu_db, &entry)) { + struct nn_dir_entry entry; + while (nn_dir_read(&camu_db, &entry)) { if (al_str_eq(&entry.name, al_str_c("users"))) { - struct aki_dir users; - if (aki_dir_open(&users, &entry.path)) { - struct aki_dir_entry user; - while (aki_dir_read(&users, &user)) { - if (user.type == AKI_ENTRY_FILE) { + struct nn_dir users; + if (nn_dir_open(&users, &entry.path)) { + struct nn_dir_entry user; + while (nn_dir_read(&users, &user)) { + if (user.type == NNWT_ENTRY_FILE) { open_user(server, &user); } - aki_dir_entry_free(&user); + nn_dir_entry_free(&user); } - aki_dir_close(&users); + nn_dir_close(&users); } } - aki_dir_entry_free(&entry); + nn_dir_entry_free(&entry); } - aki_dir_close(&camu_db); + nn_dir_close(&camu_db); return true; } diff --git a/src/server/meson.build b/src/server/meson.build index ab100e1..52c993f 100644 --- a/src/server/meson.build +++ b/src/server/meson.build @@ -1,8 +1,7 @@ server_src = [ 'server.c', - 'common.c', 'user.c', - 'db.c', + 'db.c' ] server_deps = [cache, liana_server] if get_option('portal').enabled() diff --git a/src/server/resource.h b/src/server/resource.h index 559e2b4..1cd00fc 100644 --- a/src/server/resource.h +++ b/src/server/resource.h @@ -17,6 +17,11 @@ struct camu_resource_file { str path; }; +struct camu_resource_http { + struct camu_resource r; + str url; +}; + struct camu_resource_portal { struct camu_resource r; struct camu_post *post; diff --git a/src/server/server.c b/src/server/server.c index 4ec61da..7d0b2fe 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -3,7 +3,9 @@ #include "../cache/handlers/file.h" #include "../cache/handlers/http.h" +#ifdef CACHE_HAVE_CDIO #include "../cache/handlers/cdio.h" +#endif #include "../libclient/common.h" #include "../libsink/common.h" #ifdef CAMU_HAVE_PORTAL @@ -25,7 +27,7 @@ static struct camu_user *get_user_by_username(struct camu_server *server, str *u return NULL; } -static struct camu_server_client *get_client_by_connection(struct camu_server *server, struct aki_rpc_connection *conn) +static struct camu_server_client *get_client_by_connection(struct camu_server *server, struct nn_rpc_connection *conn) { struct camu_server_client *client; al_array_foreach(server->clients, i, client) { @@ -52,26 +54,30 @@ static struct camu_server_sink *get_sink_from_name(struct camu_server *server, s return NULL; } -static void write_user_state(struct camu_server *server, struct camu_user *user, struct aki_packet *packet) +static void write_user_state(struct camu_server *server, struct camu_user *user, struct nn_packet *packet) { (void)user; +#ifdef CAMU_HAVE_PORTAL + nn_packet_write_u32(packet, server->bridge.searches.size); struct camu_search *search; - aki_packet_write_u32(packet, server->bridge.searches.size); al_array_foreach(server->bridge.searches, i, search) { - aki_packet_write_s32(packet, search->id); - aki_packet_write_str(packet, &search->module); - aki_packet_write_str(packet, &search->query); + nn_packet_write_s32(packet, search->id); + nn_packet_write_str(packet, &search->module); + nn_packet_write_str(packet, &search->query); } +#else + nn_packet_write_u32(packet, 0); +#endif } static void handle_toggle_sink(struct camu_server *server, str *name, struct camu_server_sink *sink, bool enable); -static bool identify_callback(void *userdata, struct aki_rpc_connection *conn, - struct aki_packet *packet, struct aki_packet *rpacket) +static bool identify_callback(void *userdata, struct nn_rpc_connection *conn, + struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_server *server = (struct camu_server *)userdata; - u8 op = aki_packet_read_u8(packet); + u8 op = nn_packet_read_u8(packet); switch (op) { case CAMU_NODE: { struct camu_server_node *node = al_alloc_object(struct camu_server_node); @@ -84,7 +90,7 @@ static bool identify_callback(void *userdata, struct aki_rpc_connection *conn, struct camu_server_client *client = al_alloc_object(struct camu_server_client); client->conn = conn; str username; - aki_packet_read_str(packet, &username); + nn_packet_read_str(packet, &username); struct camu_user *user = get_user_by_username(server, &username); if (!user) { user = al_alloc_object(struct camu_user); @@ -101,7 +107,7 @@ static bool identify_callback(void *userdata, struct aki_rpc_connection *conn, struct camu_server_sink *sink = al_alloc_object(struct camu_server_sink); sink->conn = conn; str name; - aki_packet_read_str(packet, &name); + nn_packet_read_str(packet, &name); al_str_clone(&sink->name, &name); sink->server = server; al_array_push(server->sinks, sink); @@ -111,21 +117,10 @@ static bool identify_callback(void *userdata, struct aki_rpc_connection *conn, } } - aki_packet_free(packet); + nn_packet_free(packet); return true; } -static bool client_still_connected(struct camu_server *server, struct aki_rpc_connection *conn) -{ - struct camu_server_client *client; - al_array_foreach(server->clients, i, client) { - if (client->conn == conn) { - return true; - } - } - return false; -} - static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing) { struct camu_server_sink *sink = (struct camu_server_sink *)userdata; @@ -134,44 +129,44 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent case LIANA_SINK_BUFFER: case LIANA_SINK_BUFFER_AND_QUEUE: { struct camu_resource *resource = (struct camu_resource *)entry->opaque; - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); - aki_packet_write_u8(packet, op); - aki_packet_write_str(packet, &sink->server->addr); - aki_packet_write_u16(packet, CAMU_PORT); - aki_packet_write_u32(packet, resource->node->id); - aki_packet_write_u32(packet, entry->id); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, timing->at); + struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); + nn_packet_write_u8(packet, op); + nn_packet_write_str(packet, &sink->server->addr); + nn_packet_write_u16(packet, CAMU_PORT); + nn_packet_write_u32(packet, resource->node->id); + nn_packet_write_u32(packet, entry->id); + nn_packet_write_s32(packet, sequence); + nn_packet_write_u64(packet, timing->at); al_assert(timing->seek_pos <= INT64_MAX); - aki_packet_write_u64(packet, timing->seek_pos); - aki_packet_write_u8(packet, timing->pause); - aki_packet_write_bool(packet, timing->previous_ended); - aki_packet_write_bool(packet, timing->ended); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + nn_packet_write_u64(packet, timing->seek_pos); + nn_packet_write_u8(packet, timing->pause); + nn_packet_write_bool(packet, timing->previous_ended); + nn_packet_write_bool(packet, timing->ended); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case LIANA_SINK_UNSET: { - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); - aki_packet_write_u8(packet, op); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); + nn_packet_write_u8(packet, op); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case LIANA_SINK_PAUSE: { - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE); - aki_packet_write_u32(packet, entry->id); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, timing->at); - aki_packet_write_u8(packet, timing->pause); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE); + nn_packet_write_u32(packet, entry->id); + nn_packet_write_s32(packet, sequence); + nn_packet_write_u64(packet, timing->at); + nn_packet_write_u8(packet, timing->pause); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case LIANA_SINK_SEEK: { - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK); - aki_packet_write_u32(packet, entry->id); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, timing->at); - aki_packet_write_u64(packet, timing->seek_pos); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK); + nn_packet_write_u32(packet, entry->id); + nn_packet_write_s32(packet, sequence); + nn_packet_write_u64(packet, timing->at); + nn_packet_write_u64(packet, timing->seek_pos); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } } @@ -188,7 +183,7 @@ void handle_toggle_sink(struct camu_server *server, str *name, struct camu_serve static void process_pending(struct camu_resource *resource) { array(struct lia_list_entry *) pending; - // resource->pending may be edited in a list_pump() call. + // resource->pending may be edited during a list_pump() call. al_array_clone(pending, resource->pending); resource->pending.size = 0; struct lia_list_entry *entry; @@ -244,42 +239,54 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v } break; } +} +#ifdef CAMU_HAVE_PORTAL +static bool client_still_connected(struct camu_server *server, struct nn_rpc_connection *conn) +{ + struct camu_server_client *client; + al_array_foreach(server->clients, i, client) { + if (client->conn == conn) { + return true; + } + } + return false; } static void client_portal_callback(void *userdata0, void *userdata1, struct camu_portal_result *result) { struct camu_server *server = (struct camu_server *)userdata0; - struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata1; + struct nn_rpc_connection *conn = (struct nn_rpc_connection *)userdata1; if (!client_still_connected(server, conn)) return; - struct aki_packet *packet = aki_rpc_get_packet(conn->rpc, CAMU_CLIENT_RESULTS); - aki_packet_write_u8(packet, result->op); - aki_packet_write_s32(packet, result->id); + struct nn_packet *packet = nn_rpc_get_packet(conn->rpc, CAMU_CLIENT_RESULTS); + nn_packet_write_u8(packet, result->op); + nn_packet_write_s32(packet, result->id); switch (result->op) { case CAMU_CLIENT_CREATE_SEARCH: { break; } case CAMU_CLIENT_GET_PAGE: { struct camu_result_page *page = result->page; - aki_packet_write_u32(packet, page->num); - aki_packet_write_u32(packet, page->posts.size); + nn_packet_write_u32(packet, page->num); + nn_packet_write_u32(packet, page->posts.size); struct camu_post *post; al_array_foreach_ptr(page->posts, i, post) { - aki_packet_write_post(packet, post); + nn_packet_write_post(packet, post); } - aki_packet_write_u32(packet, page->list.size); + nn_packet_write_u32(packet, page->list.size); str *unique_id; al_array_foreach_ptr(page->list, i, unique_id) { - aki_packet_write_str(packet, unique_id); + nn_packet_write_str(packet, unique_id); } break; } } - aki_rpc_connection_command(conn, packet, NULL, NULL); + nn_rpc_connection_command(conn, packet, NULL, NULL); } +#endif -static bool client_command_callback(void *userdata, struct aki_rpc_connection *conn, - struct aki_packet *packet, struct aki_packet *rpacket) +static bool client_command_callback(void *userdata, struct nn_rpc_connection *conn, + struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_server *server = (struct camu_server *)userdata; (void)rpacket; @@ -287,11 +294,11 @@ static bool client_command_callback(void *userdata, struct aki_rpc_connection *c struct camu_server_client *client = get_client_by_connection(server, conn); if (!client) goto out; - u8 op = aki_packet_read_u8(packet); + u8 op = nn_packet_read_u8(packet); switch (op) { case CAMU_CLIENT_CREATE_LIST: { str name; - aki_packet_read_str(packet, &name); + nn_packet_read_str(packet, &name); struct lia_list *list = al_alloc_object(struct lia_list); lia_list_init(list, &name); list->callback = list_callback; @@ -301,35 +308,38 @@ static bool client_command_callback(void *userdata, struct aki_rpc_connection *c } case CAMU_CLIENT_TOGGLE_SINK: { str name; - aki_packet_read_str(packet, &name); + nn_packet_read_str(packet, &name); struct camu_server_sink *sink = get_sink_from_name(server, &name); if (!sink) goto out; - aki_packet_read_str(packet, &name); // list name. - bool enable = aki_packet_read_bool(packet); + nn_packet_read_str(packet, &name); // list name. + bool enable = nn_packet_read_bool(packet); handle_toggle_sink(server, &name, sink, enable); break; } +#ifdef CAMU_HAVE_PORTAL case CAMU_CLIENT_CREATE_SEARCH: { str module; str query; - aki_packet_read_str(packet, &module); - aki_packet_read_str(packet, &query); + nn_packet_read_str(packet, &module); + nn_packet_read_str(packet, &query); camu_portal_create_search(&server->bridge, &module, &query, client_portal_callback, conn); break; } case CAMU_CLIENT_GET_PAGE: { - s32 id = aki_packet_read_s32(packet); - u32 num = aki_packet_read_u32(packet); + s32 id = nn_packet_read_s32(packet); + u32 num = nn_packet_read_u32(packet); camu_portal_get_page(&server->bridge, id, num, client_portal_callback, conn); break; } +#endif } out: - aki_packet_free(packet); + nn_packet_free(packet); return true; } +#ifdef CAMU_HAVE_PORTAL static struct cch_entry *entry_from_post(struct camu_server *server, struct camu_post *post, u32 index) { struct cch_entry *entry = NULL; @@ -378,17 +388,18 @@ static void simple_search_portal_callback(void *userdata0, void *userdata1, stru al_log_warn("server", "Server resource failed to load."); resource->load = LIANA_ENTRY_ERRORED; } +#endif -static void handle_add_command(struct camu_server *server, struct lia_list *list, struct aki_packet *packet) +static void handle_add_command(struct camu_server *server, struct lia_list *list, struct nn_packet *packet) { - u8 op = aki_packet_read_u8(packet); + u8 op = nn_packet_read_u8(packet); struct cch_entry *entry = NULL; struct camu_resource *resource = NULL; wstr name; switch (op) { case CAMU_RESOURCE_FILE: { str path; - aki_packet_read_str(packet, &path); + nn_packet_read_str(packet, &path); entry = cch_handler_file_create(&path); if (!entry) return; struct camu_resource_file *file = al_alloc_object(struct camu_resource_file); @@ -399,12 +410,36 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list resource->load = LIANA_ENTRY_PREPARED; break; } +#ifdef NAUNET_HAS_CURL + case CAMU_RESOURCE_HTTP: { + str url; + nn_packet_read_str(packet, &url); + entry = cch_handler_http_create(&url, server->loop); + if (!entry) return; + entry->handler->maybe_spawn_worker(entry->handler, 0); + struct camu_resource_http *http = al_alloc_object(struct camu_resource_http); + al_str_clone(&http->url, &url); + al_wstr_from_str(&name, &url); + resource = (struct camu_resource *)http; + resource->type = CAMU_RESOURCE_HTTP; + resource->load = LIANA_ENTRY_PREPARED; + break; + } +#endif +#ifdef CACHE_HAVE_CDIO case CAMU_RESOURCE_CDIO: { - u32 track = aki_packet_read_u32(packet); entry = cch_handler_cdio_create(); if (!entry) return; - struct cch_chapter *chapter = &al_array_at(entry->chapters, track); - entry->handler->maybe_spawn_worker(entry->handler, chapter->start); + str url; + nn_packet_read_str(packet, &url); + u32 track = 0; + if (url.len > 7) { + s64 index = al_str_to_long(al_str_substr(&url, 7, url.len), 10); + if (index > 0 && index <= entry->chapters.size) + track = (u32)index - 1; + } + entry->chapter = &al_array_at(entry->chapters, track); + entry->handler->maybe_spawn_worker(entry->handler, entry->chapter->start); struct camu_resource_cdio *cdio = al_alloc_object(struct camu_resource_cdio); cdio->track = track; al_wstr_from_cstr(&name, "cdio"); @@ -413,10 +448,12 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list resource->load = LIANA_ENTRY_PREPARED; break; } +#endif +#ifdef CAMU_HAVE_PORTAL case CAMU_RESOURCE_PORTAL: { str unique_id; - aki_packet_read_str(packet, &unique_id); - u32 index = aki_packet_read_u32(packet); + nn_packet_read_str(packet, &unique_id); + u32 index = nn_packet_read_u32(packet); struct camu_post *post = camu_post_cache_get(&server->cache, &unique_id); entry = entry_from_post(server, post, index); if (!entry) return; @@ -430,7 +467,7 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list } case CAMU_RESOURCE_SIMPLE_SEARCH: { str search; - aki_packet_read_str(packet, &search); + nn_packet_read_str(packet, &search); struct camu_resource_portal *portal = al_alloc_object(struct camu_resource_portal); portal->post = NULL; al_wstr_from_str(&name, &search); @@ -448,6 +485,7 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list al_str_free(&query); break; } +#endif } al_assert(resource); if (entry) { @@ -462,34 +500,34 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list al_wstr_free(&name); } -static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn, - struct aki_packet *packet, struct aki_packet *rpacket) +static bool list_action_callback(void *userdata, struct nn_rpc_connection *conn, + struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_server *server = (struct camu_server *)userdata; (void)conn; (void)rpacket; str name; - aki_packet_read_str(packet, &name); + nn_packet_read_str(packet, &name); struct lia_list *list = get_list_from_name(server, &name); if (!list) goto out; - u8 op = aki_packet_read_u8(packet); + u8 op = nn_packet_read_u8(packet); switch (op) { case CAMU_LIST_ADD: { handle_add_command(server, list, packet); break; } case CAMU_LIST_SKIP: { - s32 sequence = aki_packet_read_s32(packet); - s32 n = aki_packet_read_s32(packet); + s32 sequence = nn_packet_read_s32(packet); + s32 n = nn_packet_read_s32(packet); lia_list_skip(list, sequence, n); break; } case CAMU_LIST_SKIPTO: { - s32 sequence = aki_packet_read_s32(packet); - s32 i = aki_packet_read_s32(packet); + s32 sequence = nn_packet_read_s32(packet); + s32 i = nn_packet_read_s32(packet); lia_list_skipto(list, sequence, i); break; } @@ -498,15 +536,15 @@ static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn break; } case CAMU_LIST_TOGGLE_PAUSE: { - s32 sequence = aki_packet_read_s32(packet); - f64 pts = aki_packet_read_f64(packet); + s32 sequence = nn_packet_read_s32(packet); + f64 pts = nn_packet_read_f64(packet); lia_list_toggle_pause(list, sequence, pts); break; } case CAMU_LIST_SEEK: { - s32 sequence = aki_packet_read_s32(packet); - u32 id = aki_packet_read_u32(packet); - f64 percent = aki_packet_read_f64(packet); + s32 sequence = nn_packet_read_s32(packet); + u32 id = nn_packet_read_u32(packet); + f64 percent = nn_packet_read_f64(packet); lia_list_seek(list, sequence, id, percent); break; } @@ -515,24 +553,24 @@ static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn break; } case CAMU_LIST_END: { - u32 id = aki_packet_read_u32(packet); + u32 id = nn_packet_read_u32(packet); lia_list_end(list, id); break; } } out: - aki_packet_free(packet); + nn_packet_free(packet); return false; } -static struct aki_rpc_command commands[] = { +static struct nn_rpc_command commands[] = { { .op = CAMU_SERVER_IDENTIFY, .callback = identify_callback, .userdata = NULL }, { .op = CAMU_SERVER_CLIENT_COMMAND, .callback = client_command_callback, .userdata = NULL }, { .op = CAMU_SERVER_LIST_ACTION, .callback = list_action_callback, .userdata = NULL } }; -static void connection_callback(void *userdata, struct aki_rpc_connection *conn) +static void connection_callback(void *userdata, struct nn_rpc_connection *conn) { // TODO: Cleanup zombie connections. (void)userdata; @@ -555,7 +593,7 @@ static void cleanup_sink(struct camu_server_sink *sink) al_free(sink); } -static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) +static void connection_closed_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_server *server = (struct camu_server *)userdata; @@ -594,12 +632,12 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection } } -static bool multiplex_callback(void *userdata, u8 id, struct aki_socket *sock) +static bool multiplex_callback(void *userdata, u8 id, struct nn_socket *sock) { struct camu_server *server = (struct camu_server *)userdata; switch (id) { case CAMU_MULTIPLEX_RPC: - aki_rpc_add_socket(&server->server, sock); + nn_rpc_add_socket(&server->server, sock); return true; case CAMU_MULTIPLEX_LIANA: lia_server_add_socket(&server->data.server, sock); @@ -610,7 +648,7 @@ static bool multiplex_callback(void *userdata, u8 id, struct aki_socket *sock) return false; } -bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop) +bool camu_server_init(struct camu_server *server, u8 type, struct nn_event_loop *loop) { server->loop = loop; server->addr = al_str_zero(); @@ -626,10 +664,10 @@ bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop list->userdata = server; al_array_push(server->lists, list); - aki_rpc_init(&server->server, server->loop, connection_callback, connection_closed_callback, server); + nn_rpc_init(&server->server, server->loop, connection_callback, connection_closed_callback, server); for (u32 i = 0; i < ARRAY_SIZE(commands); i++) { commands[i].userdata = server; - aki_rpc_add_command(&server->server, &commands[i]); + nn_rpc_add_command(&server->server, &commands[i]); } lia_server_init(&server->data.server, server->loop); @@ -639,14 +677,14 @@ bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop camu_portal_init(&server->bridge, &server->cache, server->loop, server); #endif - return aki_multiplex_socket_init(&server->multi, type, multiplex_callback, server); + return nn_multiplex_socket_init(&server->multi, type, multiplex_callback, server); } bool camu_server_listen(struct camu_server *server, str *addr, u16 port) { al_str_clone(&server->addr, addr); - if (server->multi.sock.type == AKI_SOCKET_TCP) addr = NULL; // any - return aki_multiplex_socket_listen(&server->multi, server->loop, addr, port); + if (server->multi.sock.type == NNWT_SOCKET_TCP) addr = NULL; // any + return nn_multiplex_socket_listen(&server->multi, server->loop, addr, port); } void camu_server_close(struct camu_server *server) @@ -655,7 +693,7 @@ void camu_server_close(struct camu_server *server) camu_portal_close(&server->bridge); #endif lia_server_close(&server->data.server); - aki_multiplex_socket_close(&server->multi); + nn_multiplex_socket_close(&server->multi); } void camu_server_free(struct camu_server *server) @@ -668,11 +706,11 @@ void camu_server_free(struct camu_server *server) //} al_array_free(server->data.resources); lia_server_free(&server->data.server); - aki_rpc_free(&server->server); + nn_rpc_free(&server->server); al_str_free(&server->addr); } -void camu_server_local_add(struct camu_server *server, struct aki_packet *packet) +void camu_server_local_add(struct camu_server *server, struct nn_packet *packet) { handle_add_command(server, al_array_last(server->lists), packet); } diff --git a/src/server/server.h b/src/server/server.h index ec0f85b..ab96e80 100644 --- a/src/server/server.h +++ b/src/server/server.h @@ -1,7 +1,7 @@ #pragma once -#include <aki/multiplex.h> -#include <aki/rpc2.h> +#include <nnwt/multiplex.h> +#include <nnwt/rpc2.h> #include "../liana/server.h" #include "../liana/list.h" @@ -13,25 +13,25 @@ #include "resource.h" struct camu_server_node { - struct aki_rpc_connection *conn; + struct nn_rpc_connection *conn; }; struct camu_server_client { - struct aki_rpc_connection *conn; + struct nn_rpc_connection *conn; struct camu_user *user; }; struct camu_server_sink { - struct aki_rpc_connection *conn; + struct nn_rpc_connection *conn; str name; struct camu_server *server; }; struct camu_server { - struct aki_event_loop *loop; + struct nn_event_loop *loop; str addr; - struct aki_multiplex_socket multi; - struct aki_rpc server; + struct nn_multiplex_socket multi; + struct nn_rpc server; array(struct camu_server_node *) nodes; array(struct camu_server_client *) clients; array(struct camu_server_sink *) sinks; @@ -49,9 +49,9 @@ struct camu_server { void *userdata; }; -bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop); +bool camu_server_init(struct camu_server *server, u8 type, struct nn_event_loop *loop); bool camu_server_listen(struct camu_server *server, str *addr, u16 port); void camu_server_close(struct camu_server *server); void camu_server_free(struct camu_server *server); -void camu_server_local_add(struct camu_server *server, struct aki_packet *packet); +void camu_server_local_add(struct camu_server *server, struct nn_packet *packet); diff --git a/src/sink/desktop.c b/src/sink/desktop.c index ccd2e11..4a1b9ca 100644 --- a/src/sink/desktop.c +++ b/src/sink/desktop.c @@ -101,7 +101,7 @@ bool camu_desktop_init(struct camu_desktop *c, const char *name) return true; } -bool camu_desktop_connect(struct camu_desktop *c, u8 type, struct aki_event_loop *loop, str *addr, u16 port) +bool camu_desktop_connect(struct camu_desktop *c, u8 type, struct nn_event_loop *loop, str *addr, u16 port) { c->loop = loop; camu_sink_init(&c->sink, c->loop, &c->mixer, c->renderer); diff --git a/src/sink/desktop.h b/src/sink/desktop.h index efc1c34..2b3ec94 100644 --- a/src/sink/desktop.h +++ b/src/sink/desktop.h @@ -2,14 +2,14 @@ #include <al/types.h> #include <al/str.h> -#include <aki/event_loop.h> +#include <nnwt/event_loop.h> #include "../libsink/sink.h" #include "../screen/screen.h" struct camu_desktop { - struct aki_event_loop *loop; + struct nn_event_loop *loop; s32 should_quit; struct camu_screen scr; struct camu_renderer *renderer; @@ -20,7 +20,7 @@ struct camu_desktop { }; bool camu_desktop_init(struct camu_desktop *c, const char *name); -bool camu_desktop_connect(struct camu_desktop *c, u8 type, struct aki_event_loop *loop, str *addr, u16 port); +bool camu_desktop_connect(struct camu_desktop *c, u8 type, struct nn_event_loop *loop, str *addr, u16 port); bool camu_desktop_tick(struct camu_desktop *c); void camu_desktop_stop(struct camu_desktop *c); void camu_desktop_free(struct camu_desktop *c); diff --git a/src/sink/input_simulator.c b/src/sink/input_simulator.c index af0105d..773ac95 100644 --- a/src/sink/input_simulator.c +++ b/src/sink/input_simulator.c @@ -1,28 +1,28 @@ #include "../screen/screen.h" #ifdef CAMU_SCREEN_DEBUG_KEY -#include <aki/thread.h> +#include <nnwt/thread.h> #include <al/random.h> #include "input_simulator.h" // This is ignoring all thread-safety. static s32 quit = 1; -static struct aki_thread thread; +static struct nn_thread thread; enum { SKIP, BACKSKIP, - SHUFFLE, TOGGLE_PAUSE, + SHUFFLE, SEEK, MARK, // count }; -static aki_thread_result AKI_THREADCALL input_simulation_thread(void *userdata) +static nn_thread_result NNWT_THREADCALL input_simulation_thread(void *userdata) { struct camu_sink *sink = (struct camu_sink *)userdata; while (!quit) { - aki_thread_sleep(AKI_TS_FROM_USEC(30000)); + nn_thread_sleep(NNWT_TS_FROM_USEC(30000)); switch (al_rand() % MARK) { case SKIP: camu_sink_skip(sink, (al_rand() % 5)); @@ -47,7 +47,7 @@ static aki_thread_result AKI_THREADCALL input_simulation_thread(void *userdata) void camu_input_simulator_run(struct camu_sink *sink) { quit = 0; - aki_thread_create(&thread, input_simulation_thread, sink); + nn_thread_create(&thread, input_simulation_thread, sink); } bool camu_input_simulator_running(void) @@ -58,6 +58,6 @@ bool camu_input_simulator_running(void) void camu_input_simulator_stop(void) { quit = 1; - aki_thread_join(&thread); + nn_thread_join(&thread); } #endif diff --git a/src/util/blocking_ring_buffer.c b/src/util/blocking_ring_buffer.c index 5cc3325..ec6cbb8 100644 --- a/src/util/blocking_ring_buffer.c +++ b/src/util/blocking_ring_buffer.c @@ -4,17 +4,17 @@ #define LOCK(n) \ buffer->req_start = n; \ - aki_cond_wait(&buffer->cond, &buffer->mutex); + nn_cond_wait(&buffer->cond, &buffer->mutex); #define SIGNAL() \ - aki_cond_signal(&buffer->cond); + nn_cond_signal(&buffer->cond); void camu_ring_buffer_init(struct camu_ring_buffer *buffer, s32 size) { buffer->size = size; buffer->data = (u8 *)al_malloc(buffer->size); - aki_mutex_init(&buffer->mutex); - aki_cond_init(&buffer->cond); + nn_mutex_init(&buffer->mutex); + nn_cond_init(&buffer->cond); camu_ring_buffer_reset(buffer); buffer->reset_unlock = false; buffer->enabled = true; @@ -23,71 +23,71 @@ void camu_ring_buffer_init(struct camu_ring_buffer *buffer, s32 size) // TODO: does locking here actually prevent this from messing something up. void camu_ring_buffer_reset(struct camu_ring_buffer *buffer) { - aki_mutex_lock(&buffer->mutex); + nn_mutex_lock(&buffer->mutex); buffer->valid_size = 0; buffer->start = 0; buffer->end = 0; // If get_chunk was waiting it doesn't matter where req_start // was, there should be room now. - if (aki_cond_is_waiting(&buffer->cond)) { + if (nn_cond_is_waiting(&buffer->cond)) { buffer->reset_unlock = true; SIGNAL(); } buffer->wrap = false; - aki_mutex_unlock(&buffer->mutex); + nn_mutex_unlock(&buffer->mutex); } void camu_ring_buffer_disable(struct camu_ring_buffer *buffer) { - aki_mutex_lock(&buffer->mutex); + nn_mutex_lock(&buffer->mutex); buffer->enabled = false; // No state on the buffer will be edited by the time a get_chunk // returns because of this. So, the next call to get_chunk after // the buffer is re-enabled should be correct. - if (aki_cond_is_waiting(&buffer->cond)) { + if (nn_cond_is_waiting(&buffer->cond)) { buffer->reset_unlock = true; SIGNAL(); } - aki_mutex_unlock(&buffer->mutex); + nn_mutex_unlock(&buffer->mutex); } void camu_ring_buffer_enable(struct camu_ring_buffer *buffer) { - aki_mutex_lock(&buffer->mutex); + nn_mutex_lock(&buffer->mutex); buffer->enabled = true; - aki_mutex_unlock(&buffer->mutex); + nn_mutex_unlock(&buffer->mutex); } s32 camu_ring_buffer_space(struct camu_ring_buffer *buffer) { s32 available = 0; - aki_mutex_lock(&buffer->mutex); + nn_mutex_lock(&buffer->mutex); if (buffer->end < buffer->start) { available = buffer->start - buffer->end; } else { available = buffer->size - buffer->end; } - aki_mutex_unlock(&buffer->mutex); + nn_mutex_unlock(&buffer->mutex); return available; } s32 camu_ring_buffer_occupied(struct camu_ring_buffer *buffer) { s32 filled = 0; - aki_mutex_lock(&buffer->mutex); + nn_mutex_lock(&buffer->mutex); if (buffer->end < buffer->start) { filled = (buffer->valid_size - buffer->start) + buffer->end; } else { // If the buffer is empty this will equate to 0. filled = buffer->end - buffer->start; } - aki_mutex_unlock(&buffer->mutex); + nn_mutex_unlock(&buffer->mutex); return filled; } u8 *camu_ring_buffer_get_chunk(struct camu_ring_buffer *buffer, s32 n) { - aki_mutex_lock(&buffer->mutex); + nn_mutex_lock(&buffer->mutex); if (!buffer->enabled) return NULL; @@ -136,7 +136,7 @@ u8 *camu_ring_buffer_get_chunk(struct camu_ring_buffer *buffer, s32 n) s32 camu_ring_buffer_read(struct camu_ring_buffer *buffer, u8 **data, s32 n) { - aki_mutex_lock(&buffer->mutex); + nn_mutex_lock(&buffer->mutex); // Set data even if we are going to read 0 bytes to avoid possible issues // with data being NULL after a call to read. @@ -164,11 +164,11 @@ s32 camu_ring_buffer_read(struct camu_ring_buffer *buffer, u8 **data, s32 n) buffer->start = 0; buffer->valid_size = 0; // We wrapped and aren't specifically waiting for a position after the wrap, so unlock. - if (aki_cond_is_waiting(&buffer->cond) && !buffer->wrap) { + if (nn_cond_is_waiting(&buffer->cond) && !buffer->wrap) { SIGNAL(); } buffer->wrap = false; - } else if (aki_cond_is_waiting(&buffer->cond) && (!buffer->wrap && buffer->start > buffer->req_start)) { + } else if (nn_cond_is_waiting(&buffer->cond) && (!buffer->wrap && buffer->start > buffer->req_start)) { // We read past the position we were waiting for, so unlock. SIGNAL(); } @@ -182,12 +182,12 @@ s32 camu_ring_buffer_read(struct camu_ring_buffer *buffer, u8 **data, s32 n) // only after the pointer has been used. void camu_ring_buffer_unlock(struct camu_ring_buffer *buffer) { - aki_mutex_unlock(&buffer->mutex); + nn_mutex_unlock(&buffer->mutex); } void camu_ring_buffer_free(struct camu_ring_buffer *buffer) { if (buffer->data) al_free(buffer->data); - aki_mutex_destroy(&buffer->mutex); - aki_cond_destroy(&buffer->cond); + nn_mutex_destroy(&buffer->mutex); + nn_cond_destroy(&buffer->cond); } diff --git a/src/util/blocking_ring_buffer.h b/src/util/blocking_ring_buffer.h index ff3521e..a5288d0 100644 --- a/src/util/blocking_ring_buffer.h +++ b/src/util/blocking_ring_buffer.h @@ -1,7 +1,7 @@ #pragma once #include <al/types.h> -#include <aki/thread.h> +#include <nnwt/thread.h> // Ring buffer with blocking write and non-blocking read. // - Write can block and always returns the exact amount requested. @@ -20,8 +20,8 @@ struct camu_ring_buffer { // True if we need to wait for the start to wrap around // to the begining before considering signaling. bool wrap; - struct aki_mutex mutex; - struct aki_cond cond; + struct nn_mutex mutex; + struct nn_cond cond; // Signal get_chunk to return null and not modifiy the buffer // if it was blocking and reset was called. bool reset_unlock; diff --git a/src/util/color_palette.c b/src/util/color_palette.c index cc7f5b5..b67f1bb 100644 --- a/src/util/color_palette.c +++ b/src/util/color_palette.c @@ -1,7 +1,7 @@ #include "color_palette.h" -#ifdef AKIYO_HAS_JSON -#include <aki/file.h> +#ifdef NAUNET_HAS_JSON +#include <nnwt/file.h> #include <jansson.h> static char *special_colors[2] = { @@ -33,14 +33,14 @@ struct camu_color_palette global_color_palette = { 0 }; bool camu_color_palette_init(str *path) { -#ifdef AKIYO_HAS_JSON - struct aki_file file; - if (!aki_file_open(&file, path, 0)) { +#ifdef NAUNET_HAS_JSON + struct nn_file file; + if (!nn_file_open(&file, path, 0)) { return false; } str s; - aki_file_read_as_str(&file, &s); - aki_file_close(&file); + nn_file_read_as_str(&file, &s); + nn_file_close(&file); json_error_t error; json_t *root = json_loadb(s.data, s.len, 0, &error); if (!root) return false; diff --git a/src/util/queue.h b/src/util/queue.h index 4f184ba..ac0dfec 100644 --- a/src/util/queue.h +++ b/src/util/queue.h @@ -1,53 +1,53 @@ #pragma once #include <al/array.h> -#include <aki/thread.h> +#include <nnwt/thread.h> #define queue(type) \ struct { \ array(type) a; \ - struct aki_mutex mutex; \ + struct nn_mutex mutex; \ } #define camu_queue_size(q, r) \ do { \ - aki_mutex_lock(&(q).mutex); \ + nn_mutex_lock(&(q).mutex); \ r = (q).a.size; \ - aki_mutex_unlock(&(q).mutex); \ + nn_mutex_unlock(&(q).mutex); \ } while (0) #define camu_queue_init(q) \ do { \ al_array_init((q).a); \ - aki_mutex_init(&(q).mutex); \ + nn_mutex_init(&(q).mutex); \ } while (0) #define camu_queue_push(q, item) \ do { \ - aki_mutex_lock(&(q).mutex); \ + nn_mutex_lock(&(q).mutex); \ al_array_push((q).a, item); \ - aki_mutex_unlock(&(q).mutex); \ + nn_mutex_unlock(&(q).mutex); \ } while (0) #define camu_queue_pop(q, r) \ do { \ - aki_mutex_lock(&(q).mutex); \ + nn_mutex_lock(&(q).mutex); \ al_array_pop_at((q).a, 0, r); \ - aki_mutex_unlock(&(q).mutex); \ + nn_mutex_unlock(&(q).mutex); \ } while (0) #define camu_queue_try_pop(q, s, r) \ do { \ - aki_mutex_lock(&(q).mutex); \ + nn_mutex_lock(&(q).mutex); \ s = (q).a.size; \ if (s > 0) { \ al_array_pop_at((q).a, 0, r); \ } \ - aki_mutex_unlock(&(q).mutex); \ + nn_mutex_unlock(&(q).mutex); \ } while (0) #define camu_queue_free(q) \ do { \ - aki_mutex_destroy(&(q).mutex); \ + nn_mutex_destroy(&(q).mutex); \ al_array_free((q).a); \ } while (0) |