Mercurial
annotate seobeo/s_sse.c @ 264:04fee26ecce0
add authenticated JRPG conversation platform
Add reusable auth/session storage, owned conversation recovery, guest quotas, admin workflows, URL-routed conversation UI, mobile frame support, and parallel browser acceptance.
Co-authored-by: Copilot <[email protected]>
| author | MrJuneJune <me@mrjunejune.com> |
|---|---|
| date | Fri, 07 Aug 2026 07:34:12 -0700 |
| parents | 609d3c6aff4e |
| 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 } |