diff --git a/serial.cpp b/serial.cpp index c3193b17..a07be331 100644 --- a/serial.cpp +++ b/serial.cpp @@ -24,114 +24,6 @@ #include "milo/dtoa_milo.h" #include "errors.hpp" -struct zwriter { - z_stream zstream; - std::string buf; - std::string zbuf; - FILE *fp; - - zwriter(FILE *out) { - zstream.zalloc = NULL; - zstream.zfree = NULL; - zstream.opaque = NULL; - if (deflateInit(&zstream, Z_DEFAULT_COMPRESSION) != Z_OK) { - fprintf(stderr, "Compression initialization failed\n"); - exit(EXIT_MEMORY); - } - - zbuf.resize(5000); - fp = out; - } - - int fwrite_check(const void *p, size_t size, size_t nmemb, bool flush, std::atomic *fpos, const char *fname) { - buf += std::string((char *) p, size * nmemb); - bool again = false; - - while (buf.size() != 0 || again) { - again = false; - - zstream.next_in = (Bytef *) buf.c_str(); - zstream.avail_in = buf.size(); - - zstream.next_out = (Bytef *) zbuf.c_str(); - zstream.avail_out = zbuf.size(); - - int d = deflate(&zstream, flush ? Z_FULL_FLUSH : Z_NO_FLUSH); - - if (zstream.next_in != (Bytef *) buf.c_str()) { - // it consumed some input - - ssize_t consumed = zstream.next_in - (Bytef *) buf.c_str(); - buf.erase(buf.begin(), buf.begin() + consumed); - } - - if (zstream.next_out != (Bytef *) zbuf.c_str()) { - // it produced some output - - ssize_t produced = zstream.next_out - (Bytef *) zbuf.c_str(); - ::fwrite_check((void *) zbuf.c_str(), sizeof(char), produced, fp, fpos, fname); - } - - if (d == Z_BUF_ERROR) { - zbuf.resize(zbuf.size() + 5000); - again = true; - } else if (d != Z_OK) { - fprintf(stderr, "impossible return %d from deflate() in fclose()\n", d); - exit(EXIT_IMPOSSIBLE); - } - - if (zstream.avail_out == 0) { - again = true; - } - } - - return nmemb; - } - - int putc(int c, std::atomic *fpos, const char *fname) { - char ch = c; - return fwrite_check(&ch, sizeof(char), 1, false, fpos, fname); - } - - int fclose(std::atomic *fpos, const char *fname) { - 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_check((void *) zbuf.c_str(), sizeof(char), produced, fp, fpos, fname); - } - - 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 #define COORD_OFFSET (4LL << 32) #define SHIFT_RIGHT(a) ((((a) + COORD_OFFSET) >> geometry_scale) - (COORD_OFFSET >> geometry_scale)) @@ -189,45 +81,6 @@ void serialize_uint(FILE *out, unsigned n, std::atomic *fpos, const c serialize_ulong_long(out, n, fpos, fname); } -// write to compression - -size_t fwrite_check(const void *ptr, size_t size, size_t nitems, zwriter *stream, std::atomic *fpos, const char *fname) { - size_t w = stream->fwrite_check(ptr, size, nitems, false, fpos, fname); - return w; -} - -void serialize_ulong_long(zwriter *out, unsigned long long zigzag, std::atomic *fpos, const char *fname) { - while (1) { - unsigned char b = zigzag & 0x7F; - if ((zigzag >> 7) != 0) { - b |= 0x80; - out->putc(b, fpos, fname); - zigzag >>= 7; - } else { - out->putc(b, fpos, fname); - break; - } - } -} - -void serialize_long_long(zwriter *out, long long n, std::atomic *fpos, const char *fname) { - unsigned long long zigzag = protozero::encode_zigzag64(n); - - serialize_ulong_long(out, zigzag, fpos, fname); -} - -void serialize_int(zwriter *out, int n, std::atomic *fpos, const char *fname) { - serialize_long_long(out, n, fpos, fname); -} - -void serialize_byte(zwriter *out, signed char n, std::atomic *fpos, const char *fname) { - out->fwrite_check(&n, sizeof(signed char), 1, false, fpos, fname); -} - -void serialize_uint(zwriter *out, unsigned n, std::atomic *fpos, const char *fname) { - serialize_ulong_long(out, n, fpos, fname); -} - // write to memory size_t fwrite_check(const void *ptr, size_t size, size_t nitems, std::string &stream) { diff --git a/tile.cpp b/tile.cpp index 9d6372db..8d3790f6 100644 --- a/tile.cpp +++ b/tile.cpp @@ -25,6 +25,7 @@ #include #include #include +#include #include #include "mvt.hpp" #include "mbtiles.hpp" @@ -40,6 +41,7 @@ #include "milo/dtoa_milo.h" #include "evaluator.hpp" #include "errors.hpp" +#include "protozero/varint.hpp" extern "C" { #include "jsonpull/jsonpull.h" @@ -47,6 +49,103 @@ extern "C" { #include "plugin.hpp" +struct compressor { + FILE *fp = NULL; + z_stream zs; + + compressor(FILE *f) { + fp = f; + } + + compressor() { } + + void begin() { + zs.zalloc = NULL; + zs.zfree = NULL; + zs.opaque = NULL; + + if (deflateInit(&zs, Z_DEFAULT_COMPRESSION) != Z_OK) { + fprintf(stderr, "initialize compression: %s\n", zs.msg); + exit(EXIT_IMPOSSIBLE); + } + } + + void end(std::atomic *fpos, const char *fname) { + std::string buf; + buf.resize(5000); + + zs.next_out = (Bytef *) buf.c_str(); + zs.avail_out = buf.size(); + + zs.next_in = zs.next_out; + zs.avail_in = 0; + + if (deflate(&zs, Z_FINISH) != Z_STREAM_END) { + fprintf(stderr, "finish compression: %s\n", zs.msg); + exit(EXIT_IMPOSSIBLE); + } + + ::fwrite_check(buf.c_str(), sizeof(char), zs.next_out - (Bytef *) buf.c_str(), fp, fpos, fname); + + if (deflateEnd(&zs) != Z_OK) { + fprintf(stderr, "end compression: %s\n", zs.msg); + exit(EXIT_IMPOSSIBLE); + } + } + + int fclose() { + return ::fclose(fp); + } + + void fwrite_check(const char *p, size_t size, size_t nmemb, std::atomic *fpos, const char *fname) { + std::string buf; + buf.resize(size * nmemb * 2 + 200); + + zs.next_in = (Bytef *) p; + zs.avail_in = size * nmemb; + + while (zs.avail_in > 0) { + zs.next_out = (Bytef *) buf.c_str(); + zs.avail_out = buf.size(); + + if (deflate(&zs, Z_NO_FLUSH) != Z_OK) { + fprintf(stderr, "finish compression: %s\n", zs.msg); + exit(EXIT_IMPOSSIBLE); + } + + ::fwrite_check(buf.c_str(), sizeof(char), zs.next_out - (Bytef *) buf.c_str(), fp, fpos, fname); + } + } + + void serialize_ulong_long(unsigned long long val, std::atomic *fpos, const char *fname) { + while (1) { + unsigned char b = val & 0x7F; + if ((val >> 7) != 0) { + b |= 0x80; + fwrite_check((const char *) &b, 1, 1, fpos, fname); + val >>= 7; + } else { + fwrite_check((const char *) &b, 1, 1, fpos, fname); + break; + } + } + } + + void serialize_long_long(long long val, std::atomic *fpos, const char *fname) { + unsigned long long zigzag = protozero::encode_zigzag64(val); + + serialize_ulong_long(zigzag, fpos, fname); + } + + void serialize_int(int val, std::atomic *fpos, const char *fname) { + serialize_long_long(val, fpos, fname); + } + + void serialize_uint(unsigned val, std::atomic *fpos, const char *fname) { + serialize_ulong_long(val, fpos, fname); + } +}; + #define CMD_BITS 3 // Offset coordinates to keep them positive @@ -324,7 +423,7 @@ struct ordercmp { } } ordercmp; -void rewrite(drawvec &geom, int z, int nextzoom, int maxzoom, long long *bbox, unsigned tx, unsigned ty, int buffer, int *within, std::atomic *geompos, FILE **geomfile, const char *fname, signed char t, int layer, signed char feature_minzoom, int child_shards, int max_zoom_increment, long long seq, int tippecanoe_minzoom, int tippecanoe_maxzoom, int segment, unsigned *initial_x, unsigned *initial_y, std::vector &metakeys, std::vector &metavals, bool has_id, unsigned long long id, unsigned long long index, unsigned long long label_point, long long extent) { +void rewrite(drawvec &geom, int z, int nextzoom, int maxzoom, long long *bbox, unsigned tx, unsigned ty, int buffer, int *within, std::atomic *geompos, compressor **geomfile, const char *fname, signed char t, int layer, signed char feature_minzoom, int child_shards, int max_zoom_increment, long long seq, int tippecanoe_minzoom, int tippecanoe_maxzoom, int segment, unsigned *initial_x, unsigned *initial_y, std::vector &metakeys, std::vector &metavals, bool has_id, unsigned long long id, unsigned long long index, unsigned long long label_point, long long extent) { if (geom.size() > 0 && (nextzoom <= maxzoom || additional[A_EXTEND_ZOOMS])) { int xo, yo; int span = 1 << (nextzoom - z); @@ -394,9 +493,10 @@ void rewrite(drawvec &geom, int z, int nextzoom, int maxzoom, long long *bbox, u { if (!within[j]) { - serialize_int(geomfile[j], nextzoom, &geompos[j], fname); - serialize_uint(geomfile[j], tx * span + xo, &geompos[j], fname); - serialize_uint(geomfile[j], ty * span + yo, &geompos[j], fname); + geomfile[j]->begin(); + geomfile[j]->serialize_int(nextzoom, &geompos[j], fname); + geomfile[j]->serialize_uint(tx * span + xo, &geompos[j], fname); + geomfile[j]->serialize_uint(ty * span + yo, &geompos[j], fname); within[j] = 1; } @@ -423,8 +523,8 @@ void rewrite(drawvec &geom, int z, int nextzoom, int maxzoom, long long *bbox, u } std::string feature = serialize_feature(&sf, SHIFT_RIGHT(initial_x[segment]), SHIFT_RIGHT(initial_y[segment])); - serialize_long_long(geomfile[j], feature.size(), &geompos[j], fname); - fwrite_check(feature.c_str(), sizeof(char), feature.size(), geomfile[j], &geompos[j], fname); + geomfile[j]->serialize_long_long(feature.size(), &geompos[j], fname); + geomfile[j]->fwrite_check(feature.c_str(), sizeof(char), feature.size(), &geompos[j], fname); } } } @@ -1288,7 +1388,7 @@ struct write_tile_args { const char *outdir = NULL; int buffer = 0; const char *fname = NULL; - FILE **geomfile = NULL; + compressor **geomfile = NULL; double todo = 0; std::atomic *along = NULL; double gamma = 0; @@ -1426,7 +1526,7 @@ void remove_attributes(serial_feature &sf, std::set const &exclude_ } } -serial_feature next_feature(FILE *geoms, std::atomic *geompos_in, int z, unsigned tx, unsigned ty, unsigned *initial_x, unsigned *initial_y, long long *original_features, long long *unclipped_features, int nextzoom, int maxzoom, int minzoom, int max_zoom_increment, size_t pass, std::atomic *along, long long alongminus, int buffer, int *within, FILE **geomfile, std::atomic *geompos, std::atomic *oprogress, double todo, const char *fname, int child_shards, struct json_object *filter, const char *stringpool, long long *pool_off, std::vector> *layer_unmaps, bool first_time) { +serial_feature next_feature(FILE *geoms, std::atomic *geompos_in, int z, unsigned tx, unsigned ty, unsigned *initial_x, unsigned *initial_y, long long *original_features, long long *unclipped_features, int nextzoom, int maxzoom, int minzoom, int max_zoom_increment, size_t pass, std::atomic *along, long long alongminus, int buffer, int *within, compressor **geomfile, std::atomic *geompos, std::atomic *oprogress, double todo, const char *fname, int child_shards, struct json_object *filter, const char *stringpool, long long *pool_off, std::vector> *layer_unmaps, bool first_time) { while (1) { serial_feature sf = deserialize_feature(geoms, geompos_in, z, tx, ty, initial_x, initial_y); if (sf.t < 0) { @@ -1577,7 +1677,7 @@ struct run_prefilter_args { long long alongminus = 0; int buffer = 0; int *within = NULL; - FILE **geomfile = NULL; + compressor **geomfile = NULL; std::atomic *geompos = NULL; std::atomic *oprogress = NULL; double todo = 0; @@ -1846,7 +1946,7 @@ void add_sample_to(std::vector &vals, T val, size_t &increment, size_t seq) { } } -long long write_tile(FILE *geoms, std::atomic *geompos_in, char *stringpool, int z, const unsigned tx, const unsigned ty, const int detail, int min_detail, sqlite3 *outdb, const char *outdir, int buffer, const char *fname, FILE **geomfile, int minzoom, int maxzoom, double todo, std::atomic *along, long long alongminus, double gamma, int child_shards, long long *pool_off, unsigned *initial_x, unsigned *initial_y, std::atomic *running, double simplification, std::vector> *layermaps, std::vector> *layer_unmaps, size_t tiling_seg, size_t pass, unsigned long long mingap, long long minextent, double fraction, const char *prefilter, const char *postfilter, struct json_object *filter, write_tile_args *arg, atomic_strategy *strategy) { +long long write_tile(FILE *geoms, std::atomic *geompos_in, char *stringpool, int z, const unsigned tx, const unsigned ty, const int detail, int min_detail, sqlite3 *outdb, const char *outdir, int buffer, const char *fname, compressor **geomfile, int minzoom, int maxzoom, double todo, std::atomic *along, long long alongminus, double gamma, int child_shards, long long *pool_off, unsigned *initial_x, unsigned *initial_y, std::atomic *running, double simplification, std::vector> *layermaps, std::vector> *layer_unmaps, size_t tiling_seg, size_t pass, unsigned long long mingap, long long minextent, double fraction, const char *prefilter, const char *postfilter, struct json_object *filter, write_tile_args *arg, atomic_strategy *strategy) { double merge_fraction = 1; double mingap_fraction = 1; double minextent_fraction = 1; @@ -2402,7 +2502,8 @@ long long write_tile(FILE *geoms, std::atomic *geompos_in, char *stri int j; for (j = 0; j < child_shards; j++) { if (within[j]) { - serialize_ulong_long(geomfile[j], 0, &geompos[j], fname); // EOF + geomfile[j]->serialize_ulong_long(0, &geompos[j], fname); // EOF + geomfile[j]->end(&geompos[j], fname); within[j] = 0; } } @@ -2925,7 +3026,8 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *stringpool, std::atomic< for (z = 0; z <= maxzoom; z++) { std::atomic most(0); - FILE *sub[TEMP_FILES]; + compressor compressors[TEMP_FILES]; + compressor *sub[TEMP_FILES]; int subfd[TEMP_FILES]; for (size_t j = 0; j < TEMP_FILES; j++) { char geomname[strlen(tmpdir) + strlen("/geom.XXXXXXXX" XSTRINGIFY(INT_MAX)) + 1]; @@ -2936,11 +3038,13 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *stringpool, std::atomic< perror(geomname); exit(EXIT_OPEN); } - sub[j] = fopen_oflag(geomname, "wb", O_WRONLY | O_CLOEXEC); - if (sub[j] == NULL) { + FILE *fp = fopen_oflag(geomname, "wb", O_WRONLY | O_CLOEXEC); + if (fp == NULL) { perror(geomname); exit(EXIT_OPEN); } + compressors[j] = compressor(fp); + sub[j] = &compressors[j]; unlink(geomname); } @@ -3175,7 +3279,7 @@ int traverse_zooms(int *geomfd, off_t *geom_size, char *stringpool, std::atomic< exit(EXIT_CLOSE); } } - if (fclose(sub[j]) != 0) { + if (sub[j]->fclose() != 0) { perror("close subfile"); exit(EXIT_CLOSE); }