annotate seobeo/s_sse.c @ 279:b3b547563ec7

Add Google connector service and agent wiki Implement the C/Seobeo Google Drive and Gmail connector with encrypted OAuth storage, Zenbu authentication, browser testing, AI tool discovery, chunked HTTP decoding, and Bazel coverage. Consolidate repository guidance into progressive wiki documentation and enforce arena-first allocation for new first-party C code. Co-authored-by: Copilot <[email protected]> Copilot-Session: 84c338fd-0939-4bb3-b7f3-1062eb213e5d
author MrJuneJune <me@mrjunejune.com>
date Mon, 17 Aug 2026 22:22:36 -0700
parents 04fee26ecce0
children
Ignore whitespace changes - Everywhere: Within whitespace: At end of lines:
rev   line source
257
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
1 #include "seobeo/seobeo.h"
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
2
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
3 #include <stdio.h>
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
4 #include <stdlib.h>
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
5 #include <string.h>
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
6
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
7 static const char *SEOBEO_SSE_HEADERS =
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
8 "HTTP/1.1 200 OK\r\n"
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
9 "Content-Type: text/event-stream\r\n"
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
10 "Cache-Control: no-cache\r\n"
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
11 "Connection: keep-alive\r\n"
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
12 "X-Accel-Buffering: no\r\n"
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
13 "\r\n";
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
14
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
15 static pthread_mutex_t g_sse_mutex = PTHREAD_MUTEX_INITIALIZER;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
16 static Seobeo_SSE_Stream **g_sse_streams = NULL;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
17
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
18 typedef struct {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
19 char *data;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
20 uint32 size;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
21 } Seobeo_SSE_Pending_Frame;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
22
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
23 #define SEOBEO_SSE_MAX_PENDING 32
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
24
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
25 static boolean sse_is_open_unlocked(const Seobeo_SSE_Stream *p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
26 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
27 return p_stream &&
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
28 p_stream->started &&
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
29 !p_stream->closed &&
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
30 p_stream->p_handle &&
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
31 p_stream->p_handle->socket >= 0 &&
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
32 !atomic_load(&p_stream->p_handle->destroyed);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
33 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
34
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
35 static void sse_clear_pending_unlocked(Seobeo_SSE_Stream *p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
36 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
37 Seobeo_SSE_Pending_Frame *p_frames = p_stream->p_pending_frames;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
38 size_t count = Dowa_Array_Length(p_frames);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
39 for (size_t i = p_stream->pending_offset; i < count; i++)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
40 Dowa_Free(p_frames[i].data);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
41 Dowa_Array_Free(p_frames);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
42 p_stream->p_pending_frames = NULL;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
43 p_stream->pending_offset = 0;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
44 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
45
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
46 static void sse_release_unlocked(Seobeo_SSE_Stream *p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
47 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
48 if (!p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
49 return;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
50 if (atomic_fetch_sub(&p_stream->ref_count, 1) != 1)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
51 return;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
52 sse_clear_pending_unlocked(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
53 Dowa_Free(p_stream->path);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
54 Dowa_Free(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
55 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
56
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
57 static int32 sse_drain_pending_unlocked(Seobeo_SSE_Stream *p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
58 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
59 if (!sse_is_open_unlocked(p_stream))
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
60 return -1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
61 int32 flush_result = Seobeo_Handle_Flush(p_stream->p_handle);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
62 if (flush_result != 0)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
63 return flush_result;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
64
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
65 Seobeo_SSE_Pending_Frame *p_frames = p_stream->p_pending_frames;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
66 size_t count = Dowa_Array_Length(p_frames);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
67 while (p_stream->pending_offset < count)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
68 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
69 Seobeo_SSE_Pending_Frame *p_frame =
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
70 &p_frames[p_stream->pending_offset];
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
71 int32 queue_result = Seobeo_Handle_Queue(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
72 p_stream->p_handle,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
73 (const uint8 *)p_frame->data,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
74 p_frame->size);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
75 if (queue_result != 0)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
76 return queue_result;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
77 Dowa_Free(p_frame->data);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
78 p_stream->pending_offset++;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
79 flush_result = Seobeo_Handle_Flush(p_stream->p_handle);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
80 if (flush_result != 0)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
81 return flush_result;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
82 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
83
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
84 Dowa_Array_Free(p_frames);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
85 p_stream->p_pending_frames = NULL;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
86 p_stream->pending_offset = 0;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
87 return 0;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
88 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
89
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
90 static int32 sse_enqueue_frame_unlocked(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
91 Seobeo_SSE_Stream *p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
92 char *frame,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
93 uint32 frame_size)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
94 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
95 Seobeo_SSE_Pending_Frame *p_frames = p_stream->p_pending_frames;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
96 size_t pending = Dowa_Array_Length(p_frames) - p_stream->pending_offset;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
97 if (pending >= SEOBEO_SSE_MAX_PENDING)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
98 return SEOBEO_SSE_BACKPRESSURE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
99 Seobeo_SSE_Pending_Frame pending_frame = {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
100 .data = frame,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
101 .size = frame_size,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
102 };
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
103 Dowa_Array_Push(p_frames, pending_frame);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
104 p_stream->p_pending_frames = p_frames;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
105 return sse_drain_pending_unlocked(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
106 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
107
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
108 static boolean sse_field_is_safe(const char *value)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
109 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
110 return !value || (strchr(value, '\r') == NULL && strchr(value, '\n') == NULL);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
111 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
112
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
113 static size_t sse_multiline_size(const char *prefix, const char *value)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
114 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
115 const char *line = value ? value : "";
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
116 size_t prefix_length = strlen(prefix);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
117 size_t total = 0;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
118
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
119 while (TRUE)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
120 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
121 const char *end = line;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
122 while (*end && *end != '\r' && *end != '\n')
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
123 end++;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
124 size_t line_length = (size_t)(end - line);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
125 total += prefix_length + (line_length ? 1 : 0) + line_length + 1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
126 if (*end == '\0')
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
127 break;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
128 if (*end == '\r' && end[1] == '\n')
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
129 end++;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
130 line = end + 1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
131 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
132 return total;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
133 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
134
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
135 static char *sse_append_multiline(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
136 char *output,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
137 const char *prefix,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
138 const char *value)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
139 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
140 const char *line = value ? value : "";
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
141 size_t prefix_length = strlen(prefix);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
142
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
143 while (TRUE)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
144 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
145 const char *end = line;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
146 while (*end && *end != '\r' && *end != '\n')
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
147 end++;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
148 size_t line_length = (size_t)(end - line);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
149 memcpy(output, prefix, prefix_length);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
150 output += prefix_length;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
151 if (line_length)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
152 *output++ = ' ';
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
153 memcpy(output, line, line_length);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
154 output += line_length;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
155 *output++ = '\n';
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
156 if (*end == '\0')
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
157 break;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
158 if (*end == '\r' && end[1] == '\n')
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
159 end++;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
160 line = end + 1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
161 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
162 return output;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
163 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
164
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
165 static int32 sse_send_frame(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
166 Seobeo_SSE_Stream *p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
167 const char *frame,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
168 size_t frame_size)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
169 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
170 if (!sse_is_open_unlocked(p_stream) || !frame)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
171 return -1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
172 if (frame_size > p_stream->p_handle->write_buffer_capacity)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
173 return SEOBEO_SSE_EVENT_TOO_LARGE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
174
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
175 int32 flush_result = Seobeo_Handle_Flush(p_stream->p_handle);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
176 if (flush_result != 0)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
177 return flush_result;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
178
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
179 int32 queue_result = Seobeo_Handle_Queue(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
180 p_stream->p_handle,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
181 (const uint8 *)frame,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
182 (uint32)frame_size);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
183 if (queue_result != 0)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
184 return queue_result;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
185 return Seobeo_Handle_Flush(p_stream->p_handle);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
186 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
187
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
188 int32 Seobeo_SSE_Start(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
189 Seobeo_SSE_Stream *p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
190 Seobeo_Handle *p_handle)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
191 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
192 if (!p_stream || !p_handle || atomic_load(&p_handle->destroyed))
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
193 return -1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
194 p_stream->p_handle = p_handle;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
195 p_stream->started = TRUE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
196 p_stream->closed = FALSE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
197 return sse_send_frame(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
198 p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
199 SEOBEO_SSE_HEADERS,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
200 strlen(SEOBEO_SSE_HEADERS));
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
201 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
202
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
203 static int32 sse_send_event_unlocked(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
204 Seobeo_SSE_Stream *p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
205 const Seobeo_SSE_Event *p_event)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
206 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
207 if (!sse_is_open_unlocked(p_stream))
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
208 return -1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
209 if (!p_event || !sse_field_is_safe(p_event->event) ||
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
210 !sse_field_is_safe(p_event->id) ||
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
211 p_event->retry_ms < SEOBEO_SSE_RETRY_NONE)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
212 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
213 return SEOBEO_SSE_INVALID_EVENT;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
214 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
215
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
216 size_t frame_size = 1 + sse_multiline_size("data:", p_event->data);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
217 if (p_event->id)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
218 frame_size += strlen("id: ") + strlen(p_event->id) + 1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
219 if (p_event->event)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
220 frame_size += strlen("event: ") + strlen(p_event->event) + 1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
221
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
222 char retry[32] = {0};
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
223 if (p_event->retry_ms != SEOBEO_SSE_RETRY_NONE)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
224 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
225 int written = snprintf(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
226 retry,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
227 sizeof(retry),
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
228 "retry: %d\n",
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
229 p_event->retry_ms);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
230 if (written < 0 || (size_t)written >= sizeof(retry))
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
231 return SEOBEO_SSE_INVALID_EVENT;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
232 frame_size += (size_t)written;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
233 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
234
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
235 if (frame_size > p_stream->p_handle->write_buffer_capacity)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
236 return SEOBEO_SSE_EVENT_TOO_LARGE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
237 char *frame = malloc(frame_size);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
238 if (!frame)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
239 return -1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
240
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
241 char *output = frame;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
242 if (p_event->id)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
243 output += sprintf(output, "id: %s\n", p_event->id);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
244 if (p_event->event)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
245 output += sprintf(output, "event: %s\n", p_event->event);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
246 if (retry[0])
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
247 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
248 size_t retry_length = strlen(retry);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
249 memcpy(output, retry, retry_length);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
250 output += retry_length;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
251 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
252 output = sse_append_multiline(output, "data:", p_event->data);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
253 *output++ = '\n';
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
254
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
255 int32 result = sse_enqueue_frame_unlocked(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
256 p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
257 frame,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
258 (uint32)(output - frame));
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
259 if (result == SEOBEO_SSE_BACKPRESSURE)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
260 Dowa_Free(frame);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
261 return result;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
262 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
263
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
264 int32 Seobeo_SSE_Send(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
265 Seobeo_SSE_Stream *p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
266 const Seobeo_SSE_Event *p_event)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
267 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
268 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
269 int32 result = sse_send_event_unlocked(p_stream, p_event);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
270 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
271 return result;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
272 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
273
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
274 int32 Seobeo_SSE_Send_Data(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
275 Seobeo_SSE_Stream *p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
276 const char *data)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
277 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
278 Seobeo_SSE_Event event = {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
279 .data = data,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
280 .retry_ms = SEOBEO_SSE_RETRY_NONE,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
281 };
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
282 return Seobeo_SSE_Send(p_stream, &event);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
283 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
284
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
285 int32 Seobeo_SSE_Send_Comment(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
286 Seobeo_SSE_Stream *p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
287 const char *comment)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
288 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
289 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
290 if (!sse_is_open_unlocked(p_stream))
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
291 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
292 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
293 return -1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
294 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
295 size_t frame_size = sse_multiline_size(":", comment) + 1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
296 if (frame_size > p_stream->p_handle->write_buffer_capacity)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
297 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
298 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
299 return SEOBEO_SSE_EVENT_TOO_LARGE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
300 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
301 char *frame = malloc(frame_size);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
302 if (!frame)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
303 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
304 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
305 return -1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
306 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
307 char *output = sse_append_multiline(frame, ":", comment);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
308 *output++ = '\n';
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
309 int32 result = sse_enqueue_frame_unlocked(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
310 p_stream,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
311 frame,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
312 (uint32)(output - frame));
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
313 if (result == SEOBEO_SSE_BACKPRESSURE)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
314 Dowa_Free(frame);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
315 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
316 return result;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
317 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
318
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
319 int32 Seobeo_SSE_Flush(Seobeo_SSE_Stream *p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
320 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
321 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
322 int32 result = sse_is_open_unlocked(p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
323 ? sse_drain_pending_unlocked(p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
324 : -1;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
325 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
326 return result;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
327 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
328
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
329 boolean Seobeo_SSE_Is_Open(const Seobeo_SSE_Stream *p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
330 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
331 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
332 boolean open = sse_is_open_unlocked(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
333 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
334 return open;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
335 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
336
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
337 boolean Seobeo_SSE_Retain(Seobeo_SSE_Stream *p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
338 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
339 if (!p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
340 return FALSE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
341 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
342 unsigned int references = atomic_load(&p_stream->ref_count);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
343 if (references == 0)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
344 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
345 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
346 return FALSE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
347 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
348 atomic_fetch_add(&p_stream->ref_count, 1);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
349 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
350 return TRUE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
351 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
352
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
353 void Seobeo_SSE_Release(Seobeo_SSE_Stream *p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
354 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
355 if (!p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
356 return;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
357 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
358 sse_release_unlocked(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
359 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
360 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
361
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
362 void Seobeo_SSE_Close(Seobeo_SSE_Stream *p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
363 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
364 if (!p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
365 return;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
366 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
367 if (p_stream->started && !p_stream->closed && p_stream->p_handle)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
368 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
369 sse_drain_pending_unlocked(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
370 if (p_stream->managed)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
371 shutdown(p_stream->p_handle->socket, SHUT_RDWR);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
372 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
373 sse_clear_pending_unlocked(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
374 p_stream->closed = TRUE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
375 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
376 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
377
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
378 Seobeo_SSE_Stream *Seobeo_SSE_Server_Attach(
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
379 Seobeo_Handle *p_handle,
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
380 const char *path)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
381 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
382 if (!p_handle || !path)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
383 return NULL;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
384 Seobeo_SSE_Stream *p_stream = malloc(sizeof(*p_stream));
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
385 if (!p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
386 return NULL;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
387 memset(p_stream, 0, sizeof(*p_stream));
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
388 p_stream->path = strdup(path);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
389 if (!p_stream->path)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
390 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
391 Dowa_Free(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
392 return NULL;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
393 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
394 atomic_init(&p_stream->ref_count, 1);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
395 if (Seobeo_SSE_Start(p_stream, p_handle) < 0)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
396 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
397 Dowa_Free(p_stream->path);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
398 Dowa_Free(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
399 return NULL;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
400 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
401 p_stream->managed = TRUE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
402 p_handle->is_sse = TRUE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
403
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
404 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
405 Dowa_Array_Push(g_sse_streams, p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
406 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
407 return p_stream;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
408 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
409
264
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
410 void Seobeo_SSE_Set_Detach_Callback(
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
411 Seobeo_SSE_Stream *p_stream,
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
412 void (*cb)(Seobeo_SSE_Stream *, void *),
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
413 void *ctx)
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
414 {
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
415 if (!p_stream)
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
416 return;
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
417 pthread_mutex_lock(&g_sse_mutex);
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
418 p_stream->detach_cb = cb;
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
419 p_stream->detach_ctx = ctx;
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
420 pthread_mutex_unlock(&g_sse_mutex);
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
421 }
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
422
257
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
423 void Seobeo_SSE_Server_Detach_Handle(Seobeo_Handle *p_handle)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
424 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
425 if (!p_handle)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
426 return;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
427 p_handle->is_sse = FALSE;
264
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
428
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
429 /* Save callback info before releasing the stream reference. */
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
430 void (*detach_cb)(Seobeo_SSE_Stream *, void *) = NULL;
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
431 void *detach_ctx = NULL;
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
432 Seobeo_SSE_Stream *cb_stream = NULL;
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
433
257
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
434 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
435 size_t count = Dowa_Array_Length(g_sse_streams);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
436 for (size_t i = 0; i < count; i++)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
437 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
438 Seobeo_SSE_Stream *p_stream = g_sse_streams[i];
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
439 if (p_stream->p_handle != p_handle)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
440 continue;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
441 p_stream->closed = TRUE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
442 p_stream->p_handle = NULL;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
443 sse_clear_pending_unlocked(p_stream);
264
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
444 detach_cb = p_stream->detach_cb;
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
445 detach_ctx = p_stream->detach_ctx;
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
446 /* Keep a live reference for the callback; release afterward. */
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
447 cb_stream = p_stream;
257
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
448 g_sse_streams[i] = Dowa_Array_Pop(g_sse_streams);
264
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
449 /* Do NOT call sse_release_unlocked here; do it after mutex released
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
450 * so the callback can safely acquire its own locks. */
257
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
451 break;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
452 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
453 pthread_mutex_unlock(&g_sse_mutex);
264
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
454
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
455 /* Fire callback outside the mutex so callers may acquire other locks. */
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
456 if (detach_cb && cb_stream)
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
457 detach_cb(cb_stream, detach_ctx);
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
458
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
459 /* Drop the server's reference (the callback may hold its own reference). */
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
460 if (cb_stream)
04fee26ecce0 add authenticated JRPG conversation platform
MrJuneJune <me@mrjunejune.com>
parents: 257
diff changeset
461 Seobeo_SSE_Release(cb_stream);
257
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
462 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
463
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
464 void Seobeo_SSE_Server_Destroy(void)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
465 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
466 pthread_mutex_lock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
467 size_t count = Dowa_Array_Length(g_sse_streams);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
468 for (size_t i = 0; i < count; i++)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
469 {
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
470 Seobeo_SSE_Stream *p_stream = g_sse_streams[i];
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
471 if (!p_stream)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
472 continue;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
473 p_stream->closed = TRUE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
474 if (p_stream->p_handle)
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
475 p_stream->p_handle->is_sse = FALSE;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
476 p_stream->p_handle = NULL;
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
477 sse_clear_pending_unlocked(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
478 sse_release_unlocked(p_stream);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
479 }
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
480 Dowa_Array_Free(g_sse_streams);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
481 pthread_mutex_unlock(&g_sse_mutex);
609d3c6aff4e [seobeo] Add persistent SSE streams
MrJuneJune <me@mrjunejune.com>
parents:
diff changeset
482 }