21 if (!m_open_for_write) {
24 m_open_for_write =
false;
28 bool all_buffers_returned = m_avail_count == m_buffers.size();
33 const char *msg = m_fh->error.getErrText();
34 if (!msg || (*msg ==
'\0')) {msg =
"(no error message provided)";}
35 ss <<
"Failure when closing file handle: " << msg <<
" (code=" << m_fh->error.getErrInfo() <<
")";
36 m_error_buf = ss.str();
40 return all_buffers_returned;
47 return m_fh->stat(buf);
61 if (!m_open_for_write) {
62 if (!m_error_buf.size()) {m_error_buf =
"Logic error: writing to a buffer not opened for write";}
65 size_t bytes_accepted = 0;
66 ssize_t retval = size;
67 if (offset < m_offset) {
68 if (!m_error_buf.size()) {m_error_buf =
"Logic error: writing to a prior offset";}
74 if (offset == m_offset && (force || (size && !(size % (1024*1024))))) {
75 retval = WriteImpl(offset, buf, size);
76 bytes_accepted = retval;
84 if (m_avail_count == m_buffers.size()) {
95 ssize_t buffers_flushed;
97 bytes_accepted += AcceptIntoBuffers(offset + bytes_accepted,
99 size - bytes_accepted);
100 buffers_flushed = FlushBuffers(size == 0);
102 }
while ((buffers_flushed > 0) && (bytes_accepted != size));
104 if (bytes_accepted != size && size) {
105 Entry *avail_entry = FirstAvailableBuffer();
108 m_error_buf =
"No empty buffers available to place unordered data.";
111 if (avail_entry->Accept(offset + bytes_accepted, buf + bytes_accepted, size - bytes_accepted) != size - bytes_accepted) {
112 m_error_buf =
"Empty re-ordering buffer was unable to to accept data; internal logic error.";
122 if ((m_buffers.size() > 2) && (m_avail_count * 2 > m_buffers.size())) {
123 for (
auto &entry : m_buffers) {
124 entry->ShrinkIfUnused();
133Stream::AcceptIntoBuffers(off_t offset,
const char *buf,
size_t size)
135 size_t bytes_accepted = 0;
136 if (!size) {
return 0;}
137 for (
auto &entry : m_buffers) {
141 if (entry->Available()) {
continue;}
142 bytes_accepted += entry->Accept(offset + bytes_accepted,
143 buf + bytes_accepted,
144 size - bytes_accepted);
145 if (bytes_accepted == size) {
break;}
147 return bytes_accepted;
152Stream::FlushBuffers(
bool force)
154 ssize_t buffers_flushed = 0;
155 bool buffer_was_written;
157 size_t avail_count = 0;
158 buffer_was_written =
false;
159 for (
auto &entry : m_buffers) {
160 ssize_t retval = entry->Write(*
this, force);
162 if (!m_error_buf.size()) {m_error_buf =
"Unknown filesystem write failure.";}
166 buffer_was_written =
true;
169 if (entry->Available()) {avail_count ++;}
171 m_avail_count = avail_count;
174 }
while (buffer_was_written && (m_avail_count != m_buffers.size()));
175 return buffers_flushed;
180Stream::FirstAvailableBuffer()
182 for (
auto &entry : m_buffers) {
183 if (entry->Available()) {
return entry.get();}
189ssize_t Stream::WriteImpl(off_t offset,
const char *buf,
size_t size)
192 if (size == 0) {
return 0;}
193 retval = m_fh->write(offset, buf, size);
197 std::stringstream ss;
198 const char *msg = m_fh->error.getErrText();
199 if (!msg || (*msg ==
'\0')) {msg =
"(no error message provided)";}
200 ss << msg <<
" (code=" << m_fh->error.getErrInfo() <<
")";
201 m_error_buf = ss.str();
210 m_log.Emsg(
"Stream::DumpBuffers",
"Beginning dump of stream buffers.");
212 std::stringstream ss;
213 ss <<
"Stream offset: " << m_offset;
214 m_log.Emsg(
"Stream::DumpBuffers", ss.str().c_str());
217 for (
const auto &entry : m_buffers) {
218 std::stringstream ss;
219 ss <<
"Buffer " << idx <<
": Offset=" << entry->GetOffset() <<
", Size="
220 << entry->GetSize() <<
", Capacity=" << entry->GetCapacity();
221 m_log.Emsg(
"Stream::DumpBuffers", ss.str().c_str());
224 m_log.Emsg(
"Stream::DumpBuffers",
"Finish dump of stream buffers.");
231 return m_fh->read(offset, buf, size);
int Read(off_t offset, char *buffer, size_t size)
ssize_t Write(off_t offset, const char *buffer, size_t size, bool force)