From 4ccfbed07834e83007e66934e06055d2b548637e Mon Sep 17 00:00:00 2001 From: Cococry Date: Tue, 28 Jul 2026 17:48:51 +0200 Subject: [PATCH] feat: EVENT_BATCH envelopes & offline event delivery --- faithd/src/application/conversation.h | 9 + faithd/src/application/user.h | 9 + faithd/src/auth/device_link.c | 53 ++--- faithd/src/auth/handshake.c | 31 +-- faithd/src/auth/handshake.h | 2 +- faithd/src/auth/structs.h | 36 +++- faithd/src/client/client.c | 128 ++++++++---- faithd/src/codec/commands.h | 3 +- faithd/src/codec/envelopes.h | 3 +- faithd/src/codec/events.h | 22 +- faithd/src/codec/helpers.h | 88 +++++--- faithd/src/codec/protocol.c | 284 +++++++++++++++++++++++--- faithd/src/codec/protocol.h | 15 ++ faithd/src/commands/conversation.c | 30 ++- faithd/src/commands/dispatch.c | 4 +- faithd/src/core/core.h | 2 +- faithd/src/delivery/event_inbox.c | 70 ++++++- faithd/src/delivery/event_inbox.h | 10 +- faithd/src/delivery/events.c | 159 +++++++++++--- faithd/src/delivery/events.h | 7 +- faithd/src/delivery/routing.c | 59 +++--- faithd/src/server/client_lifecycle.c | 19 +- faithd/src/server/server.h | 5 +- faithd/src/server/sess_registry.c | 2 + faithd/src/server/sess_registry.h | 13 +- 25 files changed, 814 insertions(+), 249 deletions(-) create mode 100644 faithd/src/application/conversation.h create mode 100644 faithd/src/application/user.h diff --git a/faithd/src/application/conversation.h b/faithd/src/application/conversation.h new file mode 100644 index 0000000..3ada01d --- /dev/null +++ b/faithd/src/application/conversation.h @@ -0,0 +1,9 @@ +#pragma once + +#include + +#define FAITH_CONVERSATION_ID_SIZE 16 + +typedef struct { + uint8_t bytes[FAITH_CONVERSATION_ID_SIZE]; +} faith_conversation_id_t; diff --git a/faithd/src/application/user.h b/faithd/src/application/user.h new file mode 100644 index 0000000..a6e7a8a --- /dev/null +++ b/faithd/src/application/user.h @@ -0,0 +1,9 @@ +#pragma once + +#include "../auth/structs.h" + +typedef struct { + faith_auth_id_t auth_id; + faith_device_id_t device_id; + uint8_t public_key[FAITH_ED25519_PUBLIC_KEY_SIZE]; +} application_user_identity_t; diff --git a/faithd/src/auth/device_link.c b/faithd/src/auth/device_link.c index fb57ab2..d219c00 100644 --- a/faithd/src/auth/device_link.c +++ b/faithd/src/auth/device_link.c @@ -6,6 +6,7 @@ #include "../codec/signatures.h" #include "../delivery/routing.h" +#include "structs.h" #define NOB_IMPLEMENTATION #include "../../third_party/nob.h" @@ -20,7 +21,7 @@ static faith_status_code_t send_device_auth_response_failed(server_state_t *s, faith_envelope_t failed_envl = { .type = FAITH_ENVELOPE_DEVICE_AUTH_RESPONSE_FAILED, - .recipient_id = cl->auth_id, + .recipient_id = cl->ident.auth_id, .body = NULL, .body_size = 0, }; @@ -79,8 +80,9 @@ device_link_handle_device_response(server_state_t *s, client_conn_t *cl, if (!cl->authorized) { char auth_id_hex[33]; char device_id_hex[33]; - _FH_CHECK_RETURN(faith_id128_to_hex(cl->auth_id.bytes, auth_id_hex)); - _FH_CHECK_RETURN(faith_id128_to_hex(cl->device_id.bytes, device_id_hex)); + _FH_CHECK_RETURN(faith_id128_to_hex(cl->ident.auth_id.bytes, auth_id_hex)); + _FH_CHECK_RETURN( + faith_id128_to_hex(cl->ident.device_id.bytes, device_id_hex)); nob_log(ERROR, "[client=%" PRIu64 " fd=%i] Server got unauthorized %s" @@ -160,7 +162,7 @@ device_link_handle_device_response(server_state_t *s, client_conn_t *cl, memcpy(sign_msg.code, req->code, sizeof(req->code)); sign_msg.expires_at_ms = req->expires_at_ms; - sign_msg.device_id_responding = cl->device_id; + sign_msg.device_id_responding = cl->ident.device_id; sign_msg.type = response_envl->type == FAITH_ENVELOPE_DEVICE_AUTH_APPROVE ? FAITH_DEVICE_LINK_APPROVE @@ -186,11 +188,12 @@ device_link_handle_device_response(server_state_t *s, client_conn_t *cl, client_device_session_data_t *sess = NULL; { - _FH_CHECK( - sess_registry_get_session(&s->rt, &cl->auth_id, &cl->device_id, &sess)); + _FH_CHECK(sess_registry_get_session(&s->rt, &cl->ident.auth_id, + &cl->ident.device_id, &sess)); /* We specifically need routing_get_session() to return FAITH_OK. This is - * returned only if auth_id> is a registered client_route_user_t - * and device_id> is a registered client_route_device_t of that user. + * returned only if ident.auth_id> is a registered client_route_user_t + * and ident.device_id> is a registered client_route_device_t of that + * user. * */ if (_fh_rc != FAITH_OK) { _fh_result = _fh_rc; @@ -202,10 +205,10 @@ device_link_handle_device_response(server_state_t *s, client_conn_t *cl, } } - /* This means cl->auth_id is registered but cl->device_id is not, effectively - * telling us that the client connection is not yet authorized. Because we - * checked cl->authorized above, this should never happen with correct - * behaviour.*/ + /* This means cl->ident.auth_id is registered but cl->ident.device_id is not, + * effectively telling us that the client connection is not yet authorized. + * Because we checked cl->authorized above, this should never happen with + * correct behaviour.*/ if (!sess) { nob_log(ERROR, "[client=%" PRIu64 @@ -265,11 +268,12 @@ device_link_handle_device_response(server_state_t *s, client_conn_t *cl, faith_envelope_t ack_envl = {0}; ack_envl.type = FAITH_ENVELOPE_DEVICE_AUTH_RESPONSE_ACK; _FH_CHECK_RETURN( - delivery_route_envelope_to_auth_id(s, cl, &cl->auth_id, &ack_envl)); + delivery_route_envelope_to_auth_id(s, cl, &cl->ident.auth_id, &ack_envl)); faith_status_code_t device_loop_rc = FAITH_OK; - _FH_FOR_EACH_AUTH_DEVICE(s, &cl->auth_id, recipient, device_loop_rc, - { device_link_remove_request(cl); }); + _FH_FOR_EACH_AUTH_DEVICE_CONNECTION(s, &cl->ident.auth_id, recipient, + device_loop_rc, + { device_link_remove_request(cl); }); return device_loop_rc == FAITH_OK ? rc : device_loop_rc; @@ -352,7 +356,7 @@ device_link_new_device(server_state_t *s, client_conn_t *cl, /* Send device authorization request to every already registered device for that auth ID */ faith_status_code_t device_loop_rc = FAITH_OK; - _FH_FOR_EACH_AUTH_DEVICE( + _FH_FOR_EACH_AUTH_DEVICE_CONNECTION( s, ¶ms->sender_auth_id, authorized_cl, device_loop_rc, { faith_envl_stc_device_link_req_t *req = NULL; _FH_CHECK(device_link_queue_request( @@ -384,9 +388,8 @@ device_link_new_device(server_state_t *s, client_conn_t *cl, /* Send DEVICE_AUTH_PENDING to the connection that requested the * new device */ _FH_CHECK_RETURN(auth_queue_auth_pending(s, cl)); - /* temporarily assign auth_id> for disconnection purposes later. this - * does not mean that the client is authorized. */ - cl->auth_id = params->sender_auth_id; + + cl->pending_auth_id = params->sender_auth_id; server_set_client_state(s, cl, CLIENT_WAIT_FOR_DEVICE_LINK_RESPONSE); return FAITH_OK; @@ -405,7 +408,7 @@ device_link_queue_request_cancellation(server_state_t *s, &requesting_cl->temp_handshake_params; faith_status_code_t device_loop_rc = FAITH_OK; - _FH_FOR_EACH_AUTH_DEVICE( + _FH_FOR_EACH_AUTH_DEVICE_CONNECTION( s, ¶ms->sender_auth_id, authorized_cl, device_loop_rc, { if (authorized_cl->pending_device_link_conn != requesting_cl) continue; @@ -413,7 +416,7 @@ device_link_queue_request_cancellation(server_state_t *s, /* Send DEVICE_LINK_CANCELLED to the authorized device */ faith_envelope_t envl = {0}; envl.type = FAITH_ENVELOPE_DEVICE_LINK_CANCELLED; - envl.recipient_id = authorized_cl->auth_id; + envl.recipient_id = authorized_cl->ident.auth_id; _FH_CHECK(server_queue_envelope_or_mark_dead(s, authorized_cl, &envl)); if (_fh_rc != FAITH_OK && device_loop_rc == FAITH_OK) @@ -446,9 +449,9 @@ faith_status_code_t device_link_queue_request( char auth_id_hex[33]; char device_id_hex[33]; _FH_CHECK_RETURN( - faith_id128_to_hex(recipient_cl->auth_id.bytes, auth_id_hex)); + faith_id128_to_hex(recipient_cl->ident.auth_id.bytes, auth_id_hex)); _FH_CHECK_RETURN( - faith_id128_to_hex(recipient_cl->device_id.bytes, device_id_hex)); + faith_id128_to_hex(recipient_cl->ident.device_id.bytes, device_id_hex)); nob_log(ERROR, "Not sending device link request to " "device with device_id: %s (auth_id: %s). Client connection is " @@ -462,9 +465,9 @@ faith_status_code_t device_link_queue_request( char auth_id_hex[33]; char device_id_hex[33]; _FH_CHECK_RETURN( - faith_id128_to_hex(recipient_cl->auth_id.bytes, auth_id_hex)); + faith_id128_to_hex(recipient_cl->ident.auth_id.bytes, auth_id_hex)); _FH_CHECK_RETURN( - faith_id128_to_hex(recipient_cl->device_id.bytes, device_id_hex)); + faith_id128_to_hex(recipient_cl->ident.device_id.bytes, device_id_hex)); nob_log(ERROR, "Not sending device link request to " "device with device_id: %s (auth_id: %s). Another device link " diff --git a/faithd/src/auth/handshake.c b/faithd/src/auth/handshake.c index d182b98..f9b3a12 100644 --- a/faithd/src/auth/handshake.c +++ b/faithd/src/auth/handshake.c @@ -7,6 +7,8 @@ #include "../codec/protocol.h" #include "../codec/signatures.h" +#include "../delivery/events.h" + faith_status_code_t auth_handle_hello(server_state_t *s, client_conn_t *cl, const faith_envelope_t *hello_envl) { // HELLO { @@ -252,10 +254,13 @@ faith_status_code_t auth_handle_challenge_response( return FAITH_OK; } - _FH_CHECK(auth_authorize_client(s, cl, ¶ms->sender_auth_id, - ¶ms->device_id, verification_public_key, - sess == NULL)); - return _fh_rc; + _FH_CHECK_RETURN( + auth_authorize_client(s, cl, ¶ms->sender_auth_id, ¶ms->device_id, + verification_public_key, sess == NULL)); + + _FH_CHECK_RETURN(delivery_queue_pending_events(s, cl)); + + return FAITH_OK; reject: { /* =============================== */ @@ -279,7 +284,7 @@ reject: { faith_status_code_t auth_authorize_client( server_state_t *s, client_conn_t *cl, const faith_auth_id_t *auth_id, const faith_device_id_t *device_id, - uint8_t public_key[FAITH_ED25519_PUBLIC_KEY_SIZE], int register_session) { + uint8_t public_key[FAITH_ED25519_PUBLIC_KEY_SIZE], bool register_session) { if (!s || !cl || cl->closing || !auth_id || !device_id || !public_key) return FAITH_ERR_INVALID; @@ -291,13 +296,14 @@ faith_status_code_t auth_authorize_client( if (register_session) { /* Register client session */ _FH_CHECK_RETURN(sess_registry_register_session( - &s->rt, &cl->auth_id, &cl->device_id, cl, public_key)); + &s->rt, &cl->ident.auth_id, &cl->ident.device_id, cl, public_key)); } char cl_auth_id_hex[33]; char cl_device_id_hex[33]; - _FH_CHECK_RETURN(faith_id128_to_hex(cl->auth_id.bytes, cl_auth_id_hex)); - _FH_CHECK_RETURN(faith_id128_to_hex(cl->device_id.bytes, cl_device_id_hex)); + _FH_CHECK_RETURN(faith_id128_to_hex(cl->ident.auth_id.bytes, cl_auth_id_hex)); + _FH_CHECK_RETURN( + faith_id128_to_hex(cl->ident.device_id.bytes, cl_device_id_hex)); nob_log(INFO, "[client=%" PRIu64 " fd=%i] Client passed authorization for " @@ -324,8 +330,8 @@ auth_handshake_complete(server_state_t *s, client_conn_t *cl, _FH_CHECK_RETURN(server_queue_envelope_or_mark_dead(s, cl, &hello_ok_envl)); - cl->auth_id = *sender_id; - cl->device_id = *device_id; + cl->ident.auth_id = *sender_id; + cl->ident.device_id = *device_id; server_set_client_state(s, cl, CLIENT_OPEN); @@ -334,8 +340,9 @@ auth_handshake_complete(server_state_t *s, client_conn_t *cl, faith_status_code_t _fh_result = FAITH_OK; - _FH_CHECK_RETURN(faith_id128_to_hex(cl->auth_id.bytes, auth_id_hex)); - _FH_CHECK_RETURN(faith_id128_to_hex(cl->device_id.bytes, device_id_hex)); + _FH_CHECK_RETURN(faith_id128_to_hex(cl->ident.auth_id.bytes, auth_id_hex)); + _FH_CHECK_RETURN( + faith_id128_to_hex(cl->ident.device_id.bytes, device_id_hex)); nob_log(INFO, "[client=%" PRIu64 diff --git a/faithd/src/auth/handshake.h b/faithd/src/auth/handshake.h index e6ca5a3..2c4b078 100644 --- a/faithd/src/auth/handshake.h +++ b/faithd/src/auth/handshake.h @@ -14,7 +14,7 @@ auth_handle_challenge_response(server_state_t *s, client_conn_t *cl, faith_status_code_t auth_authorize_client( server_state_t *s, client_conn_t *cl, const faith_auth_id_t *auth_id, const faith_device_id_t *device_id, - uint8_t public_key[FAITH_ED25519_PUBLIC_KEY_SIZE], int register_session); + uint8_t public_key[FAITH_ED25519_PUBLIC_KEY_SIZE], bool register_session); faith_status_code_t auth_handshake_complete(server_state_t *s, client_conn_t *cl, diff --git a/faithd/src/auth/structs.h b/faithd/src/auth/structs.h index 096930d..149161e 100644 --- a/faithd/src/auth/structs.h +++ b/faithd/src/auth/structs.h @@ -6,7 +6,8 @@ #include "../core/core.h" #include "../core/crypto.h" -#define _FH_FOR_EACH_AUTH_DEVICE(SERVER, AUTH_ID, RECIPIENT, STATUS_OUT, BODY) \ +#define _FH_FOR_EACH_AUTH_DEVICE_CONNECTION(SERVER, AUTH_ID, RECIPIENT, \ + STATUS_OUT, BODY) \ do { \ server_state_t *_fh_iter_server = (SERVER); \ const faith_auth_id_t *_fh_iter_auth_id = (AUTH_ID); \ @@ -38,6 +39,39 @@ } \ } while (0) +#define _FH_FOR_EACH_AUTH_DEVICE_SESSION(SERVER, AUTH_ID, DEVICE_SESSION, \ + STATUS_OUT, BODY) \ + do { \ + server_state_t *_fh_iter_server = (SERVER); \ + const faith_auth_id_t *_fh_iter_auth_id = (AUTH_ID); \ + client_session_device_t *_fh_iter_devices = NULL; \ + \ + if (!_fh_iter_server || !_fh_iter_auth_id) { \ + (STATUS_OUT) = FAITH_ERR_INVALID; \ + break; \ + } \ + \ + (STATUS_OUT) = sess_registry_get_devices( \ + &_fh_iter_server->rt, _fh_iter_auth_id, &_fh_iter_devices); \ + \ + if ((STATUS_OUT) == FAITH_ERR_NOT_FOUND) { \ + (STATUS_OUT) = FAITH_OK; \ + } else if ((STATUS_OUT) == FAITH_OK) { \ + ptrdiff_t _fh_iter_count = hmlen(_fh_iter_devices); \ + \ + for (ptrdiff_t _fh_iter_i = 0; _fh_iter_i < _fh_iter_count; \ + ++_fh_iter_i) { \ + if (!_fh_iter_devices[_fh_iter_i].value) \ + continue; \ + \ + client_device_session_data_t *(DEVICE_SESSION) = \ + _fh_iter_devices[_fh_iter_i].value; \ + \ + BODY \ + } \ + } \ + } while (0) + #define FAITH_AUTH_ID_SIZE 16 #define FAITH_DEVICE_ID_SIZE 16 diff --git a/faithd/src/client/client.c b/faithd/src/client/client.c index 5269edf..2faef4f 100644 --- a/faithd/src/client/client.c +++ b/faithd/src/client/client.c @@ -66,7 +66,7 @@ typedef struct { } pending_command_entry_t; typedef struct { - bool occupied; + bool occupied; faith_envl_stc_event_t event; } buffered_event_slot_t; @@ -303,7 +303,7 @@ static faith_status_code_t read_frame_sync(SSL *ssl, faith_frame_t *out) { "Failed to read frame; Frame is too large, " "frame_size=%i MAX_FRAME_LEN=%i", (int32_t)frame_size, (int32_t)FAITH_MAX_FRAME_LEN); - return FAITH_ERR_FRAME_TOO_LARGE; + return FAITH_ERR_TOO_LARGE; } _FH_CHECK_RETURN(read_bytes_sync(ssl, hdr_buf, sizeof(hdr_buf))); @@ -322,7 +322,7 @@ static faith_status_code_t read_frame_sync(SSL *ssl, faith_frame_t *out) { "frame_size=%u MAX_FRAME_LEN=%i", out->payload_size, (int32_t)FAITH_MAX_FRAME_LEN); - return FAITH_ERR_FRAME_TOO_LARGE; + return FAITH_ERR_TOO_LARGE; } if (out->payload_size == 0) @@ -936,9 +936,8 @@ client_handle_disconnect(faith_client_t *client, const faith_envelope_t *envl) { return FAITH_OK; } -static faith_status_code_t -client_send_ack_until(faith_client_t *client, uint64_t seq_num, - faith_event_codec_type_t event_type) { +static faith_status_code_t client_send_ack_until(faith_client_t *client, + uint64_t seq_num) { if (!client || seq_num == UINT64_MAX) return FAITH_ERR_INVALID; @@ -948,7 +947,6 @@ client_send_ack_until(faith_client_t *client, uint64_t seq_num, faith_envl_cts_event_ack_t ack = {0}; ack.seq_num = seq_num; - ack.type = event_type; _FH_CHECK_RETURN( faith_encode_event_ack_body(body, &body_size, sizeof(body), &ack)); @@ -985,6 +983,8 @@ client_dispatch_event(faith_client_t *client, if (!client || !event) return FAITH_ERR_INVALID; + faith_status_code_t _fh_result = FAITH_OK; + switch (event->type) { case FAITH_EVENT_CONVERSATION_CREATED: _FH_CHECK_RETURN(client_handle_event_conversation_created(client, event)); @@ -995,28 +995,31 @@ client_dispatch_event(faith_client_t *client, client->ev_last_dispatched_seq = event->seq_num; - return FAITH_OK; + return _fh_result; } -static faith_status_code_t client_buffer_event(faith_client_t* client, const faith_envl_stc_event_t* event) { - if(!client || !event) return FAITH_ERR_INVALID; +static faith_status_code_t +client_buffer_event(faith_client_t *client, + const faith_envl_stc_event_t *event) { + if (!client || !event) + return FAITH_ERR_INVALID; - size_t index = event->seq_num % EV_REORDER_WINDOW_SIZE; - buffered_event_slot_t* slot = &client->ev_reorder_window[index]; + size_t index = event->seq_num % EV_REORDER_WINDOW_SIZE; + buffered_event_slot_t *slot = &client->ev_reorder_window[index]; - if(slot->occupied) { + if (slot->occupied) { /* duplicate buffered event */ - if(event->seq_num == slot->event.seq_num) return FAITH_OK; + if (event->seq_num == slot->event.seq_num) + return FAITH_OK; /* invalid state */ return FAITH_ERR_INVALID; } slot->occupied = true; - slot->event = *event; - uint8_t* data_copy = NULL; + uint8_t *data_copy = NULL; - if(event->data_size > 0) { + if (event->data_size > 0) { data_copy = malloc(event->data_size); if (!data_copy) return FAITH_ERR_NOMEM; @@ -1032,8 +1035,9 @@ static faith_status_code_t client_buffer_event(faith_client_t* client, const fai } static faith_status_code_t client_drain_reorder_buffer(faith_client_t *client) { - for(;;) { - if(client->ev_last_dispatched_seq == UINT64_MAX) return FAITH_ERR_OVERFLOW; + for (;;) { + if (client->ev_last_dispatched_seq == UINT64_MAX) + return FAITH_ERR_OVERFLOW; uint64_t next_seq = client->ev_last_dispatched_seq + 1; size_t index = next_seq % EV_REORDER_WINDOW_SIZE; @@ -1046,40 +1050,40 @@ static faith_status_code_t client_drain_reorder_buffer(faith_client_t *client) { /* advances of the client */ _FH_CHECK_RETURN(client_dispatch_event(client, &slot->event)); - free(slot->event.data); - slot->event.data = NULL; - slot->event.data_size = 0; - slot->occupied = false; } } static faith_status_code_t -client_ack_event(faith_client_t *client, const faith_envl_stc_event_t *event) { - if (!client || !event) +client_process_event(faith_client_t *client, + const faith_envl_stc_event_t *event, bool *o_saw_duplicate, + bool *o_made_progress) { + if (!client || !event || !o_saw_duplicate || !o_made_progress) return FAITH_ERR_INVALID; + if (event->seq_num == UINT64_MAX) { + nob_log(ERROR, "[client] Received event (%s) with invalid sequence number.", + faith_event_codec_type_name(event->type)); + return FAITH_ERR_INVALID; + } + /* duplicate event */ if (client->ev_last_dispatched_seq != UINT64_MAX && event->seq_num <= client->ev_last_dispatched_seq) { - /* resend ACKs for all prior events up until the last successfully - * processed/dispatched event */ - _FH_CHECK(client_send_ack_until(client, client->ev_last_dispatched_seq, - event->type)); - return _fh_rc; + + *o_saw_duplicate = true; + return FAITH_OK; } uint64_t expected_seq = client->ev_last_dispatched_seq == UINT64_MAX ? 0 : client->ev_last_dispatched_seq + 1; - if(event->seq_num == expected_seq) { + if (event->seq_num == expected_seq) { _FH_CHECK_RETURN(client_dispatch_event(client, event)); _FH_CHECK_RETURN(client_drain_reorder_buffer(client)); - _FH_CHECK_RETURN(client_send_ack_until( - client, client->ev_last_dispatched_seq, event->type)); - + *o_made_progress = true; return FAITH_OK; } @@ -1095,19 +1099,56 @@ static faith_status_code_t client_handle_event(faith_client_t *client, faith_envl_stc_event_t event = {0}; - if (event.seq_num == UINT64_MAX) { - nob_log(ERROR, "[client] Received event (%s) with invalid sequence number.", - faith_event_codec_type_name(event.type)); - return FAITH_ERR_INVALID; - } - _FH_CHECK_RETURN( faith_decode_event_body(envl->body, envl->body_size, &event)); - _FH_CHECK_RETURN(client_ack_event(client, &event)); + bool saw_duplicate = false; + bool made_progress = false; + _FH_CHECK_RETURN( + client_process_event(client, &event, &saw_duplicate, &made_progress)); + + if (saw_duplicate || made_progress) { + _FH_CHECK_RETURN( + client_send_ack_until(client, client->ev_last_dispatched_seq)); + } + + free(event.data); return FAITH_OK; } + +static faith_status_code_t +client_handle_event_batch(faith_client_t *client, + const faith_envelope_t *envl) { + if (!client || !envl) + return FAITH_ERR_INVALID; + + faith_status_code_t _fh_result = FAITH_OK; + faith_envl_stc_event_batch_t batch = {0}; + _FH_CHECK_DEFER( + faith_decode_event_batch_body(envl->body, envl->body_size, &batch)); + + bool saw_duplicate = false; + bool made_progress = false; + for (uint16_t i = 0; i < batch.n_events; i++) { + faith_envl_stc_event_t *event = &batch.events[i]; + _FH_CHECK_DEFER( + client_process_event(client, event, &saw_duplicate, &made_progress)); + } + if ((made_progress || saw_duplicate) && + client->ev_last_dispatched_seq != UINT64_MAX) { + _FH_CHECK_DEFER( + client_send_ack_until(client, client->ev_last_dispatched_seq)); + } + +defer: + for (size_t i = 0; i < batch.n_events; i++) { + free(batch.events[i].data); + } + arrfree(batch.events); + return _fh_result; +} + static faith_status_code_t client_handle_command_result(faith_client_t *client, const faith_envelope_t *envl) { @@ -1339,6 +1380,9 @@ static faith_status_code_t client_handle_envelope(faith_client_t *client, case FAITH_ENVELOPE_EVENT: _FH_CHECK_RETURN(client_handle_event(client, &envl)); break; + case FAITH_ENVELOPE_EVENT_BATCH: + _FH_CHECK_RETURN(client_handle_event_batch(client, &envl)); + break; default: break; } @@ -1852,8 +1896,6 @@ client_new_identity(client_side_identity_t *o_ident) { /* Generate 128 bit random device & auth identities */ - /*client_id_from_hex("9379839402f90a5aa848b418953cecd2", &o_ident->auth_id);*/ - _FH_CHECK_RETURN(faith_random_bytes(o_ident->auth_id.bytes, sizeof(o_ident->auth_id.bytes))); diff --git a/faithd/src/codec/commands.h b/faithd/src/codec/commands.h index d9c0185..08b25cd 100644 --- a/faithd/src/codec/commands.h +++ b/faithd/src/codec/commands.h @@ -31,7 +31,8 @@ typedef enum { X(FAITH_COMMAND_ERR_UNAUTHORIZED, 1) \ X(FAITH_COMMAND_ERR_BAD_COMMAND, 2) \ X(FAITH_COMMAND_ERR_TIMED_OUT, 3) \ - X(FAITH_COMMAND_ERR_INTERNAL_ERROR, 4) + X(FAITH_COMMAND_ERR_INTERNAL_ERROR, 4) \ + X(FAITH_COMMAND_ERR_USER_NOT_FOUND, 5) typedef enum { #define X(name, value) name = value, diff --git a/faithd/src/codec/envelopes.h b/faithd/src/codec/envelopes.h index 7714941..c9b1b62 100644 --- a/faithd/src/codec/envelopes.h +++ b/faithd/src/codec/envelopes.h @@ -22,7 +22,8 @@ X(FAITH_ENVELOPE_COMMAND, 16) \ X(FAITH_ENVELOPE_COMMAND_RESULT, 17) \ X(FAITH_ENVELOPE_EVENT, 18) \ - X(FAITH_ENVELOPE_EVENT_ACK, 19) + X(FAITH_ENVELOPE_EVENT_BATCH, 19) \ + X(FAITH_ENVELOPE_EVENT_ACK, 20) typedef enum { #define X(name, value) name = value, diff --git a/faithd/src/codec/events.h b/faithd/src/codec/events.h index 1f7f866..38b2893 100644 --- a/faithd/src/codec/events.h +++ b/faithd/src/codec/events.h @@ -1,5 +1,6 @@ #pragma once +#include "../application/conversation.h" #include "../auth/structs.h" #include "core.h" #include "helpers.h" @@ -10,11 +11,17 @@ sizeof(faith_body_size_t) /* data size */ \ ) +#define FAITH_ENVL_STC_EVENT_BATCH_BODY_SIZE_FIXED \ + _FAITH_BODY_SIZE(sizeof(uint16_t) /* number of events*/ + \ + sizeof(uint32_t) /* events data size */) + #define FAITH_ENVL_CTS_EVENT_ACK_BODY_SIZE \ _FAITH_BODY_SIZE(sizeof(uint64_t) /* sequence number */) #define FAITH_EVENTS_CODEC(X) X(FAITH_EVENT_CONVERSATION_CREATED, 0) +#define FAITH_EVENT_BATCH_MAX_EVENTS 256u + typedef enum { #define X(name, value) name = value, FAITH_EVENTS_CODEC(X) @@ -32,20 +39,19 @@ typedef struct { uint8_t *data; } faith_envl_stc_event_t; +typedef struct { + uint16_t n_events; + + faith_body_size_t events_data_size; + faith_envl_stc_event_t *events; +} faith_envl_stc_event_batch_t; + typedef struct { /* The server-generated sequence number * of the event that was acknowledged */ uint64_t seq_num; - - faith_event_codec_type_t type; } faith_envl_cts_event_ack_t; -#define FAITH_CONVERSATION_ID_SIZE 16 - -typedef struct { - uint8_t bytes[FAITH_CONVERSATION_ID_SIZE]; -} faith_conversation_id_t; - typedef struct { faith_conversation_id_t conversation_id; } faith_event_conversation_created_t; diff --git a/faithd/src/codec/helpers.h b/faithd/src/codec/helpers.h index 431fed9..caebd49 100644 --- a/faithd/src/codec/helpers.h +++ b/faithd/src/codec/helpers.h @@ -73,7 +73,8 @@ (off) += _fh_len; \ } while (0) -#define FAITH_DECODE_RETURN(payload, payload_size, offset, dst, size) \ +#define _FAITH_DECODE_IMPL(payload, payload_size, offset, dst, size, \ + invalid_failure, bad_frame_failure) \ do { \ const size_t _fh_payload_size = (size_t)(payload_size); \ const size_t _fh_offset = (size_t)(offset); \ @@ -88,7 +89,7 @@ " payload size : %zu", \ #payload, #dst, _fh_offset, _fh_size, \ _fh_payload_size); \ - return FAITH_ERR_INVALID; \ + invalid_failure; \ } \ \ if (_fh_size > 0 && !(dst)) { \ @@ -100,7 +101,7 @@ " length : %zu\n" \ " payload size : %zu", \ #payload, #dst, _fh_offset, _fh_size, _fh_payload_size); \ - return FAITH_ERR_INVALID; \ + invalid_failure; \ } \ \ if (_fh_offset > _fh_payload_size) { \ @@ -112,7 +113,7 @@ " payload size : %zu", \ #payload, #dst, _fh_offset, _fh_size, \ _fh_payload_size); \ - return FAITH_ERR_BAD_FRAME; \ + bad_frame_failure; \ } \ \ if (_fh_size > _fh_payload_size - _fh_offset) { \ @@ -125,7 +126,7 @@ " payload size : %zu", \ #payload, #dst, _fh_offset, _fh_size, \ _fh_payload_size - _fh_offset, _fh_payload_size); \ - return FAITH_ERR_BAD_FRAME; \ + bad_frame_failure; \ } \ \ if (_fh_size > 0) \ @@ -134,8 +135,17 @@ (offset) += _fh_size; \ } while (0) -#define _FAITH_DECODE_INT_BE_RETURN(payload, payload_size, offset, out, type, \ - read_fn) \ +#define FAITH_DECODE_RETURN(payload, payload_size, offset, dst, size) \ + _FAITH_DECODE_IMPL(payload, payload_size, offset, dst, size, \ + return FAITH_ERR_INVALID, return FAITH_ERR_BAD_FRAME) + +#define FAITH_DECODE_DEFER(payload, payload_size, offset, dst, size) \ + _FAITH_DECODE_IMPL(payload, payload_size, offset, dst, size, \ + _FH_RETURN_DEFER(FAITH_ERR_INVALID), \ + _FH_RETURN_DEFER(FAITH_ERR_BAD_FRAME)) + +#define _FAITH_DECODE_INT_BE_IMPL(payload, payload_size, offset, out, type, \ + read_fn, invalid_failure, bad_frame_failure) \ do { \ const size_t _fh_payload_size = (size_t)(payload_size); \ const size_t _fh_offset = (size_t)(offset); \ @@ -150,7 +160,7 @@ " payload size : %zu", \ #read_fn, #out, _fh_int_size, _fh_offset, \ _fh_payload_size); \ - return FAITH_ERR_INVALID; \ + invalid_failure; \ } \ \ if (_fh_offset > _fh_payload_size) { \ @@ -162,7 +172,7 @@ " payload size : %zu", \ #read_fn, #out, _fh_int_size, _fh_offset, \ _fh_payload_size); \ - return FAITH_ERR_BAD_FRAME; \ + bad_frame_failure; \ } \ \ if (_fh_int_size > _fh_payload_size - _fh_offset) { \ @@ -175,13 +185,25 @@ " payload size : %zu", \ #read_fn, #out, _fh_int_size, _fh_offset, \ _fh_payload_size - _fh_offset, _fh_payload_size); \ - return FAITH_ERR_BAD_FRAME; \ + bad_frame_failure; \ } \ \ (out) = read_fn((payload) + _fh_offset); \ (offset) += _fh_int_size; \ } while (0) +#define _FAITH_DECODE_INT_BE_RETURN(payload, payload_size, offset, out, type, \ + read_fn) \ + _FAITH_DECODE_INT_BE_IMPL(payload, payload_size, offset, out, type, read_fn, \ + return FAITH_ERR_INVALID, \ + return FAITH_ERR_BAD_FRAME) + +#define _FAITH_DECODE_INT_BE_DEFER(payload, payload_size, offset, out, type, \ + read_fn) \ + _FAITH_DECODE_INT_BE_IMPL(payload, payload_size, offset, out, type, read_fn, \ + _FH_RETURN_DEFER(FAITH_ERR_INVALID), \ + _FH_RETURN_DEFER(FAITH_ERR_BAD_FRAME)) + #define FAITH_DECODE_U16_BE_RETURN(payload, payload_size, offset, out) \ _FAITH_DECODE_INT_BE_RETURN(payload, payload_size, offset, out, uint16_t, \ faith_read_u16_be) @@ -194,6 +216,18 @@ _FAITH_DECODE_INT_BE_RETURN(payload, payload_size, offset, out, uint64_t, \ faith_read_u64_be) +#define FAITH_DECODE_U16_BE_DEFER(payload, payload_size, offset, out) \ + _FAITH_DECODE_INT_BE_DEFER(payload, payload_size, offset, out, uint16_t, \ + faith_read_u16_be) + +#define FAITH_DECODE_U32_BE_DEFER(payload, payload_size, offset, out) \ + _FAITH_DECODE_INT_BE_DEFER(payload, payload_size, offset, out, uint32_t, \ + faith_read_u32_be) + +#define FAITH_DECODE_U64_BE_DEFER(payload, payload_size, offset, out) \ + _FAITH_DECODE_INT_BE_DEFER(payload, payload_size, offset, out, uint64_t, \ + faith_read_u64_be) + #define _FAITH_ENCODE_INT_BE_RETURN(buf, buf_cap, offset, value, type, \ write_fn) \ do { \ @@ -291,40 +325,34 @@ } \ } while (0) -#define FAITH_DECODE_EPILOGUE(envl_size, sign) \ +#define _FAITH_DECODE_EPILOGUE_IMPL(envl_size, condition, on_failure) \ do { \ const size_t _fh_expected_size = (size_t)(envl_size); \ const size_t _fh_offset = (size_t)offset; \ const size_t _fh_payload_size = (size_t)payload_size; \ \ - if (_fh_offset sign _fh_expected_size || \ - _fh_offset sign _fh_payload_size) { \ + if (condition) { \ _FH_LOG_CODEC_FAILURE("Envelope decode", \ "final decode offset is inconsistent", \ " offset : %zu\n" \ " payload size : %zu\n" \ " expected size : %zu", \ _fh_offset, _fh_payload_size, _fh_expected_size); \ - return FAITH_ERR_BAD_FRAME; \ + on_failure; \ } \ } while (0) -#define FAITH_DECODE_EPILOGUE_DNY(envl_size, sign) \ - do { \ - const size_t _fh_expected_size = (size_t)(envl_size); \ - const size_t _fh_offset = (size_t)offset; \ - const size_t _fh_payload_size = (size_t)payload_size; \ - \ - if (_fh_offset sign _fh_expected_size) { \ - _FH_LOG_CODEC_FAILURE("Envelope decode", \ - "final decode offset is inconsistent", \ - " offset : %zu\n" \ - " payload size : %zu\n" \ - " expected size : %zu", \ - _fh_offset, _fh_payload_size, _fh_expected_size); \ - return FAITH_ERR_BAD_FRAME; \ - } \ - } while (0) +#define FAITH_DECODE_EPILOGUE(envl_size, sign) \ + _FAITH_DECODE_EPILOGUE_IMPL( \ + envl_size, \ + (_fh_offset sign _fh_expected_size || _fh_offset sign _fh_payload_size), \ + return FAITH_ERR_BAD_FRAME) + +#define FAITH_DECODE_EPILOGUE_DEFER(envl_size, sign) \ + _FAITH_DECODE_EPILOGUE_IMPL( \ + envl_size, \ + (_fh_offset sign _fh_expected_size || _fh_offset sign _fh_payload_size), \ + _FH_RETURN_DEFER(FAITH_ERR_BAD_FRAME)) #define FAITH_ENCODE_PROLOGUE(envl_size) \ do { \ diff --git a/faithd/src/codec/protocol.c b/faithd/src/codec/protocol.c index 83ffdc5..7fe31d2 100644 --- a/faithd/src/codec/protocol.c +++ b/faithd/src/codec/protocol.c @@ -1,7 +1,10 @@ #include "protocol.h" #include "../../third_party/nob.h" +#include "../../third_party/stb_ds.h" +#include "core.h" +#include "events.h" #include "helpers.h" faith_status_code_t faith_encode_frame(uint8_t *out_buf, size_t *out_size, @@ -44,13 +47,13 @@ faith_status_code_t faith_encode_frame(uint8_t *out_buf, size_t *out_size, "Failed to encode frame; Frame is too large, " "frame_size=%u, MAX_FRAME_LEN=%i", (uint32_t)frame_size, (int32_t)FAITH_MAX_FRAME_LEN); - return FAITH_ERR_FRAME_TOO_LARGE; + return FAITH_ERR_TOO_LARGE; payload_too_large: nob_log(ERROR, "Failed to encode frame; Payload is too large, " "payload_size=%u, MAX_PAYLOAD_SIZE=%i", (uint32_t)in->payload_size, (int32_t)FAITH_MAX_PAYLOAD_SIZE); - return FAITH_ERR_FRAME_TOO_LARGE; + return FAITH_ERR_TOO_LARGE; } faith_status_code_t faith_decode_frame(const uint8_t *payload, @@ -86,7 +89,7 @@ faith_status_code_t faith_decode_frame(const uint8_t *payload, "payload_size=%zu MAX_PAYLOAD_SIZE=%zu", payload_size, (size_t)FAITH_MAX_PAYLOAD_SIZE); - return FAITH_ERR_FRAME_TOO_LARGE; + return FAITH_ERR_TOO_LARGE; } if (frame_size > FAITH_MAX_FRAME_LEN) { @@ -94,7 +97,7 @@ faith_status_code_t faith_decode_frame(const uint8_t *payload, "Failed to parse frame from buffer; Frame is too large, " "frame_size=%i MAX_FRAME_LEN=%i", (int32_t)frame_size, (int32_t)FAITH_MAX_FRAME_LEN); - return FAITH_ERR_FRAME_TOO_LARGE; + return FAITH_ERR_TOO_LARGE; } size_t total_frame_size = FAITH_FRAME_LENGTH_SIZE + frame_size; @@ -112,22 +115,29 @@ faith_status_code_t faith_decode_frame(const uint8_t *payload, if (proto_ver != FAITH_PROTO_VERSION) return FAITH_ERR_UNSUPPORTED_VER; - FAITH_DECODE_EPILOGUE_DNY(FAITH_FRAME_HEADER_SIZE, <); - out->proto_ver = proto_ver; out->msg_type = msg_type; out->frame_size = frame_size; out->payload_size = frame_size - FAITH_FRAME_METADATA_SIZE; + faith_status_code_t _fh_result = FAITH_OK; + if (out->payload_size > 0) { out->payload = malloc(out->payload_size); if (!out->payload) return FAITH_ERR_NOMEM; - memcpy(out->payload, payload + FAITH_FRAME_HEADER_SIZE, out->payload_size); + FAITH_DECODE_DEFER(payload, payload_size, offset, out->payload, + out->payload_size); } + FAITH_DECODE_EPILOGUE_DEFER(FAITH_FRAME_HEADER_SIZE + out->payload_size, <); + return FAITH_OK; +defer: + free(out->payload); + out->payload = NULL; + return _fh_result; } void faith_free_frame(faith_frame_t *frame) { @@ -261,30 +271,28 @@ faith_status_code_t faith_decode_envelope(const uint8_t *payload, FAITH_DECODE_U32_BE_RETURN(payload, payload_size, offset, out->body_size); - FAITH_DECODE_EPILOGUE_DNY(FAITH_ENVL_HEADER_SIZE, <); - if (out->body_size != payload_size - offset) { return FAITH_ERR_BAD_FRAME; } - uint8_t *body = NULL; + faith_status_code_t _fh_result = FAITH_OK; + if (out->body_size != 0) { - body = malloc(out->body_size); - if (!body) + out->body = malloc(out->body_size); + if (!out->body) return FAITH_ERR_NOMEM; - memcpy(body, payload + offset, out->body_size); - offset += out->body_size; - } - - if (offset != payload_size) { - free(body); - return FAITH_ERR_BAD_FRAME; + FAITH_DECODE_DEFER(payload, payload_size, offset, out->body, + out->body_size); } - out->body = body; + FAITH_DECODE_EPILOGUE_DEFER(FAITH_ENVL_HEADER_SIZE + out->body_size, !=); return FAITH_OK; +defer: + free(out->body); + out->body = NULL; + return _fh_result; } faith_status_code_t @@ -619,17 +627,26 @@ faith_status_code_t faith_decode_command_body(const uint8_t *payload, if (out->payload_size > FAITH_COMMAND_PAYLOAD_SIZE_MAX) return FAITH_ERR_BAD_ENVELOPE; + faith_status_code_t _fh_result = FAITH_OK; + if (out->payload_size > 0) { out->payload = malloc(out->payload_size); if (!out->payload) return FAITH_ERR_NOMEM; - memcpy(out->payload, payload + offset, out->payload_size); + FAITH_DECODE_DEFER(payload, payload_size, offset, out->payload, + out->payload_size); } - FAITH_DECODE_EPILOGUE_DNY(FAITH_ENVL_CTS_COMMAND_BODY_SIZE_FIXED, <); + FAITH_DECODE_EPILOGUE_DEFER( + FAITH_ENVL_CTS_COMMAND_BODY_SIZE_FIXED + out->payload_size, !=); return FAITH_OK; + +defer: + free(out->payload); + out->payload = NULL; + return _fh_result; } faith_status_code_t @@ -673,7 +690,7 @@ faith_decode_command_result_body(const uint8_t *payload, FAITH_DECODE_U32_BE_RETURN(payload, payload_size, offset, out->result); - FAITH_DECODE_EPILOGUE_DNY(FAITH_ENVL_STC_COMMAND_RESULT_BODY_SIZE, !=); + FAITH_DECODE_EPILOGUE(FAITH_ENVL_STC_COMMAND_RESULT_BODY_SIZE, !=); return FAITH_OK; } @@ -685,6 +702,10 @@ faith_status_code_t faith_encode_event_body(uint8_t *out_buf, FAITH_ENCODE_PROLOGUE(FAITH_ENVL_STC_EVENT_BODY_SIZE_FIXED + in->data_size); + if (in->data_size != 0 && !in->data) { + return FAITH_ERR_INVALID; + } + size_t offset = 0; FAITH_ENCODE_U64_BE_RETURN(out_buf, buf_cap_in_bytes, offset, in->seq_num); @@ -714,17 +735,230 @@ faith_status_code_t faith_decode_event_body(const uint8_t *payload, FAITH_DECODE_U64_BE_RETURN(payload, payload_size, offset, out->seq_num); FAITH_DECODE_U32_BE_RETURN(payload, payload_size, offset, out->type); FAITH_DECODE_U32_BE_RETURN(payload, payload_size, offset, out->data_size); + + faith_status_code_t _fh_result = FAITH_OK; + if (out->data_size > 0) { out->data = malloc(out->data_size); if (!out->data) return FAITH_ERR_NOMEM; - memcpy(out->data, payload + offset, out->data_size); + FAITH_DECODE_DEFER(payload, payload_size, offset, out->data, + out->data_size); + } + + FAITH_DECODE_EPILOGUE_DEFER( + FAITH_ENVL_STC_EVENT_BODY_SIZE_FIXED + out->data_size, !=); + + return FAITH_OK; +defer: + free(out->data); + out->data = NULL; + return _fh_result; +} + +faith_status_code_t +faith_codec_event_batch_data_size(faith_envl_stc_event_t *events, + uint16_t n_events, + faith_body_size_t *o_size) { + if (!events || !o_size) + return FAITH_ERR_INVALID; + + *o_size = 0; + + if (n_events > FAITH_EVENT_BATCH_MAX_EVENTS) { + nob_log(ERROR, + "Invalid EVENT_BATCH body: event count %" PRIu16 + " exceeds maximum %u.", + n_events, FAITH_EVENT_BATCH_MAX_EVENTS); + return FAITH_ERR_INVALID; + } + + size_t total_data_size = 0; + + if (n_events != 0) { + for (uint16_t i = 0; i < n_events; i++) { + size_t body_size = (size_t)FAITH_ENVL_STC_EVENT_BODY_SIZE_FIXED + + (size_t)events[i].data_size; + + if (body_size > UINT16_MAX) { + nob_log(ERROR, + "Cannot encode EVENT_BATCH: event %" PRIu16 + " body size (%zu bytes) exceeds UINT16_MAX.", + i, body_size); + return FAITH_ERR_TOO_LARGE; + } + + size_t encoded_elem_size = sizeof(uint16_t) + body_size; + + if (encoded_elem_size > UINT32_MAX - total_data_size) { + nob_log(ERROR, + "Cannot encode EVENT_BATCH: adding event %" PRIu16 + " would make the total event data exceed UINT32_MAX.", + i); + return FAITH_ERR_TOO_LARGE; + } + + total_data_size += encoded_elem_size; + } + } + + if (total_data_size > SIZE_MAX - FAITH_ENVL_STC_EVENT_BATCH_BODY_SIZE_FIXED) + return FAITH_ERR_TOO_LARGE; + + *o_size = total_data_size; + + return FAITH_OK; +} + +faith_status_code_t +faith_encode_event_batch_body(uint8_t *out_buf, faith_body_size_t *out_size, + size_t buf_cap_in_bytes, + const faith_envl_stc_event_batch_t *in) { + if (!in || (!in->events && in->n_events != 0)) + return FAITH_ERR_INVALID; + + size_t encoded_size = + FAITH_ENVL_STC_EVENT_BATCH_BODY_SIZE_FIXED + in->events_data_size; + + FAITH_ENCODE_PROLOGUE(encoded_size); + + size_t offset = 0; + + FAITH_ENCODE_U16_BE_RETURN(out_buf, buf_cap_in_bytes, offset, in->n_events); + FAITH_ENCODE_U32_BE_RETURN(out_buf, buf_cap_in_bytes, offset, + (uint32_t)in->events_data_size); + + for (uint16_t i = 0; i < in->n_events; i++) { + size_t body_size = (size_t)FAITH_ENVL_STC_EVENT_BODY_SIZE_FIXED + + (size_t)in->events[i].data_size; + + /* encode element size */ + FAITH_ENCODE_U16_BE_RETURN(out_buf, buf_cap_in_bytes, offset, + (uint16_t)body_size); + + faith_body_size_t returned_body_size = 0; + /* encode element body*/ + _FH_CHECK_RETURN( + faith_encode_event_body(out_buf + offset, &returned_body_size, + buf_cap_in_bytes - offset, &in->events[i])); + + if ((size_t)returned_body_size != body_size) { + nob_log(ERROR, + "EVENT encoder returned an unexpected body size: " + "expected %zu bytes, got %" PRIu32 " bytes.", + body_size, returned_body_size); + return FAITH_ERR_INVALID; + } + offset += body_size; + } + + FAITH_ENCODE_EPILOGUE(encoded_size, !=); + + return FAITH_OK; +} + +faith_status_code_t +faith_decode_event_batch_body(const uint8_t *payload, + faith_body_size_t payload_size, + faith_envl_stc_event_batch_t *out) { + + FAITH_DECODE_PROLOGUE(FAITH_ENVL_STC_EVENT_BATCH_BODY_SIZE_FIXED, <); + + out->n_events = 0; + out->events_data_size = 0; + out->events = NULL; + + size_t offset = 0; + + FAITH_DECODE_U16_BE_RETURN(payload, payload_size, offset, out->n_events); + + if (out->n_events > FAITH_EVENT_BATCH_MAX_EVENTS) { + nob_log(ERROR, + "Invalid EVENT_BATCH body: event count %" PRIu16 + " exceeds maximum %u.", + out->n_events, FAITH_EVENT_BATCH_MAX_EVENTS); + return FAITH_ERR_BAD_ENVELOPE; + } + + FAITH_DECODE_U32_BE_RETURN(payload, payload_size, offset, + out->events_data_size); + + size_t decoded_size = (size_t)FAITH_ENVL_STC_EVENT_BATCH_BODY_SIZE_FIXED + + (size_t)out->events_data_size; + + if ((size_t)payload_size < decoded_size) { + nob_log(ERROR, + "Invalid EVENT_BATCH body: declared events data size (%" PRIu32 + " bytes) exceeds remaining payload size (%" PRIu32 " bytes).", + out->events_data_size, + payload_size - FAITH_ENVL_STC_EVENT_BATCH_BODY_SIZE_FIXED); + + return FAITH_ERR_OVERFLOW; + } + + size_t events_end = offset + (size_t)out->events_data_size; + + if (events_end != (size_t)payload_size) { + nob_log(ERROR, "Invalid EVENT_BATCH body: payload has %zu trailing bytes.", + (size_t)payload_size - events_end); + return FAITH_ERR_BAD_ENVELOPE; } - FAITH_DECODE_EPILOGUE_DNY(FAITH_ENVL_STC_EVENT_BODY_SIZE_FIXED, <); + faith_status_code_t _fh_result = FAITH_OK; + + if (out->n_events != 0) { + arrsetlen(out->events, out->n_events); + } + + uint16_t decoded_events = 0; + for (uint16_t i = 0; i < out->n_events; i++) { + /* decode element size */ + uint16_t elem_size = 0; + FAITH_DECODE_U16_BE_DEFER(payload, payload_size, offset, elem_size); + + if (elem_size < FAITH_ENVL_STC_EVENT_BODY_SIZE_FIXED) { + nob_log(ERROR, + "Invalid EVENT_BATCH body: event %" PRIu16 + " declares an invalid body size of %" PRIu16 + " bytes. The minimum required size for an event is: %" PRIu32 + " bytes.", + i, elem_size, FAITH_ENVL_STC_EVENT_BODY_SIZE_FIXED); + _FH_RETURN_DEFER(FAITH_ERR_BAD_ENVELOPE); + } + + if (offset > payload_size || payload_size - offset < (size_t)elem_size) { + nob_log(ERROR, + "Invalid EVENT_BATCH body: event %" PRIu16 " declares %" PRIu16 + " bytes, but only %zu bytes remain.", + i, elem_size, events_end - offset); + _FH_RETURN_DEFER(FAITH_ERR_BAD_FRAME); + } + + /* decode element body */ + _FH_CHECK_DEFER(faith_decode_event_body( + payload + offset, (faith_body_size_t)elem_size, &out->events[i])); + + decoded_events++; + + offset += (size_t)elem_size; + } + + FAITH_DECODE_EPILOGUE_DEFER(decoded_size, !=); return FAITH_OK; + +defer: + for (size_t i = 0; i < decoded_events; i++) { + free(out->events[i].data); + out->events[i].data = NULL; + } + arrfree(out->events); + out->events_data_size = 0; + out->n_events = 0; + out->events = NULL; + + return _fh_result; } faith_status_code_t diff --git a/faithd/src/codec/protocol.h b/faithd/src/codec/protocol.h index 2165951..d5974f5 100644 --- a/faithd/src/codec/protocol.h +++ b/faithd/src/codec/protocol.h @@ -8,6 +8,7 @@ #include "../server/envelopes.h" #include "commands.h" +#include "core.h" #include "events.h" #include "helpers.h" @@ -150,6 +151,20 @@ faith_status_code_t faith_decode_event_body(const uint8_t *payload, faith_body_size_t payload_size, faith_envl_stc_event_t *out); +faith_status_code_t +faith_codec_event_batch_data_size(faith_envl_stc_event_t *events, + uint16_t n_events, faith_body_size_t *o_size); + +faith_status_code_t +faith_encode_event_batch_body(uint8_t *out_buf, faith_body_size_t *out_size, + size_t buf_cap_in_bytes, + const faith_envl_stc_event_batch_t *in); + +faith_status_code_t +faith_decode_event_batch_body(const uint8_t *payload, + faith_body_size_t payload_size, + faith_envl_stc_event_batch_t *out); + faith_status_code_t faith_encode_event_ack_body(uint8_t *out_buf, faith_body_size_t *out_size, size_t buf_cap_in_bytes, diff --git a/faithd/src/commands/conversation.c b/faithd/src/commands/conversation.c index c202e86..113a8f5 100644 --- a/faithd/src/commands/conversation.c +++ b/faithd/src/commands/conversation.c @@ -2,20 +2,33 @@ #include "../delivery/events.h" +#include "../../third_party/stb_ds.h" + static faith_status_code_t send_conversation_created( server_state_t *s, client_conn_t *cl, - const faith_event_conversation_created_t *conv_created) { + const faith_event_conversation_created_t *conv_created, + const faith_auth_id_t *conservant_id) { if (!s || !cl || !conv_created) return FAITH_ERR_INVALID; + uint8_t data[FAITH_EVENT_CONVERSATION_CREATED_DATA_SIZE] = {0}; faith_body_size_t data_size = 0; _FH_CHECK_RETURN(faith_encode_event_conversation_created( data, &data_size, sizeof(data), conv_created)); - _FH_CHECK_RETURN(delivery_queue_event(s, cl, FAITH_EVENT_CONVERSATION_CREATED, - data, data_size)); + _FH_CHECK_RETURN( + delivery_queue_event(s, &cl->ident.auth_id, &cl->ident.device_id, + FAITH_EVENT_CONVERSATION_CREATED, data, data_size)); + + faith_status_code_t device_loop_rc = FAITH_OK; + _FH_FOR_EACH_AUTH_DEVICE_SESSION( + s, conservant_id, recipient_sess, device_loop_rc, { + _FH_CHECK_RETURN(delivery_queue_event( + s, conservant_id, &recipient_sess->ident.device_id, + FAITH_EVENT_CONVERSATION_CREATED, data, data_size)); + }); return FAITH_OK; } @@ -38,6 +51,12 @@ faith_status_code_t conv_handle_create_conversation( cmd->payload, cmd->payload_size, &create_conv_cmd)); } + if (!sess_registry_auth_id_registered(&s->rt, + &create_conv_cmd.conversant_id)) { + potential_err = FAITH_COMMAND_ERR_USER_NOT_FOUND; + _FH_RETURN_DEFER(FAITH_ERR_NOT_FOUND); + } + faith_conversation_id_t conv_id = {0}; potential_err = FAITH_COMMAND_ERR_INTERNAL_ERROR; @@ -46,9 +65,8 @@ faith_status_code_t conv_handle_create_conversation( faith_event_conversation_created_t conv_created = {.conversation_id = conv_id}; - for (size_t i = 0; i < 100; i++) { - _FH_CHECK_RETURN(send_conversation_created(s, cl, &conv_created)); - } + _FH_CHECK_DEFER(send_conversation_created(s, cl, &conv_created, + &create_conv_cmd.conversant_id)); *o_result = FAITH_COMMAND_RESULT_ACCEPTED; diff --git a/faithd/src/commands/dispatch.c b/faithd/src/commands/dispatch.c index ead5f30..eb7d076 100644 --- a/faithd/src/commands/dispatch.c +++ b/faithd/src/commands/dispatch.c @@ -32,7 +32,7 @@ queue_command_result(server_state_t *s, client_conn_t *cl, envl.body = body; envl.body_size = body_size; envl.type = FAITH_ENVELOPE_COMMAND_RESULT; - envl.recipient_id = cl->auth_id; + envl.recipient_id = cl->ident.auth_id; _FH_CHECK_RETURN(server_queue_envelope_or_mark_dead(s, cl, &envl)); @@ -57,7 +57,7 @@ faith_status_code_t command_dispatch(server_state_t *s, client_conn_t *cl, return FAITH_ERR_UNAUTHORIZED; } - if (!faith_client_id_equal(cl->auth_id, envl->sender_id)) { + if (!faith_client_id_equal(cl->ident.auth_id, envl->sender_id)) { return FAITH_ERR_NOT_EQUAL; } diff --git a/faithd/src/core/core.h b/faithd/src/core/core.h index 971b38f..8425a19 100644 --- a/faithd/src/core/core.h +++ b/faithd/src/core/core.h @@ -74,7 +74,7 @@ X(FAITH_ERR_OVERFLOW, 4) \ X(FAITH_ERR_UNDERFLOW, 5) \ X(FAITH_ERR_IO, 6) \ - X(FAITH_ERR_FRAME_TOO_LARGE, 7) \ + X(FAITH_ERR_TOO_LARGE, 7) \ X(FAITH_ERR_BAD_FRAME, 8) \ X(FAITH_ERR_CLOSED, 9) \ X(FAITH_ERR_UNSUPPORTED_VER, 10) \ diff --git a/faithd/src/delivery/event_inbox.c b/faithd/src/delivery/event_inbox.c index e38a007..0377cdc 100644 --- a/faithd/src/delivery/event_inbox.c +++ b/faithd/src/delivery/event_inbox.c @@ -9,15 +9,45 @@ void device_event_inbox_init(device_event_inbox_t *o_inbox) { o_inbox->events = NULL; o_inbox->last_acked_seq = UINT64_MAX; + o_inbox->last_sent_seq = UINT64_MAX; o_inbox->next_seq = 0; } faith_status_code_t -device_event_inbox_queue_event(device_event_inbox_t *inbox, - const faith_envl_stc_event_t *event) { +device_event_inbox_push_event(device_event_inbox_t *inbox, + const faith_envl_stc_event_t *event) { if (!inbox || !event) return FAITH_ERR_INVALID; - arrput(inbox->events, *event); + + /* Dont allow UINT64_MAX */ + if (event->seq_num == UINT64_MAX) + return FAITH_ERR_OVERFLOW; + + if (inbox->next_seq == UINT64_MAX) + return FAITH_ERR_OVERFLOW; + + if (inbox->last_sent_seq != UINT64_MAX && event->seq_num != inbox->next_seq) { + nob_log(ERROR, "Tried to push event with invalid sequence order.\n"); + return FAITH_ERR_INVALID; + } + + faith_envl_stc_event_t entry = *event; + uint8_t *data_copy = NULL; + + if (event->data_size > 0) { + data_copy = malloc(event->data_size); + if (!data_copy) + return FAITH_ERR_NOMEM; + + memcpy(data_copy, event->data, event->data_size); + } + + entry.data = data_copy; + + arrput(inbox->events, entry); + + inbox->last_sent_seq = inbox->next_seq; + inbox->next_seq++; return FAITH_OK; } @@ -28,7 +58,7 @@ device_event_inbox_advance_seq(device_event_inbox_t *inbox) { return FAITH_ERR_INVALID; /* Dont allow UINT64_MAX */ - if (inbox->next_seq >= UINT64_MAX - 1) + if (inbox->next_seq == UINT64_MAX) return FAITH_ERR_OVERFLOW; inbox->last_sent_seq = inbox->next_seq; @@ -37,19 +67,37 @@ device_event_inbox_advance_seq(device_event_inbox_t *inbox) { return FAITH_OK; } -faith_status_code_t device_event_inbox_ack_seq(device_event_inbox_t *inbox, - uint64_t acked_seq) { +faith_status_code_t device_event_inbox_remove_until(device_event_inbox_t *inbox, + uint64_t until_seq) { if (!inbox) return FAITH_ERR_INVALID; - /* Duplicate or stale ACK */ - if (acked_seq <= inbox->last_acked_seq) + if (inbox->last_sent_seq == UINT64_MAX) + return FAITH_ERR_BAD_ENVELOPE; + + if (until_seq > inbox->last_sent_seq) + return FAITH_ERR_BAD_ENVELOPE; + + if (inbox->last_acked_seq != UINT64_MAX && + until_seq <= inbox->last_acked_seq) { return FAITH_OK; + } - if (acked_seq >= inbox->next_seq) - return FAITH_ERR_INVALID; + size_t event_count = arrlen(inbox->events); + + size_t i = 0; + for (; i < event_count; i++) { + if (inbox->events[i].seq_num > until_seq) + break; + + free(inbox->events[i].data); + inbox->events[i].data = NULL; + } + + if (i > 0) + arrdeln(inbox->events, 0, i); - inbox->last_acked_seq = acked_seq; + inbox->last_acked_seq = until_seq; return FAITH_OK; } diff --git a/faithd/src/delivery/event_inbox.h b/faithd/src/delivery/event_inbox.h index 854633d..c32ea08 100644 --- a/faithd/src/delivery/event_inbox.h +++ b/faithd/src/delivery/event_inbox.h @@ -1,6 +1,7 @@ #pragma once #include "../codec/events.h" + #include typedef struct { @@ -14,11 +15,12 @@ typedef struct { void device_event_inbox_init(device_event_inbox_t *o_inbox); faith_status_code_t -device_event_inbox_queue_event(device_event_inbox_t *inbox, - const faith_envl_stc_event_t *event); +device_event_inbox_push_event(device_event_inbox_t *inbox, + const faith_envl_stc_event_t *event); faith_status_code_t device_event_inbox_advance_seq(device_event_inbox_t *inbox); -faith_status_code_t device_event_inbox_ack_seq(device_event_inbox_t *inbox, - uint64_t acked_seq); + +faith_status_code_t device_event_inbox_remove_until(device_event_inbox_t *inbox, + uint64_t until_seq); faith_status_code_t device_event_inbox_destroy(device_event_inbox_t *inbox); diff --git a/faithd/src/delivery/events.c b/faithd/src/delivery/events.c index 31573d6..e974522 100644 --- a/faithd/src/delivery/events.c +++ b/faithd/src/delivery/events.c @@ -4,21 +4,100 @@ #include "../server/client_io.h" #include "event_inbox.h" -faith_status_code_t delivery_queue_event(server_state_t *s, client_conn_t *cl, - faith_event_codec_type_t type, - uint8_t *data, - faith_body_size_t data_size) { - if (!cl || (!data && data_size != 0) || (data && data_size == 0)) +#include "../../third_party/stb_ds.h" + +#define DIV_UP(x, y) (((x) + (y) - 1) / (y)) + +static faith_status_code_t +queue_event_online_user(server_state_t *s, struct client_conn_t *cl, + device_event_inbox_t *inbox, + const faith_envl_stc_event_t *event) { + if (!s || !cl || !inbox || !event) return FAITH_ERR_INVALID; - if (!cl->authorized) - return FAITH_ERR_UNAUTHORIZED; + faith_status_code_t _fh_result = FAITH_OK; + + size_t cap = FAITH_ENVL_STC_EVENT_BODY_SIZE_FIXED + event->data_size; + uint8_t *body = malloc(cap); + faith_body_size_t body_size = 0; + _FH_CHECK_DEFER(faith_encode_event_body(body, &body_size, cap, event)); + + faith_envelope_t envl = {0}; + envl.type = FAITH_ENVELOPE_EVENT; + envl.recipient_id = cl->ident.auth_id; + envl.body = body; + envl.body_size = body_size; + + _FH_CHECK_DEFER(server_queue_envelope_or_mark_dead(s, cl, &envl)); + + _FH_CHECK_DEFER(device_event_inbox_advance_seq(inbox)); + +defer: + free(body); + return _fh_result; +} + +static faith_status_code_t +queue_event_offline_user(device_event_inbox_t *inbox, + const faith_envl_stc_event_t *event) { + if (!inbox || !event) + return FAITH_ERR_INVALID; + + _FH_CHECK_RETURN(device_event_inbox_push_event(inbox, event)); + return FAITH_OK; +} + +static faith_status_code_t +queue_event_batch_envl(server_state_t *s, struct client_conn_t *cl, + const faith_envl_stc_event_batch_t *batch_envl) { + if (!batch_envl) + return FAITH_ERR_INVALID; + + size_t batch_data_cap = + FAITH_ENVL_STC_EVENT_BATCH_BODY_SIZE_FIXED + batch_envl->events_data_size; + uint8_t *body = malloc(batch_data_cap); + if (!body) { + return FAITH_ERR_NOMEM; + } + + faith_status_code_t _fh_result = FAITH_OK; + faith_body_size_t body_size = 0; + _FH_CHECK_DEFER(faith_encode_event_batch_body(body, &body_size, + batch_data_cap, batch_envl)); + + faith_envelope_t envl = {0}; + envl.type = FAITH_ENVELOPE_EVENT_BATCH; + envl.recipient_id = cl->ident.auth_id; + envl.body = body; + envl.body_size = body_size; + + _FH_CHECK_DEFER(server_queue_envelope_or_mark_dead(s, cl, &envl)); + +defer: + free(body); + return _fh_result; +} + +faith_status_code_t delivery_queue_event(server_state_t *s, + const faith_auth_id_t *auth_id, + const faith_device_id_t *device_id, + const faith_event_codec_type_t type, + uint8_t *data, + faith_body_size_t data_size) { + if (!auth_id || !device_id || (!data && data_size != 0) || + (data && data_size == 0)) + return FAITH_ERR_INVALID; client_device_session_data_t *sess = NULL; _FH_CHECK_RETURN( - sess_registry_get_session(&s->rt, &cl->auth_id, &cl->device_id, &sess)); + sess_registry_get_session(&s->rt, auth_id, device_id, &sess)); if (!sess) + return FAITH_ERR_NOT_FOUND; + + bool online = sess->conn != NULL; + + if (online && !sess->conn->authorized) return FAITH_ERR_UNAUTHORIZED; faith_envl_stc_event_t event = {0}; @@ -27,26 +106,52 @@ faith_status_code_t delivery_queue_event(server_state_t *s, client_conn_t *cl, event.data_size = data_size; event.seq_num = sess->inbox.next_seq; - faith_status_code_t _fh_result = FAITH_OK; + if (online) { + _FH_CHECK_RETURN( + queue_event_online_user(s, sess->conn, &sess->inbox, &event)); + return FAITH_OK; + } - size_t cap = FAITH_ENVL_STC_EVENT_BODY_SIZE_FIXED + data_size; - uint8_t *body = malloc(cap); - faith_body_size_t body_size = 0; - _FH_CHECK_DEFER(faith_encode_event_body(body, &body_size, cap, &event)); + _FH_CHECK_RETURN(queue_event_offline_user(&sess->inbox, &event)); + return FAITH_OK; +} - faith_envelope_t envl = {0}; - envl.type = FAITH_ENVELOPE_EVENT; - envl.recipient_id = cl->auth_id; - envl.body = body; - envl.body_size = body_size; +faith_status_code_t delivery_queue_pending_events(server_state_t *s, + client_conn_t *cl) { + if (!s || !cl) { + return FAITH_ERR_INVALID; + } + client_device_session_data_t *sess = NULL; + _FH_CHECK_RETURN(sess_registry_get_session(&s->rt, &cl->ident.auth_id, + &cl->ident.device_id, &sess)); - _FH_CHECK_DEFER(server_queue_envelope_or_mark_dead(s, cl, &envl)); + if (!sess) + return FAITH_ERR_NOT_FOUND; - _FH_CHECK_DEFER(device_event_inbox_advance_seq(&sess->inbox)); + size_t n_events = arrlen(sess->inbox.events); -defer: - free(body); - return _fh_result; + for (size_t offset = 0; offset < n_events; offset += 256) { + size_t remaining = n_events - offset; + + faith_envl_stc_event_t *batch = &sess->inbox.events[offset]; + + size_t in_batch = remaining < FAITH_EVENT_BATCH_MAX_EVENTS + ? remaining + : FAITH_EVENT_BATCH_MAX_EVENTS; + + faith_body_size_t events_data_size = 0; + _FH_CHECK_RETURN( + faith_codec_event_batch_data_size(batch, in_batch, &events_data_size)); + + faith_envl_stc_event_batch_t batch_envl = {0}; + batch_envl.events_data_size = events_data_size; + batch_envl.events = batch; + batch_envl.n_events = (uint16_t)in_batch; + + _FH_CHECK_RETURN(queue_event_batch_envl(s, cl, &batch_envl)); + } + + return FAITH_OK; } faith_status_code_t delivery_handle_event_acked(server_state_t *s, @@ -62,17 +167,17 @@ faith_status_code_t delivery_handle_event_acked(server_state_t *s, return FAITH_ERR_UNAUTHORIZED; client_device_session_data_t *sess = NULL; - _FH_CHECK_RETURN( - sess_registry_get_session(&s->rt, &cl->auth_id, &cl->device_id, &sess)); + _FH_CHECK_RETURN(sess_registry_get_session(&s->rt, &cl->ident.auth_id, + &cl->ident.device_id, &sess)); if (!sess) - return FAITH_ERR_UNAUTHORIZED; + return FAITH_ERR_NOT_FOUND; faith_envl_cts_event_ack_t ack = {0}; _FH_CHECK_RETURN( faith_decode_event_ack_body(envl->body, envl->body_size, &ack)); - device_event_inbox_ack_seq(&sess->inbox, ack.seq_num); + _FH_CHECK_RETURN(device_event_inbox_remove_until(&sess->inbox, ack.seq_num)); return FAITH_OK; } diff --git a/faithd/src/delivery/events.h b/faithd/src/delivery/events.h index a0006b0..b59a09c 100644 --- a/faithd/src/delivery/events.h +++ b/faithd/src/delivery/events.h @@ -4,11 +4,16 @@ #include "../core/core.h" #include "../server/server.h" -faith_status_code_t delivery_queue_event(server_state_t *s, client_conn_t *cl, +faith_status_code_t delivery_queue_event(server_state_t *s, + const faith_auth_id_t *auth_id, + const faith_device_id_t *device_id, faith_event_codec_type_t type, uint8_t *data, faith_body_size_t data_size); +faith_status_code_t delivery_queue_pending_events(server_state_t *s, + client_conn_t *cl); + faith_status_code_t delivery_handle_event_acked(server_state_t *s, client_conn_t *cl, faith_envelope_t *envl); diff --git a/faithd/src/delivery/routing.c b/faithd/src/delivery/routing.c index b3f8941..185722f 100644 --- a/faithd/src/delivery/routing.c +++ b/faithd/src/delivery/routing.c @@ -44,34 +44,37 @@ delivery_route_envelope_to_auth_id(server_state_t *s, client_conn_t *cl_sender, faith_id128_to_hex(recipient_auth_id->bytes, recipient_auth_id_hex)); faith_status_code_t device_loop_rc = FAITH_OK; - _FH_FOR_EACH_AUTH_DEVICE(s, recipient_auth_id, recipient, device_loop_rc, { - faith_envelope_t routing_envl = *envl; - routing_envl.sender_id = cl_sender->auth_id; - routing_envl.recipient_id = *recipient_auth_id; - _FH_CHECK(server_queue_envelope_or_mark_dead(s, recipient, &routing_envl)); - - char recipient_device_id_hex[33]; - _FH_CHECK_RETURN(faith_id128_to_hex(recipient->device_id.bytes, - recipient_device_id_hex)); - - if (_fh_rc != FAITH_OK) { - nob_log(ERROR, - "[client=%" PRIu64 - " fd=%i] Envelope %s: Failed to route envelope to " - "recipient device (auth_id: %s, device_id: %s).", - cl_sender->conn.id, cl_sender->conn.fd, - faith_envelope_name(envl->type), recipient_auth_id_hex, - recipient_device_id_hex); - continue; - } - - nob_log( - INFO, - "[client=%" PRIu64 " fd=%i] Envelope %s: Routed envelope to recipient " - "device (auth_id: %s, device_id: %s).", - cl_sender->conn.id, cl_sender->conn.fd, faith_envelope_name(envl->type), - recipient_auth_id_hex, recipient_device_id_hex); - }); + _FH_FOR_EACH_AUTH_DEVICE_CONNECTION( + s, recipient_auth_id, recipient_cl, device_loop_rc, { + faith_envelope_t routing_envl = *envl; + routing_envl.sender_id = cl_sender->ident.auth_id; + routing_envl.recipient_id = *recipient_auth_id; + _FH_CHECK( + server_queue_envelope_or_mark_dead(s, recipient_cl, &routing_envl)); + + char recipient_device_id_hex[33]; + _FH_CHECK_RETURN(faith_id128_to_hex(recipient_cl->ident.device_id.bytes, + recipient_device_id_hex)); + + if (_fh_rc != FAITH_OK) { + nob_log(ERROR, + "[client=%" PRIu64 + " fd=%i] Envelope %s: Failed to route envelope to " + "recipient device (auth_id: %s, device_id: %s).", + cl_sender->conn.id, cl_sender->conn.fd, + faith_envelope_name(envl->type), recipient_auth_id_hex, + recipient_device_id_hex); + continue; + } + + nob_log(INFO, + "[client=%" PRIu64 + " fd=%i] Envelope %s: Routed envelope to recipient " + "device (auth_id: %s, device_id: %s).", + cl_sender->conn.id, cl_sender->conn.fd, + faith_envelope_name(envl->type), recipient_auth_id_hex, + recipient_device_id_hex); + }); return FAITH_OK; } diff --git a/faithd/src/server/client_lifecycle.c b/faithd/src/server/client_lifecycle.c index 98afc77..5ddb19b 100644 --- a/faithd/src/server/client_lifecycle.c +++ b/faithd/src/server/client_lifecycle.c @@ -81,7 +81,7 @@ faith_status_code_t server_close_client(server_state_t *s, client_session_device_t *devices = NULL; faith_status_code_t rc = - sess_registry_get_devices(&s->rt, &cl->auth_id, &devices); + sess_registry_get_devices(&s->rt, &cl->ident.auth_id, &devices); if (rc != FAITH_OK || devices == NULL) { nob_log(ERROR, "[client=%" PRIu64 @@ -105,15 +105,12 @@ faith_status_code_t server_close_client(server_state_t *s, } if (cl->authorized) { - // TODO: persistent sessions - faith_status_code_t rc = - sess_registry_unregister_session(&s->rt, &cl->auth_id, &cl->device_id); + client_device_session_data_t *sess = NULL; + _FH_CHECK_RETURN(sess_registry_get_session(&s->rt, &cl->ident.auth_id, + &cl->ident.device_id, &sess)); - if (rc != FAITH_OK) { - nob_log(ERROR, "[client=%" PRIu64 "] routing unregister failed: %s (%d)", - cl->conn.id, faith_status_code_name(rc), (int)rc); - - result = rc; + if (sess) { + sess->conn = NULL; } } @@ -205,7 +202,9 @@ server_client_queue_disconnect(server_state_t *s, struct client_conn_t *cl, faith_envelope_t envl = {0}; envl.type = FAITH_ENVELOPE_CLIENT_DISCONNECT; - envl.recipient_id = cl->auth_id; + envl.recipient_id = cl->state == CLIENT_WAIT_FOR_DEVICE_LINK_RESPONSE + ? cl->pending_auth_id + : cl->ident.auth_id; envl.body = body; envl.body_size = body_size; diff --git a/faithd/src/server/server.h b/faithd/src/server/server.h index 0c7efff..6eab6c7 100644 --- a/faithd/src/server/server.h +++ b/faithd/src/server/server.h @@ -24,9 +24,8 @@ typedef enum { } client_state_t; typedef struct client_conn_t { - // connection id - faith_auth_id_t auth_id; - faith_device_id_t device_id; + application_user_identity_t ident; + faith_auth_id_t pending_auth_id; transport_conn_t conn; reactor_source_t reactor_source; diff --git a/faithd/src/server/sess_registry.c b/faithd/src/server/sess_registry.c index b8f9425..a01862c 100644 --- a/faithd/src/server/sess_registry.c +++ b/faithd/src/server/sess_registry.c @@ -69,6 +69,8 @@ faith_status_code_t sess_registry_register_session( device_event_inbox_init(&sess->inbox); /* assign public key to new session */ + sess->ident.auth_id = cl->ident.auth_id; + sess->ident.device_id = cl->ident.device_id; memcpy(sess->ident.public_key, public_key, FAITH_ED25519_PUBLIC_KEY_SIZE); /* insert session data at */ diff --git a/faithd/src/server/sess_registry.h b/faithd/src/server/sess_registry.h index dbdac63..ae300b7 100644 --- a/faithd/src/server/sess_registry.h +++ b/faithd/src/server/sess_registry.h @@ -1,19 +1,14 @@ #pragma once +#include "../application/user.h" #include "../auth/structs.h" #include "../core/core.h" #include "../delivery/event_inbox.h" typedef struct { - uint8_t public_key[FAITH_ED25519_PUBLIC_KEY_SIZE]; -} client_identity_t; - -struct client_conn_t; - -typedef struct { - struct client_conn_t *conn; - client_identity_t ident; - device_event_inbox_t inbox; + struct client_conn_t *conn; + application_user_identity_t ident; + device_event_inbox_t inbox; } client_device_session_data_t; typedef struct {