Changeset: 58028b589089 for MonetDB
URL: https://dev.monetdb.org/hg/MonetDB?cmd=changeset;node=58028b589089
Modified Files:
common/stream/bs.c
common/stream/bs2.c
common/stream/iconv_stream.c
common/stream/memio.c
common/stream/pump.c
common/stream/stdio_stream.c
common/stream/stream.c
common/stream/stream_internal.h
Branch: makelibstreamgreatagain
Log Message:
Introduce flush levels in stream internals
pump.c maps MNSTR_FLUSH_ALL to PUMP_FLUSH_ALL which is picked up by the
compressed streams.
diffs (239 lines):
diff --git a/common/stream/bs.c b/common/stream/bs.c
--- a/common/stream/bs.c
+++ b/common/stream/bs.c
@@ -103,7 +103,7 @@ bs_write(stream *restrict ss, const void
* flushed.
*/
static int
-bs_flush(stream *ss)
+bs_flush(stream *ss, mnstr_flush_level flush_level)
{
uint16_t blksize;
bs *s;
@@ -135,7 +135,7 @@ bs_flush(stream *ss)
/* indicate that this is the last buffer of a block by
* setting the low-order bit */
blksize |= 1;
- /* allways flush (even empty blocks) needed for the protocol) */
+ /* always flush (even empty blocks) needed for the protocol) */
if ((!mnstr_writeSht(ss->inner, (int16_t) blksize) ||
(s->nr > 0 &&
ss->inner->write(ss->inner, s->buf, 1, s->nr) !=
(ssize_t) s->nr))) {
@@ -143,6 +143,8 @@ bs_flush(stream *ss)
s->nr = 0; /* data is lost due to error */
return -1;
}
+ // shouldn't we flush ss->inner too?
+ (void) flush_level;
s->blks++;
s->nr = 0;
}
@@ -300,7 +302,7 @@ bs_close(stream *ss)
if (s == NULL)
return;
if (!ss->readonly && s->nr > 0)
- bs_flush(ss);
+ bs_flush(ss, MNSTR_FLUSH_DATA);
mnstr_close(ss->inner);
}
diff --git a/common/stream/bs2.c b/common/stream/bs2.c
--- a/common/stream/bs2.c
+++ b/common/stream/bs2.c
@@ -234,7 +234,7 @@ bs2_write(stream *restrict ss, const voi
* flushed.
*/
static int
-bs2_flush(stream *ss)
+bs2_flush(stream *ss, mnstr_flush_level flush_level)
{
int64_t blksize;
bs2 *s;
@@ -291,6 +291,8 @@ bs2_flush(stream *ss)
return -1;
}
s->nr = 0;
+ // shouldn't we flush s->s too?
+ (void) flush_level;
}
return 0;
}
@@ -558,7 +560,7 @@ bs2_close(stream *ss)
if (s == NULL)
return;
if (!ss->readonly && s->nr > 0)
- bs2_flush(ss);
+ bs2_flush(ss, MNSTR_FLUSH_DATA);
assert(s->s);
if (s->s)
s->s->close(s->s);
diff --git a/common/stream/iconv_stream.c b/common/stream/iconv_stream.c
--- a/common/stream/iconv_stream.c
+++ b/common/stream/iconv_stream.c
@@ -202,7 +202,7 @@ ic_read(stream *restrict s, void *restri
}
static int
-ic_flush(stream *s)
+ic_flush(stream *s, mnstr_flush_level flush_level)
{
struct icstream *ic = (struct icstream *) s->stream_data.p;
char *outbuf;
@@ -221,7 +221,7 @@ ic_flush(stream *s)
mnstr_copy_error(s, s->inner);
return -1;
}
- return mnstr_flush(s->inner, MNSTR_FLUSH_DATA);
+ return mnstr_flush(s->inner, flush_level);
}
static void
@@ -231,7 +231,7 @@ ic_close(stream *s)
if (ic) {
if (!s->readonly)
- ic_flush(s);
+ ic_flush(s, MNSTR_FLUSH_DATA);
iconv_close(ic->cd);
close_stream(s->inner);
free(s->stream_data.p);
diff --git a/common/stream/memio.c b/common/stream/memio.c
--- a/common/stream/memio.c
+++ b/common/stream/memio.c
@@ -135,7 +135,7 @@ buffer_close(stream *s)
}
static int
-buffer_flush(stream *s)
+buffer_flush(stream *s, mnstr_flush_level flush_level)
{
buffer *b;
@@ -144,6 +144,7 @@ buffer_flush(stream *s)
if (b == NULL)
return -1;
b->pos = 0;
+ (void) flush_level;
return 0;
}
diff --git a/common/stream/pump.c b/common/stream/pump.c
--- a/common/stream/pump.c
+++ b/common/stream/pump.c
@@ -20,7 +20,7 @@ static pump_result pump_out(stream *rest
static ssize_t pump_read(stream *restrict s, void *restrict buf, size_t
elmsize, size_t cnt);
static ssize_t pump_write(stream *restrict s, const void *restrict buf, size_t
elmsize, size_t cnt);
-static int pump_flush(stream *s);
+static int pump_flush(stream *s, mnstr_flush_level flush_level);
static void pump_close(stream *s);
static void pump_destroy(stream *s);
@@ -110,17 +110,31 @@ pump_write(stream *restrict s, const voi
}
-static int pump_flush(stream *s)
+static int pump_flush(stream *s, mnstr_flush_level flush_level)
{
pump_state *state = (pump_state*) s->stream_data.p;
inner_state_t *inner_state = state->inner_state;
+ pump_action action;
+
+ switch (flush_level) {
+ case MNSTR_FLUSH_DATA:
+ action = PUMP_FLUSH_DATA;
+ break;
+ case MNSTR_FLUSH_ALL:
+ action = PUMP_FLUSH_ALL;
+ break;
+ default:
+ assert(0 /* unknown flush_level */);
+ action = PUMP_FLUSH_DATA;
+ break;
+ }
state->set_src_win(inner_state, (pump_buffer){ .start = NULL, .count =
0 });
- ssize_t nwritten = pump_out(s, PUMP_FLUSH_DATA);
+ ssize_t nwritten = pump_out(s, action);
if (nwritten < 0)
return -1;
else
- return mnstr_flush(s->inner, MNSTR_FLUSH_DATA);
+ return mnstr_flush(s->inner, action);
}
diff --git a/common/stream/stdio_stream.c b/common/stream/stdio_stream.c
--- a/common/stream/stdio_stream.c
+++ b/common/stream/stdio_stream.c
@@ -102,7 +102,7 @@ file_clrerr(stream *s)
static int
-file_flush(stream *s)
+file_flush(stream *s, mnstr_flush_level flush_level)
{
FILE *fp = (FILE *) s->stream_data.p;
@@ -110,6 +110,7 @@ file_flush(stream *s)
mnstr_set_error_errno(s, MNSTR_WRITE_ERROR, "flush
error");
return -1;
}
+ (void) flush_level;
return 0;
}
@@ -414,7 +415,7 @@ stdin_rastream(void)
#endif
// Make an attempt to skip a BOM marker.
// It would be nice to integrate this with with the BOM removal code
- // in text_stream.c but that is complicated. In text_stream,
+ // in text_stream.c but that is complicated. In text_stream,
do {
struct stat stb;
if (fstat(fileno(stdin), &stb) < 0)
diff --git a/common/stream/stream.c b/common/stream/stream.c
--- a/common/stream/stream.c
+++ b/common/stream/stream.c
@@ -492,7 +492,6 @@ mnstr_error_kind_description(mnstr_error
int
mnstr_flush(stream *s, mnstr_flush_level flush_level)
{
- (void) flush_level;
if (s == NULL)
return -1;
#ifdef STREAM_DEBUG
@@ -502,7 +501,7 @@ mnstr_flush(stream *s, mnstr_flush_level
if (s->errkind != MNSTR_NO__ERROR)
return -1;
if (s->flush)
- return s->flush(s);
+ return s->flush(s, flush_level);
return 0;
}
@@ -754,9 +753,9 @@ wrapper_destroy(stream *s)
static int
-wrapper_flush(stream *s)
+wrapper_flush(stream *s, mnstr_flush_level flush_level)
{
- return s->inner->flush(s->inner);
+ return s->inner->flush(s->inner, flush_level);
}
diff --git a/common/stream/stream_internal.h b/common/stream/stream_internal.h
--- a/common/stream/stream_internal.h
+++ b/common/stream/stream_internal.h
@@ -163,7 +163,7 @@ struct stream {
void (*close)(stream *s);
void (*clrerr)(stream *s);
void (*destroy)(stream *s);
- int (*flush)(stream *s);
+ int (*flush)(stream *s, mnstr_flush_level flush_level);
int (*fsync)(stream *s);
int (*fgetpos)(stream *restrict s, fpos_t *restrict p);
int (*fsetpos)(stream *restrict s, fpos_t *restrict p);
_______________________________________________
checkin-list mailing list
[email protected]
https://www.monetdb.org/mailman/listinfo/checkin-list