-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathpbf-reader.js
More file actions
248 lines (229 loc) · 6.01 KB
/
Copy pathpbf-reader.js
File metadata and controls
248 lines (229 loc) · 6.01 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
'use strict';
/**
* Streaming OSM PBF blob reader (fileformat layer).
* Yields { type, data, fileOffset, blobIndex, headerSize, blobSize }.
*/
var fs = require('fs');
var zlib = require('zlib');
var fileformat = require('./fileformat.proto.js');
function readUInt32BE(buf, offset) {
return buf.readUInt32BE(offset);
}
/**
* Read next blob from fd starting at fileOffset.
* @returns {object|null} null at EOF
*/
function readNextBlob(fd, fileOffset, fileSize) {
if (fileOffset + 4 > fileSize) return null;
var sizeBuf = Buffer.alloc(4);
var n = fs.readSync(fd, sizeBuf, 0, 4, fileOffset);
if (n < 4) return null;
var headerLen = readUInt32BE(sizeBuf, 0);
if (headerLen <= 0 || headerLen > 64 * 1024) {
throw new Error(
'Invalid BlobHeader size ' + headerLen + ' at offset ' + fileOffset
);
}
var headerStart = fileOffset + 4;
if (headerStart + headerLen > fileSize) {
throw new Error('Truncated BlobHeader at offset ' + fileOffset);
}
var headerBuf = Buffer.alloc(headerLen);
fs.readSync(fd, headerBuf, 0, headerLen, headerStart);
var header = fileformat.readBlobHeader(headerBuf);
var dataSize = header.datasize | 0;
if (dataSize < 0 || dataSize > 32 * 1024 * 1024) {
throw new Error(
'Invalid Blob datasize ' + dataSize + ' type=' + header.type
);
}
var blobStart = headerStart + headerLen;
if (blobStart + dataSize > fileSize) {
throw new Error('Truncated Blob at offset ' + blobStart);
}
var blobBuf = Buffer.alloc(dataSize);
fs.readSync(fd, blobBuf, 0, dataSize, blobStart);
var blob = fileformat.readBlob(blobBuf);
var data;
if (blob.raw && blob.raw.length) {
data = Buffer.from(blob.raw);
} else if (blob.zlib_data && blob.zlib_data.length) {
data = zlib.inflateSync(Buffer.from(blob.zlib_data));
} else if (blob.lz4_data || blob.zstd_data || blob.lzma_data) {
throw new Error(
'Unsupported blob compression for type ' +
header.type +
' (only raw/zlib supported)'
);
} else {
throw new Error('Empty blob data for type ' + header.type);
}
var nextOffset = blobStart + dataSize;
return {
type: header.type,
data: data,
fileOffset: fileOffset,
nextOffset: nextOffset,
headerSize: headerLen,
blobSize: dataSize,
rawSize: blob.raw_size || data.length
};
}
/**
* Open path and iterate all blobs (síncrono — bloqueia o event loop).
* @param {string} filePath
* @param {object} [options]
* @param {number} [options.startOffset=0]
* @param {function} onBlob - (blobInfo) => void | 'stop'
* @returns {{ blobsRead: number, nextOffset: number, stopped: boolean }}
*/
function forEachBlob(filePath, options, onBlob) {
if (typeof options === 'function') {
onBlob = options;
options = {};
}
options = options || {};
var startOffset = options.startOffset || 0;
var st = fs.statSync(filePath);
var fd = fs.openSync(filePath, 'r');
var offset = startOffset;
var blobIndex = options.startBlobIndex || 0;
var stopped = false;
try {
for (;;) {
var info = readNextBlob(fd, offset, st.size);
if (!info) break;
info.blobIndex = blobIndex;
info.fileSize = st.size;
var ret = onBlob(info);
blobIndex++;
offset = info.nextOffset;
if (ret === 'stop') {
stopped = true;
break;
}
}
} finally {
fs.closeSync(fd);
}
return {
blobsRead: blobIndex - (options.startBlobIndex || 0),
nextOffset: offset,
blobIndex: blobIndex,
stopped: stopped,
fileSize: st.size
};
}
/**
* Like forEachBlob, but yields to the event loop every `yieldEvery` blobs
* (default 1) so SIGINT/Ctrl+C handlers can run mid-scan.
*
* Crucial for long extracts/waves on multi-GB PBFs: the sync loop never
* processes signals until it finishes.
*
* @param {string} filePath
* @param {object} [options]
* @param {number} [options.startOffset=0]
* @param {number} [options.startBlobIndex=0]
* @param {number} [options.yieldEvery=1]
* @param {function} onBlob - (blobInfo) => void | 'stop' | Promise
* @returns {Promise<{ blobsRead, nextOffset, blobIndex, stopped, fileSize }>}
*/
function forEachBlobAsync(filePath, options, onBlob) {
if (typeof options === 'function') {
onBlob = options;
options = {};
}
options = options || {};
var startOffset = options.startOffset || 0;
var startBlobIndex = options.startBlobIndex || 0;
var yieldEvery = options.yieldEvery == null ? 1 : options.yieldEvery | 0;
if (yieldEvery < 0) yieldEvery = 0;
return new Promise(function (resolve, reject) {
var st;
var fd;
try {
st = fs.statSync(filePath);
fd = fs.openSync(filePath, 'r');
} catch (e) {
return reject(e);
}
var offset = startOffset;
var blobIndex = startBlobIndex;
var stopped = false;
var sinceYield = 0;
var closed = false;
function closeFd() {
if (closed) return;
closed = true;
try {
fs.closeSync(fd);
} catch (e) {
/* ignore */
}
}
function done() {
closeFd();
resolve({
blobsRead: blobIndex - startBlobIndex,
nextOffset: offset,
blobIndex: blobIndex,
stopped: stopped,
fileSize: st.size
});
}
function fail(err) {
closeFd();
reject(err);
}
function step() {
try {
for (;;) {
var info = readNextBlob(fd, offset, st.size);
if (!info) {
done();
return;
}
info.blobIndex = blobIndex;
info.fileSize = st.size;
var ret = onBlob(info);
// Support async onBlob without forcing all callers to async
if (ret && typeof ret.then === 'function') {
ret.then(function (asyncRet) {
blobIndex++;
offset = info.nextOffset;
if (asyncRet === 'stop') {
stopped = true;
done();
return;
}
setImmediate(step);
}, fail);
return;
}
blobIndex++;
offset = info.nextOffset;
if (ret === 'stop') {
stopped = true;
done();
return;
}
sinceYield++;
if (yieldEvery > 0 && sinceYield >= yieldEvery) {
sinceYield = 0;
setImmediate(step);
return;
}
}
} catch (err) {
fail(err);
}
}
step();
});
}
module.exports = {
readNextBlob: readNextBlob,
forEachBlob: forEachBlob,
forEachBlobAsync: forEachBlobAsync
};