summaryrefslogtreecommitdiff
path: root/src/server
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2026-09-14 08:57:42 -0400
committerAndrew Opalach <andrew@akon.city> 2026-09-14 08:57:42 -0400
commit8f208c26b6fa1a9f3372679c047cab559c06e26b (patch)
tree323d894d6ff8e1ed1445c40cb1e2f5d3cee5e8e8 /src/server
parentc66c7c64ebd16287b892f8a780cffcabafba3799 (diff)
downloadcamu-8f208c26b6fa1a9f3372679c047cab559c06e26b.tar.gz
camu-8f208c26b6fa1a9f3372679c047cab559c06e26b.tar.bz2
camu-8f208c26b6fa1a9f3372679c047cab559c06e26b.zip
Server-side fixes from DIRECT_MODE testing
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/server')
-rw-r--r--src/server/server.c101
-rw-r--r--src/server/server.h4
2 files changed, 68 insertions, 37 deletions
diff --git a/src/server/server.c b/src/server/server.c
index 513e22e..9ed6eab 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -114,10 +114,10 @@ static void handle_toggle_sink(struct camu_server *server, str *name, struct cam
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;
-
// @TODO: Return server-side IDs to the clients. This is important, for example, to
// identify which sink a list command is coming from.
+ struct camu_server *server = (struct camu_server *)userdata;
+
u8 op = nn_packet_read_u8(packet);
switch (op) {
case CAMU_NODE: {
@@ -163,6 +163,19 @@ static bool identify_callback(void *userdata, struct nn_rpc_connection *conn,
return true;
}
+static inline u64 adjust_ts_for_skew(struct camu_server_sink *sink, u64 ts)
+{
+ struct nn_skew *skew = &sink->conn->skew.s;
+ if (skew->ts == LIANA_TIMESTAMP_INVALID) {
+ return ts;
+ }
+ if (skew->direction == NNWT_SKEW_POSITIVE) {
+ return ts + skew->ts;
+ } else {
+ return ts - skew->ts;
+ }
+}
+
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;
@@ -179,7 +192,7 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
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);
+ nn_packet_write_u64(packet, adjust_ts_for_skew(sink, timing->at));
al_assert(timing->pos <= INT64_MAX); // FFmpeg.
nn_packet_write_u64(packet, timing->pos);
nn_packet_write_u8(packet, timing->pause);
@@ -188,8 +201,14 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
break;
}
case LIANA_SINK_UNSET: {
- struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET);
- nn_packet_write_u8(packet, op);
+ struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_UNSET);
+ nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
+ break;
+ }
+ case LIANA_SINK_SEQUENCE: {
+ struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEQUENCE);
+ nn_packet_write_u32(packet, entry->id);
+ nn_packet_write_s32(packet, sequence);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
@@ -197,7 +216,7 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
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_u64(packet, adjust_ts_for_skew(sink, timing->at));
nn_packet_write_u8(packet, timing->pause);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
@@ -206,7 +225,7 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
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, adjust_ts_for_skew(sink, timing->at));
nn_packet_write_u64(packet, timing->pos);
nn_packet_write_u32(packet, entry->reset_token);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
@@ -426,7 +445,7 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v
case LIANA_ENTRY_LOADING:
if (resource->load == LIANA_ENTRY_PREPARED) {
resource->load = LIANA_ENTRY_LOADING;
- lia_node_get_duration(resource->node);
+ lia_node_probe_duration(resource->node);
}
// resource->pending will only ever be read from the event loop thread.
al_array_push(resource->pending, entry);
@@ -441,7 +460,9 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v
*(u64 *)opaque = resource->duration;
break;
case LIANA_REF_ENTRY:
- resource->ref = RESOURCE_MAX_AGE;
+ if (!resource->ref) {
+ resource->ref = RESOURCE_MAX_AGE;
+ }
break;
case LIANA_UNREF_ENTRY:
if (!resource->ref || --resource->ref) {
@@ -575,9 +596,10 @@ static bool client_command_callback(void *userdata, struct nn_rpc_connection *co
nn_packet_read_str(packet, &name);
struct camu_server_sink *sink = get_sink_by_name(server, &name);
if (!sink) goto out;
- nn_packet_read_str(packet, &name); // list name.
+ str list_name;
+ nn_packet_read_str(packet, &list_name);
bool enable = nn_packet_read_bool(packet);
- handle_toggle_sink(server, &name, sink, enable);
+ handle_toggle_sink(server, &list_name, sink, enable);
break;
}
#ifdef CAMU_HAVE_PORTAL
@@ -661,21 +683,19 @@ static u8 parse_resource_type(str *line)
}
}
-static void handle_add_command(struct camu_server *server, struct lia_list *list, struct nn_packet *packet)
+static void handle_add_command(struct camu_server *server, struct lia_list *list, str *line)
{
struct camu_resource *resource = al_alloc_object(struct camu_resource);
- str line;
- nn_packet_read_str(packet, &line);
- switch (parse_resource_type(&line)) {
+ switch (parse_resource_type(line)) {
case CAMU_RESOURCE_FILE: {
- al_str_clone(&resource->uri, &line);
+ al_str_clone(&resource->uri, line);
resource->type = CAMU_RESOURCE_FILE;
resource->load = LIANA_ENTRY_UNLOADED;
break;
}
#ifdef NAUNET_HAS_CURL
case CAMU_RESOURCE_HTTP: {
- al_str_clone(&resource->uri, &line);
+ al_str_clone(&resource->uri, line);
resource->type = CAMU_RESOURCE_HTTP;
resource->load = LIANA_ENTRY_UNLOADED;
break;
@@ -684,9 +704,9 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
#ifdef CACHE_HAVE_CDIO
case CAMU_RESOURCE_CDIO: {
u32 track = 0;
- if (line.length > 7) {
+ if (line->length > 7) {
bool error;
- s64 index = al_str_to_long(&al_str_substr(&line, 7, line.length), 10, &error);
+ s64 index = al_str_to_long(&al_str_substr(line, 7, line->length), 10, &error);
if (!error && index > 0) {
track = (u32)index - 1;
}
@@ -703,11 +723,11 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
resource->type = CAMU_RESOURCE_SIMPLE_SEARCH;
resource->load = LIANA_ENTRY_PREPARING;
str query;
- if (al_str_at(&line, 0) == ';') { // search.
- al_str_clone(&query, &al_str_substr(&line, 1, line.length));
+ if (al_str_at(line, 0) == ';') { // search.
+ al_str_clone(&query, &al_str_substr(line, 1, line->length));
} else {
al_str_from(&query, "link:");
- al_str_cat(&query, &line);
+ al_str_cat(&query, line);
}
log_info("Processing search request: %.*s.", al_str_x(&query));
camu_portal_create_search(&server->bridge, &al_str_c("youtube"),
@@ -723,15 +743,14 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
resource->node = NULL;
resource->duration = LIANA_TIMESTAMP_INVALID;
al_array_init(resource->pending);
- lia_list_add(list, &line, resource, resource->duration, resource->load);
+ lia_list_add(list, line, resource, resource->duration, resource->load);
}
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;
// @TODO: Don't accept commands from non-identified connections.
- (void)conn;
+ struct camu_server *server = (struct camu_server *)userdata;
(void)rpacket;
str name;
@@ -743,7 +762,9 @@ static bool list_action_callback(void *userdata, struct nn_rpc_connection *conn,
u8 op = nn_packet_read_u8(packet);
switch (op) {
case CAMU_LIST_ADD: {
- handle_add_command(server, list, packet);
+ str line;
+ nn_packet_read_str(packet, &line);
+ handle_add_command(server, list, &line);
break;
}
case CAMU_LIST_SKIP: {
@@ -799,10 +820,18 @@ static struct nn_rpc_command commands[] = {
{ .op = CAMU_SERVER_LIST_ACTION, .callback = list_action_callback, .userdata = NULL }
};
+// @TODO: Cleanup dormant connections.
static void connection_callback(void *userdata, struct nn_rpc_connection *conn)
{
- // @TODO: Cleanup dormant connections.
- (void)userdata;
+ struct camu_server *server = (struct camu_server *)userdata;
+ (void)server;
+ (void)conn;
+}
+
+static void ready_callback(void *userdata, struct nn_rpc_connection *conn)
+{
+ struct camu_server *server = (struct camu_server *)userdata;
+ (void)server;
(void)conn;
}
@@ -879,7 +908,7 @@ static bool multiplex_callback(void *userdata, u8 id, struct nn_packet_stream *s
return true;
}
-void camu_server_init(struct camu_server *server, struct nn_event_loop *loop)
+void camu_server_init(struct camu_server *server, struct nn_event_loop *loop, bool local)
{
server->loop = loop;
server->addr = AL_STR_EMPTY;
@@ -895,11 +924,15 @@ void camu_server_init(struct camu_server *server, struct nn_event_loop *loop)
list->userdata = server;
al_array_push(server->lists, list);
- nn_rpc_init(&server->server, server->loop, connection_callback, connection_closed_callback, server);
+ nn_rpc_init(&server->server, server->loop, connection_callback, ready_callback, connection_closed_callback, server);
for (u32 i = 0; i < ARRAY_SIZE(commands); i++) {
commands[i].userdata = server;
nn_rpc_add_command(&server->server, &commands[i]);
}
+ // Set the rpc server to query the difference between this server's clock and
+ // any clients that connect to the rpc's clock (skew). Then, offset the timestamps
+ // that come from the list with adjust_ts_for_skew().
+ nn_rpc_query_clock_skews(&server->server, !local);
lia_server_init(&server->data.server, server->loop);
al_array_init(server->data.resources);
@@ -953,9 +986,7 @@ void camu_server_close(struct camu_server *server)
camu_portal_close(&server->bridge);
#endif
-#ifdef CAMU_DIRECT_MODE
- nn_multiplex_direct_close();
-#else
+#ifndef CAMU_DIRECT_MODE
nn_multiplex_socket_close(&server->multi);
#endif
}
@@ -1021,9 +1052,9 @@ void camu_server_free(struct camu_server *server)
al_str_free(&server->addr);
}
-void camu_server_local_add(struct camu_server *server, struct nn_packet *packet)
+void camu_server_local_add(struct camu_server *server, str *line)
{
- handle_add_command(server, al_array_last(server->lists), packet);
+ handle_add_command(server, al_array_last(server->lists), line);
}
void camu_server_local_start_at(struct camu_server *server, s32 index)
diff --git a/src/server/server.h b/src/server/server.h
index 3112679..d170ab6 100644
--- a/src/server/server.h
+++ b/src/server/server.h
@@ -49,12 +49,12 @@ struct camu_server {
void *userdata;
};
-void camu_server_init(struct camu_server *server, struct nn_event_loop *loop);
+void camu_server_init(struct camu_server *server, struct nn_event_loop *loop, bool local);
bool camu_server_listen(struct camu_server *server, u8 type, str *addr, u16 port);
void camu_server_bind_direct(struct camu_server *server);
void camu_server_close(struct camu_server *server);
void camu_server_free(struct camu_server *server);
// Local compat.
-void camu_server_local_add(struct camu_server *server, struct nn_packet *packet);
+void camu_server_local_add(struct camu_server *server, str *line);
void camu_server_local_start_at(struct camu_server *server, s32 index);