diff options
| -rw-r--r-- | flake.lock | 18 | ||||
| -rw-r--r-- | src/cache/backing.h | 6 | ||||
| -rw-r--r-- | src/cache/backings/file.c | 60 | ||||
| -rw-r--r-- | src/cache/backings/file.h | 1 | ||||
| -rw-r--r-- | src/cache/backings/file_common.h | 28 | ||||
| -rw-r--r-- | src/cache/backings/file_mapped.c | 41 | ||||
| -rw-r--r-- | src/cache/backings/memory.c | 65 | ||||
| -rw-r--r-- | src/cache/backings/memory.h | 4 | ||||
| -rw-r--r-- | src/cache/entry.c | 18 | ||||
| -rw-r--r-- | src/cache/entry.h | 22 | ||||
| -rw-r--r-- | src/cache/handle.c | 15 | ||||
| -rw-r--r-- | src/cache/handlers/cdio.c | 5 | ||||
| -rw-r--r-- | src/cache/handlers/file.c | 1 | ||||
| -rw-r--r-- | src/cache/handlers/http.c | 21 | ||||
| -rw-r--r-- | src/cache/threaded_waits.c | 12 | ||||
| -rw-r--r-- | src/cache/threaded_waits.h | 8 | ||||
| -rw-r--r-- | src/codec/stb_image/server.c | 3 | ||||
| -rw-r--r-- | src/liana/handlers/codec_client.c | 2 | ||||
| -rw-r--r-- | src/liana/server.c | 11 | ||||
| -rw-r--r-- | src/liana/vcr.c | 11 |
20 files changed, 198 insertions, 154 deletions
@@ -23,11 +23,11 @@ ] }, "locked": { - "lastModified": 1745782215, - "narHash": "sha256-mx27J2HYQT+nGXTyUWKrUuxRzpr1FVVr59ZH4oNzOyw=", + "lastModified": 1745987135, + "narHash": "sha256-8Up4QPuMZEJBU0eefAY+nUe7DYKQQzvaHnMpNdwRgKA=", "owner": "nix-community", "repo": "home-manager", - "rev": "7b2aae3fb39928aecc5e41c10a9c87c4881614d5", + "rev": "d2b3e6c83d457aa0e7f9344c61c3fed32bad0f7e", "type": "github" }, "original": { @@ -39,11 +39,11 @@ }, "nixos-hardware": { "locked": { - "lastModified": 1745503349, - "narHash": "sha256-bUGjvaPVsOfQeTz9/rLTNLDyqbzhl0CQtJJlhFPhIYw=", + "lastModified": 1745955289, + "narHash": "sha256-mmV2oPhQN+YF2wmnJzXX8tqgYmUYXUj3uUUBSTmYN5o=", "owner": "NixOS", "repo": "nixos-hardware", - "rev": "f7bee55a5e551bd8e7b5b82c9bc559bc50d868d1", + "rev": "72081c9fbbef63765ae82bff9727ea79cc86bd5b", "type": "github" }, "original": { @@ -91,11 +91,11 @@ }, "nixpkgs_2": { "locked": { - "lastModified": 1745526057, - "narHash": "sha256-ITSpPDwvLBZBnPRS2bUcHY3gZSwis/uTe255QgMtTLA=", + "lastModified": 1745930157, + "narHash": "sha256-y3h3NLnzRSiUkYpnfvnS669zWZLoqqI6NprtLQ+5dck=", "owner": "NixOS", "repo": "nixpkgs", - "rev": "f771eb401a46846c1aebd20552521b233dd7e18b", + "rev": "46e634be05ce9dc6d4db8e664515ba10b78151ae", "type": "github" }, "original": { diff --git a/src/cache/backing.h b/src/cache/backing.h index d9d8bad..f501dff 100644 --- a/src/cache/backing.h +++ b/src/cache/backing.h @@ -15,13 +15,15 @@ struct cch_range { struct cch_backing { u8 mode; + off_t size; array(struct cch_range) available; void (*lock)(struct cch_backing *); void (*write)(struct cch_backing *, u8 *, off_t, size_t *); void (*read)(struct cch_backing *, u8 *, off_t, size_t *); u8 *(*get_ptr)(struct cch_backing *, off_t, size_t *); void (*unlock)(struct cch_backing *); - off_t (*get_size_estimate)(struct cch_backing *); - void (*resize)(struct cch_backing *, size_t); + void (*set_size)(struct cch_backing *, off_t, bool); + off_t (*get_size_if_known)(struct cch_backing *); + void (*finalize)(struct cch_backing *); void (*free)(struct cch_backing **); }; diff --git a/src/cache/backings/file.c b/src/cache/backings/file.c index 500e16f..06cc600 100644 --- a/src/cache/backings/file.c +++ b/src/cache/backings/file.c @@ -7,40 +7,48 @@ static void file_backing_lock(struct cch_backing *backing) nn_mutex_lock(&file->mutex); } -static void file_backing_write(struct cch_backing *backing, u8 *buf, off_t index, size_t *size) +static void maybe_seek_to_index(struct cch_backing_file *file, off_t index) { - struct cch_backing_file *file = (struct cch_backing_file *)backing; - nn_mutex_lock(&file->mutex); if (file->u.pointer != index) { file->u.pointer = nn_file_seek(&file->file, index, SEEK_SET); al_assert(file->u.pointer == index); } - off_t n = (off_t)*size; - if (file->size >= 0 && index + n >= file->size) { - n = MAX(file->size - index, (off_t)0); +} + +static void file_backing_write(struct cch_backing *backing, u8 *buf, off_t index, size_t *count) +{ + struct cch_backing_file *file = (struct cch_backing_file *)backing; + off_t n = (off_t)*count; + nn_mutex_lock(&file->mutex); + maybe_seek_to_index(file, index); + off_t size = file->backing.size; + if (size > 0 && index + n >= size) { + n = MAX(size - index, (off_t)0); + } + if (n > 0) { + nn_file_write(&file->file, buf, n); + file->u.pointer += n; + cch_backing_fill_range(&file->backing, index, n); } - nn_file_write(&file->file, buf, n); - file->u.pointer += n; - cch_backing_fill_range(&file->backing, index, n); - *size = (size_t)n; + *count = (size_t)n; nn_mutex_unlock(&file->mutex); } -static void file_backing_read(struct cch_backing *backing, u8 *buf, off_t index, size_t *size) +static void file_backing_read(struct cch_backing *backing, u8 *buf, off_t index, size_t *count) { struct cch_backing_file *file = (struct cch_backing_file *)backing; + off_t n = (off_t)*count; nn_mutex_lock(&file->mutex); - if (file->u.pointer != index) { - file->u.pointer = nn_file_seek(&file->file, index, SEEK_SET); - al_assert(file->u.pointer == index); + maybe_seek_to_index(file, index); + off_t size = file->backing.size; + if (size > 0 && index + n >= size) { + n = MAX(size - index, (off_t)0); } - off_t n = (off_t)*size; - if (file->size >= 0 && index + n >= file->size) { - n = MAX(file->size - index, (off_t)0); + if (n > 0) { + nn_file_read(&file->file, buf, n); + file->u.pointer += n; } - nn_file_read(&file->file, buf, n); - file->u.pointer += n; - *size = (size_t)n; + *count = (size_t)n; nn_mutex_unlock(&file->mutex); } @@ -50,12 +58,12 @@ static void file_backing_unlock(struct cch_backing *backing) nn_mutex_unlock(&file->mutex); } -static void file_backing_resize(struct cch_backing *backing, size_t size) +static void file_backing_set_size(struct cch_backing *backing, off_t size, bool expand_to_size) { struct cch_backing_file *file = (struct cch_backing_file *)backing; nn_mutex_lock(&file->mutex); - file->size = size; - if (size > file->file.size) { + file->backing.size = size; + if (expand_to_size && (size_t)size > file->file.size) { nn_file_truncate(&file->file, size); } nn_mutex_unlock(&file->mutex); @@ -74,13 +82,15 @@ struct cch_backing *cch_backing_file_create(str *path, size_t size) { struct cch_backing_file *file = al_alloc_object(struct cch_backing_file); file->backing.mode = CACHE_BACKING_READ; + file->backing.size = -1; file->backing.lock = file_backing_lock; file->backing.write = file_backing_write; file->backing.read = file_backing_read; file->backing.get_ptr = NULL; file->backing.unlock = file_backing_unlock; - file->backing.get_size_estimate = file_backing_get_size_estimate; - file->backing.resize = file_backing_resize; + file->backing.set_size = file_backing_set_size; + file->backing.get_size_if_known = file_backing_get_size_if_known; + file->backing.finalize = file_backing_finalize; file->backing.free = file_backing_free; al_array_init(file->backing.available); if (!file_open_internal(file, path, size)) { diff --git a/src/cache/backings/file.h b/src/cache/backings/file.h index 06daac9..361aad6 100644 --- a/src/cache/backings/file.h +++ b/src/cache/backings/file.h @@ -10,7 +10,6 @@ struct cch_backing_file { struct cch_backing backing; struct nn_file file; - off_t size; union { off_t pointer; // CACHE_BACKING_READ void *map; // CACHE_BACKING_MAPPED diff --git a/src/cache/backings/file_common.h b/src/cache/backings/file_common.h index 3f4ed3b..936babe 100644 --- a/src/cache/backings/file_common.h +++ b/src/cache/backings/file_common.h @@ -7,29 +7,33 @@ AL_IGNORE_WARNING("-Wunused-function") -static off_t file_backing_get_size_estimate(struct cch_backing *backing) +static off_t file_backing_get_size_if_known(struct cch_backing *backing) { struct cch_backing_file *file = (struct cch_backing_file *)backing; nn_mutex_lock(&file->mutex); - off_t size = file->size; + off_t size = file->backing.size; nn_mutex_unlock(&file->mutex); return size; } +static void file_backing_finalize(struct cch_backing *backing) +{ + struct cch_backing_file *file = (struct cch_backing_file *)backing; + nn_mutex_lock(&file->mutex); + file->backing.size = file->file.size; + nn_mutex_unlock(&file->mutex); +} + static bool file_open_internal(struct cch_backing_file *file, str *path, size_t size) { s32 flags = size != 0 ? NNWT_FILE_CREATE : NNWT_FILE_READONLY; - if (!nn_file_open(&file->file, path, flags)) { - return false; - } + if (!nn_file_open(&file->file, path, flags)) return false; + size_t filesize = file->file.size; if (!size) { - file->size = file->file.size; - cch_backing_fill_range(&file->backing, 0, file->size); - } else { - file->size = -1; - if (file->file.size < size) { - nn_file_truncate(&file->file, size); - } + file->backing.size = filesize; + cch_backing_fill_range(&file->backing, 0, filesize); + } else if (filesize < 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 adb840c..7657a21 100644 --- a/src/cache/backings/file_mapped.c +++ b/src/cache/backings/file_mapped.c @@ -7,29 +7,33 @@ static void file_backing_lock(struct cch_backing *backing) nn_mutex_lock(&file->mutex); } -static void file_backing_write(struct cch_backing *backing, u8 *buf, off_t index, size_t *size) +static void file_backing_write(struct cch_backing *backing, u8 *buf, off_t index, size_t *count) { struct cch_backing_file *file = (struct cch_backing_file *)backing; + off_t n = (off_t)*count; nn_mutex_lock(&file->mutex); - off_t n = (off_t)*size; - if (file->size >= 0 && index + n >= file->size) { - n = MAX(file->size - index, (off_t)0); + off_t size = file->backing.size; + if (size > 0 && index + n >= size) { + n = MAX(size - index, (off_t)0); } - al_memcpy(file->u.map + index, buf, n); - cch_backing_fill_range(&file->backing, index, n); - *size = (size_t)n; + if (n > 0) { + al_memcpy(file->u.map + index, buf, n); + cch_backing_fill_range(&file->backing, index, n); + } + *count = (size_t)n; nn_mutex_unlock(&file->mutex); } -static u8 *file_backing_get_ptr(struct cch_backing *backing, off_t index, size_t *size) +static u8 *file_backing_get_ptr(struct cch_backing *backing, off_t index, size_t *count) { struct cch_backing_file *file = (struct cch_backing_file *)backing; + off_t n = (off_t)*count; nn_mutex_lock(&file->mutex); - off_t n = (off_t)*size; - if (file->size >= 0 && index + n >= file->size) { - n = MAX(file->size - index, (off_t)0); + off_t size = file->backing.size; + if (size > 0 && index + n >= size) { + n = MAX(size - index, (off_t)0); } - *size = (size_t)n; + *count = (size_t)n; return (u8 *)(file->u.map + index); } @@ -39,12 +43,13 @@ static void file_backing_unlock(struct cch_backing *backing) nn_mutex_unlock(&file->mutex); } -static void file_backing_resize(struct cch_backing *backing, size_t size) +static void file_backing_set_size(struct cch_backing *backing, off_t size, bool expand_to_size) { struct cch_backing_file *file = (struct cch_backing_file *)backing; + al_assert(expand_to_size || size == 0); nn_mutex_lock(&file->mutex); - file->size = size; - if (size > file->file.size) { + file->backing.size = size; + if ((size_t)size > file->file.size) { nn_file_munmap(&file->file, file->u.map); nn_file_truncate(&file->file, size); file->u.map = nn_file_mmap(&file->file); @@ -67,13 +72,15 @@ struct cch_backing *cch_backing_file_create(str *path, size_t size) { struct cch_backing_file *file = al_alloc_object(struct cch_backing_file); file->backing.mode = CACHE_BACKING_MAPPED; + file->backing.size = -1; file->backing.lock = file_backing_lock; file->backing.write = file_backing_write; file->backing.read = NULL; file->backing.get_ptr = file_backing_get_ptr; file->backing.unlock = file_backing_unlock; - file->backing.get_size_estimate = file_backing_get_size_estimate; - file->backing.resize = file_backing_resize; + file->backing.set_size = file_backing_set_size; + file->backing.get_size_if_known = file_backing_get_size_if_known; + file->backing.finalize = file_backing_finalize; file->backing.free = file_backing_free; al_array_init(file->backing.available); if (!file_open_internal(file, path, size)) { diff --git a/src/cache/backings/memory.c b/src/cache/backings/memory.c index 8e0306a..6b3da62 100644 --- a/src/cache/backings/memory.c +++ b/src/cache/backings/memory.c @@ -2,10 +2,11 @@ #include "memory.h" -static void ensure_alloced(struct cch_backing_memory *mem, size_t size) +static void ensure_allocated(struct cch_backing_memory *mem, off_t amount) { - if (size > mem->alloc) { - mem->alloc = al_next_power_of_two(size); + if (amount > mem->expanse) mem->expanse = amount; + if (amount > mem->alloc) { + mem->alloc = al_next_power_of_two(amount); mem->data = al_realloc(mem->data, mem->alloc); } } @@ -16,26 +17,35 @@ static void memory_backing_lock(struct cch_backing *backing) nn_mutex_lock(&mem->mutex); } -static void memory_backing_write(struct cch_backing *backing, u8 *buf, off_t index, size_t *size) +static void memory_backing_write(struct cch_backing *backing, u8 *buf, off_t index, size_t *count) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; + off_t n = (off_t)*count; nn_mutex_lock(&mem->mutex); - if (mem->size >= 0 && index + (off_t)*size >= mem->size) { - *size = MAX(mem->size - index, (off_t)0); + off_t size = mem->backing.size; + if (size > 0 && index + n >= size) { + n = MAX(size - index, (off_t)0); } - ensure_alloced(mem, index + *size); - al_memcpy(mem->data + index, buf, *size); - cch_backing_fill_range(&mem->backing, index, *size); + if (n > 0) { + ensure_allocated(mem, index + n); + al_memcpy(mem->data + index, buf, n); + cch_backing_fill_range(&mem->backing, index, n); + } + *count = (size_t)n; nn_mutex_unlock(&mem->mutex); } -static u8 *memory_backing_get_ptr(struct cch_backing *backing, off_t index, size_t *size) +static u8 *memory_backing_get_ptr(struct cch_backing *backing, off_t index, size_t *count) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; + off_t n = (off_t)*count; nn_mutex_lock(&mem->mutex); - if (mem->size >= 0 && index + (off_t)*size >= mem->size) { - *size = MAX(mem->size - index, (off_t)0); + off_t size = mem->backing.size; + if (size > 0 && index + n >= size) { + n = MAX(size - index, (off_t)0); } + *count = (size_t)n; + al_assert(mem->alloc >= index + n); return (u8 *)(mem->data + index); } @@ -45,21 +55,32 @@ static void memory_backing_unlock(struct cch_backing *backing) nn_mutex_unlock(&mem->mutex); } -static off_t memory_backing_get_size_estimate(struct cch_backing *backing) +static void memory_backing_set_size(struct cch_backing *backing, off_t size, bool expand_to_size) +{ + struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; + nn_mutex_lock(&mem->mutex); + mem->backing.size = size; + if (expand_to_size) { + al_assert(size > 0); + ensure_allocated(mem, size); + } + nn_mutex_unlock(&mem->mutex); +} + +static off_t memory_backing_get_size_if_known(struct cch_backing *backing) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; nn_mutex_lock(&mem->mutex); - off_t size = mem->size; + off_t size = mem->backing.size; nn_mutex_unlock(&mem->mutex); return size; } -static void memory_backing_resize(struct cch_backing *backing, size_t size) +static void memory_backing_finalize(struct cch_backing *backing) { struct cch_backing_memory *mem = (struct cch_backing_memory *)backing; nn_mutex_lock(&mem->mutex); - ensure_alloced(mem, size); - mem->size = size; + mem->backing.size = mem->expanse; nn_mutex_unlock(&mem->mutex); } @@ -77,16 +98,18 @@ struct cch_backing *cch_backing_memory_create(size_t size) { struct cch_backing_memory *mem = al_alloc_object(struct cch_backing_memory); mem->backing.mode = CACHE_BACKING_MAPPED; + mem->backing.size = -1; mem->backing.lock = memory_backing_lock; mem->backing.write = memory_backing_write; mem->backing.get_ptr = memory_backing_get_ptr; mem->backing.unlock = memory_backing_unlock; - mem->backing.get_size_estimate = memory_backing_get_size_estimate; - mem->backing.resize = memory_backing_resize; + mem->backing.set_size = memory_backing_set_size; + mem->backing.get_size_if_known = memory_backing_get_size_if_known; + mem->backing.finalize = memory_backing_finalize; mem->backing.free = memory_backing_free; - mem->data = al_malloc(size); + if (size > 0) mem->data = al_malloc(size); + else mem->data = NULL; mem->alloc = size; - mem->size = -1; 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 e290b48..4737eba 100644 --- a/src/cache/backings/memory.h +++ b/src/cache/backings/memory.h @@ -8,8 +8,8 @@ struct cch_backing_memory { struct cch_backing backing; u8 *data; - size_t alloc; - off_t size; + off_t alloc; + off_t expanse; struct nn_mutex mutex; }; diff --git a/src/cache/entry.c b/src/cache/entry.c index a3ad9e9..0bace31 100644 --- a/src/cache/entry.c +++ b/src/cache/entry.c @@ -19,24 +19,6 @@ str *cch_entry_get_liana(struct cch_entry *entry) return &entry->handler->liana; } -void cch_entry_set_size(struct cch_entry *entry, off_t filesize) -{ - if (filesize == 0) { - // The backing will keep size = -1 internally. - entry->unknown_size = true; - return; - } - entry->backing->resize(entry->backing, filesize); -} - -off_t cch_entry_get_size(struct cch_entry *entry) -{ - if (entry->unknown_size) { - return 0; - } - return entry->backing->get_size_estimate(entry->backing); -} - void cch_entry_return_handle(struct cch_entry *entry, struct cch_handle *handle) { al_assert(handle->entry == entry); diff --git a/src/cache/entry.h b/src/cache/entry.h index 161eaab..eb98c2f 100644 --- a/src/cache/entry.h +++ b/src/cache/entry.h @@ -29,12 +29,19 @@ - Is a liana "node handler" necessary? Cache/liana refactor: - - On creation, the handler can create multiple backings and chapters. - - Chapters include index into array of backings. - - get_handle() takes chapter as a arugment. - - get/set_size() no longer make sense. - - Remove liana field from handler. - - Method of requesting a unique file backing. + [x] get/set_size() no longer make sense. + [x] // @TODO: Unknown size is unhandled in backings. + - Handle keeps reading but will eventually have a wait cut short after a backing finalize. + [ ] Rethink liana field on handler. + - Data format hint on backing (codec or raw+description). + - Codec hints (priority)? + - Cue handling. + - How to specify what type of file we are working with on cch_handler_file_create(). + - Scan for file(s) in backing or server? + [ ] get_handle() takes chapter as a arugment. + [ ] On creation, the handler can create multiple backings and chapters. + [ ] Chapters include index into array of backings. + [ ] Method of requesting a unique file backing. */ struct cch_chapter { @@ -43,7 +50,6 @@ struct cch_chapter { }; struct cch_entry { - bool unknown_size; array(struct cch_chapter) chapters; struct cch_chapter *chapter; struct cch_backing *backing; @@ -53,7 +59,5 @@ struct cch_entry { bool cch_entry_get_handle(struct cch_entry *entry, struct cch_handle *handle); str *cch_entry_get_liana(struct cch_entry *entry); -void cch_entry_set_size(struct cch_entry *entry, off_t filesize); -off_t cch_entry_get_size(struct cch_entry *entry); void cch_entry_return_handle(struct cch_entry *entry, struct cch_handle *handle); void cch_entry_free(struct cch_entry **entry); diff --git a/src/cache/handle.c b/src/cache/handle.c index dc974d3..1ec0f47 100644 --- a/src/cache/handle.c +++ b/src/cache/handle.c @@ -15,17 +15,18 @@ static off_t wait_for_size(struct cch_handle *handle, struct cch_handler *handle log_trace("Wait for size failed."); return -1; } - return cch_entry_get_size(handler->entry); + struct cch_backing *backing = handle->entry->backing; + return backing->get_size_if_known(backing); } s32 cch_handle_read(struct cch_handle *handle, u8 *buf, s32 size) { struct cch_handler *handler = handle->entry->handler; - // @TODO: Unknown size is unhandled in backings. + struct cch_backing *backing = handle->entry->backing; // filesize < 0: Size is yet to be evaluated and we can wait on it. // filesize = 0: Size is explicitly unknown. - off_t filesize = cch_entry_get_size(handle->entry); - if (size < 0) filesize = wait_for_size(handle, handler); + off_t filesize = backing->get_size_if_known(backing); + if (filesize < 0) filesize = wait_for_size(handle, handler); if (filesize > 0) { if (handle->pointer >= filesize) { return CAMU_ERR_EOF; @@ -33,14 +34,13 @@ s32 cch_handle_read(struct cch_handle *handle, u8 *buf, s32 size) size = filesize - handle->pointer; } } + al_assert(size >= 0); handle->wait.start = handle->pointer; handle->wait.end = handle->pointer + size; if (!handler->wait_for_range(handler, &handle->wait)) { log_trace("Wait for range (%zd-%zd) failed.", handle->wait.start, handle->wait.end); return CAMU_ERR_EOF; } - struct cch_backing *backing = handle->entry->backing; - al_assert(size >= 0); size_t available = (size_t)size; if (backing->mode == CACHE_BACKING_MAPPED) { u8 *ptr = backing->get_ptr(backing, handle->pointer, &available); @@ -59,7 +59,8 @@ s32 cch_handle_read(struct cch_handle *handle, u8 *buf, s32 size) off_t cch_handle_seek(struct cch_handle *handle, off_t offset, s32 whence) { struct cch_handler *handler = handle->entry->handler; - off_t filesize = cch_entry_get_size(handle->entry); + struct cch_backing *backing = handle->entry->backing; + off_t filesize = backing->get_size_if_known(backing); if (filesize < 0) filesize = wait_for_size(handle, handler); if (filesize <= 0) return -1; if (whence == CAMU_SEEK_SIZE) { diff --git a/src/cache/handlers/cdio.c b/src/cache/handlers/cdio.c index e319f74..9a5e698 100644 --- a/src/cache/handlers/cdio.c +++ b/src/cache/handlers/cdio.c @@ -102,7 +102,7 @@ static void handler_cdio_maybe_spawn_worker(struct cch_handler *handler, size_t static bool handler_cdio_wait_for_range(struct cch_handler *handler, struct cch_handler_wait *wait) { struct cch_handler_cdio *cdio = (struct cch_handler_cdio *)handler; - return cch_threaded_wait_for_range(&cdio->waits, cdio->handler.entry, cdio->backing, wait); + return cch_threaded_wait_for_range(&cdio->waits, cdio->backing, wait); } static void handler_cdio_free(struct cch_handler **handler) @@ -172,7 +172,7 @@ static bool open_cd_drive(struct cch_handler_cdio *cdio) offset = track_end + 1; } // end is a valid index so add one framesize. - cch_entry_set_size(cdio->handler.entry, ((end - start) + 1) * CDIO_CD_FRAMESIZE_RAW); + cdio->backing->set_size(cdio->backing, ((end - start) + 1) * CDIO_CD_FRAMESIZE_RAW, true); cdio->start = -1; cdio->end = end; @@ -192,7 +192,6 @@ struct cch_entry *cch_handler_cdio_create(void) struct cch_entry *entry = al_alloc_object(struct cch_entry); entry->backing = backing; entry->handler = (struct cch_handler *)cdio; - entry->unknown_size = false; nn_mutex_init(&entry->mutex); entry->handler->can_seek = handler_cdio_can_seek; entry->handler->maybe_spawn_worker = handler_cdio_maybe_spawn_worker; diff --git a/src/cache/handlers/file.c b/src/cache/handlers/file.c index 9529772..cd78a2c 100644 --- a/src/cache/handlers/file.c +++ b/src/cache/handlers/file.c @@ -39,7 +39,6 @@ struct cch_entry *cch_handler_file_create(str *path, str *liana) struct cch_entry *entry = al_alloc_object(struct cch_entry); entry->backing = backing; entry->handler = (struct cch_handler *)file; - entry->unknown_size = false; nn_mutex_init(&entry->mutex); entry->handler->can_seek = handler_file_can_seek; entry->handler->maybe_spawn_worker = handler_file_maybe_spawn_worker; diff --git a/src/cache/handlers/http.c b/src/cache/handlers/http.c index ed84bdc..02f656b 100644 --- a/src/cache/handlers/http.c +++ b/src/cache/handlers/http.c @@ -25,8 +25,8 @@ static void http_callback(void *userdata, u8 op, u8 *buf, void *opaque) case NNWT_HTTP_WRITE: { size_t n = *(size_t *)opaque; // 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); + if (http->backing->get_size_if_known(http->backing) < 0) { + http->backing->set_size(http->backing, 0, false); } http->backing->write(http->backing, buf, http->pointer, &n); http->pointer += n; @@ -42,8 +42,10 @@ static void http_callback(void *userdata, u8 op, u8 *buf, void *opaque) case NNWT_HTTP_CONTENT_LENGTH: { curl_off_t length = *(curl_off_t *)opaque; log_debug("Content-Length: %"CURL_FORMAT_CURL_OFF_T".", length); - cch_entry_set_size(http->handler.entry, length); - cch_threaded_waits_signal_any(&http->waits); + if (length > 0) { // HTTP 204 No Content. finalize() with 0 size. + http->backing->set_size(http->backing, length, true); + } + cch_threaded_waits_signal_anys(&http->waits); break; } case NNWT_HTTP_REDIRECT: { @@ -52,6 +54,14 @@ static void http_callback(void *userdata, u8 op, u8 *buf, void *opaque) break; } case NNWT_HTTP_FINISHED: + http->backing->finalize(http->backing); + if (http->backing->get_size_if_known(http->backing) == 0) { + // The handle needs to consider that it could wake up from any + // wait with a filesize of 0. + cch_threaded_waits_disable_all(&http->waits); + } else { + cch_threaded_waits_evaluate(&http->waits, http->backing); + } log_debug("Transfer finished."); break; case NNWT_HTTP_ERROR: @@ -78,7 +88,7 @@ static void handler_http_maybe_spawn_worker(struct cch_handler *handler, size_t static bool handler_http_wait_for_range(struct cch_handler *handler, struct cch_handler_wait *wait) { struct cch_handler_http *http = (struct cch_handler_http *)handler; - return cch_threaded_wait_for_range(&http->waits, http->handler.entry, http->backing, wait); + return cch_threaded_wait_for_range(&http->waits, http->backing, wait); } static void handler_http_free(struct cch_handler **handler) @@ -106,7 +116,6 @@ struct cch_entry *cch_handler_http_create(str *url, struct nn_event_loop *loop) struct cch_entry *entry = al_alloc_object(struct cch_entry); entry->backing = backing; entry->handler = (struct cch_handler *)http; - entry->unknown_size = false; nn_mutex_init(&entry->mutex); entry->handler->can_seek = handler_http_can_seek; entry->handler->maybe_spawn_worker = handler_http_maybe_spawn_worker; diff --git a/src/cache/threaded_waits.c b/src/cache/threaded_waits.c index 6d8a5d8..3acbef1 100644 --- a/src/cache/threaded_waits.c +++ b/src/cache/threaded_waits.c @@ -9,7 +9,9 @@ void cch_threaded_waits_init(struct cch_handler_waits *thw) static bool wait_range_satisfied(struct cch_backing *backing, struct cch_handler_wait *wait) { - if (wait->end < 0) return true; + if (wait->end < 0 || (backing->size > 0 && wait->end > backing->size)) { + return true; + } struct cch_range *range; al_array_foreach_ptr(backing->available, i, range) { if (wait->start >= range->start && range->end >= wait->end) { @@ -19,10 +21,10 @@ static bool wait_range_satisfied(struct cch_backing *backing, struct cch_handler return false; } -bool cch_threaded_wait_for_range(struct cch_handler_waits *thw, struct cch_entry *entry, - struct cch_backing *backing, struct cch_handler_wait *wait) +bool cch_threaded_wait_for_range(struct cch_handler_waits *thw, struct cch_backing *backing, + struct cch_handler_wait *wait) { - off_t size = cch_entry_get_size(entry); + off_t size = backing->get_size_if_known(backing); // Lock before checking wait_range_satisfied() in case the result changes // before wait is possibly added to active. @@ -57,7 +59,7 @@ bool cch_threaded_wait_for_range(struct cch_handler_waits *thw, struct cch_entry return !canceled; } -void cch_threaded_waits_signal_any(struct cch_handler_waits *thw) +void cch_threaded_waits_signal_anys(struct cch_handler_waits *thw) { nn_mutex_lock(&thw->lock); struct cch_handler_wait *wait; diff --git a/src/cache/threaded_waits.h b/src/cache/threaded_waits.h index ce0f19c..1f41df2 100644 --- a/src/cache/threaded_waits.h +++ b/src/cache/threaded_waits.h @@ -4,7 +4,7 @@ #include <nnwt/thread.h> #include "backing.h" -#include "entry.h" +#include "wait.h" struct cch_handler_waits { bool disabled; @@ -13,9 +13,9 @@ struct cch_handler_waits { }; void cch_threaded_waits_init(struct cch_handler_waits *thw); -bool cch_threaded_wait_for_range(struct cch_handler_waits *thw, struct cch_entry *entry, - struct cch_backing *backing, struct cch_handler_wait *wait); -void cch_threaded_waits_signal_any(struct cch_handler_waits *thw); +bool cch_threaded_wait_for_range(struct cch_handler_waits *thw, struct cch_backing *backing, + struct cch_handler_wait *wait); +void cch_threaded_waits_signal_anys(struct cch_handler_waits *thw); void cch_threaded_waits_evaluate(struct cch_handler_waits *thw, struct cch_backing *backing); void cch_threaded_wait_disable(struct cch_handler_wait *wait); void cch_threaded_wait_enable(struct cch_handler_wait *wait); diff --git a/src/codec/stb_image/server.c b/src/codec/stb_image/server.c index 79327d0..e2d26e8 100644 --- a/src/codec/stb_image/server.c +++ b/src/codec/stb_image/server.c @@ -13,7 +13,8 @@ static bool stbi_demuxer_init(struct camu_demuxer *demux, struct cch_handle *han stb->eof = false; cch_handle_seek(stb->handle, 0, SEEK_SET); - off_t size = cch_entry_get_size(stb->handle->entry); + struct cch_backing *backing = stb->handle->entry->backing; + off_t size = backing->get_size_if_known(backing); nn_buffer_ensure_space(&stb->buffer, size); cch_handle_read(stb->handle, nn_buffer_get_ptr(&stb->buffer, 0), size); stb->buffer.size = size; diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c index 39bb2d1..fd2569b 100644 --- a/src/liana/handlers/codec_client.c +++ b/src/liana/handlers/codec_client.c @@ -141,7 +141,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc packet->rindex = rindex; if (!success) { - // Forcing in EOF on an error is not necessary but should exhibit less erratic behavior sink-side. + // Forcing in an EOF on an error is not necessary but behaves better in the sink. codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL); return false; } diff --git a/src/liana/server.c b/src/liana/server.c index 3670607..acb5f31 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -152,12 +152,9 @@ static void data_connection_closed_callback(void *userdata, struct nn_packet_str { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; (void)stream; - // We will never be here if init_thread() blocks or fails. disable_connection_and_wait(conn); - free_connection_stream(conn); - if (conn->disconnected) { free_connection(conn); } else { @@ -448,8 +445,11 @@ static void duration_signal_callback(void *userdata) node->handler->free(&node->handler); bool got_visual_data = false; struct lia_visual_data visual_data; - //cch_handle_seek(&node->handle, 0, SEEK_SET); - //got_visual_data = lia_prepare_visual_data(&node->handle, &visual_data); +#if 0 // @TODO: This is obviously bad because it blocks the event loop. + // Evaluate this when getting around to http/hls stream fixes. + cch_handle_seek(&node->handle, 0, SEEK_SET); + got_visual_data = lia_prepare_visual_data(&node->handle, &visual_data); +#endif cch_entry_return_handle(node->entry, &node->handle); if (should_free_node(node)) { free_node(node); @@ -516,7 +516,6 @@ void lia_server_free(struct lia_server *server) al_free(node); } al_array_free(server->nodes); - al_assert(!server->dormant_connections.count); al_array_free(server->dormant_connections); } diff --git a/src/liana/vcr.c b/src/liana/vcr.c index 4e176e7..4500681 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -104,6 +104,7 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) // NULL packet means flush. bool success = track->client->handle_packet(track->client, packet); if (!success) { + // In the case of codec_client, an EOF will be sent on an error. log_error("Error handling packet, exiting track thread."); return_entire_cache(track); track->cache.disabled = true; @@ -290,16 +291,18 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) break; } case LIANA_PACKET_EOF: + case LIANA_PACKET_ERROR: al_array_foreach(vcr->tracks, i, track) { nn_packet_cache_send_packet(&track->cache, NULL); } #ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); #endif - log_debug("Received EOF."); - break; - case LIANA_PACKET_ERROR: - log_warn("Unhandled error packet."); + if (op == LIANA_PACKET_EOF) { + log_debug("Received EOF."); + } else if (op == LIANA_PACKET_ERROR) { + log_warn("Forcing EOF because we got an error packet."); + } break; default: log_warn("Erroneous packet."); |