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

Reply via email to