From a409792fb548027b77ad547b6ffa397eb24d190a Mon Sep 17 00:00:00 2001 From: Erica Fischer Date: Wed, 25 Jan 2023 13:09:34 -0800 Subject: [PATCH] Track file position within fwrite_check() --- main.cpp | 21 +++++++++-------- serial.cpp | 67 +++++++++++++++++++++++++++++++++++++++++++++--------- serial.hpp | 2 +- 3 files changed, 68 insertions(+), 22 deletions(-) diff --git a/main.cpp b/main.cpp index 09e831bc..08af49de 100644 --- a/main.cpp +++ b/main.cpp @@ -331,8 +331,7 @@ static void merge(struct mergelist *merges, size_t nmerges, unsigned char *map, // MAGIC: This knows that the feature minzoom is the last byte of the serialized feature // and is writing one byte less and then adding the byte for the minzoom. - fwrite_check(geom_map + ix.start, 1, ix.end - ix.start - 1, geom_out, "merge geometry"); - *geompos += ix.end - ix.start - 1; + fwrite_check(geom_map + ix.start, 1, ix.end - ix.start - 1, geom_out, geompos, "merge geometry"); int feature_minzoom = calc_feature_minzoom(&ix, ds, maxzoom, gamma); serialize_byte(geom_out, feature_minzoom, geompos, "merge geometry"); @@ -346,7 +345,8 @@ static void merge(struct mergelist *merges, size_t nmerges, unsigned char *map, ix.start = pos; ix.end = *geompos; - fwrite_check(&ix, bytes, 1, indexfile, "merge temporary"); + std::atomic indexpos; + fwrite_check(&ix, bytes, 1, indexfile, &indexpos, "merge temporary"); head->start += bytes; struct mergelist *m = head; @@ -787,8 +787,7 @@ void radix1(int *geomfds_in, int *indexfds_in, int inputs, int prefix, int split unsigned long long which = (ix.ix << prefix) >> (64 - splitbits); long long pos = sub_geompos[which]; - fwrite_check(geommap + ix.start, ix.end - ix.start, 1, geomfiles[which], "geom"); - sub_geompos[which] += ix.end - ix.start; + fwrite_check(geommap + ix.start, ix.end - ix.start, 1, geomfiles[which], &sub_geompos[which], "geom"); // Count this as a 25%-accomplishment, since we will copy again *progress += (ix.end - ix.start) / 4; @@ -801,7 +800,8 @@ void radix1(int *geomfds_in, int *indexfds_in, int inputs, int prefix, int split ix.start = pos; ix.end = sub_geompos[which]; - fwrite_check(&ix, sizeof(struct index), 1, indexfiles[which], "index"); + std::atomic indexpos; + fwrite_check(&ix, sizeof(struct index), 1, indexfiles[which], &indexpos, "index"); } madvise(indexmap, indexst.st_size, MADV_DONTNEED); @@ -958,8 +958,7 @@ void radix1(int *geomfds_in, int *indexfds_in, int inputs, int prefix, int split struct index ix = indexmap[a]; long long pos = *geompos_out; - fwrite_check(geommap + ix.start, ix.end - ix.start, 1, geomfile, "geom"); - *geompos_out += ix.end - ix.start; + fwrite_check(geommap + ix.start, ix.end - ix.start, 1, geomfile, geompos_out, "geom"); int feature_minzoom = calc_feature_minzoom(&ix, ds, maxzoom, gamma); serialize_byte(geomfile, feature_minzoom, geompos_out, "merge geometry"); @@ -973,7 +972,8 @@ void radix1(int *geomfds_in, int *indexfds_in, int inputs, int prefix, int split ix.start = pos; ix.end = *geompos_out; - fwrite_check(&ix, sizeof(struct index), 1, indexfile, "index"); + std::atomic indexpos; + fwrite_check(&ix, sizeof(struct index), 1, indexfile, &indexpos, "index"); } madvise(indexmap, indexst.st_size, MADV_DONTNEED); @@ -1665,7 +1665,8 @@ std::pair read_input(std::vector &sources, char *fname, i int n; while ((n = fp->read(buf, READ_BUF)) > 0) { - fwrite_check(buf, sizeof(char), n, readfp, reading.c_str()); + std::atomic readingpos; + fwrite_check(buf, sizeof(char), n, readfp, &readingpos, reading.c_str()); ahead += n; if (buf[n - 1] == read_parallel_this && ahead > PARSE_MIN) { diff --git a/serial.cpp b/serial.cpp index 5e40c145..4917a58f 100644 --- a/serial.cpp +++ b/serial.cpp @@ -45,8 +45,11 @@ struct zwriter { int fwrite(void *p, size_t size, size_t nmemb, bool flush) { buf += std::string((char *) p, size * nmemb); + bool again = false; + + while (buf.size() != 0 || again) { + again = false; - while (true) { zstream.next_in = (Bytef *) buf.c_str(); zstream.avail_in = buf.size(); @@ -66,23 +69,67 @@ struct zwriter { // it produced some output ssize_t produced = zstream.next_out - (Bytef *) zbuf.c_str(); - fwrite((void *) zbuf.c_str(), sizeof(char), produced, fp); + ::fwrite((void *) zbuf.c_str(), sizeof(char), produced, fp); } if (d == Z_BUF_ERROR) { zbuf.resize(zbuf.size() + 5000); + again = true; } else if (d != Z_OK) { fprintf(stderr, "impossible return %d from deflate()\n", d); exit(EXIT_IMPOSSIBLE); } - if (buf.size() == 0) { - break; + if (zstream.avail_out == 0) { + again = true; } } return nmemb; } + + int putc(int c) { + char ch = c; + return fwrite(&ch, sizeof(char), 1, false); + } + + int fclose() { + while (true) { + zstream.next_in = (Bytef *) buf.c_str(); + zstream.avail_in = 0; + + zstream.next_out = (Bytef *) zbuf.c_str(); + zstream.avail_out = zbuf.size(); + + int d = deflate(&zstream, Z_FINISH); + + if (zstream.next_out != (Bytef *) zbuf.c_str()) { + // it produced some output + + ssize_t produced = zstream.next_out - (Bytef *) zbuf.c_str(); + fwrite((void *) zbuf.c_str(), sizeof(char), produced, fp); + } + + if (d == Z_BUF_ERROR || d == Z_OK) { + // still more to be written + zbuf.resize(zbuf.size() + 5000); + continue; + } else if (d == Z_STREAM_END) { + // done + break; + } else { + fprintf(stderr, "impossible return %d from deflate()\n", d); + exit(EXIT_IMPOSSIBLE); + } + } + + if (deflateEnd(&zstream) != Z_OK) { + fprintf(stderr, "Error closing compression scheme (%s)\n", zstream.msg); + exit(EXIT_WRITE); + } + + return ::fclose(fp); + } }; // Offset coordinates to keep them positive @@ -92,12 +139,13 @@ struct zwriter { // write to file -size_t fwrite_check(const void *ptr, size_t size, size_t nitems, FILE *stream, const char *fname) { +size_t fwrite_check(const void *ptr, size_t size, size_t nitems, FILE *stream, std::atomic *fpos, const char *fname) { size_t w = fwrite(ptr, size, nitems, stream); if (w != nitems) { fprintf(stderr, "%s: Write to temporary file failed: %s\n", fname, strerror(errno)); exit(EXIT_WRITE); } + *fpos += size * nitems; return w; } @@ -134,8 +182,7 @@ void serialize_ulong_long(FILE *out, unsigned long long zigzag, std::atomic *fpos, const char *fname) { - fwrite_check(&n, sizeof(signed char), 1, out, fname); - *fpos += sizeof(signed char); + fwrite_check(&n, sizeof(signed char), 1, out, fpos, fname); } void serialize_uint(FILE *out, unsigned n, std::atomic *fpos, const char *fname) { @@ -361,8 +408,7 @@ void serialize_feature(FILE *geomfile, serial_feature *sf, std::atomicfeature_minzoom); serialize_long_long(geomfile, s.size(), geompos, fname); - fwrite_check(s.c_str(), sizeof(char), s.size(), geomfile, fname); - *geompos += s.size(); + fwrite_check(s.c_str(), sizeof(char), s.size(), geomfile, geompos, fname); } serial_feature deserialize_feature(FILE *geoms, std::atomic *geompos_in, unsigned z, unsigned tx, unsigned ty, unsigned *initial_x, unsigned *initial_y) { @@ -842,8 +888,7 @@ int serialize_feature(struct serialization_state *sst, serial_feature &sf) { index.t = sf.t; index.ix = bbox_index; - fwrite_check(&index, sizeof(struct index), 1, r->indexfile, sst->fname); - r->indexpos += sizeof(struct index); + fwrite_check(&index, sizeof(struct index), 1, r->indexfile, &r->indexpos, sst->fname); for (size_t i = 0; i < 2; i++) { if (sf.bbox[i] < r->file_bbox[i]) { diff --git a/serial.hpp b/serial.hpp index 25396188..d708cbee 100644 --- a/serial.hpp +++ b/serial.hpp @@ -11,7 +11,7 @@ #include "mbtiles.hpp" #include "jsonpull/jsonpull.h" -size_t fwrite_check(const void *ptr, size_t size, size_t nitems, FILE *stream, const char *fname); +size_t fwrite_check(const void *ptr, size_t size, size_t nitems, FILE *stream, std::atomic *fpos, const char *fname); void serialize_int(FILE *out, int n, std::atomic *fpos, const char *fname); void serialize_long_long(FILE *out, long long n, std::atomic *fpos, const char *fname);