Mercurial
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 |
| 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 } |