ascii-chat 0.11.33
Video chat in your terminal
Loading...
Searching...
No Matches
server/render.c
Go to the documentation of this file.
1
151#include <stdio.h>
152#include <string.h>
153#include <time.h>
154#include <errno.h>
155#include <math.h>
156
157#include "render.h"
158#include "client.h"
159#include "main.h"
160#include "stream.h"
161#include "protocol.h"
162#include <ascii-chat/common.h>
165#include <ascii-chat/options/rcu.h> // For RCU-based options access
170#include <ascii-chat/util/time.h>
175#include <ascii-chat/util/fps.h>
176
177/* ============================================================================
178 * Cross-Platform Utility Functions
179 * ============================================================================
180 */
181
232// Removed interruptible_usleep - using regular platform_sleep_us instead
233// Sleep interruption isn't needed for small delays and isn't truly possible anyway
234
235/* ============================================================================
236 * Per-Client Video Rendering Implementation
237 * ============================================================================
238 */
239
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}
671
672/* ============================================================================
673 * Per-Client Audio Rendering Implementation
674 * ============================================================================
675 */
676
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}
1093
1094/* ============================================================================
1095 * Thread Lifecycle Management Functions
1096 * ============================================================================
1097 */
1098
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}
1240
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}
Platform abstraction layer umbrella header providing unified cross-platform API.
Application-level callbacks for library code.
#define APP_CALLBACK_VOID(callback_name)
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
⚙️ Common definitions, error codes, macros, and types shared throughout the application
Cross-platform error handling utilities.
⏱️ FPS tracking utility for monitoring frame throughput across all threads
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
#define SAFE_STRERROR(errnum)
Definition common.h:465
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
#define SET_ERRNO(code, context_msg,...)
Set error code with custom context message and log it, returning the error code.
void asciichat_errno_destroy(void)
Cleanup error system resources.
asciichat_error_t
Error and exit codes - unified status values (0-255)
Definition error_codes.h:49
@ ASCIICHAT_OK
Definition error_codes.h:51
@ ERROR_INVALID_PARAM
#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_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_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 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
#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
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
#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
int safe_snprintf(char *buffer, size_t buffer_size, const char *format,...)
Safe formatted string printing to buffer.
Definition system.c:148
void asciichat_thread_init(asciichat_thread_t *thread)
Initialize a thread handle to an uninitialized state.
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
bool asciichat_thread_is_initialized(asciichat_thread_t *thread)
Check if a thread handle has been initialized.
void platform_sleep_ms(unsigned int ms)
Sleep for a specified number of milliseconds.
#define mutex_unlock(mutex)
Unlock a mutex (with debug tracking in debug builds)
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
🔊 Audio Capture and Playback Interface for ascii-chat
ACIP server-side protocol API.
⚙️ Unified options parsing system for ascii-chat with builder pattern and lock-free access
Platform initialization and static synchronization helpers.
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)
#define asciichat_thread_join(thread, timeout_ms)
#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
🔢 Mathematical Utility Functions
Multi-Source Audio Mixing and Processing System.
Opus audio codec wrapper for real-time encoding/decoding.
mixer_t *volatile g_audio_mixer
Global audio mixer instance for multi-client audio processing.
ascii-chat Server Mode Entry Point Header
#define OPUS_FRAME_SAMPLES
void stop_client_render_threads(client_info_t *client)
Stop and cleanup per-client rendering threads.
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.
Per-client rendering threads with rate limiting.
#define VIDEO_RENDER_FPS
#define AUDIO_PACKET_FPS
Client-side H.265 media capture and encoding pipeline.
H.265 codec transport protocol definitions.
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
Multi-client video mixing and ASCII frame generation.
Per-client state structure for server-side client management.
terminal_capabilities_t terminal_caps
char display_name[MAX_DISPLAY_NAME_LEN]
asciichat_thread_t audio_render_thread
video_frame_buffer_t * outgoing_video_buffer
asciichat_thread_t video_render_thread
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
Server context - encapsulates all server state.
tcp_server_t * tcp_server
TCP server managing connections.
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)
⏱️ High-precision timing utilities using sokol_time.h and uthash
📊 String Formatting Utilities