mirror of
https://github.com/wassname/ray.git
synced 2026-08-16 11:27:09 +08:00
[Streaming] DataWriter use event driven model. (#7043)
* streaming writer use event driven model. * minor changes according reviewer comments * Fix according to reviewer's comments * fix bazel lint * code polished * Add more comments * rename Stop & Start of EventQueue to Freeze and Unfreeze.
This commit is contained in:
@@ -23,11 +23,6 @@ bool StreamingRingBuffer::Push(const StreamingMessagePtr &msg) {
|
||||
return true;
|
||||
}
|
||||
|
||||
bool StreamingRingBuffer::Push(StreamingMessagePtr &&msg) {
|
||||
message_buffer_->Push(std::forward<StreamingMessagePtr>(msg));
|
||||
return true;
|
||||
}
|
||||
|
||||
StreamingMessagePtr &StreamingRingBuffer::Front() {
|
||||
STREAMING_CHECK(!message_buffer_->Empty());
|
||||
return message_buffer_->Front();
|
||||
@@ -38,11 +33,11 @@ void StreamingRingBuffer::Pop() {
|
||||
message_buffer_->Pop();
|
||||
}
|
||||
|
||||
bool StreamingRingBuffer::IsFull() { return message_buffer_->Full(); }
|
||||
bool StreamingRingBuffer::IsFull() const { return message_buffer_->Full(); }
|
||||
|
||||
bool StreamingRingBuffer::IsEmpty() { return message_buffer_->Empty(); }
|
||||
bool StreamingRingBuffer::IsEmpty() const { return message_buffer_->Empty(); }
|
||||
|
||||
size_t StreamingRingBuffer::Size() { return message_buffer_->Size(); };
|
||||
size_t StreamingRingBuffer::Size() const { return message_buffer_->Size(); }
|
||||
|
||||
size_t StreamingRingBuffer::Capacity() const { return message_buffer_->Capacity(); }
|
||||
|
||||
|
||||
Reference in New Issue
Block a user