annotate seobeo/s_sse.c @ 257:609d3c6aff4e

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