Fix SSE parsing of CRLF field names and data-less events

A field line without a colon kept the \r of a CRLF line ending in its
name, so "data\r\n" was not recognized. Strip the \r once per line
before parsing instead of from each value.

An event without a data field was neither dispatched nor cleared, so
its event type leaked into the next event and its id only reached
last_event_id after a later event was dispatched. Reset the message on
every blank line and record the id even when nothing is dispatched.

The parsing tests exercised a copy of parse_sse_line that had drifted
from the real one. Run them through SSEClient against a local server
instead, and add regression tests for the cases above.
This commit is contained in:
yhirose
2026-10-02 21:38:39 -04:00
parent dd71728110
commit 3d40dfc727
2 changed files with 195 additions and 233 deletions
+12 -8
View File
@@ -4889,10 +4889,10 @@ inline void SSEClient::stop() {
inline bool SSEClient::parse_sse_line(const std::string &line, SSEMessage &msg,
int &retry_ms, bool &has_data) {
// Blank line signals end of event
if (line.empty() || line == "\r") { return true; }
if (line.empty()) { return true; }
// Lines starting with ':' are comments (ignored)
if (!line.empty() && line[0] == ':') { return false; }
if (line[0] == ':') { return false; }
// Find the colon separator
auto colon_pos = line.find(':');
@@ -4904,8 +4904,6 @@ inline bool SSEClient::parse_sse_line(const std::string &line, SSEMessage &msg,
auto value_start = colon_pos + 1;
if (line[value_start] == ' ') { value_start++; }
value = line.substr(value_start);
// Remove trailing \r if present
if (!value.empty() && value.back() == '\r') { value.pop_back(); }
}
// Handle known fields
@@ -5003,17 +5001,23 @@ inline void SSEClient::run_event_loop() {
auto line = buffer.substr(line_start, newline_pos - line_start);
line_start = newline_pos + 1;
// Strip the \r of a CRLF line ending so that every field, including
// one without a colon, sees the same line
if (!line.empty() && line.back() == '\r') { line.pop_back(); }
// Parse the line and check if event is complete
auto event_complete =
parse_sse_line(line, current_msg, reconnect_interval_ms_, has_data);
if (event_complete && has_data) {
// Update last_event_id for reconnection
if (event_complete) {
// Update last_event_id for reconnection, even for an event that
// has no data
if (!current_msg.id.empty()) { last_event_id_ = current_msg.id; }
// Dispatch event to appropriate handler
dispatch_event(current_msg);
// An event without a data field is not dispatched
if (has_data) { dispatch_event(current_msg); }
// Reset the message for the next event either way
current_msg.clear();
has_data = false;
}
+183 -225
View File
@@ -22091,231 +22091,6 @@ TEST(ETagTest, NegativeFileModificationTime) {
std::remove(fname);
}
//==============================================================================
// SSE Parsing Tests
//==============================================================================
class SSEParsingTest : public ::testing::Test {
protected:
// Test helper that mimics SSE parsing behavior
static bool parse_sse_line(const std::string &line, sse::SSEMessage &msg,
int &retry_ms) {
// Blank line signals end of event
if (line.empty() || line == "\r") { return true; }
// Lines starting with ':' are comments (ignored)
if (!line.empty() && line[0] == ':') { return false; }
// 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;
}
std::string field = line.substr(0, colon_pos);
std::string value;
// Value starts after colon, skip optional single space
if (colon_pos + 1 < line.size()) {
size_t value_start = colon_pos + 1;
if (line[value_start] == ' ') { value_start++; }
value = line.substr(value_start);
// Remove trailing \r if present
if (!value.empty() && value.back() == '\r') { value.pop_back(); }
}
// Handle known fields
if (field == "event") {
msg.event = value;
} else if (field == "data") {
// Multiple data lines are concatenated with newlines
if (!msg.data.empty()) { msg.data += "\n"; }
msg.data += value;
} else if (field == "id") {
// Empty id is valid (clears the last event ID)
msg.id = value;
} else if (field == "retry") {
// Parse retry interval in milliseconds
{
int v = 0;
auto res =
detail::from_chars(value.data(), value.data() + value.size(), v);
if (res.ec == std::errc{}) { retry_ms = v; }
}
}
// Unknown fields are ignored per SSE spec
return false;
}
};
// Test: Single-line data
TEST_F(SSEParsingTest, SingleLineData) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("data: hello", msg, retry_ms));
EXPECT_EQ(msg.data, "hello");
EXPECT_EQ(msg.event, "message");
// Blank line ends event
EXPECT_TRUE(parse_sse_line("", msg, retry_ms));
}
// Test: Multi-line data
TEST_F(SSEParsingTest, MultiLineData) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("data: line1", msg, retry_ms));
EXPECT_FALSE(parse_sse_line("data: line2", msg, retry_ms));
EXPECT_FALSE(parse_sse_line("data: line3", msg, retry_ms));
EXPECT_EQ(msg.data, "line1\nline2\nline3");
}
// Test: Custom event types
TEST_F(SSEParsingTest, CustomEventType) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("event: update", msg, retry_ms));
EXPECT_FALSE(parse_sse_line("data: payload", msg, retry_ms));
EXPECT_EQ(msg.event, "update");
EXPECT_EQ(msg.data, "payload");
}
// Test: Event ID handling
TEST_F(SSEParsingTest, EventIdHandling) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("id: 12345", msg, retry_ms));
EXPECT_FALSE(parse_sse_line("data: test", msg, retry_ms));
EXPECT_EQ(msg.id, "12345");
}
// Test: Empty event ID (clears last event ID)
TEST_F(SSEParsingTest, EmptyEventId) {
sse::SSEMessage msg;
msg.id = "previous";
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("id:", msg, retry_ms));
EXPECT_EQ(msg.id, "");
}
// Test: Retry field parsing
TEST_F(SSEParsingTest, RetryFieldParsing) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("retry: 5000", msg, retry_ms));
EXPECT_EQ(retry_ms, 5000);
}
// Test: Invalid retry value
TEST_F(SSEParsingTest, InvalidRetryValue) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("retry: invalid", msg, retry_ms));
EXPECT_EQ(retry_ms, 3000); // Unchanged
}
// Test: Comments (lines starting with :)
TEST_F(SSEParsingTest, CommentsIgnored) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line(": this is a comment", msg, retry_ms));
EXPECT_EQ(msg.data, "");
EXPECT_EQ(msg.event, "message");
}
// Test: Colon in value
TEST_F(SSEParsingTest, ColonInValue) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("data: hello:world:test", msg, retry_ms));
EXPECT_EQ(msg.data, "hello:world:test");
}
// Test: Line with no colon (field name only)
TEST_F(SSEParsingTest, FieldNameOnly) {
sse::SSEMessage msg;
int retry_ms = 3000;
// According to SSE spec, this is treated as field name with empty value
EXPECT_FALSE(parse_sse_line("data", msg, retry_ms));
// Since we don't recognize "data" without colon, data should be empty
EXPECT_EQ(msg.data, "");
}
// Test: Trailing \r handling
TEST_F(SSEParsingTest, TrailingCarriageReturn) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("data: hello\r", msg, retry_ms));
EXPECT_EQ(msg.data, "hello");
}
// Test: Unknown fields ignored
TEST_F(SSEParsingTest, UnknownFieldsIgnored) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("unknown: value", msg, retry_ms));
EXPECT_EQ(msg.data, "");
EXPECT_EQ(msg.event, "message");
}
// Test: Space after colon is optional
TEST_F(SSEParsingTest, SpaceAfterColonOptional) {
sse::SSEMessage msg1, msg2;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("data: hello", msg1, retry_ms));
EXPECT_FALSE(parse_sse_line("data:hello", msg2, retry_ms));
EXPECT_EQ(msg1.data, "hello");
EXPECT_EQ(msg2.data, "hello");
}
// Test: SSEMessage clear
TEST_F(SSEParsingTest, MessageClear) {
sse::SSEMessage msg;
msg.event = "custom";
msg.data = "some data";
msg.id = "123";
msg.clear();
EXPECT_EQ(msg.event, "message");
EXPECT_EQ(msg.data, "");
EXPECT_EQ(msg.id, "");
}
// Test: Complete event parsing
TEST_F(SSEParsingTest, CompleteEventParsing) {
sse::SSEMessage msg;
int retry_ms = 3000;
EXPECT_FALSE(parse_sse_line("event: notification", msg, retry_ms));
EXPECT_FALSE(parse_sse_line("id: evt-42", msg, retry_ms));
EXPECT_FALSE(parse_sse_line("data: {\"type\":\"alert\"}", msg, retry_ms));
EXPECT_FALSE(parse_sse_line("retry: 1000", msg, retry_ms));
// Blank line ends event
EXPECT_TRUE(parse_sse_line("", msg, retry_ms));
EXPECT_EQ(msg.event, "notification");
EXPECT_EQ(msg.id, "evt-42");
EXPECT_EQ(msg.data, "{\"type\":\"alert\"}");
EXPECT_EQ(retry_ms, 1000);
}
//==============================================================================
// Integration Tests with Server
//==============================================================================
@@ -22410,6 +22185,189 @@ protected:
int port_ = 0;
};
//==============================================================================
// SSE Parsing Tests
//==============================================================================
class SSEParsingTest : public SSEIntegrationTest {
protected:
// Feeds a raw event stream to SSEClient and returns the dispatched messages.
// A trailing sentinel event tells when the whole stream has been parsed.
std::vector<sse::SSEMessage> parse(const std::string &stream) {
auto body = stream + "event: end\ndata: end\n\n";
server_->Get("/parse", [body](const Request &, Response &res) {
res.set_content(body, "text/event-stream");
});
return collect("/parse");
}
// Runs SSEClient against path until an "end" event arrives and returns the
// messages dispatched before it
std::vector<sse::SSEMessage> collect(const std::string &path) {
Client client(HOST, get_port());
sse::SSEClient sse(client, path);
std::mutex mutex;
std::condition_variable cv;
std::vector<sse::SSEMessage> messages;
auto done = false;
sse.on_message([&](const sse::SSEMessage &msg) {
std::lock_guard<std::mutex> lock(mutex);
// Ignore a replay from a reconnect that races with stop()
if (!done) { messages.push_back(msg); }
});
sse.on_event("end", [&](const sse::SSEMessage &) {
std::lock_guard<std::mutex> lock(mutex);
done = true;
cv.notify_all();
});
sse.set_reconnect_interval(100);
sse.start_async();
{
std::unique_lock<std::mutex> lock(mutex);
cv.wait_for(lock, std::chrono::seconds(5), [&] { return done; });
}
sse.stop();
EXPECT_TRUE(done);
return messages;
}
};
TEST_F(SSEParsingTest, SingleLineData) {
auto msgs = parse("data: hello\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].data, "hello");
EXPECT_EQ(msgs[0].event, "message");
}
TEST_F(SSEParsingTest, MultiLineData) {
auto msgs = parse("data: line1\ndata: line2\ndata: line3\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].data, "line1\nline2\nline3");
}
TEST_F(SSEParsingTest, CustomEventType) {
auto msgs = parse("event: update\ndata: payload\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].event, "update");
EXPECT_EQ(msgs[0].data, "payload");
}
TEST_F(SSEParsingTest, EventIdHandling) {
auto msgs = parse("id: 12345\ndata: test\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].id, "12345");
}
TEST_F(SSEParsingTest, EmptyEventId) {
auto msgs = parse("id:\ndata: test\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].id, "");
}
TEST_F(SSEParsingTest, CommentsIgnored) {
auto msgs = parse(": this is a comment\ndata: hello\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].data, "hello");
EXPECT_EQ(msgs[0].event, "message");
}
TEST_F(SSEParsingTest, ColonInValue) {
auto msgs = parse("data: hello:world:test\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].data, "hello:world:test");
}
TEST_F(SSEParsingTest, FieldNameOnly) {
// A line without a colon is a field name with an empty value
auto msgs = parse("data\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].data, "");
}
TEST_F(SSEParsingTest, TrailingCarriageReturn) {
auto msgs = parse("event: update\r\ndata: hello\r\n\r\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].event, "update");
EXPECT_EQ(msgs[0].data, "hello");
}
TEST_F(SSEParsingTest, FieldNameOnlyWithCarriageReturn) {
auto msgs = parse("data\r\n\r\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].data, "");
}
TEST_F(SSEParsingTest, UnknownFieldsIgnored) {
auto msgs = parse("unknown: value\ndata: hello\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].data, "hello");
EXPECT_EQ(msgs[0].event, "message");
}
TEST_F(SSEParsingTest, SpaceAfterColonOptional) {
auto msgs = parse("data: hello\n\ndata:hello\n\n");
ASSERT_EQ(msgs.size(), 2u);
EXPECT_EQ(msgs[0].data, "hello");
EXPECT_EQ(msgs[1].data, "hello");
}
TEST_F(SSEParsingTest, EventWithoutDataResetsEventType) {
// An event without data is not dispatched, and its type must not leak
// into the next event
auto msgs = parse("event: update\n\ndata: hello\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].event, "message");
EXPECT_EQ(msgs[0].data, "hello");
}
TEST_F(SSEParsingTest, EventWithoutDataUpdatesLastEventId) {
// The first connection sends only an id; the second echoes back the
// Last-Event-ID it was reconnected with
std::atomic<int> connection_count{0};
server_->Get("/id-only", [&](const Request &req, Response &res) {
if (connection_count++ == 0) {
res.set_content("id: 42\n\n", "text/event-stream");
} else {
res.set_content("data: " + req.get_header_value("Last-Event-ID") +
"\n\nevent: end\ndata: end\n\n",
"text/event-stream");
}
});
auto msgs = collect("/id-only");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].data, "42");
}
TEST_F(SSEParsingTest, CompleteEventParsing) {
auto msgs = parse("event: notification\nid: evt-42\n"
"data: {\"type\":\"alert\"}\nretry: 1000\n\n");
ASSERT_EQ(msgs.size(), 1u);
EXPECT_EQ(msgs[0].event, "notification");
EXPECT_EQ(msgs[0].id, "evt-42");
EXPECT_EQ(msgs[0].data, "{\"type\":\"alert\"}");
}
TEST(SSEMessageTest, Clear) {
sse::SSEMessage msg;
msg.event = "custom";
msg.data = "some data";
msg.id = "123";
msg.clear();
EXPECT_EQ(msg.event, "message");
EXPECT_EQ(msg.data, "");
EXPECT_EQ(msg.id, "");
}
//==============================================================================
// SSE Integration Tests
//==============================================================================
// Test: Successful connection and on_open callback
TEST_F(SSEIntegrationTest, SuccessfulConnection) {
// Add a simple endpoint that sends one event and closes