ascii-chat 0.11.33
Video chat in your terminal
Loading...
Searching...
No Matches
client.h File Reference

Per-client state management and lifecycle orchestration. More...

Go to the source code of this file.

Data Structures

struct  client_manager_t
 Global client manager structure for server-side client coordination. More...
 

Typedefs

typedef struct server_context_t server_context_t
 

Functions

client_info_t * add_client (server_context_t *server_ctx, socket_t socket, const char *client_ip, int port)
 
client_info_t * add_webrtc_client (server_context_t *server_ctx, acip_transport_t *transport, const char *client_ip, bool start_threads)
 Register a WebRTC client with the server.
 
int start_webrtc_client_threads (server_context_t *server_ctx, const char *client_id)
 Start threads for a WebRTC client after crypto initialization.
 
int remove_client (server_context_t *server_ctx, const char *client_id)
 
client_info_t * find_client_by_id (const char *client_id)
 Fast O(1) client lookup by ID using hash table.
 
client_info_t * find_client_by_socket (socket_t socket)
 Find client by socket descriptor using linear search.
 
void cleanup_client_media_buffers (client_info_t *client)
 
void cleanup_client_packet_queues (client_info_t *client)
 
void * client_receive_thread (void *arg)
 
void stop_client_threads (client_info_t *client)
 
int process_encrypted_packet (client_info_t *client, packet_type_t *type, void **data, size_t *len, uint32_t *sender_id)
 
void process_decrypted_packet (client_info_t *client, packet_type_t type, void *data, size_t len)
 
void initialize_client_info (client_info_t *client)
 

Variables

client_manager_t g_client_manager
 Global client manager singleton - central coordination point.
 
rwlock_t g_client_manager_rwlock
 Reader-writer lock protecting the global client manager.
 
bool g_client_manager_rwlock_initialized
 

Detailed Description

Per-client state management and lifecycle orchestration.

This header provides server-specific client management functions. The client_info_t structure and network logging macros are defined in lib/network/client.h.

Definition in file src/server/client.h.

Typedef Documentation

◆ server_context_t

Definition at line 16 of file src/server/client.h.

Function Documentation

◆ add_client()

client_info_t * add_client ( server_context_t *  server_ctx,
socket_t  socket,
const char *  client_ip,
int  port 
)

Definition at line 656 of file src/server/client.c.

656 {
657 log_info("[TCP_ADD_CLIENT] ENTER: client_ip=%s, port=%d", client_ip, port);
658 // Find empty slot WITHOUT holding the global lock
659 // We'll re-verify under lock after allocations complete
660 int slot = -1;
661 int existing_count = 0;
662 for (int i = 0; i < MAX_CLIENTS; i++) {
663 if (slot == -1 && g_client_manager.clients[i].client_id[0] == '\0') {
664 slot = i; // Take first available slot
665 }
667 existing_count++;
668 }
669 }
670
671 // Quick pre-check before expensive allocations
672 log_info("[TCP_DBG] SLOT_CHECK: existing_count=%d, max_clients=%d, slot=%d", existing_count, GET_OPTION(max_clients),
673 slot);
674 if (existing_count >= GET_OPTION(max_clients) || slot == -1) {
675 const char *reject_msg = "SERVER_FULL: Maximum client limit reached\n";
676 ssize_t send_result = socket_send(socket, reject_msg, strlen(reject_msg), 0);
677 if (send_result < 0) {
678 log_warn("Failed to send rejection message to client: %s", SAFE_STRERROR(errno));
679 }
680 return NULL;
681 }
682
683 // NOW acquire the lock early for generating unique client name
684 log_info("[TCP_DBG] LOCK_ACQUIRE_START");
686 log_info("[TCP_DBG] LOCK_ACQUIRED");
687
688 // Generate unique client ID (noun only, without counter or port)
689 char new_client_id[MAX_CLIENT_ID_LEN];
690 log_info("[TCP_DBG] GENERATE_ID_START");
691 if (generate_client_id(new_client_id, sizeof(new_client_id)) != 0) {
693 log_error("Failed to generate unique client ID");
694 return NULL;
695 }
696 log_info("[TCP_DBG] GENERATE_ID_DONE: id=%s", new_client_id);
697
698 log_info("[TCP_DBG] LOCK_RELEASE_START");
700 log_info("[TCP_DBG] LOCK_RELEASED");
701
702 // DO EXPENSIVE ALLOCATIONS OUTSIDE THE LOCK
703 // This prevents blocking frame processing during client initialization
704 log_info("[TCP_DBG] VIDEO_CREATE_START: About to create video frame buffer for client_id=%s", new_client_id);
705 video_frame_buffer_t *incoming_video_buffer = video_frame_buffer_create(new_client_id);
706 log_info("[TCP_DBG] VIDEO_CREATE_DONE: Video buffer creation returned %p", (void *)incoming_video_buffer);
707 if (!incoming_video_buffer) {
708 SET_ERRNO(ERROR_MEMORY, "Failed to create video buffer for client %s", new_client_id);
709 log_error("Failed to create video buffer for client %s", new_client_id);
710 return NULL;
711 }
712
713 log_info("[TCP_DBG] VIDEO_DONE: About to create audio buffer");
715 log_info("[TCP_DBG] AUDIO_DONE: Audio buffer created");
716 if (!incoming_audio_buffer) {
717 SET_ERRNO(ERROR_MEMORY, "Failed to create audio buffer for client %s", new_client_id);
718 log_error("Failed to create audio buffer for client %s", new_client_id);
719 video_frame_buffer_destroy(incoming_video_buffer);
720 return NULL;
721 }
722
723 log_info("[TCP_DBG] AUDIO_QUEUE_START: Creating audio queue");
724 // Keep the per-client queue allocation independent. Registering a 1000-node
725 // pool for every client can block a second client while the first client's
726 // debug registrations are being retired during disconnect.
727 packet_queue_t *audio_queue = packet_queue_create(500);
728 log_info("[TCP_DBG] AUDIO_QUEUE_DONE");
729 if (!audio_queue) {
730 LOG_ERRNO_IF_SET("Failed to create audio queue for client");
731 audio_ring_buffer_destroy(incoming_audio_buffer);
732 video_frame_buffer_destroy(incoming_video_buffer);
733 return NULL;
734 }
735
736 log_info("[TCP_DBG] OUTGOING_VID_START: Creating outgoing video buffer");
737 video_frame_buffer_t *outgoing_video_buffer = video_frame_buffer_create(new_client_id);
738 log_info("[TCP_DBG] OUTGOING_VID_DONE");
739 if (!outgoing_video_buffer) {
740 LOG_ERRNO_IF_SET("Failed to create outgoing video buffer for client");
741 packet_queue_destroy(audio_queue);
742 audio_ring_buffer_destroy(incoming_audio_buffer);
743 video_frame_buffer_destroy(incoming_video_buffer);
744 return NULL;
745 }
746
747 void *send_buffer = SAFE_MALLOC_ALIGNED(MAX_FRAME_BUFFER_SIZE, 64, void *);
748 if (!send_buffer) {
749 log_error("Failed to allocate send buffer for client %s", new_client_id);
750 video_frame_buffer_destroy(outgoing_video_buffer);
751 packet_queue_destroy(audio_queue);
752 audio_ring_buffer_destroy(incoming_audio_buffer);
753 video_frame_buffer_destroy(incoming_video_buffer);
754 return NULL;
755 }
756
757 // NOW acquire the lock for the critical section: slot assignment + registration
758 log_info("[TCP_DBG] LOCK_START: About to acquire client manager write lock");
760 log_info("[TCP_DBG] LOCK_ACQUIRED: Client manager lock acquired");
761 log_info("[TCP_DBG] SLOT_RECHECK: About to check if slot is still available");
762
763 // Re-check slot availability under lock (another thread might have taken it)
764 if (g_client_manager.clients[slot].client_id[0] != '\0') {
766 SAFE_FREE(send_buffer);
767 video_frame_buffer_destroy(outgoing_video_buffer);
768 packet_queue_destroy(audio_queue);
769 audio_ring_buffer_destroy(incoming_audio_buffer);
770 video_frame_buffer_destroy(incoming_video_buffer);
771
772 const char *reject_msg = "SERVER_FULL: Slot reassigned, try again\n";
773 socket_send(socket, reject_msg, strlen(reject_msg), 0);
774 return NULL;
775 }
776
777 // Now we have exclusive access to the slot - do the actual registration
778 client_info_t *client = &g_client_manager.clients[slot];
779 memset(client, 0, sizeof(client_info_t));
780
781 // Initialize codec capabilities to zero (no codecs supported until CLIENT_CAPABILITIES received)
782 client->codec_capabilities_video = 0;
783 client->codec_capabilities_audio = 0;
784
785 client->socket = socket;
786 client->is_tcp_client = true;
787 SAFE_STRNCPY(client->client_id, new_client_id, sizeof(client->client_id) - 1);
788 SAFE_STRNCPY(client->client_ip, client_ip, sizeof(client->client_ip) - 1);
789 client->port = port;
790 atomic_store_bool(&client->active, true);
791 client->server_ctx = server_ctx; // Store server context for cleanup
792 atomic_store_bool(&client->shutting_down, false);
795 client->connected_at = time(NULL);
796
797 memset(&client->crypto_handshake_ctx, 0, sizeof(client->crypto_handshake_ctx));
798 client->crypto_initialized = false;
799 client->transport_encrypted = false;
800
801 client->pending_packet_type = 0;
802 client->pending_packet_payload = NULL;
803 client->pending_packet_length = 0;
804
805 // Assign pre-allocated buffers
806 client->incoming_video_buffer = incoming_video_buffer;
807 client->incoming_audio_buffer = incoming_audio_buffer;
808 client->audio_queue = audio_queue;
809 client->outgoing_video_buffer = outgoing_video_buffer;
810 client->send_buffer = send_buffer;
812
813 // Generate unique noun-based client name with transport type and port
814 // Note: We're holding the write lock here, so it's safe to iterate existing clients
815 log_info("[TCP_DBG] BEFORE_GENERATE_CLIENT_NAME");
817 true /* is_tcp */) != 0) {
818 // Fallback to numeric name if generation fails
819 safe_snprintf(client->display_name, sizeof(client->display_name), "client_%u (tcp:%d)", new_client_id, port);
820 }
821 log_info("[TCP_DBG] AFTER_GENERATE_CLIENT_NAME: display_name=%s", client->display_name);
822
823 // Register client with named debug system using the generated name
824 (void)NAMED_REGISTER_CLIENT(client, client->display_name, NULL);
825
826 log_info("Added new client %s from %s:%d (socket=%d, slot=%d)", new_client_id, client_ip, port, socket, slot);
827 log_debug("Client slot assigned: client_id=%s assigned to slot %d, socket=%d", new_client_id, slot, socket);
828
829 // Register socket with tcp_server
830 asciichat_error_t reg_result = tcp_server_add_client(server_ctx->tcp_server, socket, client);
831 if (reg_result != ASCIICHAT_OK) {
832 SET_ERRNO(ERROR_INTERNAL, "Failed to register client socket with tcp_server");
833 log_error("Failed to register client %s socket with tcp_server", new_client_id);
834 // Don't unlock here - error_cleanup will do it
835 goto error_cleanup;
836 }
837
839 log_debug("Client count updated: now %d clients (added client_id=%s to slot %d)", g_client_manager.client_count,
840 new_client_id, slot);
841
842 HASH_ADD_STR(g_client_manager.clients_by_id, client_id, client);
843 log_debug("Added client %s to uthash table", new_client_id);
844
846
847 // Configure socket OUTSIDE lock
848 configure_client_socket(socket, new_client_id);
849
850 // Initialize mutexes OUTSIDE lock BEFORE registering them
851 // register_client_info_atomics needs these mutexes to be initialized first
852 if (mutex_init(&client->client_state_mutex, "client_state") != 0) {
853 log_error("Failed to initialize client state mutex for client %s", new_client_id);
854 remove_client(server_ctx, new_client_id);
855 return NULL;
856 }
857
858 if (mutex_init(&client->send_mutex, "client_send") != 0) {
859 log_error("Failed to initialize send mutex for client %s", new_client_id);
860 remove_client(server_ctx, new_client_id);
861 return NULL;
862 }
863
864 // Register all atomic fields for debug tracking - DISABLED to prevent deadlock
865 // register_client_info_atomics() and audio_ring_buffer_register_atomics() call
866 // NAMED_REGISTER which tries to acquire WRITE lock on named_registry.entries_lock
867 // The debug_sync thread holds READ lock while printing, causing reader-writer deadlock
868 // These are debug-only features not critical for functionality
869 // log_info("[TCP_DBG] BEFORE_REGISTER_CLIENT_ATOMICS");
870 // register_client_info_atomics(client);
871 // log_info("[TCP_DBG] AFTER_REGISTER_CLIENT_ATOMICS");
872
873 // Audio buffer atomics registration disabled - same deadlock reason
874 // log_info("[TCP_DBG] AUDIO_BUFFER_REGISTER_ATOMICS_START: client=%p", (void *)client);
875 // if (client->incoming_audio_buffer) {
876 // audio_ring_buffer_register_atomics(client->incoming_audio_buffer, new_client_id);
877 // log_info("[TCP_DBG] AUDIO_BUFFER_REGISTER_ATOMICS_DONE");
878 // } else {
879 // log_info("[TCP_DBG] AUDIO_BUFFER_REGISTER_ATOMICS_SKIPPED: incoming_audio_buffer is NULL");
880 // }
881
882 // Register with audio mixer OUTSIDE lock
883 // CRITICAL: Do this BEFORE crypto handshake to avoid deadlock with mixer thread
884 log_info("[TCP_DBG] MIXER_ADD_SOURCE_START: g_audio_mixer=%p, client->incoming_audio_buffer=%p",
885 (void *)g_audio_mixer, (void *)client->incoming_audio_buffer);
886 if (g_audio_mixer && client->incoming_audio_buffer) {
887 log_info("[TCP_DBG] MIXER_ADD_SOURCE_CALLING: client_id=%s", new_client_id);
888 if (mixer_add_source(g_audio_mixer, new_client_id, client->incoming_audio_buffer) < 0) {
889 log_warn("Failed to add client %s to audio mixer", new_client_id);
890 } else {
891 log_info("[TCP_DBG] MIXER_ADD_SOURCE_SUCCESS: client_id=%s", new_client_id);
892#ifdef DEBUG_AUDIO
893 log_debug("Added client %s to audio mixer", new_client_id);
894#endif
895 }
896 } else {
897 log_info("[TCP_DBG] MIXER_ADD_SOURCE_SKIPPED: conditions not met");
898 }
899 log_info("[TCP_DBG] MIXER_ADD_SOURCE_DONE");
900
901 // Create the client's persistent transport early (before handshake).
902 // It starts with NULL crypto; server_crypto_handshake() sets crypto_ctx after success.
903 client->transport = acip_tcp_transport_create(new_client_id, socket, NULL);
904 if (!client->transport) {
905 log_error("Failed to create ACIP transport for client %s", new_client_id);
906 if (remove_client(server_ctx, new_client_id) != 0) {
907 log_error("Failed to remove client after transport creation failure");
908 }
909 return NULL;
910 }
911 log_debug("Created ACIP transport for client %s (crypto pending handshake)", new_client_id);
912
913 // Perform crypto handshake before starting threads.
914 // This ensures the handshake uses the socket directly without interference from receive thread.
915 // The handshake uses client->transport and sets its crypto_ctx on success.
916 log_info("[TCP_DBG] CRYPTO_INIT_START: About to call server_crypto_init()");
917 if (server_crypto_init() == 0) {
918 log_info("[TCP_DBG] CRYPTO_INIT_DONE: server_crypto_init() returned 0");
919 // Set timeout for crypto handshake to prevent indefinite blocking
920 // This prevents clients from connecting but never completing the handshake
921 const uint64_t HANDSHAKE_TIMEOUT_NS = 30ULL * NS_PER_SEC_INT;
922 asciichat_error_t timeout_result = set_socket_timeout(socket, HANDSHAKE_TIMEOUT_NS);
923 if (timeout_result != ASCIICHAT_OK) {
924 log_warn("Failed to set handshake timeout for client %s: %s", new_client_id,
925 asciichat_error_string(timeout_result));
926 // Continue anyway - timeout is a safety feature, not critical
927 }
928
929 log_info("[TCP_DBG] SERVER_CRYPTO_HANDSHAKE_START: About to call server_crypto_handshake()");
930 int crypto_result = server_crypto_handshake(client);
931 log_info("[TCP_DBG] SERVER_CRYPTO_HANDSHAKE_DONE: result=%d", crypto_result);
932 if (crypto_result != 0) {
933 log_error("Crypto handshake failed for client %s: %s", new_client_id, network_error_string());
934 if (remove_client(server_ctx, new_client_id) != 0) {
935 log_error("Failed to remove client after crypto handshake failure");
936 }
937 return NULL;
938 }
939
940 // Clear socket timeout after handshake completes successfully
941 // This allows normal operation without timeouts on data transfer
942 log_info("[TCP_DBG] CLEAR_SOCKET_TIMEOUT_START");
943 asciichat_error_t clear_timeout_result = set_socket_timeout(socket, 0);
944 if (clear_timeout_result != ASCIICHAT_OK) {
945 log_warn("Failed to clear handshake timeout for client %s: %s", new_client_id,
946 asciichat_error_string(clear_timeout_result));
947 // Continue anyway - we can still communicate even with timeout set
948 }
949
950 log_debug("Crypto handshake completed successfully for client %s", new_client_id);
951
952 // After handshake completes, the client immediately sends PACKET_TYPE_CLIENT_CAPABILITIES
953 // We must read and process this packet BEFORE starting the receive thread to avoid a race condition
954 // where the packet arrives but no thread is listening for it.
955 //
956 // SPECIAL CASE: If client used --no-encrypt, we already received this packet during the handshake
957 // attempt and stored it in pending_packet_*. Use that instead of receiving a new one.
958 packet_envelope_t envelope;
959 bool used_pending_packet = false;
960
961 if (client->pending_packet_payload) {
962 // Client used --no-encrypt mode - use the packet we already received
963 log_info("Client %s using --no-encrypt mode - processing pending packet type %u", client->client_id,
964 client->pending_packet_type);
965 envelope.type = client->pending_packet_type;
966 envelope.data = client->pending_packet_payload;
967 envelope.len = client->pending_packet_length;
968 envelope.allocated_buffer = client->pending_packet_payload; // Will be freed below
969 envelope.allocated_size = client->pending_packet_length;
970 used_pending_packet = true;
971
972 // Clear pending packet fields
973 client->pending_packet_type = 0;
974 client->pending_packet_payload = NULL;
975 client->pending_packet_length = 0;
976 } else {
977 // Normal encrypted mode - receive capabilities packet
978 log_debug("Waiting for initial capabilities packet from client %s", client->client_id);
979
980 // Protect crypto context access with client state mutex
982 const crypto_context_t *crypto_ctx = crypto_server_get_context(client->client_id);
983
984 // Use per-client crypto state to determine enforcement
985 // At this point, handshake is complete, so crypto_initialized=true and handshake is ready
986 bool enforce_encryption = !GET_OPTION(no_encrypt) && client->crypto_initialized &&
988
989 packet_recv_result_t result = receive_packet_secure(socket, (void *)crypto_ctx, enforce_encryption, &envelope);
991
992 if (result != PACKET_RECV_SUCCESS) {
993 log_error("Failed to receive initial capabilities packet from client %s: result=%d", client->client_id, result);
994 if (envelope.allocated_buffer) {
995 buffer_pool_free(NULL, envelope.allocated_buffer, envelope.allocated_size);
996 }
997 if (remove_client(server_ctx, client->client_id) != 0) {
998 log_error("Failed to remove client after crypto handshake failure");
999 }
1000 return NULL;
1001 }
1002 }
1003
1004 if (envelope.type != PACKET_TYPE_CLIENT_CAPABILITIES) {
1005 log_error("Expected PACKET_TYPE_CLIENT_CAPABILITIES but got packet type %d from client %s", envelope.type,
1006 client->client_id);
1007 if (envelope.allocated_buffer) {
1008 buffer_pool_free(NULL, envelope.allocated_buffer, envelope.allocated_size);
1009 }
1010 if (remove_client(server_ctx, client->client_id) != 0) {
1011 log_error("Failed to remove client after crypto handshake failure");
1012 }
1013 return NULL;
1014 }
1015
1016 // The unencrypted socket receive path retains the wire header in its
1017 // allocation. Packet handlers always consume payload-only data, matching
1018 // the encrypted and ACIP-dispatch paths.
1019 if (envelope.data == envelope.allocated_buffer) {
1020 const packet_header_t *header = (const packet_header_t *)envelope.allocated_buffer;
1021 envelope.type = (packet_type_t)NET_TO_HOST_U16(header->type);
1022 envelope.data = (uint8_t *)envelope.allocated_buffer + sizeof(*header);
1023 envelope.len = NET_TO_HOST_U32(header->length);
1024 }
1025
1026 // Process the capabilities packet directly
1027 log_debug("Processing initial capabilities packet from client %s (from %s)", client->client_id,
1028 used_pending_packet ? "pending packet" : "network");
1029 log_info("INITIAL_CAPABILITIES: client_id=%s source=%s type=%u payload_len=%zu allocated_size=%zu data=%p",
1030 client->client_id, used_pending_packet ? "pending" : "network", envelope.type, envelope.len,
1031 envelope.allocated_size, envelope.data);
1032 handle_client_capabilities_packet(client, envelope.data, envelope.len);
1033
1034 // Free the packet data
1035 if (envelope.allocated_buffer) {
1036 buffer_pool_free(NULL, envelope.allocated_buffer, envelope.allocated_size);
1037 }
1038 log_debug("Successfully received and processed initial capabilities for client %s", client->client_id);
1039 }
1040
1041 // Start all client threads in the correct order (unified path for TCP and WebRTC)
1042 // This creates: receive thread -> render threads -> send thread
1043 // The render threads MUST be created before send thread to avoid the race condition
1044 // where send thread reads empty frames before render thread generates the first real frame
1045 char client_id_snapshot[MAX_CLIENT_ID_LEN];
1046 SAFE_STRNCPY(client_id_snapshot, client->client_id, sizeof(client_id_snapshot) - 1);
1047 if (start_client_threads(server_ctx, client, true) != 0) {
1048 log_error("Failed to start threads for TCP client %s", client_id_snapshot);
1049 // Client is already in hash table - use remove_client for proper cleanup
1050 remove_client(server_ctx, client_id_snapshot);
1051 return NULL;
1052 }
1053 log_debug("Successfully created render threads for client %s", client_id_snapshot);
1054
1055 // Register client with session_host (for discovery mode support)
1056 if (server_ctx->session_host) {
1057 client->session_client_id = session_host_add_client(server_ctx->session_host, socket, client_ip, port);
1058 if (client->session_client_id == 0) {
1059 log_warn("Failed to register client %s with session_host", client_id_snapshot);
1060 } else {
1061 log_debug("Client %s registered with session_host as %u", client_id_snapshot, client->session_client_id);
1062 }
1063 }
1064
1065 // Broadcast server state to ALL clients AFTER the new client is fully set up
1066 // This notifies all clients (including the new one) about the updated grid
1067 log_info("[TCP_DBG] BROADCAST_SERVER_STATE_START: About to broadcast to all clients");
1069 log_info("[TCP_DBG] BROADCAST_SERVER_STATE_DONE");
1070
1071 // Register client atomics for debug sync inspection
1072 register_client_info_atomics(client);
1073
1074 log_info("[TCP_DBG] ADD_CLIENT_RETURNING: client_id=%s, client=%p", client->client_id, (void *)client);
1075 return client;
1076
1077error_cleanup:
1078 // Clean up all partially allocated resources
1079 // NOTE: This label is reached when allocation or initialization fails BEFORE
1080 // the client is added to the hash table. Don't call remove_client() here.
1081 cleanup_client_all_buffers(client);
1083 return NULL;
1084}
#define LOG_ERRNO_IF_SET(message)
Check if any error occurred and log it if so.
void atomic_store_bool(atomic_t *a, bool value)
Atomically store a boolean value.
Definition atomic.c:177
bool atomic_load_bool(atomic_t *a)
Atomically load a boolean value.
Definition atomic.c:169
void atomic_store_u64(atomic_t *a, uint64_t value)
Atomically store a uint64_t value.
Definition atomic.c:241
bool crypto_handshake_is_ready(const crypto_handshake_context_t *ctx)
Check if handshake is complete and encryption is ready.
#define NET_TO_HOST_U16(val)
Definition endian.h:111
#define NET_TO_HOST_U32(val)
Definition endian.h:81
void audio_ring_buffer_destroy(audio_ring_buffer_t *rb)
Destroy an audio ring buffer.
audio_ring_buffer_t * audio_ring_buffer_create_for_capture(void)
Create a new audio ring buffer for capture (without jitter buffering)
int mixer_add_source(mixer_t *mixer, const char *client_id, audio_ring_buffer_t *buffer)
Add an audio source to the mixer.
Definition mixer.c:424
void buffer_pool_free(buffer_pool_t *pool, const void *data, size_t size)
Free a buffer back to the pool (lock-free)
#define MAX_FRAME_BUFFER_SIZE
Maximum frame buffer size including headers and compression.
#define SAFE_MALLOC_ALIGNED(size, alignment, cast)
Definition common.h:349
#define SAFE_STRNCPY(dst, src, size)
Definition common.h:414
#define SAFE_FREE(ptr)
Definition common.h:376
#define SAFE_STRERROR(errnum)
Definition common.h:465
unsigned long long uint64_t
Definition common.h:59
unsigned char uint8_t
Definition common.h:56
#define NAMED_REGISTER_CLIENT(client, name, parent_ptr)
Register a client connection with automatic format specifier.
#define SET_ERRNO(code, context_msg,...)
Set error code with custom context message and log it, returning the error code.
asciichat_error_t
Error and exit codes - unified status values (0-255)
Definition error_codes.h:49
@ ERROR_MEMORY
Definition error_codes.h:56
@ ASCIICHAT_OK
Definition error_codes.h:51
@ ERROR_INTERNAL
Definition error_codes.h:92
#define MAX_CLIENTS
Maximum possible clients (static array size) - actual runtime limit set by –max-clients (1-32)
Definition limits.h:26
#define MAX_CLIENT_ID_LEN
Maximum client ID length - format: "noun.N (transport:port)" e.g., "mountain.0 (tcp:48036)".
Definition limits.h:23
#define log_warn(...)
Log a WARN message.
Definition log/log.h:574
#define log_error(...)
Log an ERROR message.
Definition log/log.h:587
#define log_info(...)
Log an INFO message.
Definition log/log.h:561
#define log_debug(...)
Log a DEBUG message.
Definition log/log.h:548
#define NS_PER_SEC_INT
Definition time.h:157
asciichat_error_t set_socket_timeout(socket_t sockfd, uint64_t timeout_ns)
Set socket timeout.
packet_recv_result_t receive_packet_secure(socket_t sockfd, void *crypto_ctx, bool enforce_encryption, packet_envelope_t *envelope)
Receive a packet with decryption and decompression support.
Definition packet.c:569
const char * network_error_string()
Get human-readable error string for network errors.
packet_recv_result_t
Packet reception result codes.
Definition packet.h:1152
@ PACKET_RECV_SUCCESS
Packet received successfully.
Definition packet.h:1154
#define GET_OPTION(field)
Safely get a specific option field (lock-free read)
packet_queue_t * packet_queue_create(size_t max_size)
Create a new packet queue.
Definition queue.c:125
void packet_queue_destroy(packet_queue_t *queue)
Signal queue shutdown (causes dequeue to return NULL)
Definition queue.c:173
packet_type_t
Network protocol packet type enumeration.
Definition packet.h:286
@ PACKET_TYPE_CLIENT_CAPABILITIES
Client reports terminal capabilities.
Definition packet.h:381
#define rwlock_wrunlock(lock)
Release a write lock (with debug tracking in debug builds)
Definition rwlock.h:349
int safe_snprintf(char *buffer, size_t buffer_size, const char *format,...)
Safe formatted string printing to buffer.
Definition system.c:148
#define mutex_lock(mutex)
Lock a mutex (with debug tracking in debug builds)
int mutex_init(mutex_t *mutex, const char *name)
Initialize a mutex with a name.
Definition threading.c:16
ssize_t socket_send(socket_t sock, const void *buf, size_t len, int flags)
Send data on a socket.
#define rwlock_wrlock(lock)
Acquire a write lock (with debug tracking in debug builds)
Definition rwlock.h:313
int errno
#define mutex_unlock(mutex)
Unlock a mutex (with debug tracking in debug builds)
uint32_t session_host_add_client(session_host_t *host, socket_t socket, const char *ip, int port)
Add a client from an accepted socket.
Definition host.c:1157
void video_frame_buffer_destroy(video_frame_buffer_t *vfb)
Destroy frame buffer and free all resources.
video_frame_buffer_t * video_frame_buffer_create(const char *client_id)
Create a double-buffered video frame manager.
Definition video_frame.c:17
asciichat_error_t tcp_server_add_client(tcp_server_t *server, socket_t socket, void *client_data)
Add client to registry.
int generate_client_id(char *buffer, size_t buffer_size)
Generate a unique client ID (noun only, no counter or port)
Definition nouns.c:780
int generate_client_name(char *buffer, size_t buffer_size, void *existing_clients_hash, int port, bool is_tcp)
Generate a unique client name from the nouns wordlist.
Definition nouns.c:809
mixer_t *volatile g_audio_mixer
Global audio mixer instance for multi-client audio processing.
void handle_client_capabilities_packet(client_info_t *client, const void *data, size_t len)
Process CLIENT_CAPABILITIES packet - configure client-specific rendering.
rwlock_t g_client_manager_rwlock
Reader-writer lock protecting the global client manager.
client_manager_t g_client_manager
Global client manager singleton - central coordination point.
int remove_client(server_context_t *server_ctx, const char *client_id)
void broadcast_server_state_to_all_clients(void)
Notify all clients of state changes.
const crypto_context_t * crypto_server_get_context(const char *client_id)
int server_crypto_init(void)
int server_crypto_handshake(client_info_t *client)
Audio ring buffer for real-time audio streaming.
Definition ringbuffer.h:205
Per-client state structure for server-side client management.
bool transport_encrypted
TLS or DTLS already encrypts this client's transport.
char display_name[MAX_DISPLAY_NAME_LEN]
video_frame_buffer_t * outgoing_video_buffer
crypto_handshake_context_t crypto_handshake_ctx
video_frame_buffer_t * incoming_video_buffer
audio_ring_buffer_t * incoming_audio_buffer
char client_ip[INET_ADDRSTRLEN]
char client_id[MAX_CLIENT_ID_LEN]
client_info_t * clients_by_id
uthash head pointer for O(1) client_id -> client_info_t* lookups
client_info_t clients[MAX_CLIENTS]
Array of client_info_t structures (backing storage)
int client_count
Current number of active clients.
Cryptographic context structure.
Packet envelope containing received packet data.
Definition packet.h:1129
void * allocated_buffer
Buffer that needs to be freed by caller (may be NULL if not allocated)
Definition packet.h:1139
size_t allocated_size
Size of allocated buffer in bytes.
Definition packet.h:1141
size_t len
Length of payload data in bytes.
Definition packet.h:1135
void * data
Packet payload data (decrypted and decompressed if applicable)
Definition packet.h:1133
packet_type_t type
Packet type (from packet_types.h)
Definition packet.h:1131
Network packet header structure.
Definition packet.h:598
Thread-safe packet queue for producer-consumer communication.
Definition queue.h:211
tcp_server_t * tcp_server
TCP server managing connections.
session_host_t * session_host
Session host for discovery mode support.
Video frame buffer manager.
acip_transport_t * acip_tcp_transport_create(const char *name, socket_t sockfd, crypto_context_t *crypto_ctx)
Create TCP transport from existing socket.

References acip_tcp_transport_create(), client_info::active, packet_envelope_t::allocated_buffer, packet_envelope_t::allocated_size, ASCIICHAT_OK, atomic_load_bool(), atomic_store_bool(), atomic_store_u64(), client_info::audio_queue, audio_ring_buffer_create_for_capture(), audio_ring_buffer_destroy(), broadcast_server_state_to_all_clients(), buffer_pool_free(), client_manager_t::client_count, client_info::client_id, client_info::client_ip, client_info::client_state_mutex, client_manager_t::clients, client_manager_t::clients_by_id, client_info::codec_capabilities_audio, client_info::codec_capabilities_video, client_info::connected_at, client_info::crypto_handshake_ctx, crypto_handshake_is_ready(), client_info::crypto_initialized, crypto_server_get_context(), packet_envelope_t::data, client_info::display_name, errno, ERROR_INTERNAL, ERROR_MEMORY, g_audio_mixer, g_client_manager, g_client_manager_rwlock, generate_client_id(), generate_client_name(), GET_OPTION, handle_client_capabilities_packet(), client_info::incoming_audio_buffer, client_info::incoming_video_buffer, client_info::is_tcp_client, client_info::last_rendered_grid_sources, client_info::last_sent_grid_sources, packet_envelope_t::len, packet_header_t::length, log_debug, LOG_ERRNO_IF_SET, log_error, log_info, log_warn, MAX_CLIENT_ID_LEN, MAX_CLIENTS, MAX_FRAME_BUFFER_SIZE, mixer_add_source(), mutex_init(), mutex_lock, mutex_unlock, NAMED_REGISTER_CLIENT, NET_TO_HOST_U16, NET_TO_HOST_U32, network_error_string(), NS_PER_SEC_INT, client_info::outgoing_video_buffer, packet_queue_create(), packet_queue_destroy(), PACKET_RECV_SUCCESS, PACKET_TYPE_CLIENT_CAPABILITIES, client_info::pending_packet_length, client_info::pending_packet_payload, client_info::pending_packet_type, client_info::port, receive_packet_secure(), remove_client(), rwlock_wrlock, rwlock_wrunlock, SAFE_FREE, SAFE_MALLOC_ALIGNED, safe_snprintf(), SAFE_STRERROR, SAFE_STRNCPY, client_info::send_buffer, client_info::send_buffer_size, client_info::send_mutex, server_crypto_handshake(), server_crypto_init(), client_info::server_ctx, client_info::session_client_id, server_context_t::session_host, session_host_add_client(), SET_ERRNO, set_socket_timeout(), client_info::shutting_down, client_info::socket, socket_send(), server_context_t::tcp_server, tcp_server_add_client(), client_info::transport, client_info::transport_encrypted, packet_header_t::type, packet_envelope_t::type, video_frame_buffer_create(), and video_frame_buffer_destroy().

◆ add_webrtc_client()

client_info_t * add_webrtc_client ( server_context_t *  server_ctx,
acip_transport_t *  transport,
const char *  client_ip,
bool  start_threads 
)

Register a WebRTC client with the server.

Registers a client that connected via WebRTC data channel instead of TCP socket. This function reuses most of add_client() logic but skips:

  • Crypto handshake (already done via ACDS signaling)
  • Socket-specific configuration
  • TCP thread pool registration

DIFFERENCES FROM add_client():

  • Takes an already-created acip_transport_t* instead of socket
  • No crypto handshake (WebRTC signaling handled authentication)
  • No socket configuration (WebRTC handles buffering)
  • Uses generic thread spawning instead of tcp_server thread pool
Parameters
server_ctxServer context
transportWebRTC transport (already created and connected)
client_ipClient IP address for logging (may be empty for P2P)
Returns
Pointer to client_info_t on success, NULL on failure
Note
The transport must be fully initialized and ready to send/receive
Client capabilities are still expected as first packet

Definition at line 1109 of file src/server/client.c.

1110 {
1111 if (!server_ctx || !transport || !client_ip) {
1112 SET_ERRNO(ERROR_INVALID_PARAM, "Invalid parameters to add_webrtc_client");
1113 return NULL;
1114 }
1115
1117
1118 // Find empty slot - this is the authoritative check
1119 int slot = -1;
1120 int existing_count = 0;
1121 for (int i = 0; i < MAX_CLIENTS; i++) {
1122 if (slot == -1 && g_client_manager.clients[i].client_id[0] == '\0') {
1123 slot = i; // Take first available slot
1124 }
1125 // Count only active clients
1127 existing_count++;
1128 }
1129 }
1130
1131 // Check if we've hit the configured max-clients limit (not the array size)
1132 if (existing_count >= GET_OPTION(max_clients)) {
1134 SET_ERRNO(ERROR_RESOURCE_EXHAUSTED, "Maximum client limit reached (%d/%d active clients)", existing_count,
1135 GET_OPTION(max_clients));
1136 log_error("Maximum client limit reached (%d/%d active clients)", existing_count, GET_OPTION(max_clients));
1137 return NULL;
1138 }
1139
1140 if (slot == -1) {
1142 SET_ERRNO(ERROR_RESOURCE_EXHAUSTED, "No available client slots (all %d array slots are in use)", MAX_CLIENTS);
1143 log_error("No available client slots (all %d array slots are in use)", MAX_CLIENTS);
1144 return NULL;
1145 }
1146
1147 // Update client_count to match actual count before adding new client
1148 g_client_manager.client_count = existing_count;
1149
1150 // Generate unique client ID for WebRTC client (noun only, without counter or port)
1151 char new_client_id[MAX_CLIENT_ID_LEN];
1152 if (generate_client_id(new_client_id, sizeof(new_client_id)) != 0) {
1154 log_error("Failed to generate unique client ID for WebRTC client");
1155 return NULL;
1156 }
1157
1158 // Initialize client
1159 client_info_t *client = &g_client_manager.clients[slot];
1160
1161 // Free any existing buffers from previous client in this slot
1162 if (client->incoming_video_buffer) {
1164 client->incoming_video_buffer = NULL;
1165 }
1166 if (client->outgoing_video_buffer) {
1168 client->outgoing_video_buffer = NULL;
1169 }
1170 if (client->incoming_audio_buffer) {
1172 client->incoming_audio_buffer = NULL;
1173 }
1174 if (client->send_buffer) {
1175 SAFE_FREE(client->send_buffer);
1176 client->send_buffer = NULL;
1177 }
1178
1179 memset(client, 0, sizeof(client_info_t));
1180
1181 // Set up WebRTC-specific fields
1182 client->socket = INVALID_SOCKET_VALUE; // WebRTC has no traditional socket
1183 client->is_tcp_client = false; // WebRTC client - threads managed directly
1184 client->transport_encrypted = acip_transport_get_type(transport) == ACIP_TRANSPORT_WEBRTC;
1185 client->transport = transport; // Use provided transport
1186 SAFE_STRNCPY(client->client_id, new_client_id, sizeof(client->client_id) - 1);
1187 SAFE_STRNCPY(client->client_ip, client_ip, sizeof(client->client_ip) - 1);
1188 client->port = 0; // WebRTC doesn't use port numbers
1189 atomic_store_bool(&client->active, true);
1190 client->server_ctx = server_ctx; // Store server context for receive thread cleanup
1191 log_info("Added new WebRTC client %s from %s (transport=%p, slot=%d)", new_client_id, client_ip, transport, slot);
1192 atomic_store_bool(&client->shutting_down, false);
1193 atomic_store_u64(&client->last_rendered_grid_sources, 0); // Render thread updates this
1194 atomic_store_u64(&client->last_sent_grid_sources, 0); // Send thread updates this
1195 log_debug("WebRTC client slot assigned: client_id=%s assigned to slot %d", client->client_id, slot);
1196 client->connected_at = time(NULL);
1197
1198 // Initialize crypto context for this client
1199 memset(&client->crypto_handshake_ctx, 0, sizeof(client->crypto_handshake_ctx));
1200 log_debug("[ADD_WEBRTC_CLIENT] Setting crypto_initialized=false (will be set to true after handshake completes)");
1201 // Do NOT set crypto_initialized = true here! The handshake must complete first.
1202 // For WebSocket clients: websocket_client_handler will perform the handshake
1203 // For WebRTC clients: ACDS signaling may have already done crypto (future enhancement)
1204 client->crypto_initialized = false;
1205
1206 // Initialize pending packet storage (unused for WebRTC, but keep for consistency)
1207 client->pending_packet_type = 0;
1208 client->pending_packet_payload = NULL;
1209 client->pending_packet_length = 0;
1210
1211 // Use noun-based client name for display_name as well
1212 SAFE_STRNCPY(client->display_name, new_client_id, sizeof(client->display_name) - 1);
1213
1214 // Create individual video buffer for this client using modern double-buffering
1215 client->incoming_video_buffer = video_frame_buffer_create(new_client_id);
1216 if (!client->incoming_video_buffer) {
1217 SET_ERRNO(ERROR_MEMORY, "Failed to create video buffer for WebRTC client %s", new_client_id);
1218 log_error("Failed to create video buffer for WebRTC client %s", new_client_id);
1219 goto error_cleanup_webrtc;
1220 }
1221
1222 // Create individual audio buffer for this client
1224 if (!client->incoming_audio_buffer) {
1225 SET_ERRNO(ERROR_MEMORY, "Failed to create audio buffer for WebRTC client %s", new_client_id);
1226 log_error("Failed to create audio buffer for WebRTC client %s", new_client_id);
1227 goto error_cleanup_webrtc;
1228 }
1229
1230 // Create packet queues for outgoing data
1231 client->audio_queue = packet_queue_create_with_pools(500, 1000, false);
1232 if (!client->audio_queue) {
1233 LOG_ERRNO_IF_SET("Failed to create audio queue for WebRTC client");
1234 goto error_cleanup_webrtc;
1235 }
1236
1237 // Create packet queue for async dispatch: received packets waiting to be processed
1238 // Capacity of 100 packets allows buffering of ~10-50 frames (depending on fragmentation)
1239 client->received_packet_queue = packet_queue_create_with_pools(100, 500, false);
1240 if (!client->received_packet_queue) {
1241 LOG_ERRNO_IF_SET("Failed to create received packet queue for client");
1242 goto error_cleanup_webrtc;
1243 }
1244
1245 // Create outgoing video buffer for ASCII frames (double buffered, no dropping)
1246 client->outgoing_video_buffer = video_frame_buffer_create(new_client_id);
1247 if (!client->outgoing_video_buffer) {
1248 LOG_ERRNO_IF_SET("Failed to create outgoing video buffer for WebRTC client");
1249 goto error_cleanup_webrtc;
1250 }
1251
1252 // Pre-allocate send buffer to avoid malloc/free in send thread (prevents deadlocks)
1253 client->send_buffer_size = MAX_FRAME_BUFFER_SIZE; // 2MB should handle largest frames
1254 client->send_buffer = SAFE_MALLOC_ALIGNED(client->send_buffer_size, 64, void *);
1255 if (!client->send_buffer) {
1256 log_error("Failed to allocate send buffer for WebRTC client %s", new_client_id);
1257 goto error_cleanup_webrtc;
1258 }
1259
1260 g_client_manager.client_count = existing_count + 1; // We just added a client
1261 log_debug("Client count updated: now %d clients (added WebRTC client_id=%s to slot %d)",
1262 g_client_manager.client_count, new_client_id, slot);
1263
1264 // Add client to uthash table for O(1) lookup
1265 HASH_ADD_STR(g_client_manager.clients_by_id, client_id, client);
1266 log_debug("Added WebRTC client %s to uthash table", new_client_id);
1267
1268 // Release the write lock IMMEDIATELY after adding to hash table
1269 // All subsequent operations (mutex init, mixer registration) don't need global client manager lock
1271
1272 // Register this client's audio buffer with the mixer
1273 if (g_audio_mixer && client->incoming_audio_buffer) {
1274 if (mixer_add_source(g_audio_mixer, new_client_id, client->incoming_audio_buffer) < 0) {
1275 log_warn("Failed to add WebRTC client %s to audio mixer", new_client_id);
1276 } else {
1277#ifdef DEBUG_AUDIO
1278 log_debug("Added WebRTC client %s to audio mixer", new_client_id);
1279#endif
1280 }
1281 }
1282
1283 // Initialize mutexes BEFORE creating any threads to prevent race conditions
1284 if (mutex_init(&client->client_state_mutex, "client_state") != 0) {
1285 log_error("Failed to initialize client state mutex for WebRTC client %s", new_client_id);
1286 // Client is already in hash table - use remove_client for proper cleanup
1287 remove_client(server_ctx, new_client_id);
1288 return NULL;
1289 }
1290
1291 initialize_default_render_capabilities(client);
1292
1293 // Initialize send mutex to protect concurrent socket writes
1294 if (mutex_init(&client->send_mutex, "client_send") != 0) {
1295 log_error("Failed to initialize send mutex for WebRTC client %s", new_client_id);
1296 // Client is already in hash table - use remove_client for proper cleanup
1297 remove_client(server_ctx, client->client_id);
1298 return NULL;
1299 }
1300
1301 // Initialize condition variable for dispatch thread wake-up
1302 if (cond_init(&client->dispatch_queue_cond, "dispatch_queue_cond") != 0) {
1303 log_error("Failed to initialize dispatch queue condition variable for WebRTC client %s", new_client_id);
1304 // Client is already in hash table - use remove_client for proper cleanup
1305 remove_client(server_ctx, client->client_id);
1306 return NULL;
1307 }
1308
1309 // For WebRTC clients, the capabilities packet will be received by the receive thread
1310 // when it starts. Unlike TCP clients where we handle it synchronously in add_client(),
1311 // WebRTC uses the transport abstraction which handles packet reception automatically.
1312 log_debug("WebRTC client %s initialized - receive thread will process capabilities", new_client_id);
1313
1314 // Send initial server state to the new client
1315 if (send_server_state_to_client(client) != 0) {
1316 log_warn("Failed to send initial server state to WebRTC client %s", new_client_id);
1317 } else {
1318#ifdef DEBUG_NETWORK
1319 log_info("Sent initial server state to WebRTC client %s", new_client_id);
1320#endif
1321 }
1322
1323 // Send initial server state via ACIP transport
1325 uint32_t connected_count = g_client_manager.client_count;
1327
1329 state.connected_client_count = connected_count;
1330 state.active_client_count = 0; // Will be updated by broadcast thread
1331 memset(state.reserved, 0, sizeof(state.reserved));
1332
1333 // Convert to network byte order
1334 server_state_packet_t net_state;
1337 memset(net_state.reserved, 0, sizeof(net_state.reserved));
1338
1339 asciichat_error_t packet_send_result = acip_send_server_state(client->transport, &net_state);
1340 if (packet_send_result != ASCIICHAT_OK) {
1341 log_warn("Failed to send initial server state to WebRTC client %s: %s", new_client_id,
1342 asciichat_error_string(packet_send_result));
1343 } else {
1344 log_debug("Sent initial server state to WebRTC client %s: %u connected clients", new_client_id,
1346 }
1347
1348 // Register client with session_host (for discovery mode support)
1349 // WebRTC clients use INVALID_SOCKET_VALUE since they don't have a TCP socket
1350 if (server_ctx->session_host) {
1351 client->session_client_id = session_host_add_client(server_ctx->session_host, INVALID_SOCKET_VALUE, client_ip, 0);
1352 if (client->session_client_id == 0) {
1353 log_warn("Failed to register WebRTC client %s with session_host", new_client_id);
1354 } else {
1355 log_debug("WebRTC client %s registered with session_host as %u", new_client_id, client->session_client_id);
1356 }
1357 }
1358
1359 // Broadcast server state to ALL clients AFTER the new client is fully set up
1360 // This notifies all clients (including the new one) about the updated grid
1362
1363 // The browser can send CLIENT_CAPABILITIES and STREAM_START immediately when
1364 // its DataChannel opens. Start media threads only after this client is fully
1365 // registered and all initial control traffic is complete; the WebRTC
1366 // transport queues those incoming packets until the receive thread begins.
1367 // WebSocket callers still defer this until their crypto setup is complete.
1368 if (start_threads) {
1369 log_debug("[ADD_WEBRTC_CLIENT] Starting client threads after initialization for %s", new_client_id);
1370 if (start_client_threads(server_ctx, client, false) != 0) {
1371 log_error("Failed to start threads for WebRTC client %s", new_client_id);
1372 return NULL;
1373 }
1374 } else {
1375 log_debug("[ADD_WEBRTC_CLIENT] Deferring client thread startup for %s until crypto setup", new_client_id);
1376 }
1377
1378 return client;
1379
1380error_cleanup_webrtc:
1381 // Clean up all partially allocated resources for WebRTC client
1382 // NOTE: This label is reached when allocation or initialization fails BEFORE
1383 // the client is added to the hash table. Don't call remove_client() here.
1384 cleanup_client_all_buffers(client);
1386 return NULL;
1387}
#define HOST_TO_NET_U32(val)
Definition endian.h:66
unsigned int uint32_t
Definition common.h:58
@ ERROR_RESOURCE_EXHAUSTED
@ ERROR_INVALID_PARAM
packet_queue_t * packet_queue_create_with_pools(size_t max_size, size_t node_pool_size, bool use_buffer_pool)
Create a packet queue with both node and buffer pools.
Definition queue.c:133
#define INVALID_SOCKET_VALUE
Invalid socket value (POSIX: -1)
Definition socket.h:278
#define rwlock_rdlock(lock)
Acquire a read lock (with debug tracking in debug builds)
Definition rwlock.h:294
int cond_init(cond_t *cond, const char *name)
Initialize a condition variable with a name.
#define rwlock_rdunlock(lock)
Release a read lock (with debug tracking in debug builds)
Definition rwlock.h:331
asciichat_error_t acip_send_server_state(acip_transport_t *transport, const server_state_packet_t *state)
Send server state update to client (server → client)
int send_server_state_to_client(client_info_t *client)
Send current server state to a specific client.
packet_queue_t * received_packet_queue
Server state packet structure.
Definition packet.h:706
uint32_t reserved[6]
Reserved fields for future use (must be zero)
Definition packet.h:712
uint32_t active_client_count
Number of clients actively sending video/audio streams.
Definition packet.h:710
uint32_t connected_client_count
Total number of currently connected clients.
Definition packet.h:708
@ ACIP_TRANSPORT_WEBRTC
WebRTC DataChannel (P2P)
Definition transport.h:96

References acip_send_server_state(), ACIP_TRANSPORT_WEBRTC, client_info::active, server_state_packet_t::active_client_count, ASCIICHAT_OK, atomic_load_bool(), atomic_store_bool(), atomic_store_u64(), client_info::audio_queue, audio_ring_buffer_create_for_capture(), audio_ring_buffer_destroy(), broadcast_server_state_to_all_clients(), client_manager_t::client_count, client_info::client_id, client_info::client_ip, client_info::client_state_mutex, client_manager_t::clients, client_manager_t::clients_by_id, cond_init(), client_info::connected_at, server_state_packet_t::connected_client_count, client_info::crypto_handshake_ctx, client_info::crypto_initialized, client_info::dispatch_queue_cond, client_info::display_name, ERROR_INVALID_PARAM, ERROR_MEMORY, ERROR_RESOURCE_EXHAUSTED, g_audio_mixer, g_client_manager, g_client_manager_rwlock, generate_client_id(), GET_OPTION, HOST_TO_NET_U32, client_info::incoming_audio_buffer, client_info::incoming_video_buffer, INVALID_SOCKET_VALUE, client_info::is_tcp_client, client_info::last_rendered_grid_sources, client_info::last_sent_grid_sources, log_debug, LOG_ERRNO_IF_SET, log_error, log_info, log_warn, MAX_CLIENT_ID_LEN, MAX_CLIENTS, MAX_FRAME_BUFFER_SIZE, mixer_add_source(), mutex_init(), client_info::outgoing_video_buffer, packet_queue_create_with_pools(), client_info::pending_packet_length, client_info::pending_packet_payload, client_info::pending_packet_type, client_info::port, client_info::received_packet_queue, remove_client(), server_state_packet_t::reserved, rwlock_rdlock, rwlock_rdunlock, rwlock_wrlock, rwlock_wrunlock, SAFE_FREE, SAFE_MALLOC_ALIGNED, SAFE_STRNCPY, client_info::send_buffer, client_info::send_buffer_size, client_info::send_mutex, send_server_state_to_client(), client_info::server_ctx, client_info::session_client_id, server_context_t::session_host, session_host_add_client(), SET_ERRNO, client_info::shutting_down, client_info::socket, client_info::transport, client_info::transport_encrypted, video_frame_buffer_create(), and video_frame_buffer_destroy().

◆ cleanup_client_media_buffers()

void cleanup_client_media_buffers ( client_info_t *  client)

Definition at line 2873 of file src/server/client.c.

2873 {
2874 if (!client) {
2875 SET_ERRNO(ERROR_INVALID_PARAM, "Client is NULL");
2876 return;
2877 }
2878
2879 if (client->digital_rain) {
2881 client->digital_rain = NULL;
2882 }
2883
2884 if (client->incoming_video_buffer) {
2886 client->incoming_video_buffer = NULL;
2887 }
2888
2889 // Clean up outgoing video buffer (for ASCII frames)
2890 if (client->outgoing_video_buffer) {
2892 client->outgoing_video_buffer = NULL;
2893 }
2894
2895 // Clean up pre-allocated send buffer
2896 if (client->send_buffer) {
2897 SAFE_FREE(client->send_buffer);
2898 client->send_buffer = NULL;
2899 client->send_buffer_size = 0;
2900 }
2901
2902 if (client->incoming_audio_buffer) {
2904 client->incoming_audio_buffer = NULL;
2905 }
2906
2907 // Clean up Opus decoder
2908 if (client->opus_decoder) {
2910 client->opus_decoder = NULL;
2911 }
2912
2913 // Clean up H.265 decoder (persistent per-client decoder)
2914 if (client->h265_decoder_ctx) {
2915 avcodec_free_context((AVCodecContext **)&client->h265_decoder_ctx);
2916 client->h265_decoder_ctx = NULL;
2917 }
2918 if (client->h265_decode_frame) {
2919 av_frame_free((AVFrame **)&client->h265_decode_frame);
2920 client->h265_decode_frame = NULL;
2921 }
2922 if (client->h265_rgb_frame) {
2923 av_frame_free((AVFrame **)&client->h265_rgb_frame);
2924 client->h265_rgb_frame = NULL;
2925 }
2926 if (client->h265_sws_ctx) {
2927 sws_freeContext(client->h265_sws_ctx);
2928 client->h265_sws_ctx = NULL;
2929 }
2930}
void opus_codec_destroy(opus_codec_t *codec)
Destroy an Opus codec instance.
Definition opus.c:230
void digital_rain_destroy(digital_rain_t *rain)
Destroy digital rain effect context.
Opus codec context for encoding or decoding.
Definition opus.h:95

References audio_ring_buffer_destroy(), client_info::digital_rain, digital_rain_destroy(), ERROR_INVALID_PARAM, client_info::h265_decode_frame, client_info::h265_decoder_ctx, client_info::h265_rgb_frame, client_info::h265_sws_ctx, client_info::incoming_audio_buffer, client_info::incoming_video_buffer, opus_codec_destroy(), client_info::opus_decoder, client_info::outgoing_video_buffer, SAFE_FREE, client_info::send_buffer, client_info::send_buffer_size, SET_ERRNO, and video_frame_buffer_destroy().

◆ cleanup_client_packet_queues()

void cleanup_client_packet_queues ( client_info_t *  client)

Definition at line 2932 of file src/server/client.c.

2932 {
2933 if (!client)
2934 return;
2935
2936 if (client->audio_queue) {
2938 client->audio_queue = NULL;
2939 }
2940
2941 // Async dispatch: clean up received packet queue
2942 if (client->received_packet_queue) {
2944 client->received_packet_queue = NULL;
2945 }
2946
2947 // Video now uses double buffer, cleaned up in cleanup_client_media_buffers
2948}

References client_info::audio_queue, packet_queue_destroy(), and client_info::received_packet_queue.

◆ client_receive_thread()

void * client_receive_thread ( void *  arg)

Definition at line 1885 of file src/server/client.c.

1885 {
1886 // Log thread startup
1887 log_debug("RECV_THREAD: Thread function entered, arg=%p", arg);
1888
1889 client_info_t *client = (client_info_t *)arg;
1890
1891 // Validate client pointer immediately before any access.
1892 // This prevents crashes if remove_client() has zeroed the client struct
1893 // while the thread was still starting at RtlUserThreadStart.
1894 if (!client) {
1895 log_error("Invalid client info in receive thread (NULL pointer)");
1896 return NULL;
1897 }
1898
1899 // Save this thread's ID so remove_client() can detect self-joins
1901
1902 log_debug("RECV_THREAD_DEBUG: Thread started, client=%p, client_id=%s, is_tcp=%d", (void *)client, client->client_id,
1903 client->is_tcp_client);
1904
1906 log_debug("Receive thread for client %s exiting before start (protocol disconnect requested)", client->client_id);
1907 return NULL;
1908 }
1909
1910 // Check if client_id is empty (client struct has been zeroed by remove_client)
1911 // This must be checked BEFORE accessing any client fields
1912 if (client->client_id[0] == 0) {
1913 log_debug("Receive thread: client_id is empty, client struct may have been zeroed, exiting");
1914 return NULL;
1915 }
1916
1917 // Additional validation: check socket is valid
1918 // For TCP clients, validate socket. WebRTC clients use DataChannel (no socket)
1919 if (client->is_tcp_client && client->socket == INVALID_SOCKET_VALUE) {
1920 log_error("Invalid client socket in receive thread");
1921 return NULL;
1922 }
1923
1924 // NOTE: Atomic registration disabled
1925 // register_client_info_atomics() and audio_ring_buffer_register_atomics() attempt
1926 // to acquire WRITE lock on named_registry.entries_lock, which conflicts with debug_sync
1927 // thread holding READ lock during printing, causing deadlock.
1928 // These are debug-only features, not critical for functionality.
1929 // The atomics work fine without being registered in the debug system.
1930
1931 // Enable thread cancellation for clean shutdown
1932 // Thread cancellation not available in platform abstraction
1933 // Threads should exit when g_should_exit is set
1934
1935 log_debug("Started receive thread for client %s (%s)", client->client_id, client->display_name);
1936
1937 // Main receive loop - processes packets from transport
1938 // For TCP clients: receives from socket
1939 // For WebRTC clients: receives from transport ringbuffer (via ACDS signaling)
1940
1941 log_info("RECV_THREAD_LOOP_START: client_id=%s, is_tcp=%d, transport=%p, active=%d", client->client_id,
1942 client->is_tcp_client, (void *)client->transport, atomic_load_bool(&client->active));
1943
1944 while (!atomic_load_bool(&g_should_exit) && atomic_load_bool(&client->active)) {
1945 // For TCP clients, check socket validity
1946 // For WebRTC clients, continue even if no socket (transport handles everything)
1947 if (client->is_tcp_client && client->socket == INVALID_SOCKET_VALUE) {
1948 log_debug("TCP client %s has invalid socket, exiting receive thread", client->client_id);
1949 break;
1950 }
1951
1952 // Check client_id is still valid before accessing transport.
1953 // This prevents accessing freed memory if remove_client() has zeroed the client struct.
1954 if (client->client_id[0] == 0) {
1955 log_debug("Client client_id reset, exiting receive thread");
1956 break;
1957 }
1958
1959 // Receive packet (without dispatching) - decouple from dispatch for async processing
1960 // For WebRTC clients with async dispatch, we queue packets instead of processing immediately
1961 // This prevents backpressure on the network socket
1962
1963 if (client->is_tcp_client) {
1964 // TCP clients: use original synchronous dispatch
1965 asciichat_error_t acip_result =
1966 acip_server_receive_and_dispatch(client->transport, client, &g_acip_server_callbacks);
1967
1968 // Check if shutdown was requested during the network call
1970 log_debug("RECV_EXIT: Server shutdown requested, breaking loop");
1971 break;
1972 }
1973
1974 // Handle receive errors
1975 if (acip_result != ASCIICHAT_OK) {
1977 if (HAS_ERRNO(&err_ctx)) {
1978 log_error("🔴 ACIP error for client %s: code=%u msg=%s", client->client_id, err_ctx.code,
1979 err_ctx.context_message);
1980 if (err_ctx.code == ERROR_NETWORK) {
1981 log_debug("Client %s disconnected (network error): %s", client->client_id, err_ctx.context_message);
1982 break;
1983 } else if (err_ctx.code == ERROR_CRYPTO) {
1985 client, "SECURITY VIOLATION: Unencrypted packet when encryption required - terminating connection");
1987 break;
1988 }
1989 }
1990 log_warn("ACIP error for TCP client %s: %s (disconnecting)", client->client_id,
1991 asciichat_error_string(acip_result));
1992 break;
1993 }
1994 } else {
1995 // WebRTC/WebSocket clients: async dispatch - receive packet and queue for async processing
1996 void *packet_data = NULL;
1997 void *allocated_buffer = NULL;
1998 size_t packet_len = 0;
1999
2000 // Snapshot transport pointer to prevent use-after-free
2001 // Transport can be destroyed by WebSocket callback while we're receiving
2002 acip_transport_t *transport_snapshot = NULL;
2003 mutex_lock(&client->send_mutex);
2004 if (!client->transport || atomic_load_bool(&client->shutting_down)) {
2005 mutex_unlock(&client->send_mutex);
2006 log_debug("RECV_THREAD[%s]: Transport destroyed or client shutting down, exiting receive loop",
2007 client->client_id);
2008 break;
2009 }
2010 transport_snapshot = client->transport;
2011 mutex_unlock(&client->send_mutex);
2012
2013 log_debug("🔍 RECV_THREAD[%s]: About to call transport->recv() (transport=%p)", client->client_id,
2014 (void *)transport_snapshot);
2015
2016 asciichat_error_t recv_result =
2017 transport_snapshot->methods->recv(transport_snapshot, &packet_data, &packet_len, &allocated_buffer);
2018
2019 // Check error BEFORE logging - logging system may clear thread-local errno
2020 if (recv_result != ASCIICHAT_OK) {
2022 if (HAS_ERRNO(&err_ctx)) {
2023 // Check for reassembly timeout (fragments arriving slowly)
2024 // This is NOT a connection failure - safe to retry
2025 if ((err_ctx.code == ERROR_NETWORK) && err_ctx.context_message &&
2026 strstr(err_ctx.context_message, "reassembly timeout")) {
2027 // Fragments are arriving slowly - this is normal, retry without disconnecting
2028 log_dev_every(100000, "Client %s: fragment reassembly timeout, retrying in 10ms", client->client_id);
2029 APP_CALLBACK_VOID(platform_pump_events);
2030 platform_sleep_ms(10); // Sleep 10ms to allow fragments to arrive
2031 continue; // Retry without disconnecting
2032 }
2033
2034 if (err_ctx.code == ERROR_NETWORK) {
2035 log_debug("Client %s disconnected (network error): %s", client->client_id, err_ctx.context_message);
2036 break;
2037 }
2038 }
2039 log_warn("Receive failed for WebRTC client %s: %s (disconnecting)", client->client_id,
2040 asciichat_error_string(recv_result));
2041 break;
2042 }
2043
2044 log_dev("RECV_THREAD[%s]: recv result=%d packet_len=%zu", client->client_id, recv_result, packet_len);
2045
2046 // Validate received packet before queueing
2047 if (packet_len < sizeof(packet_header_t)) {
2048 log_warn("RECV_THREAD[%s]: Received packet too small (%zu < %zu), dropping", client->client_id, packet_len,
2049 sizeof(packet_header_t));
2050 if (allocated_buffer) {
2051 buffer_pool_free(NULL, allocated_buffer, packet_len);
2052 }
2053 continue;
2054 }
2055
2056 // Queue the received packet for async dispatch
2057 // This prevents the receive thread from blocking on dispatch
2058 log_dev("RECV_THREAD[%s]: queueing %zu-byte packet", client->client_id, packet_len);
2059
2060 // Extract packet type from the header to preserve it when queueing
2061 const packet_header_t *pkt_header = (const packet_header_t *)allocated_buffer;
2062 packet_type_t pkt_type = (packet_type_t)NET_TO_HOST_U16(pkt_header->type);
2063 uint32_t payload_len = NET_TO_HOST_U32(pkt_header->length);
2064
2065 log_dev("RECV_THREAD[%s]: type=%d payload_len=%u total_len=%zu", client->client_id, pkt_type, payload_len,
2066 packet_len);
2067
2068 // Build a complete packet to queue (header + payload)
2069 // The entire buffer (allocated_buffer) contains the full packet
2070 // The queue copies the packet; release the transport buffer afterward.
2071 uint32_t client_id_hash = fnv1a_hash_string(client->client_id);
2072 // Pair queue publication and notification with the dispatch wait lock.
2074 int enqueue_result = packet_queue_enqueue(client->received_packet_queue, pkt_type, allocated_buffer, packet_len,
2075 client_id_hash, true);
2076 if (enqueue_result >= 0)
2079 buffer_pool_free(NULL, allocated_buffer, packet_len);
2080
2081 if (enqueue_result < 0) {
2082 log_error("🔴 RECV_THREAD[%s]: Failed to queue received packet (queue full?) - DROPPING FRAME",
2083 client->client_id);
2084 } else {
2085 log_dev("RECV_THREAD[%s]: queued type=%d len=%zu", client->client_id, pkt_type, packet_len);
2086 }
2087 }
2088 }
2089
2090 // Mark client as inactive and stop all threads
2091 // Must stop render threads when client disconnects.
2092 // OPTIMIZED: Use atomic operations for thread control flags (lock-free)
2093 char client_id_snapshot[MAX_CLIENT_ID_LEN] = {0};
2094 SAFE_STRNCPY(client_id_snapshot, client->client_id, sizeof(client_id_snapshot) - 1);
2095 bool is_tcp_client = client->is_tcp_client;
2096 server_context_t *server_ctx = (server_context_t *)client->server_ctx;
2097 log_debug("Setting active=false in receive_thread_fn (client_id=%s, exiting receive loop)", client_id_snapshot);
2098 atomic_store_bool(&client->active, false);
2099 atomic_store_bool(&client->send_thread_running, false);
2102
2103 // Call remove_client() to trigger cleanup
2104 // Safe to call from receive thread now: remove_client() detects self-join via thread IDs
2105 // and skips the receive thread join when called from the receive thread itself
2106 log_debug("Receive thread for client %s calling remove_client() for cleanup", client_id_snapshot);
2107 if (server_ctx && !is_tcp_client) {
2108 if (remove_client(server_ctx, client_id_snapshot) != 0) {
2109 log_warn("Failed to remove client %s from receive thread cleanup", client_id_snapshot);
2110 }
2111 } else if (!server_ctx) {
2112 log_error("Receive thread for client %s: server_ctx is NULL, cannot call remove_client()", client_id_snapshot);
2113 }
2114
2115 log_debug("Receive thread for client %s terminated", client_id_snapshot);
2116
2117 // Clean up thread-local error context before exit
2119
2120 return NULL;
2121}
#define APP_CALLBACK_VOID(callback_name)
atomic_t g_should_exit
Global application exit flag (shared across all modes)
Definition globals.c:34
#define HAS_ERRNO(var)
Check if an error occurred and get full context.
void asciichat_errno_destroy(void)
Cleanup error system resources.
@ ERROR_NETWORK
Definition error_codes.h:77
@ ERROR_CRYPTO
Definition error_codes.h:96
#define log_dev(...)
Log a DEV message (most verbose, development only)
Definition log/log.h:534
#define log_error_client(client, fmt,...)
Server sends ERROR log message to client.
Definition network/log.h:75
int packet_queue_enqueue(packet_queue_t *queue, packet_type_t type, const void *data, size_t data_len, uint32_t client_id, bool copy_data)
Enqueue a packet into the queue.
Definition queue.c:209
int cond_signal(cond_t *cond)
Signal a condition variable (wake one waiting thread)
void platform_sleep_ms(unsigned int ms)
Sleep for a specified number of milliseconds.
asciichat_error_t acip_server_receive_and_dispatch(acip_transport_t *transport, void *client_ctx, const acip_server_callbacks_t *callbacks)
Receive packet from client and dispatch to callbacks.
#define asciichat_thread_self()
#define log_dev_every(interval_us, fmt,...)
Rate-limited DEV logging.
Definition log/log.h:699
asciichat_error_t(* recv)(acip_transport_t *transport, void **buffer, size_t *out_len, void **out_allocated_buffer)
Receive data from this transport.
Definition transport.h:145
Transport instance structure.
Definition transport.h:214
const acip_transport_methods_t * methods
Method table (virtual functions)
Definition transport.h:215
Error context structure.
char * context_message
Optional custom message (dynamically allocated, owned by system)
asciichat_error_t code
Error code (asciichat_error_t enum value)
uint32_t length
Payload data length in bytes (0 for header-only packets)
Definition packet.h:604
uint16_t type
Packet type (packet_type_t enumeration)
Definition packet.h:602
Server context - encapsulates all server state.

References acip_server_receive_and_dispatch(), client_info::active, APP_CALLBACK_VOID, asciichat_errno_destroy(), ASCIICHAT_OK, asciichat_thread_self, atomic_load_bool(), atomic_store_bool(), client_info::audio_render_thread_running, buffer_pool_free(), client_info::client_id, client_info::client_state_mutex, asciichat_error_context_t::code, cond_signal(), asciichat_error_context_t::context_message, client_info::dispatch_queue_cond, client_info::display_name, ERROR_CRYPTO, ERROR_NETWORK, g_should_exit, HAS_ERRNO, INVALID_SOCKET_VALUE, client_info::is_tcp_client, packet_header_t::length, log_debug, log_dev, log_dev_every, log_error, log_error_client, log_info, log_warn, MAX_CLIENT_ID_LEN, acip_transport::methods, mutex_lock, mutex_unlock, NET_TO_HOST_U16, NET_TO_HOST_U32, packet_queue_enqueue(), platform_sleep_ms(), client_info::protocol_disconnect_requested, client_info::receive_thread_id, client_info::received_packet_queue, acip_transport_methods_t::recv, remove_client(), SAFE_STRNCPY, client_info::send_mutex, client_info::send_thread_running, client_info::server_ctx, client_info::shutting_down, client_info::socket, client_info::transport, packet_header_t::type, and client_info::video_render_thread_running.

◆ find_client_by_id()

client_info_t * find_client_by_id ( const char *  client_id)

Fast O(1) client lookup by ID using hash table.

This is the primary method for locating clients throughout the server. It uses a hash table for constant-time lookups regardless of client count, making it suitable for high-performance operations like rendering and stats.

PERFORMANCE CHARACTERISTICS:

  • Time Complexity: O(1) average case, O(n) worst case (hash collision)
  • Space Complexity: O(1)
  • Thread Safety: Hash table is internally thread-safe for lookups

USAGE PATTERNS:

  • Called by render threads to find target clients for frame generation
  • Used by protocol handlers to locate clients for packet processing
  • Stats collection for per-client performance monitoring
Parameters
client_idUnique identifier for the client (0 is invalid)
Returns
Pointer to client_info_t if found, NULL if not found or invalid ID
Note
Does not require external locking - hash table provides thread safety
Returns direct pointer to client struct - caller should use snapshot pattern

Definition at line 429 of file src/server/client.c.

429 {
430 if (!client_id || client_id[0] == '\0') {
431 SET_ERRNO(ERROR_INVALID_PARAM, "Invalid client ID");
432 return NULL;
433 }
434
435 // Protect uthash lookup with read lock to prevent concurrent access issues
437
438 client_info_t *result = NULL;
439 HASH_FIND_STR(g_client_manager.clients_by_id, client_id, result);
440
442
443 if (!result) {
444 log_warn("Client not found for ID %s", client_id);
445 }
446
447 return result;
448}

References client_manager_t::clients_by_id, ERROR_INVALID_PARAM, g_client_manager, g_client_manager_rwlock, log_warn, rwlock_rdlock, rwlock_rdunlock, and SET_ERRNO.

Referenced by broadcast_server_state_to_all_clients(), crypto_server_cleanup_client(), crypto_server_decrypt_packet(), crypto_server_encrypt_packet(), crypto_server_get_context(), crypto_server_is_ready(), and start_webrtc_client_threads().

◆ find_client_by_socket()

client_info_t * find_client_by_socket ( socket_t  socket)

Find client by socket descriptor using linear search.

This function provides socket-based client lookup, primarily used during connection establishment before client IDs are assigned. Less efficient than find_client_by_id() but necessary for socket-based operations.

PERFORMANCE CHARACTERISTICS:

  • Time Complexity: O(n) where n = number of active clients
  • Space Complexity: O(1)
  • Thread Safety: Internally acquires read lock on g_client_manager_rwlock

USAGE PATTERNS:

  • Connection establishment during add_client() processing
  • Socket error handling and cleanup operations
  • Debugging and diagnostic functions
Parameters
socketPlatform-abstracted socket descriptor to search for
Returns
Pointer to client_info_t if found, NULL if not found
Note
Only searches active clients (avoids returning stale entries)
Caller should use snapshot pattern when accessing returned client data

Definition at line 473 of file src/server/client.c.

473 {
475
476 for (int i = 0; i < MAX_CLIENTS; i++) {
480 return client;
481 }
482 }
483
485 return NULL;
486}

References client_info::active, atomic_load_bool(), client_manager_t::clients, g_client_manager, g_client_manager_rwlock, MAX_CLIENTS, rwlock_rdlock, rwlock_rdunlock, and client_info::socket.

◆ initialize_client_info()

void initialize_client_info ( client_info_t *  client)

◆ process_decrypted_packet()

void process_decrypted_packet ( client_info_t *  client,
packet_type_t  type,
void *  data,
size_t  len 
)

Process a decrypted packet from a client

Parameters
clientClient info structure
typePacket type
dataPacket data
lenPacket length

Definition at line 3934 of file src/server/client.c.

3934 {
3935 if (type == 5000) { // CLIENT_CAPABILITIES
3936 log_debug("CLIENT: client_id=%s, data=%p, len=%zu", client->client_id, data, len);
3937 }
3938
3939 // Rate limiting: Check and record packet-specific rate limits
3940 if (g_rate_limiter) {
3941 if (!check_and_record_packet_rate_limit(g_rate_limiter, client->client_ip, client->socket, type)) {
3942 // Rate limit exceeded - error response already sent by utility function
3943 return;
3944 }
3945 }
3946
3947 // O(1) dispatch via hash table lookup
3948 int idx = client_dispatch_hash_lookup(g_client_dispatch_hash, type);
3949 if (type == 5000 || type == 3001) {
3950 log_error("DISPATCH_LOOKUP: type=%d, idx=%d (len=%zu)", type, idx, len);
3951 }
3952 if (idx < 0) {
3953 disconnect_client_for_bad_data(client, "Unknown packet type: %d (len=%zu)", type, len);
3954 return;
3955 }
3956
3957 if (type == 5000 || type == 3001) {
3958 log_error("DISPATCH_HANDLER: type=%d, calling handler[%d]...", type, idx);
3959 }
3960 g_client_dispatch_handlers[idx](client, data, len);
3961 if (type == 5000 || type == 3001) {
3962 log_error("DISPATCH_HANDLER: type=%d, handler returned", type);
3963 }
3964}
bool check_and_record_packet_rate_limit(rate_limiter_t *rate_limiter, const char *client_ip, socket_t client_socket, packet_type_t packet_type)
Map packet type to rate event type and check rate limit.
Definition errors.c:54
void disconnect_client_for_bad_data(client_info_t *client, const char *format,...)
rate_limiter_t * g_rate_limiter
Global rate limiter for connection attempts and packet processing.

References check_and_record_packet_rate_limit(), client_info::client_id, client_info::client_ip, disconnect_client_for_bad_data(), g_rate_limiter, log_debug, log_error, and client_info::socket.

◆ process_encrypted_packet()

int process_encrypted_packet ( client_info_t *  client,
packet_type_t *  type,
void **  data,
size_t *  len,
uint32_t *  sender_id 
)

Process an encrypted packet from a client

Parameters
clientClient info structure
typePointer to packet type (will be updated with decrypted type)
dataPointer to packet data (will be updated with decrypted data)
lenPointer to packet length (will be updated with decrypted length)
sender_idPointer to sender ID (will be updated with decrypted sender ID)
Returns
0 on success, -1 on error

Definition at line 2974 of file src/server/client.c.

2975 {
2976 if (!crypto_server_is_ready(client->client_id)) {
2977 log_error("Received encrypted packet but crypto not ready for client %s", client->client_id);
2978 buffer_pool_free(NULL, *data, *len);
2979 *data = NULL;
2980 return -1;
2981 }
2982
2983 // Store original allocation size before it gets modified
2984 size_t original_alloc_size = *len;
2985 void *decrypted_data = buffer_pool_alloc(NULL, original_alloc_size);
2986 size_t decrypted_len;
2987 int decrypt_result = crypto_server_decrypt_packet(client->client_id, (const uint8_t *)*data, *len,
2988 (uint8_t *)decrypted_data, original_alloc_size, &decrypted_len);
2989
2990 if (decrypt_result != 0) {
2991 SET_ERRNO(ERROR_CRYPTO, "Failed to process encrypted packet from client %s (result=%d)", client->client_id,
2992 decrypt_result);
2993 buffer_pool_free(NULL, *data, original_alloc_size);
2994 buffer_pool_free(NULL, decrypted_data, original_alloc_size);
2995 *data = NULL;
2996 return -1;
2997 }
2998
2999 // Replace encrypted data with decrypted data
3000 // Use original allocation size for freeing the encrypted buffer
3001 buffer_pool_free(NULL, *data, original_alloc_size);
3002
3003 *data = decrypted_data;
3004 *len = decrypted_len;
3005
3006 // Now process the decrypted packet by parsing its header
3007 if (*len < sizeof(packet_header_t)) {
3008 SET_ERRNO(ERROR_CRYPTO, "Decrypted packet too small for header from client %s", client->client_id);
3009 buffer_pool_free(NULL, *data, *len);
3010 *data = NULL;
3011 return -1;
3012 }
3013
3014 packet_header_t *header = (packet_header_t *)*data;
3015 *type = (packet_type_t)NET_TO_HOST_U16(header->type);
3016 *sender_id = NET_TO_HOST_U32(header->client_id);
3017
3018 // Adjust data pointer to skip header
3019 *data = (uint8_t *)*data + sizeof(packet_header_t);
3020 *len -= sizeof(packet_header_t);
3021
3022 return 0;
3023}
void * buffer_pool_alloc(buffer_pool_t *pool, size_t size)
Allocate a buffer from the pool (lock-free fast path)
bool crypto_server_is_ready(const char *client_id)
int crypto_server_decrypt_packet(const char *client_id, const uint8_t *ciphertext, size_t ciphertext_len, uint8_t *plaintext, size_t plaintext_size, size_t *plaintext_len)
uint32_t client_id
Client ID (0 = server, >0 = client identifier)
Definition packet.h:608

References buffer_pool_alloc(), buffer_pool_free(), client_info::client_id, packet_header_t::client_id, crypto_server_decrypt_packet(), crypto_server_is_ready(), ERROR_CRYPTO, log_error, NET_TO_HOST_U16, NET_TO_HOST_U32, SET_ERRNO, and packet_header_t::type.

◆ remove_client()

int remove_client ( server_context_t *  server_ctx,
const char *  client_id 
)

Definition at line 1389 of file src/server/client.c.

1389 {
1390 if (!server_ctx) {
1391 SET_ERRNO(ERROR_INVALID_PARAM, "Cannot remove client %s: NULL server_ctx", client_id);
1392 return -1;
1393 }
1394
1395 // Phase 1: Mark client inactive and prepare for cleanup while holding write lock
1396 client_info_t *target_client = NULL;
1397 char display_name_copy[MAX_DISPLAY_NAME_LEN];
1398 socket_t client_socket = INVALID_SOCKET_VALUE; // Save socket for thread cleanup
1399
1400 log_debug("SOCKET_DEBUG: Attempting to remove client %s", client_id);
1402
1403 for (int i = 0; i < MAX_CLIENTS; i++) {
1405 if (strcmp(client->client_id, client_id) == 0 && client->client_id[0] != '\0') {
1406 // Check if already being removed by another thread
1407 // This prevents double-free and use-after-free crashes during concurrent cleanup
1408 if (atomic_load_bool(&client->shutting_down)) {
1410 log_debug("Client %s already being removed by another thread, skipping", client_id);
1411 return 0; // Return success - removal is in progress
1412 }
1413 // Mark as shutting down and inactive immediately to stop new operations
1414 log_debug("Setting active=false in remove_client (client_id=%s, socket=%d)", client_id, client->socket);
1415 log_info("Removing client %s (socket=%d) - marking inactive and clearing video flags", client_id, client->socket);
1416 atomic_store_bool(&client->shutting_down, true);
1417 atomic_store_bool(&client->active, false);
1418 atomic_store_bool(&client->is_sending_video, false);
1419 atomic_store_bool(&client->is_sending_audio, false);
1420 target_client = client;
1421
1422 // Store display name before clearing
1423 SAFE_STRNCPY(display_name_copy, client->display_name, MAX_DISPLAY_NAME_LEN - 1);
1424
1425 // Save socket for tcp_server_stop_client_threads(). Keep the descriptor
1426 // open until thread cleanup completes so the OS cannot reuse its fd for
1427 // a new client while this client's cleanup is still in progress.
1429 client_socket = client->socket; // Save socket for thread cleanup
1431
1432 // Wake blocked receive/send operations without closing the descriptor.
1433 // Closing here would allow a concurrently accepted client to reuse the
1434 // same descriptor before the old client's worker threads finish.
1435 if (client_socket != INVALID_SOCKET_VALUE) {
1436 (void)socket_shutdown(client_socket, SHUT_RDWR);
1437 }
1438
1439 // Shutdown packet queues to unblock send thread
1440 if (client->audio_queue) {
1442 }
1443 // Video now uses double buffer, no queue to shutdown
1444
1445 break;
1446 }
1447 }
1448
1449 // If client not found, unlock and return
1450 if (!target_client) {
1452 log_warn("Cannot remove client %s: not found", client_id);
1453 return -1;
1454 }
1455
1456 // Unregister client from session_host (for discovery mode support)
1457 // NOTE: Client may not be registered if crypto handshake failed before session_host registration
1458 if (server_ctx->session_host && target_client->session_client_id != 0) {
1459 asciichat_error_t session_result =
1460 session_host_remove_client(server_ctx->session_host, target_client->session_client_id);
1461 if (session_result != ASCIICHAT_OK) {
1462 // ERROR_NOT_FOUND (91) is expected if client failed crypto before being registered with session_host
1463 if (session_result == ERROR_NOT_FOUND) {
1464 log_debug("Client %s not found in session_host (likely failed crypto before registration)", client_id);
1465 } else {
1466 log_warn("Failed to unregister client %s from session_host: %s", client_id,
1467 asciichat_error_string(session_result));
1468 }
1469 } else {
1470 log_debug("Client %s unregistered from session_host", client_id);
1471 }
1472 }
1473
1474 // Release write lock before joining threads.
1475 // This prevents deadlock with render threads that need read locks.
1477
1478 // Phase 2: Stop all client threads
1479 // For TCP clients: use tcp_server thread pool management
1480 // For WebRTC clients: manually join threads (no socket-based thread pool)
1481 // Use is_tcp_client flag, not socket value - socket may already be INVALID_SOCKET_VALUE
1482 // even for TCP clients if it was closed earlier during cleanup.
1483 bool receive_thread_is_current = false;
1484 log_debug("Stopping all threads for client %s (socket %d, is_tcp=%d)", client_id, client_socket,
1485 target_client ? target_client->is_tcp_client : -1);
1486
1487 if (target_client && target_client->is_tcp_client) {
1488 // TCP client: use tcp_server thread pool
1489 // This joins threads in stop_id order: receive(1), render(2), send(3)
1490 // Use saved client_socket for lookup (tcp_server needs original socket as key)
1491 if (client_socket != INVALID_SOCKET_VALUE) {
1492 asciichat_error_t stop_result = tcp_server_stop_client_threads(server_ctx->tcp_server, client_socket);
1493 if (stop_result != ASCIICHAT_OK) {
1494 log_warn("Failed to stop threads for TCP client %s: error %d", client_id, stop_result);
1495 // Continue with cleanup even if thread stopping failed
1496 }
1497 } else {
1498 log_debug("TCP client %s socket already closed, threads should have already exited", client_id);
1499 }
1500 } else if (target_client) {
1501 // WebRTC client: manually join threads
1502 log_debug("Stopping WebRTC client %s threads (receive and send)", client_id);
1503
1504 // Join receive thread (but skip if called from the receive thread itself to avoid deadlock)
1505 thread_id_t current_thread_id = asciichat_thread_self();
1506 receive_thread_is_current = asciichat_thread_equal(current_thread_id, target_client->receive_thread_id);
1507 if (receive_thread_is_current) {
1508 log_debug("remove_client() called from receive thread for client %s, skipping self-join", client_id);
1509 } else {
1510 void *recv_result = NULL;
1511 asciichat_error_t recv_join_result = asciichat_thread_join(&target_client->receive_thread, &recv_result);
1512 if (recv_join_result != ASCIICHAT_OK) {
1513 log_warn("Failed to join receive thread for WebRTC client %s: error %d", client_id, recv_join_result);
1514 } else {
1515 log_debug("Joined receive thread for WebRTC client %s", client_id);
1516 }
1517 }
1518
1519 // Join dispatch thread (BEFORE destroying packet queues)
1520 if (asciichat_thread_is_initialized(&target_client->dispatch_thread)) {
1521 atomic_store_bool(&target_client->dispatch_thread_running, false);
1522 void *dispatch_result = NULL;
1523 asciichat_error_t dispatch_join_result = asciichat_thread_join(&target_client->dispatch_thread, &dispatch_result);
1524 if (dispatch_join_result != ASCIICHAT_OK) {
1525 log_warn("Failed to join dispatch thread for WebRTC client %s: error %d", client_id, dispatch_join_result);
1526 } else {
1527 log_debug("Joined dispatch thread for WebRTC client %s", client_id);
1528 }
1529 }
1530
1531 // Join send thread
1532 void *send_result = NULL;
1533 asciichat_error_t send_join_result = asciichat_thread_join(&target_client->send_thread, &send_result);
1534 if (send_join_result != ASCIICHAT_OK) {
1535 log_warn("Failed to join send thread for WebRTC client %s: error %d", client_id, send_join_result);
1536 } else {
1537 log_debug("Joined send thread for WebRTC client %s", client_id);
1538 }
1539
1540 // Render workers access the client buffers and synchronization primitives.
1541 // Join them before those resources are destroyed or the client slot is reused.
1542 stop_client_render_threads(target_client);
1543 }
1544
1545 // Destroy ACIP transport before closing socket
1546 // For WebSocket clients: LWS_CALLBACK_CLOSED already closed and destroyed the transport
1547 // Trying to destroy it again causes heap-use-after-free since LWS callbacks might still fire
1548 // For TCP clients: transport is ours to clean up
1549 if (target_client && target_client->transport && target_client->is_tcp_client) {
1550 acip_transport_destroy(target_client->transport);
1551 target_client->transport = NULL;
1552 log_debug("Destroyed ACIP transport for TCP client %s", client_id);
1553 } else if (target_client && target_client->transport && !target_client->is_tcp_client) {
1554 // WebSocket client - just NULL it out, LWS_CALLBACK_CLOSED already destroyed it
1555 target_client->transport = NULL;
1556 log_debug("Skipped transport destruction for WebSocket client %s (LWS already destroyed)", client_id);
1557 }
1558
1559 // Now safe to close the socket (threads are stopped)
1560 if (client_socket != INVALID_SOCKET_VALUE) {
1561 log_debug("SOCKET_DEBUG: Closing socket %d for client %s after thread cleanup", client_socket, client_id);
1562 socket_close(client_socket);
1563 }
1564
1565 // Phase 3: Clean up resources with write lock
1567
1568 // Re-validate target_client pointer after reacquiring lock.
1569 // Another thread might have invalidated the pointer while we had the lock released.
1570 if (target_client) {
1571 // Verify client_id still matches and client is still in shutting_down state
1572 bool still_shutting_down = atomic_load_bool(&target_client->shutting_down);
1573 if (strcmp(target_client->client_id, client_id) != 0 || !still_shutting_down) {
1574 log_warn("Client %s pointer invalidated during thread cleanup (id=%s, shutting_down=%d)", client_id,
1575 target_client->client_id, still_shutting_down);
1577 return 0; // Another thread completed the cleanup
1578 }
1579 }
1580
1581 // Mark socket as closed in client structure
1582 if (target_client && target_client->socket != INVALID_SOCKET_VALUE) {
1583 mutex_lock(&target_client->client_state_mutex);
1584 target_client->socket = INVALID_SOCKET_VALUE;
1585 mutex_unlock(&target_client->client_state_mutex);
1586 log_debug("SOCKET_DEBUG: Client %s socket set to INVALID", target_client->client_id);
1587 }
1588
1589 // Use the dedicated cleanup function to ensure all resources are freed
1590 cleanup_client_all_buffers(target_client);
1591
1592 // Remove from audio mixer
1593 if (g_audio_mixer) {
1595#ifdef DEBUG_AUDIO
1596 log_debug("Removed client %s from audio mixer", client_id);
1597#endif
1598 }
1599
1600 // Remove from uthash table
1601 // Verify client is actually in the hash table before deleting.
1602 // Another thread might have already removed it.
1603 if (target_client) {
1604 client_info_t *hash_entry = NULL;
1605 HASH_FIND_STR(g_client_manager.clients_by_id, target_client->client_id, hash_entry);
1606 if (hash_entry == target_client) {
1607 HASH_DELETE(hh, g_client_manager.clients_by_id, target_client);
1608 log_debug("Removed client %s from uthash table", client_id);
1609 } else {
1610 log_warn("Client %s already removed from hash table by another thread (found=%p, expected=%p)", client_id,
1611 (void *)hash_entry, (void *)target_client);
1612 }
1613 } else {
1614 log_warn("Failed to remove client %s from hash table (client not found)", client_id);
1615 }
1616
1617 // Cleanup crypto context for this client
1618 if (target_client->crypto_initialized) {
1620 target_client->crypto_initialized = false;
1621 log_debug("Crypto context cleaned up for client %s", client_id);
1622 }
1623
1624 // Verify all threads have actually exited before resetting client_id.
1625 // Threads that are still starting (at RtlUserThreadStart) haven't checked client_id yet.
1626 // We must ensure threads are fully joined before zeroing the client struct.
1627 // Use exponential backoff for thread termination verification
1628 int retry_count = 0;
1629 const int max_retries = 5;
1630 while (retry_count < max_retries && (asciichat_thread_is_initialized(&target_client->send_thread) ||
1631 (!receive_thread_is_current &&
1635 // Exponential backoff: 10ms, 20ms, 40ms, 80ms, 160ms
1636 uint32_t delay_ms = 10 * (1 << retry_count);
1637 log_warn("Client %s: Some threads still appear initialized (attempt %d/%d), waiting %ums", client_id,
1638 retry_count + 1, max_retries, delay_ms);
1639 platform_sleep_us(delay_ms * 1000);
1640 retry_count++;
1641 }
1642
1643 if (retry_count == max_retries) {
1644 log_error("Client %s: Threads did not terminate after %d retries, proceeding with cleanup anyway", client_id,
1645 max_retries);
1646 }
1647
1648 // Only reset client_id to 0 AFTER confirming threads are joined
1649 // This prevents threads that are starting from accessing a zeroed client struct
1650 // Reset client_id to NULL before destroying mutexes to prevent race conditions.
1651 // This ensures worker threads can detect shutdown and exit before the mutex is destroyed.
1652 // If we destroy the mutex first, threads might try to access a destroyed mutex.
1653 target_client->client_id[0] = '\0'; // Clear the client_id string
1654
1655 // Wait for threads to observe the client_id reset
1656 // Use sufficient delay for memory visibility across all CPU cores
1657 platform_sleep_us(5 * US_PER_MS_INT); // 5ms delay for memory barrier propagation
1658
1659 // Destroy mutexes and condition variables
1660 // IMPORTANT: Always destroy these even if threads didn't join properly
1661 // to prevent issues when the slot is reused
1662 mutex_destroy(&target_client->client_state_mutex);
1663 mutex_destroy(&target_client->send_mutex);
1664 cond_destroy(&target_client->dispatch_queue_cond);
1665
1666 // Clear client structure
1667 // NOTE: After memset, the mutex handles are zeroed but the OS resources
1668 // have been released by the destroy calls above
1669 memset(target_client, 0, sizeof(client_info_t));
1670
1671 // Recalculate client count
1672 int remaining_count = 0;
1673 for (int j = 0; j < MAX_CLIENTS; j++) {
1674 if (g_client_manager.clients[j].client_id[0] != '\0') {
1675 remaining_count++;
1676 }
1677 }
1678 g_client_manager.client_count = remaining_count;
1679
1680 log_debug("Client removed: client_id=%s (%s) removed, remaining clients: %d", client_id, display_name_copy,
1681 remaining_count);
1682
1684
1685 // Broadcast updated state
1687
1688 return 0;
1689}
void crypto_handshake_destroy(crypto_handshake_context_t *ctx)
Cleanup crypto handshake context with secure memory wiping.
void mixer_remove_source(mixer_t *mixer, const char *client_id)
Remove an audio source from the mixer.
Definition mixer.c:477
@ ERROR_NOT_FOUND
#define MAX_DISPLAY_NAME_LEN
Maximum display name length in characters.
Definition limits.h:20
#define US_PER_MS_INT
Definition time.h:160
void packet_queue_stop(packet_queue_t *queue)
Destroy a packet queue and free all resources.
Definition queue.c:616
int socket_shutdown(socket_t sock, int how)
Shutdown socket I/O.
int socket_close(socket_t sock)
Close a socket.
bool asciichat_thread_is_initialized(asciichat_thread_t *thread)
Check if a thread handle has been initialized.
void platform_sleep_us(unsigned int us)
High-precision sleep function with microsecond precision.
int cond_destroy(cond_t *cond)
Destroy a condition variable.
int mutex_destroy(mutex_t *mutex)
Destroy a mutex.
Definition threading.c:22
asciichat_error_t session_host_remove_client(session_host_t *host, uint32_t client_id)
Remove a client by ID.
Definition host.c:1391
int socket_t
asciichat_error_t tcp_server_stop_client_threads(tcp_server_t *server, socket_t client_socket)
Stop all threads for a client in stop_id order.
#define asciichat_thread_equal(t1, t2)
#define asciichat_thread_join(thread, timeout_ms)
void stop_client_render_threads(client_info_t *client)
Stop and cleanup per-client rendering threads.
asciichat_thread_t audio_render_thread
asciichat_thread_t dispatch_thread
asciichat_thread_t receive_thread
asciichat_thread_t video_render_thread
void acip_transport_destroy(acip_transport_t *transport)
Destroy transport and free all resources.

References acip_transport_destroy(), client_info::active, ASCIICHAT_OK, asciichat_thread_equal, asciichat_thread_is_initialized(), asciichat_thread_join, asciichat_thread_self, atomic_load_bool(), atomic_store_bool(), client_info::audio_queue, client_info::audio_render_thread, broadcast_server_state_to_all_clients(), client_manager_t::client_count, client_info::client_id, client_info::client_state_mutex, client_manager_t::clients, client_manager_t::clients_by_id, cond_destroy(), client_info::crypto_handshake_ctx, crypto_handshake_destroy(), client_info::crypto_initialized, client_info::dispatch_queue_cond, client_info::dispatch_thread, client_info::dispatch_thread_running, client_info::display_name, ERROR_INVALID_PARAM, ERROR_NOT_FOUND, g_audio_mixer, g_client_manager, g_client_manager_rwlock, INVALID_SOCKET_VALUE, client_info::is_sending_audio, client_info::is_sending_video, client_info::is_tcp_client, log_debug, log_error, log_info, log_warn, MAX_CLIENTS, MAX_DISPLAY_NAME_LEN, mixer_remove_source(), mutex_destroy(), mutex_lock, mutex_unlock, packet_queue_stop(), platform_sleep_us(), client_info::receive_thread, client_info::receive_thread_id, rwlock_wrlock, rwlock_wrunlock, SAFE_STRNCPY, client_info::send_mutex, client_info::send_thread, client_info::session_client_id, server_context_t::session_host, session_host_remove_client(), SET_ERRNO, client_info::shutting_down, client_info::socket, socket_close(), socket_shutdown(), stop_client_render_threads(), server_context_t::tcp_server, tcp_server_stop_client_threads(), client_info::transport, US_PER_MS_INT, and client_info::video_render_thread.

Referenced by add_client(), add_webrtc_client(), client_receive_thread(), and client_send_thread_func().

◆ start_webrtc_client_threads()

int start_webrtc_client_threads ( server_context_t *  server_ctx,
const char *  client_id 
)

Start threads for a WebRTC client after crypto initialization.

This is called by WebSocket handler after crypto handshake is initialized. It ensures receive thread doesn't try to process packets before crypto context exists.

Parameters
server_ctxServer context
client_idClient ID to start threads for
Returns
0 on success, -1 on failure

Definition at line 2832 of file src/server/client.c.

2832 {
2833 if (!server_ctx) {
2834 SET_ERRNO(ERROR_INVALID_PARAM, "Server context is NULL");
2835 return -1;
2836 }
2837
2838 client_info_t *client = find_client_by_id(client_id);
2839 if (!client) {
2840 SET_ERRNO(ERROR_NOT_FOUND, "Client %s not found", client_id);
2841 return -1;
2842 }
2843
2844 log_debug("Starting threads for WebRTC client %s...", client_id);
2845 return start_client_threads(server_ctx, client, false);
2846}
client_info_t * find_client_by_id(const char *client_id)
Fast O(1) client lookup by ID using hash table.

References ERROR_INVALID_PARAM, ERROR_NOT_FOUND, find_client_by_id(), log_debug, and SET_ERRNO.

◆ stop_client_threads()

void stop_client_threads ( client_info_t *  client)

Definition at line 2848 of file src/server/client.c.

2848 {
2849 if (!client) {
2850 SET_ERRNO(ERROR_INVALID_PARAM, "Client is NULL");
2851 return;
2852 }
2853
2854 // Signal threads to stop
2855 log_debug("Setting active=false in stop_client_threads (client_id=%s)", client->client_id);
2856 atomic_store_bool(&client->active, false);
2857 atomic_store_bool(&client->send_thread_running, false);
2858
2859 // Wait for threads to finish
2861 asciichat_thread_join(&client->send_thread, NULL);
2862 }
2864 asciichat_thread_join(&client->receive_thread, NULL);
2865 }
2866 // For async dispatch: stop dispatch thread if running
2869 asciichat_thread_join(&client->dispatch_thread, NULL);
2870 }
2871}

References client_info::active, asciichat_thread_is_initialized(), asciichat_thread_join, atomic_store_bool(), client_info::client_id, client_info::dispatch_thread, client_info::dispatch_thread_running, ERROR_INVALID_PARAM, log_debug, client_info::receive_thread, client_info::send_thread, client_info::send_thread_running, and SET_ERRNO.

Variable Documentation

◆ g_client_manager

client_manager_t g_client_manager
extern

Global client manager singleton - central coordination point.

This is the primary data structure for managing all connected clients. It serves as the bridge between main.c's connection accept loop and the per-client threading architecture.

STRUCTURE COMPONENTS:

  • clients[]: Array backing storage for client_info_t structs
  • client_hashtable: O(1) lookup table for client_id -> client_info_t*
  • client_count: Current number of active clients
  • mutex: Legacy mutex (mostly replaced by rwlock)
  • next_client_id: Monotonic counter for unique client identification

THREAD SAFETY: Protected by g_client_manager_rwlock for concurrent access

Definition at line 347 of file src/server/client.c.

◆ g_client_manager_rwlock

rwlock_t g_client_manager_rwlock
extern

Reader-writer lock protecting the global client manager.

This lock enables high-performance concurrent access patterns:

  • Multiple threads can read client data simultaneously (stats, rendering)
  • Only one thread can modify client data at a time (add/remove operations)
  • Eliminates contention between read-heavy operations

USAGE PATTERN:

Definition at line 362 of file src/server/client.c.

362{0};

◆ g_client_manager_rwlock_initialized

bool g_client_manager_rwlock_initialized
extern

Track whether g_client_manager_rwlock was successfully initialized.

Definition at line 365 of file src/server/client.c.