XRootD
Loading...
Searching...
No Matches
XrdHttpTpcStream.cc
Go to the documentation of this file.
1
2#include <sstream>
3
4#include "XrdHttpTpcStream.hh"
5
8
9using namespace TPC;
10
12{
13 m_fh->close();
14}
15
16
17bool
19{
20 // Do not close twice
21 if (!m_open_for_write) {
22 return false;
23 }
24 m_open_for_write = false;
25
26 // If there are outstanding buffers to reorder, finalization failed; the
27 // check has to happen before the buffers are released.
28 bool all_buffers_returned = m_avail_count == m_buffers.size();
29 m_buffers.clear();
30
31 if (m_fh->close() == SFS_ERROR) {
32 std::stringstream ss;
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();
37 return false;
38 }
39
40 return all_buffers_returned;
41}
42
43
44int
45Stream::Stat(struct stat* buf)
46{
47 return m_fh->stat(buf);
48}
49
50ssize_t
51Stream::Write(off_t offset, const char *buf, size_t size, bool force)
52{
53/*
54 * NOTE: these lines are useful for debuggin the state of the buffer
55 * management code; too expensive to compile in and have a runtime switch.
56 std::stringstream ss;
57 ss << "Offset=" << offset << ", Size=" << size << ", force=" << force;
58 m_log.Emsg("Stream::Write", ss.str().c_str());
59 DumpBuffers();
60*/
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";}
63 return SFS_ERROR;
64 }
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";}
69 return SFS_ERROR;
70 }
71 // If this is write is appending to the stream and
72 // MB-aligned, then we write it to disk; otherwise, the
73 // data will be buffered.
74 if (offset == m_offset && (force || (size && !(size % (1024*1024))))) {
75 retval = WriteImpl(offset, buf, size);
76 bytes_accepted = retval;
77 // On failure, we don't care about flushing buffers from memory --
78 // the stream is now invalid.
79 if (retval < 0) {
80 return retval;
81 }
82 // If there are no in-use buffers, then we don't need to
83 // do any accounting.
84 if (m_avail_count == m_buffers.size()) {
85 return retval;
86 }
87 }
88 // Even if we already accepted the current data, always iterate through the
89 // buffers and try to write as much out to disk as possible.
90 //
91 // Accepting data can complete a buffer, and flushing a buffer advances
92 // m_offset, which can in turn let another buffer accept more data or become
93 // writable. Alternate between the two until neither makes progress. When
94 // size == 0 we force a flush even if things are not MB-aligned.
95 ssize_t buffers_flushed;
96 do {
97 bytes_accepted += AcceptIntoBuffers(offset + bytes_accepted,
98 buf + bytes_accepted,
99 size - bytes_accepted);
100 buffers_flushed = FlushBuffers(size == 0);
101 if (buffers_flushed == SFS_ERROR) {return SFS_ERROR;}
102 } while ((buffers_flushed > 0) && (bytes_accepted != size));
103
104 if (bytes_accepted != size && size) { // No place for this data in the buffers currently in use
105 Entry *avail_entry = FirstAvailableBuffer();
106 if (!avail_entry) { // No available buffers to allocate; logic error, should not happen.
107 DumpBuffers();
108 m_error_buf = "No empty buffers available to place unordered data.";
109 return SFS_ERROR;
110 }
111 if (avail_entry->Accept(offset + bytes_accepted, buf + bytes_accepted, size - bytes_accepted) != size - bytes_accepted) { // Empty buffer cannot accept?!?
112 m_error_buf = "Empty re-ordering buffer was unable to to accept data; internal logic error.";
113 return SFS_ERROR;
114 }
115 // The buffer we just filled may already be complete and contiguous with
116 // m_offset; flush it now instead of waiting for a later callback to
117 // notice, as every curl handle may be idle by then.
118 if (FlushBuffers(false) == SFS_ERROR) {return SFS_ERROR;}
119 }
120
121 // If we have low buffer occupancy, then release memory.
122 if ((m_buffers.size() > 2) && (m_avail_count * 2 > m_buffers.size())) {
123 for (auto &entry : m_buffers) {
124 entry->ShrinkIfUnused();
125 }
126 }
127
128 return retval;
129}
130
131
132size_t
133Stream::AcceptIntoBuffers(off_t offset, const char *buf, size_t size)
134{
135 size_t bytes_accepted = 0;
136 if (!size) {return 0;}
137 for (auto &entry : m_buffers) {
138 // Empty buffers are deliberately skipped here: they are only handed out
139 // as a last resort by Write() so that buffer occupancy keeps tracking
140 // the number of transfers in flight.
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;}
146 }
147 return bytes_accepted;
148}
149
150
151ssize_t
152Stream::FlushBuffers(bool force)
153{
154 ssize_t buffers_flushed = 0;
155 bool buffer_was_written;
156 do {
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);
161 if (retval == SFS_ERROR) {
162 if (!m_error_buf.size()) {m_error_buf = "Unknown filesystem write failure.";}
163 return SFS_ERROR;
164 }
165 if (retval > 0) {
166 buffer_was_written = true;
167 buffers_flushed ++;
168 }
169 if (entry->Available()) {avail_count ++;}
170 }
171 m_avail_count = avail_count;
172 // Writing a buffer advances m_offset, which may have made a buffer we
173 // already walked past contiguous with the stream; go around again.
174 } while (buffer_was_written && (m_avail_count != m_buffers.size()));
175 return buffers_flushed;
176}
177
178
179Stream::Entry *
180Stream::FirstAvailableBuffer()
181{
182 for (auto &entry : m_buffers) {
183 if (entry->Available()) {return entry.get();}
184 }
185 return nullptr;
186}
187
188
189ssize_t Stream::WriteImpl(off_t offset, const char *buf, size_t size)
190{
191 ssize_t retval;
192 if (size == 0) {return 0;}
193 retval = m_fh->write(offset, buf, size);
194 if (retval != SFS_ERROR) {
195 m_offset += retval;
196 } else {
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();
202 }
203 return retval;
204}
205
206
207void
209{
210 m_log.Emsg("Stream::DumpBuffers", "Beginning dump of stream buffers.");
211 {
212 std::stringstream ss;
213 ss << "Stream offset: " << m_offset;
214 m_log.Emsg("Stream::DumpBuffers", ss.str().c_str());
215 }
216 size_t idx = 0;
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());
222 idx ++;
223 }
224 m_log.Emsg("Stream::DumpBuffers", "Finish dump of stream buffers.");
225}
226
227
228int
229Stream::Read(off_t offset, char *buf, size_t size)
230{
231 return m_fh->read(offset, buf, size);
232}
#define stat(a, b)
Definition XrdPosix.hh:101
#define SFS_ERROR
int Read(off_t offset, char *buffer, size_t size)
ssize_t Write(off_t offset, const char *buffer, size_t size, bool force)
void DumpBuffers() const
int Stat(struct stat *)