diff options
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 2 | ||||
| -rw-r--r-- | src/liana/handlers/cdio_server.c | 8 | ||||
| -rw-r--r-- | src/liana/list.c | 2 | ||||
| -rw-r--r-- | src/liana/server.c | 21 | ||||
| -rw-r--r-- | src/liana/server.h | 4 | ||||
| -rw-r--r-- | src/liana/vcr.c | 14 |
6 files changed, 35 insertions, 16 deletions
diff --git a/src/liana/client.c b/src/liana/client.c index 2b62612..1e873c6 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -138,7 +138,7 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe nn_packet_read_str(packet, &handler); client->duration = nn_packet_read_u64(packet); collect_streams(client, packet); - if (client->streams.count == 0) { + if (!client->streams.count) { log_warn("Resource has no streams."); goto out; } diff --git a/src/liana/handlers/cdio_server.c b/src/liana/handlers/cdio_server.c index c67330d..65bd990 100644 --- a/src/liana/handlers/cdio_server.c +++ b/src/liana/handlers/cdio_server.c @@ -6,7 +6,7 @@ #include "cdio.h" -#define SECTORS_PER_PACKET 8 +#define SAMPLES_PER_PACKET 128 static bool cdio_server_init(struct lia_server_handler *handler, struct cch_handle *handle) { @@ -75,7 +75,8 @@ static void cdio_server_step(struct lia_server_handler *handler) 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; + size_t bps = camu_audio_format_bytes_per_sample(&cdio->fmt) * cdio->fmt.channel_count; + off_t size = bps * SAMPLES_PER_PACKET; 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) { @@ -87,7 +88,8 @@ static void cdio_server_write_packet(struct lia_server_handler *handler, struct 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); + al_assert(ret % bps == 0); + s32 sample_count = ret / bps; nn_packet_write_s32(packet, sample_count); cdio->buffer.size = ret; nn_packet_write_buffer(packet, &cdio->buffer); diff --git a/src/liana/list.c b/src/liana/list.c index d12caab..01c2ff3 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -712,7 +712,7 @@ static void handle_clear(struct lia_list *list) static void run_queue(struct lia_list *list) { if (!list->cmd) { - if (list->command_queue.count) { + if (list->command_queue.count > 0) { al_array_pop_at(list->command_queue, 0, list->cmd); } else { return; diff --git a/src/liana/server.c b/src/liana/server.c index 5065425..791fe0e 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -124,7 +124,7 @@ static void free_connection(struct lia_node_connection *conn) struct lia_node *node = conn->node; al_assert(conn->handler); conn->handler->free(&conn->handler); - cch_entry_return_handle(conn->node->entry, &conn->handle); + cch_entry_return_handle(node->entry, &conn->handle); nn_packet_pool_free(&conn->pool); bool removed = al_array_remove(node->connections, conn); al_assert(removed); @@ -272,31 +272,32 @@ static void signal_callback(void *userdata) nn_thread_join(&conn->thread); nn_signal_stop(&conn->signal); - al_array_remove(node->requests, conn); struct nn_packet_stream *stream = conn->stream; struct nn_packet *packet = conn->packet; conn->packet = NULL; + if (packet) { + al_array_remove(node->requests, conn); + } + if (!packet || conn->errored) { conn->handler->free(&conn->handler); cch_entry_return_handle(node->entry, &conn->handle); } - if (!packet) { - // Connection was closed before init was done. + if (!packet) { // Connection was closed before init was done. + al_free(conn); if (should_free_node(node)) { free_node(node); } nn_packet_stream_free(stream); al_free(stream); - al_free(conn); } else if (conn->errored) { // We must return the packet before disconnecting. nn_packet_stream_return_packet(stream, packet); - al_free(conn); - conn = NULL; demote_and_disconnect_stream(server, stream); + al_free(conn); } else { conn->id = get_incremental_id(server); al_array_push(node->connections, conn); @@ -337,10 +338,12 @@ static struct lia_node_connection *get_connection_from_id(struct lia_node *node, static void pre_init_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + struct lia_node *node = conn->node; struct nn_packet *packet = conn->packet; // Checked in signal_callback and will signal to cleanup the connection. conn->packet = NULL; nn_packet_stream_return_packet(stream, packet); + al_array_remove(node->requests, conn); } static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) @@ -439,6 +442,10 @@ static nn_thread_result NNWT_THREADCALL init_duration_thread(void *userdata) { struct lia_node *node = (struct lia_node *)userdata; nn_thread_set_name("liana_init_dur"); +#if 0 // This is really bad for large resources and is just for testing. + lia_prepare_visual_data(&node->handle, &node->visual); + cch_handle_seek(&node->handle, 0, SEEK_SET); +#endif if (!node->handler->init(node->handler, &node->handle)) { node->errored = true; } else { diff --git a/src/liana/server.h b/src/liana/server.h index bafe95b..ec6f82e 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -7,6 +7,8 @@ #include "../cache/entry.h" +#include "process.h" + //#define LIANA_SERVER_LOOP struct lia_node_connection { @@ -46,6 +48,8 @@ struct lia_node { struct cch_handle handle; struct nn_thread thread; struct nn_signal signal; + // @TODO: This should probably be on camu_resource instead. + struct lia_visual_data visual; }; struct lia_server { diff --git a/src/liana/vcr.c b/src/liana/vcr.c index f9dad7b..3a61e32 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -26,6 +26,7 @@ enum { enum { VCR_TRACK_RUNNING = 0, VCR_TRACK_STOPPED, + VCR_TRACK_ERRORED, VCR_TRACK_CLOSED }; @@ -131,8 +132,8 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) nn_thread_set_priority(NNWT_THREAD_SCHED_FIFO, 32); struct lia_vcr_track *track = (struct lia_vcr_track *)userdata; struct lia_vcr *vcr = track->vcr; - const char thread_name[16] = "\0"; // 16 = limit. - al_snprintf((char *)thread_name, sizeof(thread_name), "vcr:%hu_%d", vcr->node_id, track->stream->index); + char thread_name[16] = "\0"; // 16 = limit. + al_snprintf(thread_name, sizeof(thread_name), "vcr:%hu_%d", vcr->node_id, track->stream->index); nn_thread_set_name(thread_name); bool corked; @@ -177,12 +178,13 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) return_entire_cache(track); track->cache.disabled = true; nn_packet_cache_unlock(&track->cache); + // Let flush()/close() know we forcefully exited the thread. + track->state = VCR_TRACK_ERRORED; nn_mutex_unlock(&track->lock); return 0; } - if (!packet) { - // Wait on EOF. + if (!packet) { // Wait on EOF. track->state = VCR_TRACK_STOPPED; } @@ -410,6 +412,7 @@ void lia_vcr_set_buffered(struct lia_vcr_track *track) // Meaning track->lock will be held. void lia_vcr_cork(struct lia_vcr_track *track) { + al_assert(track->state != VCR_TRACK_ERRORED); track->state = VCR_TRACK_STOPPED; } @@ -422,6 +425,9 @@ void lia_vcr_uncork(struct lia_vcr_track *track) // being held by flush(). flush() sets the state to CLOSED and signals the // cond. Now nn_cond_is_waiting() is false at the point uncork() acquires the lock. nn_mutex_lock(&track->lock); + // If a buffer is running under MARK_LOW, it will be continuously trying to uncork() + // the track. Meaning an erroring handle_packet() could happen at the same time as + // an uncork(). Causing track->state to be ERRORED after we acquire the lock here. if (track->state != VCR_TRACK_STOPPED) { nn_mutex_unlock(&track->lock); return; |