Mercurial
view mrjunejune/inference_bridge.c @ 264:04fee26ecce0
add authenticated JRPG conversation platform
Add reusable auth/session storage, owned conversation recovery, guest quotas, admin workflows, URL-routed conversation UI, mobile frame support, and parallel browser acceptance.
Co-authored-by: Copilot <[email protected]>
| author | MrJuneJune <me@mrjunejune.com> |
|---|---|
| date | Fri, 07 Aug 2026 07:34:12 -0700 |
| parents | b401627fc49e |
| children | 056790c4fb0d |
line wrap: on
line source
#include "mrjunejune/inference_bridge.h" #include <errno.h> #include <pthread.h> #include <signal.h> #include <stdatomic.h> #include <stdio.h> #include <stdlib.h> #include <string.h> #include <sys/wait.h> #include <time.h> #include <unistd.h> #define INFERENCE_EVENT_MAX (1024 * 1024 + 4096) #define INFERENCE_READY_TIMEOUT_MS 15000 struct Inference_Bridge { char *sidecar_path; char *copilot_cli_path; FILE *p_commands; FILE *p_events; pid_t child_pid; pthread_t reader_thread; pthread_mutex_t write_mutex; pthread_mutex_t ready_mutex; pthread_cond_t ready_condition; boolean write_mutex_initialized; boolean ready_mutex_initialized; boolean ready_condition_initialized; boolean reader_started; _Atomic boolean running; _Atomic boolean ready; Inference_Event_Handler handler; void *p_user_data; }; static int64 Inference_Monotonic_Milliseconds(void) { struct timespec now; clock_gettime(CLOCK_MONOTONIC, &now); return (int64)now.tv_sec * 1000 + now.tv_nsec / 1000000; } static const char *Inference_JSON_String( Dowa_JSON_Entry *object, const char *key) { char *value = Dowa_JSON_Get_String(object, key); return value ? value : ""; } static boolean Inference_JSON_Boolean( Dowa_JSON_Entry *object, const char *key) { Dowa_JSON_Value *p_value = Dowa_JSON_Get(object, key); return p_value && p_value->type == DOWA_JSON_BOOL ? p_value->bool_val : FALSE; } static int64 Inference_JSON_Integer( Dowa_JSON_Entry *object, const char *key) { Dowa_JSON_Value *p_value = Dowa_JSON_Get(object, key); return p_value && p_value->type == DOWA_JSON_NUMBER ? (int64)p_value->num_val : 0; } static void Inference_Bridge_Handle_Line( Inference_Bridge *p_bridge, const char *line, size_t length) { Dowa_Arena *p_arena = Dowa_Arena_Create(length * 3 + 4096); if (!p_arena) return; Dowa_JSON_Value parsed = Dowa_JSON_Parse(line, (int32)length, p_arena); if (parsed.type != DOWA_JSON_OBJECT) { Dowa_Arena_Free(p_arena); return; } Dowa_JSON_Entry *object = parsed.object_val; const char *type = Inference_JSON_String(object, "type"); if (strcmp(type, "ready") == 0) { atomic_store(&p_bridge->ready, TRUE); pthread_mutex_lock(&p_bridge->ready_mutex); pthread_cond_broadcast(&p_bridge->ready_condition); pthread_mutex_unlock(&p_bridge->ready_mutex); } Inference_Event event = { .type = type, .request_id = Inference_JSON_String(object, "request_id"), .conversation_id = Inference_JSON_String(object, "conversation_id"), .delta = Inference_JSON_String(object, "delta"), .content = Inference_JSON_String(object, "content"), .raw_json = line, .raw_json_length = length, .failed = Inference_JSON_Boolean(object, "failed"), .aborted = Inference_JSON_Boolean(object, "aborted"), }; Dowa_JSON_Value *p_error = Dowa_JSON_Get(object, "error"); if (p_error && p_error->type == DOWA_JSON_OBJECT) { Dowa_JSON_Entry *error = p_error->object_val; event.error_code = Inference_JSON_String(error, "code"); event.error_message = Inference_JSON_String(error, "message"); } else { event.error_code = ""; event.error_message = ""; } Dowa_JSON_Value *p_usage = Dowa_JSON_Get(object, "usage"); if (p_usage && p_usage->type == DOWA_JSON_OBJECT) { Dowa_JSON_Entry *usage = p_usage->object_val; event.input_tokens = Inference_JSON_Integer(usage, "input_tokens"); event.output_tokens = Inference_JSON_Integer(usage, "output_tokens"); } if (p_bridge->handler) p_bridge->handler(&event, p_bridge->p_user_data); Dowa_Arena_Free(p_arena); } static void *Inference_Bridge_Read_Events(void *p_context) { Inference_Bridge *p_bridge = p_context; char *line = NULL; size_t capacity = 0; while (atomic_load(&p_bridge->running)) { ssize_t amount = getline(&line, &capacity, p_bridge->p_events); if (amount < 0) break; if ((size_t)amount > INFERENCE_EVENT_MAX) continue; size_t length = (size_t)amount; while (length > 0 && (line[length - 1] == '\n' || line[length - 1] == '\r')) length--; if (length > 0) { line[length] = '\0'; Inference_Bridge_Handle_Line(p_bridge, line, length); } } free(line); atomic_store(&p_bridge->ready, FALSE); atomic_store(&p_bridge->running, FALSE); if (p_bridge->handler) { Inference_Event closed = { .type = "bridge.closed", .error_code = "sidecar_closed", .error_message = "Inference sidecar connection closed", .failed = TRUE, }; p_bridge->handler(&closed, p_bridge->p_user_data); } pthread_mutex_lock(&p_bridge->ready_mutex); pthread_cond_broadcast(&p_bridge->ready_condition); pthread_mutex_unlock(&p_bridge->ready_mutex); return NULL; } static boolean Inference_Bridge_Write( Inference_Bridge *p_bridge, const char *payload) { if (!p_bridge || !payload || !atomic_load(&p_bridge->running)) return FALSE; pthread_mutex_lock(&p_bridge->write_mutex); boolean success = fputs(payload, p_bridge->p_commands) >= 0 && fputc('\n', p_bridge->p_commands) != EOF && fflush(p_bridge->p_commands) == 0; pthread_mutex_unlock(&p_bridge->write_mutex); return success; } static boolean Inference_Bridge_Command( Inference_Bridge *p_bridge, const char *command, const char *request_id, const char *conversation_id, const char *prompt) { if (!command || !request_id) return FALSE; size_t input_length = strlen(command) + strlen(request_id) + strlen(conversation_id ? conversation_id : "") + strlen(prompt ? prompt : ""); if (input_length > (((size_t)-1) - 4096) / 12) return FALSE; Dowa_Arena *p_arena = Dowa_Arena_Create(input_length * 12 + 4096); if (!p_arena) return FALSE; char *escaped_command = Dowa_JSON_Escape_String(command, 0, p_arena); char *escaped_request = Dowa_JSON_Escape_String(request_id, 0, p_arena); char *escaped_conversation = Dowa_JSON_Escape_String( conversation_id ? conversation_id : "", 0, p_arena); char *escaped_prompt = prompt ? Dowa_JSON_Escape_String(prompt, 0, p_arena) : NULL; if (!escaped_command || !escaped_request || !escaped_conversation || (prompt && !escaped_prompt)) { Dowa_Arena_Free(p_arena); return FALSE; } size_t capacity = strlen(escaped_command) + strlen(escaped_request) + strlen(escaped_conversation) + (escaped_prompt ? strlen(escaped_prompt) : 0) + 160; char *payload = Dowa_Arena_Allocate(p_arena, capacity); if (!payload) { Dowa_Arena_Free(p_arena); return FALSE; } if (escaped_prompt) { snprintf( payload, capacity, "{\"command\":\"%s\",\"request_id\":\"%s\"," "\"conversation_id\":\"%s\",\"prompt\":\"%s\"}", escaped_command, escaped_request, escaped_conversation, escaped_prompt); } else { snprintf( payload, capacity, "{\"command\":\"%s\",\"request_id\":\"%s\"," "\"conversation_id\":\"%s\"}", escaped_command, escaped_request, escaped_conversation); } boolean success = Inference_Bridge_Write(p_bridge, payload); Dowa_Arena_Free(p_arena); return success; } Inference_Bridge *Inference_Bridge_Create( const char *sidecar_path, const char *copilot_cli_path, Inference_Event_Handler handler, void *p_user_data) { if (!sidecar_path || !copilot_cli_path) return NULL; Inference_Bridge *p_bridge = calloc(1, sizeof(*p_bridge)); if (!p_bridge) return NULL; p_bridge->sidecar_path = strdup(sidecar_path); p_bridge->copilot_cli_path = strdup(copilot_cli_path); p_bridge->handler = handler; p_bridge->p_user_data = p_user_data; p_bridge->child_pid = -1; if (!p_bridge->sidecar_path || !p_bridge->copilot_cli_path) { Inference_Bridge_Destroy(p_bridge); return NULL; } if (pthread_mutex_init(&p_bridge->write_mutex, NULL) != 0) { Inference_Bridge_Destroy(p_bridge); return NULL; } p_bridge->write_mutex_initialized = TRUE; if (pthread_mutex_init(&p_bridge->ready_mutex, NULL) != 0) { Inference_Bridge_Destroy(p_bridge); return NULL; } p_bridge->ready_mutex_initialized = TRUE; if (pthread_cond_init(&p_bridge->ready_condition, NULL) != 0) { Inference_Bridge_Destroy(p_bridge); return NULL; } p_bridge->ready_condition_initialized = TRUE; return p_bridge; } static void Inference_Bridge_Stop_Process(Inference_Bridge *p_bridge) { if (!p_bridge) return; if (atomic_load(&p_bridge->running) && p_bridge->p_commands) Inference_Bridge_Command( p_bridge, "shutdown", "server-shutdown", "", NULL); if (p_bridge->p_commands) { fclose(p_bridge->p_commands); p_bridge->p_commands = NULL; } if (p_bridge->child_pid > 0) { int status = 0; pid_t result = 0; for (int attempt = 0; attempt < 20; attempt++) { result = waitpid(p_bridge->child_pid, &status, WNOHANG); if (result != 0) break; usleep(50000); } if (result == 0) { kill(p_bridge->child_pid, SIGTERM); waitpid(p_bridge->child_pid, &status, 0); } p_bridge->child_pid = -1; } atomic_store(&p_bridge->running, FALSE); atomic_store(&p_bridge->ready, FALSE); if (p_bridge->reader_started) { pthread_join(p_bridge->reader_thread, NULL); p_bridge->reader_started = FALSE; } if (p_bridge->p_events) { fclose(p_bridge->p_events); p_bridge->p_events = NULL; } } boolean Inference_Bridge_Start(Inference_Bridge *p_bridge) { if (!p_bridge || atomic_load(&p_bridge->running)) return FALSE; int commands[2]; int events[2]; if (pipe(commands) != 0) return FALSE; if (pipe(events) != 0) { close(commands[0]); close(commands[1]); return FALSE; } pid_t child = fork(); if (child < 0) { close(commands[0]); close(commands[1]); close(events[0]); close(events[1]); return FALSE; } if (child == 0) { dup2(commands[0], STDIN_FILENO); dup2(events[1], STDOUT_FILENO); close(commands[0]); close(commands[1]); close(events[0]); close(events[1]); execl( p_bridge->sidecar_path, p_bridge->sidecar_path, p_bridge->copilot_cli_path, NULL); _exit(127); } close(commands[0]); close(events[1]); p_bridge->p_commands = fdopen(commands[1], "w"); p_bridge->p_events = fdopen(events[0], "r"); if (!p_bridge->p_commands || !p_bridge->p_events) { if (p_bridge->p_commands) fclose(p_bridge->p_commands); else close(commands[1]); if (p_bridge->p_events) fclose(p_bridge->p_events); else close(events[0]); kill(child, SIGTERM); waitpid(child, NULL, 0); return FALSE; } setvbuf(p_bridge->p_commands, NULL, _IOLBF, 0); p_bridge->child_pid = child; atomic_store(&p_bridge->running, TRUE); atomic_store(&p_bridge->ready, FALSE); if (pthread_create( &p_bridge->reader_thread, NULL, Inference_Bridge_Read_Events, p_bridge) != 0) { atomic_store(&p_bridge->running, FALSE); fclose(p_bridge->p_commands); fclose(p_bridge->p_events); kill(child, SIGTERM); waitpid(child, NULL, 0); p_bridge->p_commands = NULL; p_bridge->p_events = NULL; p_bridge->child_pid = -1; return FALSE; } p_bridge->reader_started = TRUE; int64 deadline = Inference_Monotonic_Milliseconds() + INFERENCE_READY_TIMEOUT_MS; pthread_mutex_lock(&p_bridge->ready_mutex); while (atomic_load(&p_bridge->running) && !atomic_load(&p_bridge->ready)) { int64 remaining = deadline - Inference_Monotonic_Milliseconds(); if (remaining <= 0) break; struct timespec timeout; clock_gettime(CLOCK_REALTIME, &timeout); timeout.tv_sec += remaining / 1000; timeout.tv_nsec += (remaining % 1000) * 1000000; if (timeout.tv_nsec >= 1000000000) { timeout.tv_sec++; timeout.tv_nsec -= 1000000000; } pthread_cond_timedwait( &p_bridge->ready_condition, &p_bridge->ready_mutex, &timeout); } boolean ready = atomic_load(&p_bridge->ready); pthread_mutex_unlock(&p_bridge->ready_mutex); if (!ready) { Inference_Bridge_Stop_Process(p_bridge); return FALSE; } return TRUE; } boolean Inference_Bridge_Is_Ready(const Inference_Bridge *p_bridge) { return p_bridge && atomic_load(&p_bridge->running) && atomic_load(&p_bridge->ready); } boolean Inference_Bridge_Start_Turn( Inference_Bridge *p_bridge, const char *request_id, const char *conversation_id, const char *prompt) { return Inference_Bridge_Command( p_bridge, "turn.start", request_id, conversation_id, prompt); } boolean Inference_Bridge_Abort_Turn( Inference_Bridge *p_bridge, const char *request_id, const char *conversation_id) { return Inference_Bridge_Command( p_bridge, "turn.abort", request_id, conversation_id, NULL); } boolean Inference_Bridge_Delete_Conversation( Inference_Bridge *p_bridge, const char *request_id, const char *conversation_id) { return Inference_Bridge_Command( p_bridge, "conversation.delete", request_id, conversation_id, NULL); } void Inference_Bridge_Destroy(Inference_Bridge *p_bridge) { if (!p_bridge) return; Inference_Bridge_Stop_Process(p_bridge); if (p_bridge->ready_condition_initialized) pthread_cond_destroy(&p_bridge->ready_condition); if (p_bridge->ready_mutex_initialized) pthread_mutex_destroy(&p_bridge->ready_mutex); if (p_bridge->write_mutex_initialized) pthread_mutex_destroy(&p_bridge->write_mutex); free(p_bridge->sidecar_path); free(p_bridge->copilot_cli_path); free(p_bridge); }