Stroika Library 3.0d24
 
Loading...
Searching...
No Matches
BufferedInputStream.inl
1/*
2 * Copyright(c) Sophist Solutions, Inc. 1990-2026. All rights reserved
3 */
6#include "Stroika/Foundation/Streams/InternallySynchronizedInputStream.h"
8
9namespace Stroika::Foundation::Streams::BufferedInputStream {
10
11 namespace Private_ {
12
13 [[noreturn]] void ThrowCannotSeekFromEnd_ ();
14
15 // this case easy, delegate to StreamReader to do all the work
16 template <typename ELEMENT_TYPE>
17 class Rep_Seekable_FromSeekable_ : public IRep_<ELEMENT_TYPE> {
18 public:
19 Rep_Seekable_FromSeekable_ (const typename InputStream::Ptr<ELEMENT_TYPE>& realIn)
20 : fRealIn_{realIn}
21 , fReader_{realIn}
22 {
23 Require (realIn.IsSeekable ());
24 }
25 virtual bool IsSeekable () const override
26 {
27 return true;
28 }
29 virtual void CloseRead () override
30 {
31 if (fRealIn_ != nullptr) {
32 fRealIn_.Close ();
33 }
34 Ensure (not IsOpenRead ());
35 Assert (fRealIn_ == nullptr);
36 }
37 virtual bool IsOpenRead () const override
38 {
39 return fRealIn_ != nullptr;
40 }
41 virtual SeekOffsetType GetReadOffset () const override
42 {
43 Require (IsOpenRead ());
44 return fReader_.GetOffset ();
45 }
46 virtual optional<size_t> AvailableToRead () override
47 {
48 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
49 Require (IsOpenRead ());
50 return fReader_.AvailableToRead (); // since no actual buffering here yet
51 }
52 virtual optional<SeekOffsetType> RemainingLength () override
53 {
54 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
55 Require (IsOpenRead ());
56 return fReader_.RemainingLength ();
57 }
58 virtual SeekOffsetType SeekRead (Whence whence, SignedSeekOffsetType offset) override
59 {
60 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
61 Require (IsOpenRead ());
62 return fReader_.Seek (whence, offset);
63 }
64 virtual optional<span<ELEMENT_TYPE>> Read (span<ELEMENT_TYPE> intoBuffer, NoDataAvailableHandling blockFlag) override
65 {
66 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
67 Require (IsOpenRead ());
68 return fReader_.Read (intoBuffer, blockFlag);
69 }
70
71 private:
72 typename InputStream::Ptr<ELEMENT_TYPE> fRealIn_;
73 StreamReader<ELEMENT_TYPE> fReader_;
74 qStroika_ATTRIBUTE_NO_UNIQUE_ADDRESS_VCFORCE Debug::AssertExternallySynchronizedChecker fThisAssertExternallySynchronized_;
75 };
76
77 // read the source into one big buffer. Keep it all around, so seekable
78 template <typename ELEMENT_TYPE, size_t INLINE_BUF_SIZE>
79 class Rep_Seekable_FromUnSeekable_ : public IRep_<ELEMENT_TYPE> {
80 public:
81 Rep_Seekable_FromUnSeekable_ (const typename InputStream::Ptr<ELEMENT_TYPE>& realIn)
82 : fRealIn_{realIn}
83 {
84 }
85 virtual bool IsSeekable () const override
86 {
87 return true;
88 }
89 virtual void CloseRead () override
90 {
91 if (fRealIn_ != nullptr) {
92 fRealIn_.Close ();
93 }
94 Ensure (not IsOpenRead ());
95 Assert (fRealIn_ == nullptr);
96 }
97 virtual bool IsOpenRead () const override
98 {
99 return fRealIn_ != nullptr;
100 }
101 virtual SeekOffsetType GetReadOffset () const override
102 {
103 Require (IsOpenRead ());
104 return fSeekOffset_;
105 }
106 virtual optional<size_t> AvailableToRead () override
107 {
108 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
109 Require (IsOpenRead ());
110 if (fSeekOffset_ < fBufferOfAllReadDataSoFar_.size ()) [[likely]] {
111 return fBufferOfAllReadDataSoFar_.size () - static_cast<size_t> (fSeekOffset_); // don't include what we might get upstream cuz more costly to compute
112 }
113 return fRealIn_.AvailableToRead ();
114 }
115 virtual optional<SeekOffsetType> RemainingLength () override
116 {
117 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
118 Require (IsOpenRead ());
119 if (auto rl = fRealIn_.RemainingLength ()) {
120 return MapOffsetFromReal2Mine_ (*rl);
121 }
122 return nullopt;
123 }
124 virtual auto SeekRead (Whence whence, SignedSeekOffsetType offset) -> SeekOffsetType override
125 {
126 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
127 Require (IsOpenRead ());
128 // @todo - allow seek forward past fBufferOfAllReadDataSoFar_?
129 switch (whence) {
130 case Whence::eFromCurrent:
131 fSeekOffset_ += offset;
132 break;
133 case Whence::eFromStart:
134 fSeekOffset_ = offset;
135 break;
136 case Whence::eFromEnd:
137 if (auto remaining = this->RemainingLength ()) {
138 fSeekOffset_ += static_cast<SignedSeekOffsetType> (*remaining) - offset;
139 break;
140 }
141 else {
142 Private_::ThrowCannotSeekFromEnd_ ();
143 }
144 default:
146 }
147 return fSeekOffset_;
148 }
149 virtual optional<span<ELEMENT_TYPE>> Read (span<ELEMENT_TYPE> intoBuffer, NoDataAvailableHandling blockFlag) override
150 {
151 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
152 Require (IsOpenRead ());
153 Assert (fSeekOffset_ <= fBufferOfAllReadDataSoFar_.size ());
154 if (fSeekOffset_ == fBufferOfAllReadDataSoFar_.size ()) [[unlikely]] {
155 ELEMENT_TYPE buf[1024];
156 if (auto r = fRealIn_.Read (span{buf}, blockFlag)) {
157 fBufferOfAllReadDataSoFar_.push_back (*r); // continue, and fall through
158 }
159 else {
160 return nullopt; // no data pre-read, and nothing available upstream
161 }
162 }
163 if (fSeekOffset_ <= fBufferOfAllReadDataSoFar_.size ()) [[likely]] {
164 size_t n2Read = min<size_t> (intoBuffer.size (), static_cast<size_t> (fBufferOfAllReadDataSoFar_.size () - fSeekOffset_));
165 auto result = Memory::CopySpanData (span{fBufferOfAllReadDataSoFar_}.subspan (static_cast<size_t> (fSeekOffset_), n2Read), intoBuffer);
166 Assert (result.size () == n2Read);
167 fSeekOffset_ += n2Read;
168 return result;
169 }
170 return nullopt;
171 }
172
173 private:
174 nonvirtual SeekOffsetType MapOffsetFromReal2Mine_ (SeekOffsetType so) const
175 {
176 return static_cast<SeekOffsetType> (static_cast<SignedSeekOffsetType> (so) + static_cast<SignedSeekOffsetType> (fRealIn_.GetOffset ()) -
177 static_cast<SignedSeekOffsetType> (fSeekOffset_));
178 }
179
180 private:
181 typename InputStream::Ptr<ELEMENT_TYPE> fRealIn_;
182 Memory::InlineBuffer<ELEMENT_TYPE, INLINE_BUF_SIZE> fBufferOfAllReadDataSoFar_;
183 SeekOffsetType fSeekOffset_{0}; // always inside fBufferOfAllReadDataSoFar_
184 qStroika_ATTRIBUTE_NO_UNIQUE_ADDRESS_VCFORCE Debug::AssertExternallySynchronizedChecker fThisAssertExternallySynchronized_;
185 };
186
187 // pretty easy/efficient case cuz we can throw away data as we go, and since not seekable, not many cases to analyze
188 template <typename ELEMENT_TYPE, size_t INLINE_BUF_SIZE>
189 class Rep_UnSeekable_ : public IRep_<ELEMENT_TYPE> {
190 public:
191 Rep_UnSeekable_ (const typename InputStream::Ptr<ELEMENT_TYPE>& realIn)
192 : fRealIn_{realIn}
193 {
194 }
195 virtual bool IsSeekable () const override
196 {
197 return false;
198 }
199 virtual void CloseRead () override
200 {
201 if (fRealIn_ != nullptr) {
202 fRealIn_.Close ();
203 }
204 Ensure (not IsOpenRead ());
205 Assert (fRealIn_ == nullptr);
206 }
207 virtual bool IsOpenRead () const override
208 {
209 return fRealIn_ != nullptr;
210 }
211 virtual SeekOffsetType GetReadOffset () const override
212 {
213 Require (IsOpenRead ());
214 // fRealIn_ has typically pre-read PAST what the caller has consumed, so back out
215 // whatever is still sitting unread in fIntermediateBuffer_
216 Assert (fRealIn_.GetOffset () >= GetNEltsAlreadyBufferedFromUpstream_ ());
217 return fRealIn_.GetOffset () - GetNEltsAlreadyBufferedFromUpstream_ ();
218 }
219 virtual optional<size_t> AvailableToRead () override
220 {
221 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
222 Require (IsOpenRead ());
223 size_t n = GetNEltsAlreadyBufferedFromUpstream_ ();
224 if (auto o = fRealIn_.AvailableToRead ()) {
225 // if we KNOW what's available upstream, add to what we've pre-read
226 return *o + n;
227 }
228 else if (n != 0) {
229 return n; // if zero buffered, and nothing KNOWN about upstream, return nullopt
230 }
231 return nullopt; // if nothing upstream available (but not zero/eof)
232 }
233 virtual optional<SeekOffsetType> RemainingLength () override
234 {
235 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
236 Require (IsOpenRead ());
237 if (auto o = fRealIn_.RemainingLength ()) {
238 return *o + GetNEltsAlreadyBufferedFromUpstream_ ();
239 }
240 return nullopt; // if nothing upstream available (but not zero/eof)
241 }
242 virtual optional<span<ELEMENT_TYPE>> Read (span<ELEMENT_TYPE> intoBuffer, NoDataAvailableHandling blockFlag) override
243 {
244 Debug::AssertExternallySynchronizedChecker::WriteContext declareContext{fThisAssertExternallySynchronized_};
245 Require (IsOpenRead ());
246 auto n = GetNEltsAlreadyBufferedFromUpstream_ ();
247 if (n == 0) [[unlikely]] {
248 // read a big chunk upstream (not into intoBuffer, which may be tiny) - OK to
249 // overwrite what we had cuz no seeking allowed, so we never re-examine it.
250 // NB: read via a local buffer, NOT span{fIntermediateBuffer_}: an InlineBuffer's
251 // size () starts out ZERO (INLINE_BUF_SIZE is only its inline CAPACITY), so that
252 // span would be empty - which is a precondition violation in Read (). Going
253 // through a local also means fIntermediateBuffer_ is only touched once the read
254 // has actually succeeded.
255 ELEMENT_TYPE buf[INLINE_BUF_SIZE];
256 if (auto bufR = fRealIn_.Read (span{buf}, blockFlag)) {
257 // we filled buffer (possibly with zero elements)
258 fIntermediateBuffer_.resize_uninitialized (0);
259 fIntermediateBuffer_.push_back (*bufR);
260 fReadOffsetIntoIntermediateBuf_ = 0;
261 n = bufR->size ();
262 }
263 else {
264 return nullopt; // no new information, don't change state so can Read again
265 }
266 }
267 // if we get here, and n = 0, really EOF, cuz return above on NOT-AVAIL case
268 // nb: IsAtEOF () peeks - read then seek back - so it can only be asked of a
269 // seekable stream, and fRealIn_ here often is not one
270 Assert (n != 0 or not fRealIn_.IsSeekable () or fRealIn_.IsAtEOF ());
271 size_t n2Read = Math::AtMost (n, intoBuffer.size ());
272 auto t = Memory::CopySpanData (span{fIntermediateBuffer_}.subspan (fReadOffsetIntoIntermediateBuf_, n2Read), intoBuffer);
273 Assert (t.size () == n2Read);
274 fReadOffsetIntoIntermediateBuf_ += n2Read;
275 return t;
276 }
277
278 private:
279 nonvirtual size_t GetNEltsAlreadyBufferedFromUpstream_ () const
280 {
281 Assert (fReadOffsetIntoIntermediateBuf_ <= fIntermediateBuffer_.size ());
282 return fIntermediateBuffer_.size () - fReadOffsetIntoIntermediateBuf_;
283 }
284
285 private:
286 typename InputStream::Ptr<ELEMENT_TYPE> fRealIn_;
287 Memory::InlineBuffer<ELEMENT_TYPE, INLINE_BUF_SIZE> fIntermediateBuffer_;
288 size_t fReadOffsetIntoIntermediateBuf_{0};
289 qStroika_ATTRIBUTE_NO_UNIQUE_ADDRESS_VCFORCE Debug::AssertExternallySynchronizedChecker fThisAssertExternallySynchronized_;
290 };
291 }
292
293 /*
294 ********************************************************************************
295 ********************* Streams::BufferedInputStream::New ************************
296 ********************************************************************************
297 */
298 template <typename ELEMENT_TYPE>
299 inline auto New (const typename InputStream::Ptr<ELEMENT_TYPE>& realIn, optional<SeekableFlag> seekable) -> Ptr<ELEMENT_TYPE>
300 {
301 using PTR = Ptr<ELEMENT_TYPE>;
302 SeekableFlag srcSeekable = realIn.GetSeekability ();
303 SeekableFlag useSeekable = seekable.value_or (srcSeekable);
304 constexpr size_t INLINE_BUF_SIZE = 4 * 1024;
305 if (useSeekable == SeekableFlag::eSeekable) {
306 return (srcSeekable == SeekableFlag::eSeekable)
307 ? PTR{Memory::MakeSharedPtr<Private_::Rep_Seekable_FromSeekable_<ELEMENT_TYPE>> (realIn)}
308 : PTR{Memory::MakeSharedPtr<Private_::Rep_Seekable_FromUnSeekable_<ELEMENT_TYPE, INLINE_BUF_SIZE>> (realIn)};
309 }
310 else {
311 return PTR{Memory::MakeSharedPtr<Private_::Rep_UnSeekable_<ELEMENT_TYPE, INLINE_BUF_SIZE>> (realIn)};
312 }
313 }
314 template <typename ELEMENT_TYPE>
315 inline auto New (Execution::InternallySynchronized internallySynchronized, const typename InputStream::Ptr<ELEMENT_TYPE>& realIn,
316 optional<SeekableFlag> seekable) -> Ptr<ELEMENT_TYPE>
317 {
318 constexpr size_t INLINE_BUF_SIZE = 4 * 1024;
319 switch (internallySynchronized) {
320 case Execution::eInternallySynchronized: {
321 SeekableFlag srcSeekable = realIn.GetSeekability ();
322 SeekableFlag useSeekable = seekable.value_or (srcSeekable);
323 if (useSeekable == SeekableFlag::eSeekable) {
324 return (srcSeekable == SeekableFlag::eSeekable)
325 ? InternallySynchronizedInputStream::New<Private_::Rep_Seekable_FromSeekable_<ELEMENT_TYPE>> ({}, realIn)
326 : InternallySynchronizedInputStream::New<Private_::Rep_Seekable_FromUnSeekable_<ELEMENT_TYPE, INLINE_BUF_SIZE>> ({}, realIn);
327 }
328 else {
329 return InternallySynchronizedInputStream::New<Private_::Rep_UnSeekable_<ELEMENT_TYPE, INLINE_BUF_SIZE>> ({}, realIn);
330 }
331 }
332 case Execution::eNotKnownInternallySynchronized:
333 return New<ELEMENT_TYPE> (realIn, seekable);
334 default:
336 return nullptr;
337 }
338 }
339
340 /*
341 ********************************************************************************
342 ******************* BufferedInputStream::Ptr<ELEMENT_TYPE> *********************
343 ********************************************************************************
344 */
345 template <typename ELEMENT_TYPE>
346 inline Ptr<ELEMENT_TYPE>::Ptr (const shared_ptr<Private_::IRep_<ELEMENT_TYPE>>& from)
347 : inherited{from}
348 {
349 }
350
351 template <typename ELEMENT_TYPE>
352 [[deprecated ("Since Stroika v3.0d19 use Seekability overload")]] Ptr<ELEMENT_TYPE> New (const typename InputStream::Ptr<ELEMENT_TYPE>& realIn,
353 optional<bool> seekable)
354 {
355 optional<SeekableFlag> sf;
356 if (seekable) {
357 sf = *seekable ? SeekableFlag::eSeekable : SeekableFlag::eNotSeekable;
358 }
359 return New (realIn, sf);
360 }
361 template <typename ELEMENT_TYPE>
362 [[deprecated ("Since Stroika v3.0d19 use Seekability overload")]] Ptr<ELEMENT_TYPE>
363 New (Execution::InternallySynchronized internallySynchronized, const typename InputStream::Ptr<ELEMENT_TYPE>& realIn, optional<bool> seekable = {})
364 {
365 optional<SeekableFlag> sf;
366 if (seekable) {
367 sf = *seekable ? SeekableFlag::eSeekable : SeekableFlag::eNotSeekable;
368 }
369 return New (internallySynchronized, realIn, sf);
370 }
371}
#define RequireNotReached()
Definition Assertions.h:386
#define qStroika_ATTRIBUTE_NO_UNIQUE_ADDRESS_VCFORCE
[[msvc::no_unique_address]] isn't always broken in MSVC. Annotate with this on things where its not b...
Definition StdCompat.h:443
unique_lock< AssertExternallySynchronizedChecker > WriteContext
Instantiate AssertExternallySynchronizedChecker::WriteContext to designate an area of code where prot...