Track streams
This commit is contained in:
2
firecgi
2
firecgi
Submodule firecgi updated: ec99454756...0ba446bacb
13
server.cc
13
server.cc
@@ -26,6 +26,19 @@ void Server::OnRequest(firecgi::Request* request) {
|
||||
request->WriteHeader("X-Accel-Buffering", "no");
|
||||
request->WriteBody("");
|
||||
auto stream = new Stream(request);
|
||||
|
||||
{
|
||||
std::lock_guard l(mu_);
|
||||
streams_.insert(stream);
|
||||
|
||||
request->OnClose([this, stream]() {
|
||||
std::lock_guard l(mu_);
|
||||
stream->Close();
|
||||
streams_.erase(stream);
|
||||
delete stream;
|
||||
});
|
||||
}
|
||||
|
||||
callback_(stream);
|
||||
}
|
||||
|
||||
|
||||
5
server.h
5
server.h
@@ -1,6 +1,6 @@
|
||||
#pragma once
|
||||
|
||||
#include <queue>
|
||||
#include <set>
|
||||
|
||||
#include "firecgi/server.h"
|
||||
#include "stream.h"
|
||||
@@ -19,6 +19,9 @@ class Server {
|
||||
|
||||
std::function<void(Stream*)> callback_;
|
||||
firecgi::Server firecgi_server_;
|
||||
|
||||
std::mutex mu_;
|
||||
std::set<Stream*, IsFresherStream> streams_;
|
||||
};
|
||||
|
||||
} // namespace firesse
|
||||
|
||||
28
stream.cc
28
stream.cc
@@ -3,19 +3,18 @@
|
||||
namespace firesse {
|
||||
|
||||
Stream::Stream(firecgi::Request* request)
|
||||
: request_(request) {
|
||||
request->OnClose([this]() {
|
||||
if (on_close_) {
|
||||
on_close_();
|
||||
}
|
||||
});
|
||||
}
|
||||
: request_(request) {}
|
||||
|
||||
void Stream::OnClose(const std::function<void()>& callback) {
|
||||
on_close_ = callback;
|
||||
}
|
||||
|
||||
bool Stream::WriteEvent(const std::string& data, uint64_t id, const std::string& type) {
|
||||
{
|
||||
std::lock_guard l(mu_);
|
||||
last_message_time_ = std::chrono::steady_clock::now();
|
||||
}
|
||||
|
||||
return request_->InTransaction<bool>([=]() {
|
||||
if (id) {
|
||||
request_->WriteBody("id: ", std::to_string(id), "\n");
|
||||
@@ -32,4 +31,19 @@ bool Stream::End() {
|
||||
return request_->End();
|
||||
}
|
||||
|
||||
std::chrono::steady_clock::time_point Stream::LastMessageTime() {
|
||||
std::lock_guard l(mu_);
|
||||
return last_message_time_;
|
||||
}
|
||||
|
||||
void Stream::Close() {
|
||||
if (on_close_) {
|
||||
on_close_();
|
||||
}
|
||||
}
|
||||
|
||||
bool IsFresherStream::operator() (Stream* a, Stream* b) const {
|
||||
return a->LastMessageTime() > b->LastMessageTime();
|
||||
}
|
||||
|
||||
} // namespace firesse
|
||||
|
||||
11
stream.h
11
stream.h
@@ -15,10 +15,21 @@ class Stream {
|
||||
bool WriteEvent(const std::string& data, uint64_t id=0, const std::string& type="");
|
||||
bool End();
|
||||
|
||||
std::chrono::steady_clock::time_point LastMessageTime();
|
||||
void Close();
|
||||
|
||||
private:
|
||||
firecgi::Request* request_;
|
||||
|
||||
std::function<void()> on_close_;
|
||||
|
||||
std::mutex mu_;
|
||||
std::chrono::steady_clock::time_point last_message_time_;
|
||||
};
|
||||
|
||||
class IsFresherStream {
|
||||
public:
|
||||
bool operator() (Stream* a, Stream* b) const;
|
||||
};
|
||||
|
||||
} // namespace firesse
|
||||
|
||||
Reference in New Issue
Block a user