Skip to content

Commit a09050b

Browse files
committed
feat(daemon): pull from input pipes during each flush round
1 parent 9503106 commit a09050b

7 files changed

Lines changed: 65 additions & 24 deletions

File tree

daemon/buffer.cpp

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -56,12 +56,12 @@ size_t HopperBuffer::write(void *src, size_t len) {
5656
return done_len;
5757
}
5858

59-
size_t HopperBuffer::write(HopperPipe *pipe) {
59+
size_t HopperBuffer::write(HopperPipe *pipe, bool *more) {
6060
size_t max_len = max_write();
6161
size_t done_len = 0;
6262

6363
size_t next_len = std::min((m_buf.size() - m_edge), max_len);
64-
size_t res = pipe->read_pipe(&m_buf[m_edge], next_len);
64+
size_t res = pipe->read_pipe(&m_buf[m_edge], next_len, more);
6565
if (res == (size_t)-1)
6666
// -1 indicates read error
6767
return -1;
@@ -74,7 +74,7 @@ size_t HopperBuffer::write(HopperPipe *pipe) {
7474
return done_len;
7575

7676
next_len = max_len - next_len;
77-
res = pipe->read_pipe(&m_buf[m_edge], next_len);
77+
res = pipe->read_pipe(&m_buf[m_edge], next_len, more);
7878
if (res == (size_t)-1)
7979
return -1;
8080

@@ -106,13 +106,13 @@ size_t HopperBuffer::read(BufferMarker *m, void *dst, size_t len) {
106106
return done_len;
107107
}
108108

109-
size_t HopperBuffer::read(HopperPipe *pipe) {
109+
size_t HopperBuffer::read(HopperPipe *pipe, bool *more) {
110110
BufferMarker *m = pipe->marker();
111111
size_t max_len = max_read(m);
112112
size_t done_len = 0;
113113

114114
size_t next_len = std::min((m_buf.size() - m->pos()), max_len);
115-
size_t res = pipe->write_pipe(&m_buf[m->pos()], next_len);
115+
size_t res = pipe->write_pipe(&m_buf[m->pos()], next_len, more);
116116
if (res == (size_t)-1)
117117
return -1;
118118

@@ -122,7 +122,7 @@ size_t HopperBuffer::read(HopperPipe *pipe) {
122122
return done_len;
123123

124124
next_len = max_len - next_len;
125-
res = pipe->write_pipe(&m_buf[m->pos()], next_len);
125+
res = pipe->write_pipe(&m_buf[m->pos()], next_len, more);
126126
if (res == (size_t)-1)
127127
return -1;
128128

@@ -134,7 +134,7 @@ size_t HopperBuffer::read(HopperPipe *pipe) {
134134
size_t HopperBuffer::max_write() {
135135
size_t cap = m_buf.size();
136136
if (m_markers.empty())
137-
return cap;
137+
return std::max((size_t)0, cap - 1);
138138

139139
size_t min_dist = cap;
140140

@@ -148,7 +148,7 @@ size_t HopperBuffer::max_write() {
148148
min_dist = d;
149149
}
150150

151-
return min_dist;
151+
return std::max((size_t)0, min_dist - 1);
152152
}
153153

154154
size_t HopperBuffer::max_read(BufferMarker *m) {

daemon/endpoint.cpp

Lines changed: 28 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,7 @@
1+
#include <fcntl.h>
12
#include <filesystem>
3+
#include <sys/ioctl.h>
4+
#include <unistd.h>
25

36
#include "hopper/daemon/endpoint.hpp"
47
#include "hopper/daemon/pipe.hpp"
@@ -25,21 +28,41 @@ void HopperEndpoint::on_pipe_readable(uint64_t id) {
2528

2629
HopperPipe *pipe = m_inputs[id];
2730

28-
size_t res = m_buffer.write(pipe);
31+
/*
32+
* Calculation for "pipe pressure"
33+
int p_sz, p_nb;
34+
if ((p_sz = fcntl(pipe->fd(), F_GETPIPE_SZ)) != -1 &&
35+
(ioctl(pipe->fd(), FIONREAD, &p_nb) > -1)) {
36+
m_logger.trace(*pipe, ": ", (float)p_nb / (float)p_sz);
37+
}
38+
*/
39+
40+
bool more = false;
41+
size_t res = m_buffer.write(pipe, &more);
2942
if (res == (size_t)-1)
3043
throw_errno("read");
3144

32-
m_logger.trace(*pipe, " -> ", res, " bytes");
45+
if (more && !m_more.contains(id))
46+
m_more.insert(id);
47+
48+
m_logger.trace(*pipe, " -> ", res, " bytes", (more ? " ++" : ""));
3349
}
3450

3551
void HopperEndpoint::flush_pipes() {
52+
auto c_more = std::unordered_set<uint64_t>(m_more.begin(), m_more.end());
53+
for (const uint64_t id : c_more) {
54+
m_more.erase(id);
55+
on_pipe_readable(id);
56+
}
57+
3658
for (const auto &[_, pipe] : m_outputs) {
3759
if (pipe->status() == PipeStatus::INACTIVE)
3860
continue;
3961

40-
size_t res = m_buffer.read(pipe);
62+
bool more = false;
63+
size_t res = m_buffer.read(pipe, &more);
4164
if (res > 0)
42-
m_logger.trace(*pipe, " <- ", res, " bytes");
65+
m_logger.trace(*pipe, " <- ", res, " bytes", (more ? " ++" : ""));
4366
}
4467
}
4568

@@ -98,6 +121,7 @@ void HopperEndpoint::remove_by_id(uint64_t pipe_id) {
98121

99122
delete pipe;
100123
m_inputs.erase(pipe_id);
124+
m_more.erase(pipe_id);
101125
} else if (type == PipeType::OUT && m_outputs.contains(pipe_id)) {
102126
HopperPipe *pipe = m_outputs[pipe_id];
103127
m_buffer.delete_marker(pipe->marker());

daemon/pipe.cpp

Lines changed: 20 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -58,10 +58,13 @@ void HopperPipe::close_pipe() {
5858
close(m_fd);
5959
}
6060

61-
size_t HopperPipe::write_pipe(void *src, size_t len) {
61+
size_t HopperPipe::write_pipe(void *src, size_t len, bool *more) {
6262
if (m_type == PipeType::IN)
6363
return -1;
6464

65+
if (more)
66+
*more = true;
67+
6568
if (len == 0)
6669
return 0;
6770

@@ -71,9 +74,11 @@ size_t HopperPipe::write_pipe(void *src, size_t len) {
7174
ssize_t res = write(m_fd, reinterpret_cast<char *>(src) + done_len,
7275
len - done_len);
7376

74-
if (res == -1 && (errno == EWOULDBLOCK || errno == EINTR))
77+
if (res == -1 && (errno == EWOULDBLOCK || errno == EINTR)) {
78+
if (more)
79+
*more = false;
7580
return done_len;
76-
else if (res == -1) {
81+
} else if (res == -1) {
7782
perror("write");
7883
return -1;
7984
}
@@ -84,10 +89,13 @@ size_t HopperPipe::write_pipe(void *src, size_t len) {
8489
return done_len;
8590
}
8691

87-
size_t HopperPipe::read_pipe(void *dst, size_t len) {
92+
size_t HopperPipe::read_pipe(void *dst, size_t len, bool *more) {
8893
if (m_type == PipeType::OUT)
8994
return -1;
9095

96+
if (more)
97+
*more = true;
98+
9199
if (len == 0)
92100
return 0;
93101

@@ -97,13 +105,18 @@ size_t HopperPipe::read_pipe(void *dst, size_t len) {
97105
ssize_t res = read(m_fd, reinterpret_cast<char *>(dst) + done_len,
98106
len - done_len);
99107

100-
if (res == -1 && (errno == EWOULDBLOCK || errno == EINTR))
108+
if (res == -1 && (errno == EWOULDBLOCK || errno == EINTR)) {
109+
if (more)
110+
*more = false;
101111
return done_len;
102-
else if (res == -1) {
112+
} else if (res == -1) {
103113
perror("read");
104114
return -1;
105-
} else if (res == 0)
115+
} else if (res == 0) {
116+
if (more)
117+
*more = false;
106118
break;
119+
}
107120

108121
done_len += res;
109122
}

include/hopper/daemon/buffer.hpp

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,10 +25,10 @@ class HopperBuffer {
2525
void delete_marker(BufferMarker *marker);
2626

2727
size_t write(void *src, size_t len);
28-
size_t write(HopperPipe *pipe);
28+
size_t write(HopperPipe *pipe, bool *more);
2929

3030
size_t read(BufferMarker *marker, void *dst, size_t len);
31-
size_t read(HopperPipe *pipe);
31+
size_t read(HopperPipe *pipe, bool *more);
3232

3333
size_t max_write();
3434
size_t max_read(BufferMarker *marker);

include/hopper/daemon/daemon.hpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ class HopperDaemon {
3434
int m_epoll_fd = -1;
3535

3636
int m_max_events = 64;
37-
int m_timeout = 250;
37+
int m_timeout = 100;
3838

3939
std::filesystem::path m_path;
4040
Logger &m_logger;

include/hopper/daemon/endpoint.hpp

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22
#define endpoint_hpp_INCLUDED
33

44
#include <unordered_map>
5+
#include <unordered_set>
56

67
#include "hopper/daemon/buffer.hpp"
78
#include "hopper/daemon/logging.hpp"
@@ -14,6 +15,9 @@ class HopperEndpoint {
1415
std::unordered_map<uint64_t, HopperPipe *> m_inputs;
1516
std::unordered_map<uint64_t, HopperPipe *> m_outputs;
1617

18+
// inputs that may have more data
19+
std::unordered_set<uint64_t> m_more;
20+
1721
HopperBuffer m_buffer{};
1822

1923
uint64_t m_last_pipe_id = 1;

include/hopper/daemon/pipe.hpp

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -40,8 +40,8 @@ class HopperPipe {
4040

4141
int open_pipe();
4242
void close_pipe();
43-
size_t write_pipe(void *src, size_t len);
44-
size_t read_pipe(void *dst, size_t len);
43+
size_t write_pipe(void *src, size_t len, bool *more);
44+
size_t read_pipe(void *dst, size_t len, bool *more);
4545

4646
const std::filesystem::path &path() { return m_path; }
4747
BufferMarker *marker() { return m_marker; }

0 commit comments

Comments
 (0)