XRootD
XrdHttpTpcStream.cc
Go to the documentation of this file.
1 
2 #include <sstream>
3 
4 #include "XrdHttpTpcStream.hh"
5 
7 #include "XrdSys/XrdSysError.hh"
8 
9 using namespace TPC;
10 
12 {
13  m_fh->close();
14 }
15 
16 
17 bool
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 
44 int
45 Stream::Stat(struct stat* buf)
46 {
47  return m_fh->stat(buf);
48 }
49 
50 ssize_t
51 Stream::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 
132 size_t
133 Stream::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 
151 ssize_t
152 Stream::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 
179 Stream::Entry *
180 Stream::FirstAvailableBuffer()
181 {
182  for (auto &entry : m_buffers) {
183  if (entry->Available()) {return entry.get();}
184  }
185  return nullptr;
186 }
187 
188 
189 ssize_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 
207 void
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 
228 int
229 Stream::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 *)
int Emsg(const char *esfx, int ecode, const char *text1, const char *text2=0)
Definition: XrdSysError.cc:95