CX Framework
Cross-platform C utility framework
Loading...
Searching...
No Matches
streambuf.h
Go to the documentation of this file.
1
103
104#pragma once
105
106#include <cx/buffer/bufring.h>
107#include <cx/closure/closure.h>
108#include <cx/stype/stype.h>
109#include <cx/thread/condvar.h>
110#include <cx/thread/mutex.h>
111
112CX_C_BEGIN
113
114typedef struct StreamBuffer StreamBuffer;
115
132typedef size_t (*sbufPullCB)(stvlist* cvars, _Pre_valid_ StreamBuffer* sb,
133 _Out_writes_bytes_(sz) uint8* buf, size_t sz);
134
144typedef void (*sbufPushCB)(stvlist* cvars, _Pre_valid_ StreamBuffer* sb,
145 _In_reads_bytes_(sz) const uint8* buf, size_t sz);
146
147// Send callback
148// This callback is used with sbufCSend. It may be called multiple times with varying
149// offsets. The offset passed is always from the start of the available bytes in the
150// buffer.
151// Consumption is all-or-nothing across the whole sbufCSend call: if every invocation
152// returns true the bytes are consumed and removed from the buffer, and if any one of them
153// returns false the buffer keeps all of them, like the peek functions. A callback that can
154// accept some of the data but not the rest should therefore return false every time and
155// have the caller sbufCSkip() exactly what it took.
156// This callback MUST NOT call any sbuf function on the buffer it was passed, with the single
157// exception of sbufError(), which a sink that failed partway through uses to report the failure.
158typedef bool (*sbufSendCB)(_Pre_valid_ StreamBuffer* sb, _In_reads_bytes_(sz) const uint8* buf,
159 size_t off, size_t sz, _Pre_opt_valid_ void* ctx);
160
173typedef void (*sbufNotifyCB)(stvlist* cvars, _Pre_valid_ StreamBuffer* sb, size_t sz);
174
184typedef void (*sbufResumeCB)(stvlist* cvars, _Pre_valid_ StreamBuffer* sb);
185
191
193enum STREAM_BUFFER_FLAGS_ENUM {
194 SBUF_Push = 0x0001,
195 SBUF_Pull = 0x0002,
196 SBUF_Direct = 0x0010,
197 SBUF_Error = 0x0800,
198 SBUF_Closed = 0x1000,
199};
201
206 SBUF_Locked = 0x0004,
207
211 SBUF_Wait = 0x0008,
212};
213
215// Internal state, not for callers.
216enum STREAM_BUFFER_STATE_ENUM {
217 SBUF_PHeld = 0x0020, // producer is held at the high watermark
218 SBUF_PResumeOwed = 0x0040, // a write was refused; owe the producer a resume callback
219};
220
222typedef struct StreamBuffer {
223 BufRing buf;
224 size_t targetsz;
225
226 closure producerPull;
227 struct SbufResume* producerResume;
228
229 closure consumerNotify;
230 closure consumerPush;
231
232 // Closures of roles that unregistered from inside a callback, destroyed once the stack has
233 // unwound back out of the buffer.
234 closure pendingP;
235 closure pendingC;
236
237 size_t high;
238 size_t low;
239
240 Mutex lock;
241 CondVar ready;
242 CondVar flushed;
243 atomic(intptr) owner;
244 uint32 depth;
245
246 int refcount;
247 bool locked;
248 bool resumePending;
249 bool destroyPending;
250 bool walking;
251 atomic(uint32) flags;
252} StreamBuffer;
254
255// Internal function - use sbufCreate() macro instead
256_Ret_valid_ StreamBuffer* _sbufCreate(size_t targetsz, flags_t flags);
257
275#define sbufCreate(targetsz, ...) _sbufCreate(targetsz, opt_flags(__VA_ARGS__))
276
288_Ret_valid_ StreamBuffer* sbufAcquire(_Inout_ StreamBuffer* sb);
289
298_At_(*sb, _Pre_maybenull_ _Post_null_) void sbufRelease(_Inout_ StreamBuffer** sb);
299
311void sbufClose(_Inout_opt_ StreamBuffer* sb);
312
327_At_(*sb, _Pre_maybenull_ _Post_null_) void sbufFinish(_Inout_ StreamBuffer** sb);
328
341void sbufError(_Inout_ StreamBuffer* sb);
342
351void sbufClearError(_Inout_ StreamBuffer* sb);
352
368void sbufSetWatermark(_Inout_ StreamBuffer* sb, size_t high, size_t low);
369
376_meta_inline bool sbufIsLocked(_In_ StreamBuffer* sb)
377{
378 return sb->locked;
379}
380
387_meta_inline bool sbufIsPull(_In_ StreamBuffer* sb)
388{
389 return (atomicLoad(uint32, &sb->flags, Relaxed) & SBUF_Pull) != 0;
390}
391
398_meta_inline bool sbufIsPush(_In_ StreamBuffer* sb)
399{
400 return (atomicLoad(uint32, &sb->flags, Relaxed) & SBUF_Push) != 0;
401}
402
409_meta_inline bool sbufIsError(_In_ StreamBuffer* sb)
410{
411 return (atomicLoad(uint32, &sb->flags, Relaxed) & SBUF_Error) != 0;
412}
413
422_meta_inline bool sbufIsClosed(_In_ StreamBuffer* sb)
423{
424 return (atomicLoad(uint32, &sb->flags, Relaxed) & SBUF_Closed) != 0;
425}
426
438bool sbufCMore(_Inout_ StreamBuffer* sb);
439
441
447
466_Check_return_ bool sbufPRegisterPull(_Inout_ StreamBuffer* sb, _In_ closure ppull);
467
477void sbufPUnregister(_Inout_ StreamBuffer* sb);
478
485bool sbufPAttached(_Inout_ StreamBuffer* sb);
486
503void sbufPSetResume(_Inout_ StreamBuffer* sb, _In_opt_ closure resume);
504
512bool sbufPIsHeld(_Inout_ StreamBuffer* sb);
513
520size_t sbufPAvail(_Inout_ StreamBuffer* sb);
521
522// Internal function - use sbufPWrite() macro instead
523bool _sbufPWrite(_Inout_ StreamBuffer* sb, _In_reads_bytes_(sz) const uint8* buf, size_t sz,
524 flags_t flags);
525
541#define sbufPWrite(sb, buf, sz, ...) _sbufPWrite(sb, buf, sz, opt_flags(__VA_ARGS__))
542
543// Internal function - use sbufPWriteStr() macro instead
544bool _sbufPWriteStr(_Inout_ StreamBuffer* sb, _In_opt_ strref str, flags_t flags);
545
554#define sbufPWriteStr(sb, str, ...) _sbufPWriteStr(sb, str, opt_flags(__VA_ARGS__))
555
556// Internal function - use sbufPWriteLine() macro instead
557bool _sbufPWriteLine(_Inout_ StreamBuffer* sb, _In_opt_ strref str, flags_t flags);
558
569#define sbufPWriteLine(sb, str, ...) _sbufPWriteLine(sb, str, opt_flags(__VA_ARGS__))
570
571// Internal function - use sbufPWriteEOL() macro instead
572bool _sbufPWriteEOL(_Inout_ StreamBuffer* sb, flags_t flags);
573
583#define sbufPWriteEOL(sb, ...) _sbufPWriteEOL(sb, opt_flags(__VA_ARGS__))
584
600bool sbufPFlush(_Inout_ StreamBuffer* sb);
601
603
609
626_Check_return_ bool sbufCRegisterPush(_Inout_ StreamBuffer* sb, _In_ closure cnotify);
627
640_Check_return_ bool sbufCRegisterPushDirect(_Inout_ StreamBuffer* sb, _In_ closure cpush);
641
654void sbufCUnregister(_Inout_ StreamBuffer* sb);
655
662bool sbufCAttached(_Inout_ StreamBuffer* sb);
663
670size_t sbufCAvail(_Inout_ StreamBuffer* sb);
671
686_Success_(return) bool
687sbufCRead(_Inout_ StreamBuffer* sb, _Out_writes_bytes_to_(sz, *bytesread) uint8* buf, size_t sz,
688 _Out_ _Deref_out_range_(0, sz) size_t* bytesread);
689
704_Success_(return > 0) bool
705sbufCPeek(_Inout_ StreamBuffer* sb, _Out_writes_bytes_(sz) uint8* buf, size_t off, size_t sz);
706
718bool sbufCFeed(_Inout_ StreamBuffer* sb, size_t minsz);
719
737bool sbufCSend(_Inout_ StreamBuffer* sb, _In_ sbufSendCB func, size_t sz, _Inout_opt_ void* ctx);
738
748bool sbufCSkip(_Inout_ StreamBuffer* sb, size_t bytes);
749
752
753CX_C_END
Ring buffer implementation for efficient streaming I/O.
Basic closure functionality.
Condition variable synchronization primitive.
bool sbufCFeed(StreamBuffer *sb, size_t minsz)
bool sbufCPeek(StreamBuffer *sb, uint8 *buf, size_t off, size_t sz)
bool sbufCAttached(StreamBuffer *sb)
void sbufCUnregister(StreamBuffer *sb)
size_t sbufCAvail(StreamBuffer *sb)
bool sbufCSkip(StreamBuffer *sb, size_t bytes)
bool sbufCSend(StreamBuffer *sb, sbufSendCB func, size_t sz, void *ctx)
bool sbufCRegisterPushDirect(StreamBuffer *sb, closure cpush)
bool sbufCRead(StreamBuffer *sb, uint8 *buf, size_t sz, size_t *bytesread)
bool sbufCRegisterPush(StreamBuffer *sb, closure cnotify)
STREAM_BUFFER_OPT_FLAGS
Optional flags for sbufCreate() and the sbufPWrite() family.
Definition streambuf.h:203
void sbufClose(StreamBuffer *sb)
void sbufError(StreamBuffer *sb)
void sbufFinish(StreamBuffer **sb)
bool sbufIsLocked(StreamBuffer *sb)
Definition streambuf.h:376
bool sbufIsClosed(StreamBuffer *sb)
Definition streambuf.h:422
void sbufClearError(StreamBuffer *sb)
void sbufSetWatermark(StreamBuffer *sb, size_t high, size_t low)
bool sbufIsPush(StreamBuffer *sb)
Definition streambuf.h:398
StreamBuffer * sbufAcquire(StreamBuffer *sb)
void sbufRelease(StreamBuffer **sb)
bool sbufIsPull(StreamBuffer *sb)
Definition streambuf.h:387
bool sbufCMore(StreamBuffer *sb)
bool sbufIsError(StreamBuffer *sb)
Definition streambuf.h:409
@ SBUF_Locked
Definition streambuf.h:206
@ SBUF_Wait
Definition streambuf.h:211
void sbufPUnregister(StreamBuffer *sb)
size_t sbufPAvail(StreamBuffer *sb)
bool sbufPAttached(StreamBuffer *sb)
void sbufPSetResume(StreamBuffer *sb, closure resume)
bool sbufPIsHeld(StreamBuffer *sb)
bool sbufPFlush(StreamBuffer *sb)
bool sbufPRegisterPull(StreamBuffer *sb, closure ppull)
size_t(* sbufPullCB)(stvlist *cvars, StreamBuffer *sb, uint8 *buf, size_t sz)
Definition streambuf.h:132
void(* sbufResumeCB)(stvlist *cvars, StreamBuffer *sb)
Definition streambuf.h:184
void(* sbufNotifyCB)(stvlist *cvars, StreamBuffer *sb, size_t sz)
Definition streambuf.h:173
void(* sbufPushCB)(stvlist *cvars, StreamBuffer *sb, const uint8 *buf, size_t sz)
Definition streambuf.h:144
Mutex synchronization primitive.
Definition mutex.h:60
Runtime type system and type descriptor infrastructure.