master
h 186 lines 6.27 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 #ifndef NETDATA_STREAM_COMPRESSION_H
4 #define NETDATA_STREAM_COMPRESSION_H 1
5
6 #include "libnetdata/libnetdata.h"
7
8 // signature MUST end with a newline
9
10 #if COMPRESSION_MAX_MSG_SIZE >= (COMPRESSION_MAX_CHUNK - COMPRESSION_MAX_OVERHEAD)
11 #error "COMPRESSION_MAX_MSG_SIZE >= (COMPRESSION_MAX_CHUNK - COMPRESSION_MAX_OVERHEAD)"
12 #endif
13
14 typedef uint32_t stream_compression_signature_t;
15 #define STREAM_COMPRESSION_SIGNATURE ((stream_compression_signature_t)('z' | 0x80) | (0x80 << 8) | (0x80 << 16) | ('\n' << 24))
16 #define STREAM_COMPRESSION_SIGNATURE_MASK ((stream_compression_signature_t) 0xffU | (0x80U << 8) | (0x80U << 16) | (0xffU << 24))
17 #define STREAM_COMPRESSION_SIGNATURE_SIZE sizeof(stream_compression_signature_t)
18
19 static inline stream_compression_signature_t stream_compress_encode_signature(size_t compressed_data_size) {
20 stream_compression_signature_t len = ((compressed_data_size & 0x7f) | 0x80 | (((compressed_data_size & (0x7f << 7)) << 1) | 0x8000)) << 8;
21 return len | STREAM_COMPRESSION_SIGNATURE;
22 }
23
24 typedef enum {
25 COMPRESSION_ALGORITHM_NONE = 0,
26 COMPRESSION_ALGORITHM_ZSTD,
27 COMPRESSION_ALGORITHM_LZ4,
28 COMPRESSION_ALGORITHM_GZIP,
29 COMPRESSION_ALGORITHM_BROTLI,
30
31 // terminator
32 COMPRESSION_ALGORITHM_MAX,
33 } compression_algorithm_t;
34
35 // this defines the order the algorithms will be selected by the receiver (parent)
36 #define STREAM_COMPRESSION_ALGORITHMS_ORDER "zstd lz4 brotli gzip"
37
38 // ----------------------------------------------------------------------------
39
40 typedef struct simple_ring_buffer {
41 const char *data;
42 size_t size;
43 size_t read_pos;
44 size_t write_pos;
45 } SIMPLE_RING_BUFFER;
46
47 static inline void simple_ring_buffer_reset(SIMPLE_RING_BUFFER *b) {
48 b->read_pos = b->write_pos = 0;
49 }
50
51 static inline void simple_ring_buffer_make_room(SIMPLE_RING_BUFFER *b, size_t size) {
52 if(b->write_pos + size > b->size) {
53 if(!b->size)
54 b->size = COMPRESSION_MAX_CHUNK;
55 else
56 b->size *= 2;
57
58 if(b->write_pos + size > b->size)
59 b->size += size;
60
61 b->data = (const char *)reallocz((void *)b->data, b->size);
62 }
63 }
64
65 static inline void simple_ring_buffer_append_data(SIMPLE_RING_BUFFER *b, const void *data, size_t size) {
66 simple_ring_buffer_make_room(b, size);
67 memcpy((void *)(b->data + b->write_pos), data, size);
68 b->write_pos += size;
69 }
70
71 static inline void simple_ring_buffer_destroy(SIMPLE_RING_BUFFER *b) {
72 freez((void *)b->data);
73 b->data = NULL;
74 b->read_pos = b->write_pos = b->size = 0;
75 }
76
77 // ----------------------------------------------------------------------------
78
79 struct compressor_state {
80 bool initialized;
81 compression_algorithm_t algorithm;
82
83 SIMPLE_RING_BUFFER input;
84 SIMPLE_RING_BUFFER output;
85
86 int level;
87 void *stream;
88
89 struct {
90 size_t total_compressed;
91 size_t total_uncompressed;
92 size_t total_compressions;
93 } sender_locked;
94 };
95
96 void stream_compressor_init(struct compressor_state *state);
97 void stream_compressor_destroy(struct compressor_state *state);
98 size_t stream_compress(struct compressor_state *state, const char *data, size_t size, const char **out);
99
100 // ----------------------------------------------------------------------------
101
102 struct decompressor_state {
103 bool initialized;
104 compression_algorithm_t algorithm;
105 size_t signature_size;
106
107 size_t total_compressed;
108 size_t total_uncompressed;
109 size_t total_compressions;
110
111 SIMPLE_RING_BUFFER output;
112
113 void *stream;
114 };
115
116 void stream_decompressor_destroy(struct decompressor_state *state);
117 void stream_decompressor_init(struct decompressor_state *state);
118 size_t stream_decompress(struct decompressor_state *state, const char *compressed_data, size_t compressed_size);
119
120 static inline size_t stream_decompress_decode_signature(const char *data, size_t data_size) {
121 if (unlikely(!data || !data_size))
122 return 0;
123
124 if (unlikely(data_size != STREAM_COMPRESSION_SIGNATURE_SIZE))
125 return 0;
126
127 stream_compression_signature_t sign;
128 memcpy(&sign, data, sizeof(stream_compression_signature_t)); // Safe copy to aligned variable
129 // stream_compression_signature_t sign = *(stream_compression_signature_t *)data;
130
131 if (unlikely((sign & STREAM_COMPRESSION_SIGNATURE_MASK) != STREAM_COMPRESSION_SIGNATURE))
132 return 0;
133
134 size_t length = ((sign >> 8) & 0x7f) | ((sign >> 9) & (0x7f << 7));
135 return length;
136 }
137
138 static inline size_t stream_decompressor_start(struct decompressor_state *state, const char *header, size_t header_size) {
139 if(unlikely(state->output.read_pos != state->output.write_pos))
140 fatal("STREAM_DECOMPRESS: asked to decompress new data, while there are unread data in the decompression buffer!");
141
142 return stream_decompress_decode_signature(header, header_size);
143 }
144
145 static inline size_t stream_decompressed_bytes_in_buffer(struct decompressor_state *state) {
146 if(unlikely(state->output.read_pos > state->output.write_pos))
147 fatal("STREAM_DECOMPRESS: invalid read/write stream positions");
148
149 return state->output.write_pos - state->output.read_pos;
150 }
151
152 static inline size_t stream_decompressor_get(struct decompressor_state *state, char *dst, size_t size) {
153 if (unlikely(!state || !size || !dst))
154 return 0;
155
156 size_t remaining = stream_decompressed_bytes_in_buffer(state);
157
158 if(unlikely(!remaining))
159 return 0;
160
161 size_t bytes_to_return = size;
162 if(bytes_to_return > remaining)
163 bytes_to_return = remaining;
164
165 memcpy(dst, state->output.data + state->output.read_pos, bytes_to_return);
166 state->output.read_pos += bytes_to_return;
167
168 if(unlikely(state->output.read_pos > state->output.write_pos))
169 fatal("STREAM_DECOMPRESS: invalid read/write stream positions");
170
171 return bytes_to_return;
172 }
173
174 // ----------------------------------------------------------------------------
175
176 struct sender_state;
177 struct receiver_state;
178 struct stream_receiver_config;
179
180 bool stream_compression_initialize(struct sender_state *s);
181 bool stream_decompression_initialize(struct receiver_state *rpt);
182 void stream_parse_compression_order(struct stream_receiver_config *config, const char *order);
183 void stream_select_receiver_compression_algorithm(struct receiver_state *rpt);
184 void stream_compression_deactivate(struct sender_state *s);
185
186 #endif // NETDATA_STREAM_COMPRESSION_H 1