-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathtable_writer.v
More file actions
165 lines (153 loc) · 4.04 KB
/
Copy pathtable_writer.v
File metadata and controls
165 lines (153 loc) · 4.04 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
module leveldb
import os
import compress.zlib
import compress.deflate
const table_magic = u64(0xdb4775248b80fb57)
const footer_len = 48
const filter_base_lg = 11
const filter_base = 1 << filter_base_lg
struct TableWriter {
opts Options
mut:
f os.File
offset u64
data_block &BlockBuilder
index_block &BlockBuilder
filter_keys [][]u8
filter_data []u8
filter_offs []u32
bloom &BloomFilter
pending bool
pending_h BlockHandle
last_key []u8
num_entries int
closed bool
}
fn new_table_writer(path string, opts Options) !&TableWriter {
f := os.open_file(path, 'wb+', 0o644)!
return &TableWriter{
opts: opts
f: f
data_block: new_block_builder(opts.block_restart_interval)
index_block: new_block_builder(1)
bloom: new_bloom_filter(opts.bloom_bits_per_key)
}
}
fn (mut w TableWriter) add(key []u8, value []u8) ! {
if w.pending {
w.index_block.add(w.last_key, w.pending_h.encode())
w.pending = false
}
if w.opts.bloom_bits_per_key > 0 {
w.filter_keys << internal_ukey(key).clone()
}
w.data_block.add(key, value)
w.last_key = key.clone()
w.num_entries++
if w.data_block.size_estimate() >= w.opts.block_size {
w.flush_data_block()!
}
}
fn (mut w TableWriter) flush_data_block() ! {
if w.data_block.empty() {
return
}
block := w.data_block.finish()
w.pending_h = w.write_block(block, w.opts.compression)!
w.pending = true
w.data_block.reset()
w.generate_filters()
}
fn (mut w TableWriter) generate_filters() {
if w.opts.bloom_bits_per_key <= 0 {
return
}
filter_index := int(w.offset / u64(filter_base))
for w.filter_offs.len < filter_index {
w.finish_filter_slot()
}
}
fn (mut w TableWriter) finish_filter_slot() {
w.filter_offs << u32(w.filter_data.len)
if w.filter_keys.len > 0 {
f := w.bloom.create(w.filter_keys)
w.filter_data << f
w.filter_keys.clear()
}
}
fn (mut w TableWriter) write_block(block []u8, compression Compression) !BlockHandle {
mut data := unsafe { block }
mut ctype := u8(0)
match compression {
.zlib {
compressed := zlib.compress(block) or { block.clone() }
if compressed.len < block.len {
data = compressed.clone()
ctype = u8(Compression.zlib)
}
}
.raw_deflate {
compressed := deflate.compress(block) or { block.clone() }
if compressed.len < block.len {
data = compressed.clone()
ctype = u8(Compression.raw_deflate)
}
}
else {}
}
handle := BlockHandle{
offset: w.offset
size: u64(data.len)
}
write_fd_all(w.f.fd, data)!
mut trailer := []u8{cap: 5}
trailer << ctype
crc := crc32c_update(crc32c(data), [ctype])
append_u32_le(mut trailer, mask_crc(crc))
write_fd_all(w.f.fd, trailer)!
w.offset += u64(data.len) + 5
return handle
}
fn (mut w TableWriter) finish() ! {
w.flush_data_block()!
if w.pending {
w.index_block.add(w.last_key, w.pending_h.encode())
w.pending = false
}
mut metaindex := new_block_builder(w.opts.block_restart_interval)
if w.opts.bloom_bits_per_key > 0 {
w.finish_filter_slot()
mut filter_block := w.filter_data.clone()
offs_start := u32(filter_block.len)
for off in w.filter_offs {
append_u32_le(mut filter_block, off)
}
append_u32_le(mut filter_block, offs_start)
filter_block << u8(filter_base_lg)
fh := w.write_block(filter_block, .none)!
metaindex.add('filter.${bloom_filter_name}'.bytes(), fh.encode())
}
mh := w.write_block(metaindex.finish(), w.opts.compression)!
ih := w.write_block(w.index_block.finish(), w.opts.compression)!
mut footer := []u8{cap: footer_len}
footer << mh.encode()
footer << ih.encode()
for footer.len < footer_len - 8 {
footer << u8(0)
}
append_u64_le(mut footer, table_magic)
write_fd_all(w.f.fd, footer)!
w.offset += u64(footer_len)
// The table is about to be named by a manifest edit that is itself made
// durable. Get the contents to the device first or a crash between the two
// leaves durable metadata pointing at a table that was never written.
sync_file(w.f.fd)!
w.f.close()
w.closed = true
}
fn (w &TableWriter) file_size() u64 {
return w.offset
}
fn (w &TableWriter) entries() int {
return w.num_entries
}