Preserve empty SSE data fields (#2594)

SSE events may contain an empty data field, and a data field without a colon also has an empty value. Track whether a data field was seen separately from the accumulated payload so empty events are dispatched and leading empty lines are preserved. Add an integration regression test for both forms.
This commit is contained in:
DosX
2026-10-02 21:17:47 -04:00
committed by GitHub
parent 8a3abfb597
commit 7255a7e979
2 changed files with 49 additions and 11 deletions
+10 -11
View File
@@ -4376,7 +4376,8 @@ public:
void stop();
private:
bool parse_sse_line(const std::string &line, SSEMessage &msg, int &retry_ms);
bool parse_sse_line(const std::string &line, SSEMessage &msg, int &retry_ms,
bool &has_data);
void run_event_loop();
void dispatch_event(const SSEMessage &msg);
bool should_reconnect(int count) const;
@@ -4886,7 +4887,7 @@ inline void SSEClient::stop() {
}
inline bool SSEClient::parse_sse_line(const std::string &line, SSEMessage &msg,
int &retry_ms) {
int &retry_ms, bool &has_data) {
// Blank line signals end of event
if (line.empty() || line == "\r") { return true; }
@@ -4895,16 +4896,11 @@ inline bool SSEClient::parse_sse_line(const std::string &line, SSEMessage &msg,
// Find the colon separator
auto colon_pos = line.find(':');
if (colon_pos == std::string::npos) {
// Line with no colon is treated as field name with empty value
return false;
}
auto field = line.substr(0, colon_pos);
std::string value;
// Value starts after colon, skip optional single space
if (colon_pos + 1 < line.size()) {
if (colon_pos != std::string::npos && colon_pos + 1 < line.size()) {
auto value_start = colon_pos + 1;
if (line[value_start] == ' ') { value_start++; }
value = line.substr(value_start);
@@ -4917,8 +4913,9 @@ inline bool SSEClient::parse_sse_line(const std::string &line, SSEMessage &msg,
msg.event = value;
} else if (field == "data") {
// Multiple data lines are concatenated with newlines
if (!msg.data.empty()) { msg.data += "\n"; }
if (has_data) { msg.data += "\n"; }
msg.data += value;
has_data = true;
} else if (field == "id") {
// Empty id is valid (clears the last event ID)
msg.id = value;
@@ -4992,6 +4989,7 @@ inline void SSEClient::run_event_loop() {
// Event receiving loop
std::string buffer;
SSEMessage current_msg;
bool has_data = false;
while (running_.load() && result.next()) {
buffer.append(result.data(), result.size());
@@ -5007,9 +5005,9 @@ inline void SSEClient::run_event_loop() {
// Parse the line and check if event is complete
auto event_complete =
parse_sse_line(line, current_msg, reconnect_interval_ms_);
parse_sse_line(line, current_msg, reconnect_interval_ms_, has_data);
if (event_complete && !current_msg.data.empty()) {
if (event_complete && has_data) {
// Update last_event_id for reconnection
if (!current_msg.id.empty()) { last_event_id_ = current_msg.id; }
@@ -5017,6 +5015,7 @@ inline void SSEClient::run_event_loop() {
dispatch_event(current_msg);
current_msg.clear();
has_data = false;
}
}
+39
View File
@@ -22678,6 +22678,45 @@ TEST_F(SSEIntegrationTest, MultiLineDataIntegration) {
EXPECT_EQ(received_data, "line1\nline2\nline3");
}
TEST_F(SSEIntegrationTest, EmptyDataLines) {
server_->Get("/empty-data-lines", [](const Request &, Response &res) {
res.set_chunked_content_provider("text/event-stream", [](size_t offset,
DataSink &sink) {
if (offset == 0) {
const std::string events = "data:\n\ndata:\ndata: hello\n\ndata\n\n";
sink.write(events.data(), events.size());
}
return false;
});
});
Client client("localhost", get_port());
sse::SSEClient sse(client, "/empty-data-lines");
std::mutex mutex;
std::condition_variable cv;
std::vector<std::string> received;
sse.on_message([&](const sse::SSEMessage &msg) {
std::lock_guard<std::mutex> lock(mutex);
received.push_back(msg.data);
cv.notify_all();
});
sse.set_max_reconnect_attempts(1);
sse.start_async();
{
std::unique_lock<std::mutex> lock(mutex);
cv.wait_for(lock, std::chrono::seconds(2),
[&] { return received.size() >= 3; });
}
sse.stop();
ASSERT_GE(received.size(), 3u);
EXPECT_EQ(received[0], "");
EXPECT_EQ(received[1], "\nhello");
EXPECT_EQ(received[2], "");
}
// Test: Auto-reconnect after server disconnection
TEST_F(SSEIntegrationTest, AutoReconnectAfterDisconnect) {
std::atomic<int> connection_count{0};