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

Per-client rendering threads with rate limiting. More...

Go to the source code of this file.

Macros

#define VIDEO_RENDER_FPS   60
 
#define AUDIO_PACKET_FPS   50
 

Functions

void * client_video_render_thread (void *arg)
 Interruptible sleep function with platform-specific optimizations.
 
void * client_audio_render_thread (void *arg)
 Main audio rendering thread function for individual clients.
 
int create_client_render_threads (server_context_t *server_ctx, client_info_t *client)
 Create and initialize per-client rendering threads.
 
void stop_client_render_threads (client_info_t *client)
 Stop and cleanup per-client rendering threads.
 

Detailed Description

Per-client rendering threads with rate limiting.

Definition in file server/render.h.

Macro Definition Documentation

◆ AUDIO_PACKET_FPS

#define AUDIO_PACKET_FPS   50

Definition at line 26 of file server/render.h.

◆ VIDEO_RENDER_FPS

#define VIDEO_RENDER_FPS   60

Definition at line 22 of file server/render.h.

Function Documentation

◆ client_audio_render_thread()

void * client_audio_render_thread ( void *  arg)

Main audio rendering thread function for individual clients.

This is the audio processing thread that generates personalized audio mixes for a specific client at 100fps (10ms intervals, 480 samples @ 48kHz). Each client receives audio from all other clients while excluding their own audio to prevent echo and feedback.

THREAD EXECUTION FLOW:

  1. INITIALIZATION:
    • Validate client parameters and socket
    • Allocate local mixing buffer
    • Log thread startup for debugging
  2. MAIN PROCESSING LOOP:
    • Check thread running state (with mutex protection)
    • Take atomic snapshot of client state
    • Generate audio mix excluding this client's audio
    • Queue mixed audio for delivery to client
    • Sleep for precise timing (10ms intervals)
  3. CLEANUP AND EXIT:
    • Log thread termination
    • Return NULL to indicate clean exit

AUDIO MIXING STRATEGY:

ECHO PREVENTION:

  • Excludes client's own audio from their mix
  • Prevents audio feedback loops
  • Uses mixer_process_excluding_source() function
  • Maintains audio quality for other participants

PERFORMANCE CHARACTERISTICS:

TIMING PRECISION:

  • Target: 100fps (10ms intervals, 480 samples @ 48kHz)
  • Accumulates to 960 samples (20ms) for Opus encoding
  • Fixed timing (no dynamic rate adjustment)
  • Low-latency audio delivery

BUFFER MANAGEMENT:

  • Uses fixed-size local buffer (AUDIO_FRAMES_PER_BUFFER)
  • No dynamic allocation in processing loop
  • Predictable memory usage
  • Cache-friendly processing pattern

THREAD SAFETY MECHANISMS:

STATE SYNCHRONIZATION:

  • Uses client->client_state_mutex for state access
  • Implements snapshot pattern for client ID and queues
  • Prevents race conditions during client disconnect
  • Safe concurrent access with other threads

MIXER COORDINATION:

  • Global audio mixer is internally thread-safe
  • Multiple audio threads can process concurrently
  • Lock-free audio buffer operations where possible

AUDIO PIPELINE INTEGRATION:

AUDIO MIXING:

  • Uses global g_audio_mixer for multi-client processing
  • Processes AUDIO_FRAMES_PER_BUFFER samples per iteration
  • Handles varying number of active audio sources

AUDIO DELIVERY:

  • Queues audio directly in client's audio packet queue
  • Higher priority than video packets in send thread
  • Real-time delivery requirements

ERROR HANDLING STRATEGY:

  • Invalid client parameters: Log error and exit immediately
  • Missing audio mixer: Continue with polling (mixer may initialize later)
  • Queue failures: Log debug info and continue (expected under load)
  • Client disconnect: Clean exit without error spam

PERFORMANCE OPTIMIZATIONS:

  • Local buffer allocation (no malloc in loop)
  • Minimal processing per iteration
  • Efficient mixer integration
  • Low-overhead packet queuing
Parameters
argPointer to client_info_t for the target client
Returns
NULL on thread completion (always clean exit)
Note
This function runs in its own dedicated thread
Audio processing has higher priority than video
Thread lifetime matches client connection lifetime
Warning
Invalid client parameter causes immediate thread termination
Thread must be joined properly to prevent resource leaks
See also
mixer_process_excluding_source() For Audio mixing implementation
AUDIO_FRAMES_PER_BUFFER For buffer size constant

Definition at line 783 of file server/render.c.

783 {
784 client_info_t *client = (client_info_t *)arg;
785
786 if (!client) {
787 log_error("Invalid client info in audio render thread");
788 return NULL;
789 }
790
791 // Take snapshot of client ID and display name at start to avoid race conditions
792 char thread_client_id[MAX_CLIENT_ID_LEN];
793 SAFE_STRNCPY(thread_client_id, client->client_id, sizeof(thread_client_id) - 1);
794 char thread_display_name[64];
795 bool is_webrtc = (client->socket == INVALID_SOCKET_VALUE);
796 (void)is_webrtc; // May be unused in release builds
797
798 // LOCK OPTIMIZATION: Only need client_state_mutex, not global rwlock
799 // We already have a stable client pointer
801 SAFE_STRNCPY(thread_display_name, client->display_name, sizeof(thread_display_name));
803
804#ifdef DEBUG_THREADS
805 log_debug("Audio render thread started for client %s (%s), webrtc=%d", thread_client_id, thread_display_name,
806 is_webrtc);
807#endif
808
809 // Mix buffer: up to 960 samples for adaptive reading
810 // Normal: 480 samples = 10ms @ 48kHz
811 // Catchup: 960 samples = 20ms when buffers are filling up
812 float mix_buffer[960];
813
814// Opus frame accumulation buffer (960 samples = 20ms @ 48kHz)
815// Opus requires minimum 480 samples, 960 is optimal for 20ms frames
816#define OPUS_FRAME_SAMPLES 960
817 float opus_frame_buffer[OPUS_FRAME_SAMPLES];
818 int opus_frame_accumulated = 0;
819
820 // Create Opus encoder for this client's audio stream (48kHz, mono, 128kbps, AUDIO mode for music quality)
821 opus_codec_t *opus_encoder = opus_codec_create_encoder(OPUS_APPLICATION_AUDIO, 48000, 128000);
822 if (!opus_encoder) {
823 log_error("Failed to create Opus encoder for audio render thread (client %s)", thread_client_id);
824 return NULL;
825 }
826
827 // FPS tracking for audio render thread
828 fps_t audio_fps_tracker = {0};
829 fps_init(&audio_fps_tracker, AUDIO_PACKET_FPS, "SERVER AUDIO");
830
831 // Per-thread counters (NOT static - each thread instance gets its own)
832 int server_audio_frame_count = 0;
833
834 bool should_continue = true;
835 uint64_t audio_deadline = time_get_ns();
836 while (should_continue && !atomic_load_bool(&g_should_exit) && !atomic_load_bool(&client->shutting_down)) {
837 log_debug_every(LOG_RATE_SLOW, "Audio render loop iteration for client %s", thread_client_id);
838
839 // Check for immediate shutdown
841 log_debug("Audio render thread stopping for client %s (g_should_exit)", thread_client_id);
842 break;
843 }
844
845 // Check thread state before acquiring any locks to prevent use-after-destroy.
846 // If we acquire locks after client is being destroyed, we'll crash with SIGSEGV
847 should_continue = (((int)atomic_load_bool(&client->audio_render_thread_running) != 0) &&
848 ((int)atomic_load_bool(&client->active) != 0) && !atomic_load_bool(&client->shutting_down));
849
850 if (!should_continue) {
851 log_debug("Audio render thread stopping for client %s (should_continue=false)", thread_client_id);
852 break;
853 }
854
855 if (!g_audio_mixer) {
856 log_dev_every(10 * NS_PER_MS_INT, "Audio render waiting for mixer (client %s)", thread_client_id);
857 // Check shutdown flag while waiting
859 break;
861 continue;
862 }
863
864 // Optimization: No mutex needed - all fields are atomic or stable.
865 // client_id: atomic_t - use atomic_load for thread safety
866 // active: atomic_t - use atomic_load
867 // audio_queue: Assigned once at init and never changes
868 const char *client_id_snapshot = client->client_id; // Atomic read
869 bool active_snapshot = atomic_load_bool(&client->active); // Atomic read
870 packet_queue_t *audio_queue_snapshot = client->audio_queue; // Stable after init
871
872 // Check if client is still active after getting snapshot
873 if (!active_snapshot || !audio_queue_snapshot) {
874 break;
875 }
876
877 // Create mix excluding THIS client's audio using snapshot data
878 START_TIMER("mix_%s", client_id_snapshot);
879
880 // ADAPTIVE READING: Read more samples when we're behind to catch up
881 // Normal: 480 samples per 10ms iteration
882 // When behind: read up to 960 samples to catch up faster
883 // Check source buffer levels to decide
884 int samples_to_read = 480; // Default: 10ms worth
885
886 // Log latency at each stage in the server pipeline
887 // NOTE: Skipping mixer source latency logging due to potential use-after-free
888 // when source IDs are freed during concurrent client cleanup
889 if (g_audio_mixer) {
890 // Log outgoing queue latency
891 size_t queue_depth = packet_queue_size(audio_queue_snapshot);
892 float queue_latency_ms = (float)queue_depth * 20.0f; // ~20ms per Opus packet
893 log_dev_every(5 * NS_PER_MS_INT, "LATENCY: Server send queue for client %s: %.1fms (%zu packets)",
894 client_id_snapshot, queue_latency_ms, queue_depth);
895 }
896
897 int samples_mixed = 0;
898 uint32_t client_id_hash = fnv1a_hash_string(client_id_snapshot);
899 samples_mixed = mixer_process_excluding_source(g_audio_mixer, mix_buffer, samples_to_read, client_id_hash);
900
901 STOP_TIMER_AND_LOG_EVERY(dev, NS_PER_SEC_INT, 5 * NS_PER_MS_INT, "mix_%s", "Mixer for client %s: took",
902 client_id_snapshot);
903
904 // Debug logging every 100 iterations (disabled - can slow down audio rendering)
905 // log_debug_every(LOG_RATE_SLOW, "Audio render for client %u: samples_mixed=%d", client_id_snapshot,
906 // samples_mixed);
907
908 // Accumulate all samples (including 0 or partial) until we have a full Opus frame
909 // This maintains continuous stream without silence padding
910 START_TIMER("accum_%s", client_id_snapshot);
911
912 int space_available = OPUS_FRAME_SAMPLES - opus_frame_accumulated;
913 int samples_to_copy = (samples_mixed <= space_available) ? samples_mixed : space_available;
914
915 // Only copy if we have samples, otherwise just wait for next frame
916 if (samples_to_copy > 0) {
917 SAFE_MEMCPY(opus_frame_buffer + opus_frame_accumulated,
918 (OPUS_FRAME_SAMPLES - opus_frame_accumulated) * sizeof(float), mix_buffer,
919 samples_to_copy * sizeof(float));
920 opus_frame_accumulated += samples_to_copy;
921 }
922
923 STOP_TIMER_AND_LOG_EVERY(dev, NS_PER_SEC_INT, 2 * NS_PER_MS_INT, "accum_%s", "Accumulate for client %s: took",
924 client_id_snapshot);
925
926 // Only encode and send when we have accumulated a full Opus frame
927 if (opus_frame_accumulated >= OPUS_FRAME_SAMPLES) {
928 // A completed Opus frame is only 20 ms, while the outgoing audio queue
929 // is the same ordered path that precedes video. Enforce the bound for
930 // every frame: checking only once per 100 iterations lets a 200 ms
931 // backlog grow past a second before the producer reacts.
932 size_t queue_depth = packet_queue_size(audio_queue_snapshot);
933 bool apply_backpressure = queue_depth >= 10;
934 if (apply_backpressure) {
936 "Audio backpressure for client %s: queue depth %zu packets (%.1fms buffered)",
937 client_id_snapshot, queue_depth, (float)queue_depth / 50.0f * 1000.0f);
938 }
939
940 if (apply_backpressure) {
941 // Skip this packet to let the queue drain
942 // Reset accumulation buffer so fresh samples can be captured on next iteration.
943 // Without this reset, we'd loop forever with stale audio and no space for new samples
944 opus_frame_accumulated = 0;
945 // Reduced sleep from 5.8ms to 1ms to minimize audio gaps on localhost
947 continue;
948 }
949
950 // Encode accumulated Opus frame (960 samples = 20ms @ 48kHz)
951 uint8_t opus_buffer[1024]; // Max Opus frame size
952
953 START_TIMER("opus_encode_%s", client_id_snapshot);
954
955 int opus_size =
956 opus_codec_encode(opus_encoder, opus_frame_buffer, OPUS_FRAME_SAMPLES, opus_buffer, sizeof(opus_buffer));
957
958 STOP_TIMER_AND_LOG_EVERY(dev, NS_PER_SEC_INT, 10 * NS_PER_MS_INT, "opus_encode_%s",
959 "Opus encode for client %s: took", client_id_snapshot);
960
961 // DEBUG: Log mix buffer and encoding results to see audio levels being sent
962 {
963 float peak = 0.0f, rms = 0.0f;
964 for (int i = 0; i < OPUS_FRAME_SAMPLES; i++) {
965 float abs_val = fabsf(opus_frame_buffer[i]);
966 if (abs_val > peak)
967 peak = abs_val;
968 rms += opus_frame_buffer[i] * opus_frame_buffer[i];
969 }
970 rms = sqrtf(rms / OPUS_FRAME_SAMPLES);
971 // NOTE: server_audio_frame_count is now per-thread (not static), so each client thread has its own counter
972 server_audio_frame_count++;
973 if (server_audio_frame_count <= 5 || server_audio_frame_count % 20 == 0) {
974 // Log first 4 samples to verify they look like valid audio (not NaN/Inf/garbage)
976 "Server audio frame #%d for client %s: samples_mixed=%d, Peak=%.6f, RMS=%.6f, opus_size=%d, "
977 "first4=[%.4f,%.4f,%.4f,%.4f]",
978 server_audio_frame_count, client_id_snapshot, samples_mixed, peak, rms, opus_size,
979 opus_frame_buffer[0], opus_frame_buffer[1], opus_frame_buffer[2], opus_frame_buffer[3]);
980 }
981 }
982
983 // Always reset accumulation buffer after attempting to encode - we've consumed these samples
984 // If we don't reset, new audio samples would be dropped while stale data sits in the buffer
985 opus_frame_accumulated = 0;
986
987 if (opus_size <= 0) {
988 log_error("Failed to encode audio to Opus for client %s: opus_size=%d", client_id_snapshot, opus_size);
989 } else {
990 // Wrap Opus data in proper batch packet format for client parsing:
991 // Header: sample_rate (4B) + frame_duration (4B) + frame_count (4B)
992 // Frame sizes: frame_count * 2B (one uint16 per frame with Opus data size)
993 // Opus data: raw encoded audio
994
995 // Fixed values for Opus encoding (standard 48kHz, 960 samples per frame = 20ms)
996 const uint32_t sample_rate = 48000;
997 const uint32_t frame_duration = 20; // milliseconds (960 samples / 48kHz = 20ms)
998 const uint32_t frame_count = 1; // one Opus frame per packet
999 const uint16_t frame_size = (uint16_t)opus_size;
1000
1001 // Calculate total packet size: 16-byte header + frame_sizes + opus data
1002 const size_t header_size = 16; // 4 uint32s (sample_rate, frame_duration, frame_count, reserved)
1003 const size_t frame_sizes_size = (size_t)frame_count * sizeof(uint16_t);
1004 const size_t total_packet_size = header_size + frame_sizes_size + (size_t)opus_size;
1005
1006 // Allocate buffer for the complete packet
1007 uint8_t *packet_buffer = SAFE_MALLOC(total_packet_size, uint8_t *);
1008 if (!packet_buffer) {
1009 log_error("Failed to allocate Opus batch packet buffer for client %s (size=%zu)", client_id_snapshot,
1010 total_packet_size);
1011 } else {
1012 // Build header in network byte order
1013 uint8_t *ptr = packet_buffer;
1014
1015 // Write sample_rate (4 bytes, network byte order)
1016 uint32_t sr_net = htonl(sample_rate);
1017 memcpy(ptr, &sr_net, 4);
1018 ptr += 4;
1019
1020 // Write frame_duration (4 bytes, network byte order)
1021 uint32_t fd_net = htonl(frame_duration);
1022 memcpy(ptr, &fd_net, 4);
1023 ptr += 4;
1024
1025 // Write frame_count (4 bytes, network byte order)
1026 uint32_t fc_net = htonl(frame_count);
1027 memcpy(ptr, &fc_net, 4);
1028 ptr += 4;
1029
1030 // Write reserved field (4 bytes, all zeros)
1031 uint32_t reserved = 0;
1032 memcpy(ptr, &reserved, 4);
1033 ptr += 4;
1034
1035 // Write frame sizes (2 bytes per frame, network byte order)
1036 for (uint32_t i = 0; i < frame_count; i++) {
1037 uint16_t fs_net = htons(frame_size);
1038 memcpy(ptr, &fs_net, 2);
1039 ptr += 2;
1040 }
1041
1042 // Copy Opus data
1043 memcpy(ptr, opus_buffer, (size_t)opus_size);
1044
1045 // Queue the properly formatted packet
1046 START_TIMER("audio_queue_%s", client_id_snapshot);
1047
1048 int result = packet_queue_enqueue(audio_queue_snapshot, PACKET_TYPE_AUDIO_OPUS_BATCH, packet_buffer,
1049 total_packet_size, 0, true);
1050
1051 STOP_TIMER_AND_LOG_EVERY(dev, NS_PER_SEC_INT, 1 * NS_PER_MS_INT, "audio_queue_%s",
1052 "Audio queue for client %s: took", client_id_snapshot);
1053
1054 if (result < 0) {
1055 log_debug("Failed to queue Opus audio for client %s", client_id_snapshot);
1056 } else {
1057 // FPS tracking - audio packet successfully queued (handles lag detection and periodic reporting)
1058 fps_frame_ns(&audio_fps_tracker, time_get_ns(), "audio packet queued");
1059 }
1060
1061 // Free the allocated packet buffer after queuing (packet_queue_enqueue copies data)
1062 SAFE_FREE(packet_buffer);
1063 }
1064 }
1065 // NOTE: opus_frame_accumulated is already reset at line 928 after encode attempt
1066 }
1067
1068 // Audio mixing rate limiting using adaptive sleep system
1069 // Target: 10ms intervals (100 FPS) for 480 samples @ 48kHz
1070 // Use queue_depth=0 and target_depth=0 for constant-rate audio processing
1071 audio_deadline += 10 * NS_PER_MS_INT;
1072 uint64_t now = time_get_ns();
1073 if (audio_deadline > now)
1074 platform_sleep_ns(audio_deadline - now);
1075 else if (now - audio_deadline > 100 * NS_PER_MS_INT)
1076 audio_deadline = now;
1077 }
1078
1079#ifdef DEBUG_THREADS
1080 log_debug("Audio render thread stopped for client %s", thread_client_id);
1081#endif
1082
1083 // Clean up Opus encoder
1084 if (opus_encoder) {
1085 opus_codec_destroy(opus_encoder);
1086 }
1087
1088 // Clean up thread-local error context before exit
1090
1091 return NULL;
1092}
bool atomic_load_bool(atomic_t *a)
Atomically load a boolean value.
Definition atomic.c:169
atomic_t g_should_exit
Global application exit flag (shared across all modes)
Definition globals.c:34
opus_codec_t * opus_codec_create_encoder(opus_application_t application, int sample_rate, int bitrate)
Create an Opus encoder.
Definition opus.c:19
int mixer_process_excluding_source(mixer_t *mixer, float *output, int num_samples, uint32_t exclude_client_id)
Process audio from all sources except one (for per-client output)
Definition mixer.c:687
void opus_codec_destroy(opus_codec_t *codec)
Destroy an Opus codec instance.
Definition opus.c:230
size_t opus_codec_encode(opus_codec_t *codec, const float *samples, int num_samples, uint8_t *out_data, size_t out_size)
Encode audio frame with Opus.
Definition opus.c:112
@ OPUS_APPLICATION_AUDIO
General audio (optimized for music)
Definition opus.h:78
unsigned short uint16_t
Definition common.h:57
unsigned int uint32_t
Definition common.h:58
#define SAFE_STRNCPY(dst, src, size)
Definition common.h:414
#define SAFE_FREE(ptr)
Definition common.h:376
#define SAFE_MALLOC(size, cast)
Definition common.h:264
void fps_frame_ns(fps_t *tracker, uint64_t current_time_ns, const char *context)
Track a frame and detect lag conditions (nanosecond version - PRIMARY)
Definition fps.c:52
unsigned long long uint64_t
Definition common.h:59
unsigned char uint8_t
Definition common.h:56
#define SAFE_MEMCPY(dest, dest_size, src, count)
Definition common.h:468
void fps_init(fps_t *tracker, int expected_fps, const char *name)
Initialize FPS tracker.
Definition fps.c:32
void asciichat_errno_destroy(void)
Cleanup error system resources.
#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_RATE_SLOW
Log rate limit: 10 seconds (10,000,000 microseconds)
Definition log_rates.h:35
#define log_error(...)
Log an ERROR message.
Definition log/log.h:587
#define log_debug(...)
Log a DEBUG message.
Definition log/log.h:548
#define STOP_TIMER_AND_LOG_EVERY(log_level, interval_ns, threshold_ns, timer_name, msg_fmt,...)
Stop a timer and log the result with rate limiting.
Definition time.h:647
uint64_t time_get_ns(void)
Get current monotonic time in nanoseconds.
Definition util/time.c:108
#define NS_PER_SEC_INT
Definition time.h:157
#define START_TIMER(name_fmt,...)
Start a timer with formatted name.
Definition time.h:333
#define NS_PER_MS_INT
Definition time.h:156
#define US_PER_MS_INT
Definition time.h:160
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
size_t packet_queue_size(packet_queue_t *queue)
Get current number of packets in queue.
Definition queue.c:595
@ PACKET_TYPE_AUDIO_OPUS_BATCH
Batched Opus-encoded audio frames.
Definition packet.h:374
void platform_sleep_ns(uint64_t ns)
Platform-safe sleep function with nanosecond precision.
#define mutex_lock(mutex)
Lock a mutex (with debug tracking in debug builds)
#define INVALID_SOCKET_VALUE
Invalid socket value (POSIX: -1)
Definition socket.h:278
#define mutex_unlock(mutex)
Unlock a mutex (with debug tracking in debug builds)
#define log_debug_every(interval_us, fmt,...)
Rate-limited DEBUG logging.
Definition log/log.h:702
#define log_warn_every(interval_us, fmt,...)
Rate-limited WARN logging.
Definition log/log.h:708
#define log_dev_every(interval_us, fmt,...)
Rate-limited DEV logging.
Definition log/log.h:699
mixer_t *volatile g_audio_mixer
Global audio mixer instance for multi-client audio processing.
#define OPUS_FRAME_SAMPLES
#define AUDIO_PACKET_FPS
Per-client state structure for server-side client management.
char display_name[MAX_DISPLAY_NAME_LEN]
char client_id[MAX_CLIENT_ID_LEN]
FPS tracking state.
Definition fps.h:52
Opus codec context for encoding or decoding.
Definition opus.h:95
Thread-safe packet queue for producer-consumer communication.
Definition queue.h:211

References client_info::active, asciichat_errno_destroy(), atomic_load_bool(), AUDIO_PACKET_FPS, client_info::audio_queue, client_info::audio_render_thread_running, client_info::client_id, client_info::client_state_mutex, client_info::display_name, fps_frame_ns(), fps_init(), g_audio_mixer, g_should_exit, INVALID_SOCKET_VALUE, log_debug, log_debug_every, log_dev_every, log_error, LOG_RATE_SLOW, log_warn_every, MAX_CLIENT_ID_LEN, mixer_process_excluding_source(), mutex_lock, mutex_unlock, NS_PER_MS_INT, NS_PER_SEC_INT, OPUS_APPLICATION_AUDIO, opus_codec_create_encoder(), opus_codec_destroy(), opus_codec_encode(), OPUS_FRAME_SAMPLES, packet_queue_enqueue(), packet_queue_size(), PACKET_TYPE_AUDIO_OPUS_BATCH, platform_sleep_ns(), SAFE_FREE, SAFE_MALLOC, SAFE_MEMCPY, SAFE_STRNCPY, client_info::shutting_down, client_info::socket, START_TIMER, STOP_TIMER_AND_LOG_EVERY, time_get_ns(), and US_PER_MS_INT.

Referenced by create_client_render_threads().

◆ client_video_render_thread()

void * client_video_render_thread ( void *  arg)

Interruptible sleep function with platform-specific optimizations.

Provides a sleep mechanism that can be interrupted by the global shutdown signal, enabling responsive thread termination. The implementation varies by platform to optimize for responsiveness vs CPU usage.

PLATFORM-SPECIFIC BEHAVIOR:

WINDOWS IMPLEMENTATION:

  • Uses Sleep() API with millisecond precision
  • Simple approach due to Windows timer characteristics
  • Sleep(1) may sleep up to 15.6ms due to timer resolution
  • Checks shutdown flag before and after sleep

POSIX IMPLEMENTATION (Linux/macOS):

  • Uses condition variables for precise timing
  • Can be interrupted immediately by shutdown signal
  • Higher responsiveness for thread termination
  • Uses static mutex/condition variable for coordination

USAGE PATTERNS:

  • Rate limiting in render threads (maintain target FPS)
  • Backoff delays when queues are full
  • Polling intervals for state changes
  • Graceful busy-wait alternatives

PERFORMANCE CHARACTERISTICS:

  • Low CPU overhead (proper sleep, not busy wait)
  • Sub-second responsiveness to shutdown
  • Maintains timing precision for media processing
  • Scales well with many concurrent threads

ERROR HANDLING:

  • Early return if shutdown flag is set
  • Platform-specific error handling for sleep APIs
  • Logging for debugging (throttled to avoid spam)
Parameters
usecSleep duration in microseconds
Note
Function may return early if shutdown is requested
Minimum sleep time is platform-dependent (typically 1ms)
Static variables are used for call counting (debugging)
Warning
Not suitable for precise timing (use for rate limiting only)
See also
VIDEO_RENDER_FPS For Video to ASCII Conversion thread timing requirements

Main video rendering thread function for individual clients

This is the core video processing thread that generates personalized ASCII art frames for a specific client at 60fps. Each connected client gets their own dedicated video thread, providing linear performance scaling and personalized rendering based on terminal capabilities.

THREAD EXECUTION FLOW:

  1. INITIALIZATION:
    • Validate client parameters and socket
    • Initialize timing variables for rate limiting
    • Log thread startup for debugging
  2. MAIN PROCESSING LOOP:
    • Check thread running state (with mutex protection)
    • Rate limit to 60fps using high-resolution timing
    • Take atomic snapshot of client state
    • Generate personalized ASCII frame for this client
    • Queue frame for delivery to client
    • Update timing for next iteration
  3. CLEANUP AND EXIT:
    • Log thread termination
    • Return NULL to indicate clean exit

PERFORMANCE CHARACTERISTICS:

TIMING PRECISION:

  • Target: 60fps (16.67ms intervals)
  • Uses CLOCK_MONOTONIC for accurate timing
  • Only sleeps if ahead of schedule (no CPU waste)
  • Maintains consistent frame rate under varying load

CLIENT-SPECIFIC PROCESSING:

  • Generates frames customized for client's terminal
  • Respects color depth, palette, and size preferences
  • Adapts to client capability changes dynamically
  • Handles client disconnection gracefully

THREAD SAFETY MECHANISMS:

STATE SYNCHRONIZATION:

  • Uses client->client_state_mutex for all state access
  • Implements snapshot pattern (copy state, release lock)
  • Prevents race conditions with other threads
  • Safe concurrent access to client data

SHUTDOWN COORDINATION:

  • Monitors global g_should_exit flag
  • Checks per-client running flags with mutex protection
  • Responds to shutdown within one frame interval (16.67ms)

VIDEO PIPELINE INTEGRATION:

FRAME GENERATION:

FRAME DELIVERY:

  • Calls queue_ascii_frame_for_client() for packet queuing
  • Uses client's dedicated video packet queue
  • Send thread delivers frames asynchronously

ERROR HANDLING STRATEGY:

  • Invalid client parameters: Log error and exit immediately
  • Frame generation failures: Continue with debug logging
  • Queue overflow: Continue with debug logging (expected under load)
  • Client disconnect: Clean exit without error spam
  • Timing issues: Self-correcting (no accumulation)

LOGGING AND MONITORING:

  • Success logging: Throttled (every ~4 seconds)
  • Failure logging: Throttled (every ~10 seconds)
  • Performance metrics: Frame count and size tracking
  • Debug information: Timing and queue status
Parameters
argPointer to client_info_t for the target client
Returns
NULL on thread completion (always clean exit)
Note
This function runs in its own dedicated thread
Thread lifetime matches client connection lifetime
Memory allocated for frames is automatically freed
Warning
Invalid client parameter causes immediate thread termination
Thread must be joined properly to prevent resource leaks
See also
create_mixed_ascii_frame_for_client() For frame generation
queue_ascii_frame_for_client() For frame delivery
VIDEO_RENDER_FPS For timing constant definition

Definition at line 340 of file server/render.c.

340 {
341 client_info_t *client = (client_info_t *)arg;
342 if (!client) {
343 log_error("NULL client pointer in video render thread");
344 return NULL;
345 }
346
347 // Take snapshot of client ID and socket at start to avoid race conditions
348 char thread_client_id[MAX_CLIENT_ID_LEN];
349 SAFE_STRNCPY(thread_client_id, client->client_id, sizeof(thread_client_id) - 1);
350 socket_t thread_socket = client->socket;
351 bool is_webrtc = (thread_socket == INVALID_SOCKET_VALUE);
352 (void)is_webrtc; // May be unused in release builds
353
354 log_debug("Video render thread: client_id=%s, webrtc=%d", thread_client_id, is_webrtc);
355
356 log_info("[VIDEO_RENDER_THREAD_START] ★★★ LOCK STATE at thread entry for client %s", thread_client_id);
357 // debug_sync_print_state(); // Disabled: causes AddressSanitizer stack-use-after-return crash
358
359 // Wait for client to send terminal capabilities before rendering
360 // Without capabilities, convert_composite_to_ascii() will fail
361 int timeout_ms = 5000;
362 int waited_ms = 0;
364 if (waited_ms == 0) {
365 log_debug("Waiting for terminal capabilities from client %s...", thread_client_id);
366 }
367 APP_CALLBACK_VOID(platform_pump_events);
369 waited_ms += 10;
370 if (waited_ms >= timeout_ms) {
371 log_debug("Still waiting for terminal capabilities from client %s", thread_client_id);
372 waited_ms = 0;
373 }
374 }
375 if (client->has_terminal_caps) {
376 log_debug("Received terminal capabilities for client %s after %dms", thread_client_id, waited_ms);
377 }
378
379 // Get client's desired FPS from capabilities or use default
380 int client_fps = VIDEO_RENDER_FPS; // Default to 60 FPS
381 // Use snapshot pattern to avoid mutex in render thread
382 bool has_caps = client->has_terminal_caps;
383 int desired_fps = has_caps ? client->terminal_caps.desired_fps : 0;
384 if (has_caps && desired_fps > 0) {
385 client_fps = desired_fps;
386 log_debug("Client %s requested FPS: %d (has_caps=%d, desired_fps=%d)", thread_client_id, client_fps, has_caps,
387 desired_fps);
388 } else {
389 log_debug("Client %s using default FPS: %d (has_caps=%d, desired_fps=%d)", thread_client_id, client_fps, has_caps,
390 desired_fps);
391 }
392
393 int base_frame_interval_ms = 1000 / client_fps;
394 log_debug("Client %s render interval: %dms (%d FPS)", thread_client_id, base_frame_interval_ms, client_fps);
395
396 // FPS tracking for video render thread
397 fps_t video_fps_tracker = {0};
398 fps_init(&video_fps_tracker, client_fps, "SERVER VIDEO");
399
400 const uint64_t frame_interval_ns = NS_PER_SEC_INT / client_fps;
401 uint64_t next_frame_ns = time_get_ns() + frame_interval_ns;
402 bool last_has_sources = false;
403 uint32_t frame_gen_count = 0;
404 uint64_t frame_gen_start_time = 0;
405 uint32_t last_frame_hash = UINT32_MAX;
406 uint32_t commits_count = 0;
407 uint64_t commits_start_time = 0;
408 uint64_t render_rate_window_start_ns = time_get_ns();
409 uint64_t rendered_frames_in_window = 0;
410 uint64_t missed_deadlines_in_window = 0;
411
412 log_info("Video render loop STARTING for client %s", thread_client_id);
413
414 bool should_continue = true;
415 while (should_continue && !atomic_load_bool(&g_should_exit) && !atomic_load_bool(&client->shutting_down)) {
416 START_TIMER("render_iteration");
417 log_dev_every(10 * NS_PER_MS_INT, "Video render loop iteration for client %s", thread_client_id);
418
419 // Check for immediate shutdown
421 log_debug("Video render thread stopping for client %s (g_should_exit)", thread_client_id);
422 break;
423 }
424
425 bool video_running = atomic_load_bool(&client->video_render_thread_running);
426 bool active = atomic_load_bool(&client->active);
427 bool shutting_down = atomic_load_bool(&client->shutting_down);
428
429 should_continue = video_running && active && !shutting_down;
430
431 if (!should_continue) {
432 log_debug("Video render thread stopping for client %s (should_continue=false: video_running=%d, active=%d, "
433 "shutting_down=%d)",
434 thread_client_id, video_running, active, shutting_down);
435 break;
436 }
437
438 // Processing consumes part of the frame interval. Sleep only until the
439 // scheduled deadline, and skip missed deadlines without a catch-up burst.
440 uint64_t iter_start_ns = time_get_ns();
441 START_TIMER("render_frame_sleep");
442 if (iter_start_ns < next_frame_ns) {
443 platform_sleep_ns(next_frame_ns - iter_start_ns);
444 }
445 STOP_TIMER_AND_LOG(dev, 0, "render_frame_sleep", "frame deadline sleep completed");
446 uint64_t after_sleep_ns = time_get_ns();
447 uint64_t elapsed_intervals =
448 after_sleep_ns >= next_frame_ns ? (after_sleep_ns - next_frame_ns) / frame_interval_ns + 1 : 1;
449 if (elapsed_intervals > 1) {
450 missed_deadlines_in_window += elapsed_intervals - 1;
451 }
452 next_frame_ns += elapsed_intervals * frame_interval_ns;
453
454 // Capture timestamp for FPS tracking and frame timestamps
455 uint64_t current_time_ns = time_get_ns();
456
457 // Check thread state again before acquiring locks (client might have been destroyed during sleep).
458 should_continue = atomic_load_bool(&client->video_render_thread_running) && atomic_load_bool(&client->active) &&
460 if (!should_continue) {
461 break;
462 }
463
464 // Optimization: No mutex needed - all fields are atomic or stable.
465 // client_id: stable string pointer
466 // width/height: stable uint16_t values (set once, never modified)
467 // active: atomic_t - use atomic_load for thread safety
468 const char *client_id_snapshot = client->client_id; // Stable read
469 unsigned short width_snapshot = client->width; // Stable uint16_t
470 unsigned short height_snapshot = client->height; // Stable uint16_t
471 bool active_snapshot = atomic_load_bool(&client->active); // Atomic read
472
473 // Check if client is still active after getting snapshot
474 if (!active_snapshot) {
475 break;
476 }
477
478 // Phase 2 IMPLEMENTED: Generate frame specifically for THIS client using snapshot data
479 size_t frame_size = 0;
480
481 // Check if any clients are sending video
482 bool has_video_sources = any_clients_sending_video();
483
484 // DIAGNOSTIC: Track when video sources become available
485 if (has_video_sources != last_has_sources) {
486 log_warn("DIAGNOSTIC: Client %s video sources: %s", thread_client_id,
487 has_video_sources ? "AVAILABLE" : "UNAVAILABLE");
488 last_has_sources = has_video_sources;
489 }
490
492 "Video render iteration for client %s: has_video_sources=%d, width=%u, height=%u", thread_client_id,
493 has_video_sources, width_snapshot, height_snapshot);
494
495 // Use default dimensions if client dimensions not received
496 if (width_snapshot == 0 || height_snapshot == 0) {
497 log_dev_every(5 * NS_PER_MS_INT, "Using default dimensions for client %s (client reported: width=%u, height=%u)",
498 thread_client_id, width_snapshot, height_snapshot);
499 // Use 80x25 as standard terminal size fallback for snapshot/test clients
500 width_snapshot = 80;
501 height_snapshot = 25;
502 }
503
504 // Generate frames for all clients with valid dimensions, regardless of whether other clients are sending video.
505 // This ensures WebSocket snapshot clients (receive-only mode) receive frames even when no other clients are
506 // sending video.
507 int sources_count = 0; // Track number of video sources in this frame
508
509 // DIAGNOSTIC: Track every frame generation attempt
510 frame_gen_count++;
511 if (frame_gen_count == 1) {
512 frame_gen_start_time = current_time_ns;
513 }
514
515 // Log every 120 attempts (should be ~2 seconds at 60 Hz)
516 if (frame_gen_count % 120 == 0) {
517 uint64_t elapsed_ns = current_time_ns - frame_gen_start_time;
518 double gen_fps = (120.0 / (elapsed_ns / (double)NS_PER_SEC_INT));
519 log_dev("render loop: client=%s fps=%.1f", thread_client_id, gen_fps);
520 }
521
522 log_dev_every(5 * NS_PER_MS_INT, "About to call create_mixed_ascii_frame_for_client for client %s with dims %ux%u",
523 thread_client_id, width_snapshot, height_snapshot);
524
525 // debug_sync_print_state() is too expensive to call every frame (kills FPS)
526
527 uint64_t frame_create_start_ns = time_get_ns();
528 START_TIMER("render_create_frame");
529 char *ascii_frame = create_mixed_ascii_frame_for_client(client_id_snapshot, width_snapshot, height_snapshot, false,
530 &frame_size, NULL, &sources_count);
531 STOP_TIMER_AND_LOG(dev, 0, "render_create_frame", "create_mixed_ascii_frame completed");
532 uint64_t frame_create_end_ns = time_get_ns();
533
534 if (frame_gen_count % 120 == 0) {
535 char sleep_str[32], create_str[32];
536 time_pretty((uint64_t)(after_sleep_ns - iter_start_ns), -1, sleep_str, sizeof(sleep_str));
537 time_pretty((uint64_t)(frame_create_end_ns - frame_create_start_ns), -1, create_str, sizeof(create_str));
538 log_dev("render timing: adaptive_sleep=%s frame_create=%s", sleep_str, create_str);
539 }
540
541 // Always send frames at configured FPS, even if no external video sources yet
542 // This ensures single-client mode (like testing) still gets continuous frames
543 // The frame data (animated background or test pattern) is still valid content
544
545 // DEBUG: Log frame generation details
546 uint32_t current_frame_hash = 0;
547 if (ascii_frame && frame_size > 0) {
548 for (size_t i = 0; i < frame_size && i < 1000; i++) {
549 current_frame_hash = (uint32_t)((uint64_t)current_frame_hash * 31 + ((unsigned char *)ascii_frame)[i]);
550 }
551 if (current_frame_hash != last_frame_hash) {
552 log_dev("RENDER_FRAME CHANGE: Client %s frame #%zu sources=%d hash=0x%08x (prev=0x%08x)", thread_client_id,
553 frame_size, sources_count, current_frame_hash, last_frame_hash);
554 last_frame_hash = current_frame_hash;
555 } else {
556 log_dev_every(25000, "RENDER_FRAME DUPLICATE: Client %s frame #%zu sources=%d hash=0x%08x (no change)",
557 thread_client_id, frame_size, sources_count, current_frame_hash);
558 }
559 }
560
562 "create_mixed_ascii_frame_for_client returned: ascii_frame=%p, frame_size=%zu, sources_count=%d",
563 (void *)ascii_frame, frame_size, sources_count);
564
565 // Phase 2 IMPLEMENTED: Write frame to double buffer (never drops!)
566 if (ascii_frame && frame_size > 0) {
567 log_debug_every(5 * NS_PER_MS_INT, "Buffering frame for client %s (size=%zu)", thread_client_id, frame_size);
568 // GRID LAYOUT CHANGE DETECTION: Store source count with frame
569 // Send thread will compare this with last sent count to detect grid changes
570 atomic_store_u64(&client->last_rendered_grid_sources, (uint64_t)sources_count);
571
572 // Use double-buffer system which has its own internal swap_mutex
573 // No external locking needed - the double-buffer is thread-safe by design
574 video_frame_buffer_t *vfb_snapshot = client->outgoing_video_buffer;
575
576 if (vfb_snapshot) {
577 video_frame_t *write_frame = video_frame_begin_write(vfb_snapshot);
578 if (write_frame) {
579 // Copy ASCII frame data to the back buffer (NOT holding rwlock - just double-buffer's internal lock)
580 if (write_frame->data && frame_size <= vfb_snapshot->allocated_buffer_size) {
581 memcpy(write_frame->data, ascii_frame, frame_size);
582 write_frame->size = frame_size;
583 write_frame->width = width_snapshot;
584 write_frame->height = height_snapshot;
585 write_frame->capture_timestamp_ns = current_time_ns;
586
587 // Commit at the configured cadence enforced by the frame deadline.
588 // Never skip frames based on content changes.
589 {
590 uint64_t commit_start_ns = time_get_ns();
591 // Commit the frame (swaps buffers atomically using vfb->swap_mutex, NOT rwlock)
592 START_TIMER("render_video_frame_commit");
593 video_frame_commit(vfb_snapshot);
594 STOP_TIMER_AND_LOG(dev, 0, "render_video_frame_commit", "video_frame_commit completed");
595 uint64_t commit_end_ns = time_get_ns();
596 char commit_duration_str[32];
597 time_pretty((uint64_t)(commit_end_ns - commit_start_ns), -1, commit_duration_str,
598 sizeof(commit_duration_str));
599
600 commits_count++;
601 if (commits_count == 1) {
602 commits_start_time = commit_end_ns;
603 }
604 if (commits_count % 10 == 0) {
605 uint64_t elapsed_ns = commit_end_ns - commits_start_time;
606 double commit_fps = (10.0 / (elapsed_ns / (double)NS_PER_SEC_INT));
607 log_dev("render commit rate: client=%s fps=%.1f", thread_client_id, commit_fps);
608 }
609
610 log_dev("FRAME_COMMIT: Client %s took %s (hash=0x%08x)", thread_client_id,
611 commit_duration_str, current_frame_hash);
612
613 // ARCHITECTURE: Frame transmission is handled EXCLUSIVELY by the send thread
614 // The send thread reads from this buffer and applies proper 60 FPS rate limiting
615 // Do NOT send from render thread - it bypasses rate limiting and causes network congestion
616 // Render thread responsibility: Generate frame → Write to buffer → Done
617 // Send thread responsibility: Read from buffer → Apply 60 FPS pacing → Transmit
618 }
619
620 } else {
621 log_warn("Frame too large for buffer: %zu > %zu", frame_size, vfb_snapshot->allocated_buffer_size);
622 }
623
624 // FPS tracking - frame successfully generated (handles lag detection and periodic reporting)
625 fps_frame_ns(&video_fps_tracker, current_time_ns, "frame rendered");
626 rendered_frames_in_window++;
627 uint64_t render_rate_now_ns = time_get_ns();
628 uint64_t render_rate_elapsed_ns = render_rate_now_ns - render_rate_window_start_ns;
629 if (render_rate_elapsed_ns >= NS_PER_SEC_INT) {
630 double render_fps = (double)rendered_frames_in_window * NS_PER_SEC_INT / (double)render_rate_elapsed_ns;
631 if (render_fps < (double)client_fps * 0.9 || missed_deadlines_in_window > 0) {
632 log_warn("Server video render lagging client=%s fps=%.1f target=%d missed_deadlines=%llu",
633 thread_client_id, render_fps, client_fps,
634 (unsigned long long)missed_deadlines_in_window);
635 } else {
636 log_info("Server video render client=%s fps=%.1f target=%d missed_deadlines=0", thread_client_id,
637 render_fps, client_fps);
638 }
639 render_rate_window_start_ns = render_rate_now_ns;
640 rendered_frames_in_window = 0;
641 missed_deadlines_in_window = 0;
642 }
643 }
644 }
645
646 SAFE_FREE(ascii_frame);
647 } else {
648 // No frame generated (probably no video sources) - this is normal, no error logging needed
649 log_dev_every(10 * NS_PER_MS_INT, "Per-client render: No video sources available for client %s",
650 client_id_snapshot);
651 }
652
653 uint64_t iter_end_ns = time_get_ns();
654 if (frame_gen_count % 120 == 0) {
655 char total_str[32];
656 time_pretty((uint64_t)(iter_end_ns - iter_start_ns), -1, total_str, sizeof(total_str));
657 log_warn(" TIMING: TOTAL_ITERATION=%s", total_str);
658 }
659 STOP_TIMER_AND_LOG(dev, 0, "render_iteration", "render_iteration completed (total)");
660 }
661
662#ifdef DEBUG_THREADS
663 log_debug("Video render thread stopped for client %s", thread_client_id);
664#endif
665
666 // Clean up thread-local error context before exit
668
669 return NULL;
670}
#define APP_CALLBACK_VOID(callback_name)
void atomic_store_u64(atomic_t *a, uint64_t value)
Atomically store a uint64_t value.
Definition atomic.c:241
#define log_warn(...)
Log a WARN message.
Definition log/log.h:574
#define log_dev(...)
Log a DEV message (most verbose, development only)
Definition log/log.h:534
#define log_info(...)
Log an INFO message.
Definition log/log.h:561
#define STOP_TIMER_AND_LOG(log_level, threshold_ns, timer_name, msg_fmt,...)
Stop a timer and log the result with a custom message.
Definition time.h:612
int time_pretty(uint64_t nanoseconds, int decimals, char *buffer, size_t buffer_size)
Format nanoseconds as pretty duration with spaces and configurable precision.
Definition util/time.c:424
void platform_sleep_ms(unsigned int ms)
Sleep for a specified number of milliseconds.
video_frame_t * video_frame_begin_write(video_frame_buffer_t *vfb)
Writer API: Start writing a new frame.
void video_frame_commit(video_frame_buffer_t *vfb)
Writer API: Commit the frame and swap buffers.
int socket_t
#define VIDEO_RENDER_FPS
bool any_clients_sending_video(void)
Check if any connected clients are currently sending video.
Definition stream.c:1335
char * create_mixed_ascii_frame_for_client(const char *target_client_id, unsigned short width, unsigned short height, bool wants_stretch, size_t *out_size, bool *out_grid_changed, int *out_sources_count)
Generate personalized ASCII frame for a specific client.
Definition stream.c:973
terminal_capabilities_t terminal_caps
video_frame_buffer_t * outgoing_video_buffer
uint8_t desired_fps
Client's desired frame rate (1-144 FPS)
Definition terminal.h:733
Video frame buffer manager.
size_t allocated_buffer_size
Size of allocated data buffers (for cleanup)
Video frame structure.
uint32_t height
Frame height in pixels.
size_t size
Size of frame data in bytes.
uint32_t width
Frame width in pixels.
uint64_t capture_timestamp_ns
Timestamp when frame was captured (nanoseconds)
void * data
Frame data pointer (points to pre-allocated buffer)

References client_info::active, video_frame_buffer_t::allocated_buffer_size, any_clients_sending_video(), APP_CALLBACK_VOID, asciichat_errno_destroy(), atomic_load_bool(), atomic_store_u64(), video_frame_t::capture_timestamp_ns, client_info::client_id, create_mixed_ascii_frame_for_client(), video_frame_t::data, terminal_capabilities_t::desired_fps, fps_frame_ns(), fps_init(), g_should_exit, client_info::has_terminal_caps, client_info::height, video_frame_t::height, INVALID_SOCKET_VALUE, client_info::last_rendered_grid_sources, log_debug, log_debug_every, log_dev, log_dev_every, log_error, log_info, log_warn, MAX_CLIENT_ID_LEN, NS_PER_MS_INT, NS_PER_SEC_INT, client_info::outgoing_video_buffer, platform_sleep_ms(), platform_sleep_ns(), SAFE_FREE, SAFE_STRNCPY, client_info::shutting_down, video_frame_t::size, client_info::socket, START_TIMER, STOP_TIMER_AND_LOG, client_info::terminal_caps, time_get_ns(), time_pretty(), video_frame_begin_write(), video_frame_commit(), VIDEO_RENDER_FPS, client_info::video_render_thread_running, client_info::width, and video_frame_t::width.

Referenced by create_client_render_threads().

◆ create_client_render_threads()

int create_client_render_threads ( server_context_t *  server_ctx,
client_info_t *  client 
)

Create and initialize per-client rendering threads.

Sets up the complete per-client threading infrastructure including both video and audio rendering threads plus all necessary synchronization primitives. This function is called once per client during connection establishment.

INITIALIZATION SEQUENCE:

  1. MUTEX INITIALIZATION:
    • client_state_mutex: Protects client state variables
  2. THREAD CREATION:
    • Create video rendering thread (60fps)
    • Create audio rendering thread (172fps)
    • Set running flags with proper synchronization
  3. ERROR HANDLING:
    • Complete cleanup if any step fails
    • Prevent partially initialized client state
    • Ensure no resource leaks on failure

THREAD SYNCHRONIZATION SETUP:

PER-CLIENT MUTEXES: Each client gets dedicated mutexes for fine-grained locking:

  • Prevents contention between clients
  • Allows concurrent processing of multiple clients
  • Enables lock-free operations within client context

THREAD RUNNING FLAGS:

  • Protected by client_state_mutex
  • Allow graceful thread termination
  • Checked frequently by render threads

ERROR RECOVERY:

PARTIAL FAILURE HANDLING: If thread creation fails partway through:

  • Stop and join any already-created threads
  • Destroy any initialized mutexes
  • Return error code to caller
  • Leave client in clean state for retry or cleanup

RESOURCE LEAK PREVENTION:

  • All allocations have corresponding cleanup
  • Thread handles are properly managed
  • Mutex destruction is deterministic

INTEGRATION POINTS:

Parameters
clientTarget client for thread creation
Returns
0 on success, -1 on failure
Note
Function performs complete initialization or complete cleanup
Mutexes must be destroyed by stop_client_render_threads()
Thread handles must be joined before client destruction
Warning
Partial initialization is not supported - complete success or failure
See also
stop_client_render_threads() For cleanup implementation

Definition at line 1169 of file server/render.c.

1169 {
1170 log_info("★★★ create_client_render_threads() CALLED for client_id=%s", client ? client->client_id : "NULL");
1171
1172 if (!server_ctx || !client) {
1173 log_error("Cannot create render threads: NULL %s", !server_ctx ? "server_ctx" : "client");
1174 return -1;
1175 }
1176
1177#ifdef DEBUG_THREADS
1178 log_debug("Creating render threads for client %s", client->client_id);
1179#endif
1180
1181 // NOTE: Mutexes are already initialized in add_client() before any threads start
1182 // This prevents race conditions where receive thread tries to use uninitialized mutexes
1183
1184 // Initialize render thread control flags
1185 // IMPORTANT: Set to true BEFORE creating thread to avoid race condition
1186 // where thread starts and immediately exits because flag is false
1189
1190 // Create video rendering thread (stop_id=2, stop after receive thread)
1191 char thread_name[64];
1192 safe_snprintf(thread_name, sizeof(thread_name), "video_render_%s", client->client_id);
1193 asciichat_error_t video_result;
1194
1195 if (client->is_tcp_client) {
1196 video_result = tcp_server_spawn_thread(server_ctx->tcp_server, client->socket, client_video_render_thread, client,
1197 2, thread_name);
1198 } else {
1199 log_debug("THREAD_CREATE: WebRTC/WebSocket render video thread for client %s", client->client_id);
1200 video_result =
1202 }
1203
1204 if (video_result != ASCIICHAT_OK) {
1205 // Reset flag since thread creation failed
1207 // Mutexes will be destroyed by remove_client() which called us
1208 return -1;
1209 }
1210
1211 // Create audio rendering thread (stop_id=2, same priority as video)
1212 safe_snprintf(thread_name, sizeof(thread_name), "audio_render_%s", client->client_id);
1213 asciichat_error_t audio_result;
1214
1215 if (client->is_tcp_client) {
1216 audio_result = tcp_server_spawn_thread(server_ctx->tcp_server, client->socket, client_audio_render_thread, client,
1217 2, thread_name);
1218 } else {
1219 log_debug("THREAD_CREATE: WebRTC/WebSocket render audio thread for client %s", client->client_id);
1220 audio_result =
1222 }
1223 if (audio_result != ASCIICHAT_OK) {
1224 // Clean up video thread (atomic operation, no mutex needed)
1226 // Reset audio flag since thread creation failed
1228 // tcp_server_stop_client_threads() will be called by remove_client()
1229 // to clean up the video thread we just created
1230 // Mutexes will be destroyed by remove_client() which called us
1231 return -1;
1232 }
1233
1234#ifdef DEBUG_THREADS
1235 log_debug("Created render threads for client %s", client->client_id);
1236#endif
1237
1238 return 0;
1239}
void atomic_store_bool(atomic_t *a, bool value)
Atomically store a boolean value.
Definition atomic.c:177
asciichat_error_t
Error and exit codes - unified status values (0-255)
Definition error_codes.h:49
@ ASCIICHAT_OK
Definition error_codes.h:51
int safe_snprintf(char *buffer, size_t buffer_size, const char *format,...)
Safe formatted string printing to buffer.
Definition system.c:148
asciichat_error_t tcp_server_spawn_thread(tcp_server_t *server, socket_t client_socket, void *(*thread_func)(void *), void *thread_arg, int stop_id, const char *thread_name)
Spawn a worker thread for a client.
#define asciichat_thread_create(thread_ptr, attr, start_routine, arg)
void * client_video_render_thread(void *arg)
Interruptible sleep function with platform-specific optimizations.
void * client_audio_render_thread(void *arg)
Main audio rendering thread function for individual clients.
asciichat_thread_t audio_render_thread
asciichat_thread_t video_render_thread
tcp_server_t * tcp_server
TCP server managing connections.

References ASCIICHAT_OK, asciichat_thread_create, atomic_store_bool(), client_info::audio_render_thread, client_info::audio_render_thread_running, client_audio_render_thread(), client_info::client_id, client_video_render_thread(), client_info::is_tcp_client, log_debug, log_error, log_info, safe_snprintf(), client_info::socket, server_context_t::tcp_server, tcp_server_spawn_thread(), client_info::video_render_thread, and client_info::video_render_thread_running.

◆ stop_client_render_threads()

void stop_client_render_threads ( client_info_t *  client)

Stop and cleanup per-client rendering threads.

Performs graceful shutdown of both video and audio rendering threads for a specific client, including proper thread joining and resource cleanup. This function ensures deterministic cleanup without resource leaks.

SHUTDOWN SEQUENCE:

  1. SIGNAL SHUTDOWN:
    • Set thread running flags to false
    • Threads detect flag change and exit processing loops
    • Use mutex protection for atomic flag updates
  2. THREAD JOINING:
    • Wait for video render thread to complete
    • Wait for audio render thread to complete
    • Handle join failures appropriately
  3. RESOURCE CLEANUP:
    • Destroy per-client mutexes
    • Clear thread handles
    • Reset client thread state

GRACEFUL TERMINATION:

THREAD COORDINATION:

  • Threads monitor running flags in their main loops
  • Flags are checked frequently (every iteration)
  • Threads exit cleanly within one processing cycle
  • No forced thread termination (unsafe)

DETERMINISTIC CLEANUP:

  • asciichat_thread_join() waits for complete thread exit
  • All thread resources are properly released
  • No zombie threads or resource leaks

ERROR HANDLING:

THREAD JOIN FAILURES:

  • Logged as errors but don't prevent cleanup
  • Continue with resource cleanup regardless
  • Platform-specific error reporting

NULL CLIENT HANDLING:

  • Safe to call with NULL client pointer
  • Error logged and function returns safely
  • No undefined behavior or crashes

RESOURCE MANAGEMENT:

MUTEX CLEANUP:

  • Destroys all per-client mutexes
  • Prevents resource leaks on client disconnect
  • Platform-independent destruction

THREAD HANDLE MANAGEMENT:

  • Clears thread handles after joining
  • Prevents accidental reuse of stale handles
  • Memory zeroing for safety

INTEGRATION REQUIREMENTS:

  • Called by remove_client() in client.c
  • Must be called before freeing client structure
  • Coordinates with other cleanup functions
Parameters
clientTarget client for thread cleanup
Note
Function is safe to call multiple times (idempotent)
All resources are cleaned up regardless of errors
Thread joining may take up to one processing cycle
Warning
Must be called before client structure deallocation
See also
create_client_render_threads() For thread creation
asciichat_thread_join() For platform-specific joining

Definition at line 1322 of file server/render.c.

1322 {
1323 if (!client) {
1324 SET_ERRNO(ERROR_INVALID_PARAM, "Client is NULL");
1325 return;
1326 }
1327
1328 log_debug("Stopping render threads for client %s", client->client_id);
1329
1330 // Signal threads to stop (atomic operations, no mutex needed)
1333
1334 // Wait for threads to finish (deterministic cleanup)
1335 // During shutdown, don't wait forever for threads to join
1336 bool is_shutting_down = atomic_load_bool(&g_should_exit);
1337
1339 log_debug("Joining video render thread for client %s", client->client_id);
1340 int result;
1341 if (is_shutting_down) {
1342 // During shutdown, don't timeout - wait for thread to exit
1343 // Timeouts mask the real problem: threads that are still running
1344 log_debug("Shutdown mode: joining video render thread for client %s (no timeout)", client->client_id);
1345 result = asciichat_thread_join(&client->video_render_thread, NULL);
1346 if (result != 0) {
1347 log_warn("Video render thread for client %s failed to join during shutdown: %s", client->client_id,
1348 SAFE_STRERROR(result));
1349 }
1350 } else {
1351 log_debug("Calling asciichat_thread_join for video thread of client %s", client->client_id);
1352 result = asciichat_thread_join(&client->video_render_thread, NULL);
1353 log_debug("asciichat_thread_join returned %d for video thread of client %s", result, client->client_id);
1354 }
1355
1356 if (result == 0) {
1357#ifdef DEBUG_THREADS
1358 log_debug("Video render thread joined for client %s", client->client_id);
1359#endif
1360 } else if (result != -2) { // Don't log timeout errors again
1361 if (is_shutting_down) {
1362 log_warn("Failed to join video render thread for client %s during shutdown (continuing): %s", client->client_id,
1363 SAFE_STRERROR(result));
1364 } else {
1365 log_error("Failed to join video render thread for client %s: %s", client->client_id, SAFE_STRERROR(result));
1366 }
1367 }
1368 if (result == 0) {
1370 }
1371 }
1372
1374 int result;
1375 if (is_shutting_down) {
1376 // During shutdown, don't timeout - wait for thread to exit
1377 // Timeouts mask the real problem: threads that are still running
1378 log_debug("Shutdown mode: joining audio render thread for client %s (no timeout)", client->client_id);
1379 result = asciichat_thread_join(&client->audio_render_thread, NULL);
1380 if (result != 0) {
1381 log_warn("Audio render thread for client %s failed to join during shutdown: %s", client->client_id,
1382 SAFE_STRERROR(result));
1383 }
1384 } else {
1385 result = asciichat_thread_join(&client->audio_render_thread, NULL);
1386 }
1387
1388 if (result == 0) {
1389#ifdef DEBUG_THREADS
1390 log_debug("Audio render thread joined for client %s", client->client_id);
1391#endif
1392 } else if (result != -2) { // Don't log timeout errors again
1393 if (is_shutting_down) {
1394 log_warn("Failed to join audio render thread for client %s during shutdown (continuing): %s", client->client_id,
1395 SAFE_STRERROR(result));
1396 } else {
1397 log_error("Failed to join audio render thread for client %s: %s", client->client_id, SAFE_STRERROR(result));
1398 }
1399 }
1400 if (result == 0) {
1402 }
1403 }
1404
1405 // DO NOT destroy the mutex here - client.c will handle it
1406 // mutex_destroy(&client->client_state_mutex);
1407
1408#ifdef DEBUG_THREADS
1409 log_debug("Successfully destroyed render threads for client %s", client->client_id);
1410#endif
1411}
#define SAFE_STRERROR(errnum)
Definition common.h:465
#define SET_ERRNO(code, context_msg,...)
Set error code with custom context message and log it, returning the error code.
@ ERROR_INVALID_PARAM
void asciichat_thread_init(asciichat_thread_t *thread)
Initialize a thread handle to an uninitialized state.
bool asciichat_thread_is_initialized(asciichat_thread_t *thread)
Check if a thread handle has been initialized.
#define asciichat_thread_join(thread, timeout_ms)

References asciichat_thread_init(), asciichat_thread_is_initialized(), asciichat_thread_join, atomic_load_bool(), atomic_store_bool(), client_info::audio_render_thread, client_info::audio_render_thread_running, client_info::client_id, ERROR_INVALID_PARAM, g_should_exit, log_debug, log_error, log_warn, SAFE_STRERROR, SET_ERRNO, client_info::video_render_thread, and client_info::video_render_thread_running.

Referenced by remove_client().