Raw
1 /*
2 * Copyright (c) 2011, Google Inc.
3 */
4
5 #include "git-compat-util.h"
6 #include "convert.h"
7 #include "environment.h"
8 #include "repository.h"
9 #include "odb.h"
10 #include "odb/source.h"
11 #include "odb/streaming.h"
12 #include "replace-object.h"
13
14 #define FILTER_BUFFER (1024*16)
15
16 /*****************************************************************
17 *
18 * Filtered stream
19 *
20 *****************************************************************/
21
22 struct odb_filtered_read_stream {
23 struct odb_read_stream base;
24 struct odb_read_stream *upstream;
25 struct stream_filter *filter;
26 char ibuf[FILTER_BUFFER];
27 char obuf[FILTER_BUFFER];
28 int i_end, i_ptr;
29 int o_end, o_ptr;
30 int input_finished;
31 };
32
33 static int close_istream_filtered(struct odb_read_stream *_fs)
34 {
35 struct odb_filtered_read_stream *fs = (struct odb_filtered_read_stream *)_fs;
36 free_stream_filter(fs->filter);
37 return odb_read_stream_close(fs->upstream);
38 }
39
40 static ssize_t read_istream_filtered(struct odb_read_stream *_fs, char *buf,
41 size_t sz)
42 {
43 struct odb_filtered_read_stream *fs = (struct odb_filtered_read_stream *)_fs;
44 size_t filled = 0;
45
46 while (sz) {
47 /* do we already have filtered output? */
48 if (fs->o_ptr < fs->o_end) {
49 size_t to_move = fs->o_end - fs->o_ptr;
50 if (sz < to_move)
51 to_move = sz;
52 memcpy(buf + filled, fs->obuf + fs->o_ptr, to_move);
53 fs->o_ptr += to_move;
54 sz -= to_move;
55 filled += to_move;
56 continue;
57 }
58 fs->o_end = fs->o_ptr = 0;
59
60 /* do we have anything to feed the filter with? */
61 if (fs->i_ptr < fs->i_end) {
62 size_t to_feed = fs->i_end - fs->i_ptr;
63 size_t to_receive = FILTER_BUFFER;
64 if (stream_filter(fs->filter,
65 fs->ibuf + fs->i_ptr, &to_feed,
66 fs->obuf, &to_receive))
67 return -1;
68 fs->i_ptr = fs->i_end - to_feed;
69 fs->o_end = FILTER_BUFFER - to_receive;
70 continue;
71 }
72
73 /* tell the filter to drain upon no more input */
74 if (fs->input_finished) {
75 size_t to_receive = FILTER_BUFFER;
76 if (stream_filter(fs->filter,
77 NULL, NULL,
78 fs->obuf, &to_receive))
79 return -1;
80 fs->o_end = FILTER_BUFFER - to_receive;
81 if (!fs->o_end)
82 break;
83 continue;
84 }
85 fs->i_end = fs->i_ptr = 0;
86
87 /* refill the input from the upstream */
88 if (!fs->input_finished) {
89 fs->i_end = odb_read_stream_read(fs->upstream, fs->ibuf, FILTER_BUFFER);
90 if (fs->i_end < 0)
91 return -1;
92 if (fs->i_end)
93 continue;
94 }
95 fs->input_finished = 1;
96 }
97 return filled;
98 }
99
100 static struct odb_read_stream *attach_stream_filter(struct odb_read_stream *st,
101 struct stream_filter *filter)
102 {
103 struct odb_filtered_read_stream *fs;
104
105 CALLOC_ARRAY(fs, 1);
106 fs->base.close = close_istream_filtered;
107 fs->base.read = read_istream_filtered;
108 fs->upstream = st;
109 fs->filter = filter;
110 fs->base.size = -1; /* unknown */
111 fs->base.type = st->type;
112
113 return &fs->base;
114 }
115
116 /*****************************************************************
117 *
118 * In-core stream
119 *
120 *****************************************************************/
121
122 struct odb_incore_read_stream {
123 struct odb_read_stream base;
124 char *buf; /* from odb_read_object_info_extended() */
125 unsigned long read_ptr;
126 };
127
128 static int close_istream_incore(struct odb_read_stream *_st)
129 {
130 struct odb_incore_read_stream *st = (struct odb_incore_read_stream *)_st;
131 free(st->buf);
132 return 0;
133 }
134
135 static ssize_t read_istream_incore(struct odb_read_stream *_st, char *buf, size_t sz)
136 {
137 struct odb_incore_read_stream *st = (struct odb_incore_read_stream *)_st;
138 size_t read_size = sz;
139 size_t remainder = st->base.size - st->read_ptr;
140
141 if (remainder <= read_size)
142 read_size = remainder;
143 if (read_size) {
144 memcpy(buf, st->buf + st->read_ptr, read_size);
145 st->read_ptr += read_size;
146 }
147 return read_size;
148 }
149
150 static int open_istream_incore(struct odb_read_stream **out,
151 struct object_database *odb,
152 const struct object_id *oid)
153 {
154 struct object_info oi = OBJECT_INFO_INIT;
155 struct odb_incore_read_stream stream = {
156 .base.close = close_istream_incore,
157 .base.read = read_istream_incore,
158 };
159 struct odb_incore_read_stream *st;
160 unsigned long size_ul;
161 int ret;
162
163 oi.typep = &stream.base.type;
164 /*
165 * object_info.sizep is unsigned long* (32-bit on Windows), but
166 * stream.base.size is size_t (64-bit). We use a temporary variable
167 * because the types are incompatible. Note: this path still truncates
168 * for >4GB objects, but large objects should use pack streaming
169 * (packfile_store_read_object_stream) which handles size_t properly.
170 * This incore fallback is only used for small objects or when pack
171 * streaming is unavailable.
172 */
173 oi.sizep = &size_ul;
174 oi.contentp = (void **)&stream.buf;
175 ret = odb_read_object_info_extended(odb, oid, &oi,
176 OBJECT_INFO_DIE_IF_CORRUPT);
177 if (ret)
178 return ret;
179 stream.base.size = size_ul;
180
181 CALLOC_ARRAY(st, 1);
182 *st = stream;
183 *out = &st->base;
184
185 return 0;
186 }
187
188 /*****************************************************************************
189 * static helpers variables and functions for users of streaming interface
190 *****************************************************************************/
191
192 static int istream_source(struct odb_read_stream **out,
193 struct object_database *odb,
194 const struct object_id *oid)
195 {
196 struct odb_source *source;
197
198 odb_prepare_alternates(odb);
199 for (source = odb->sources; source; source = source->next)
200 if (!odb_source_read_object_stream(out, source, oid))
201 return 0;
202
203 return open_istream_incore(out, odb, oid);
204 }
205
206 /****************************************************************
207 * Users of streaming interface
208 ****************************************************************/
209
210 int odb_read_stream_close(struct odb_read_stream *st)
211 {
212 int r = st->close(st);
213 free(st);
214 return r;
215 }
216
217 ssize_t odb_read_stream_read(struct odb_read_stream *st, void *buf, size_t sz)
218 {
219 return st->read(st, buf, sz);
220 }
221
222 struct odb_read_stream *odb_read_stream_open(struct object_database *odb,
223 const struct object_id *oid,
224 struct stream_filter *filter)
225 {
226 struct odb_read_stream *st;
227 const struct object_id *real = lookup_replace_object(odb->repo, oid);
228 int ret = istream_source(&st, odb, real);
229
230 if (ret)
231 return NULL;
232
233 if (filter) {
234 /* Add "&& !is_null_stream_filter(filter)" for performance */
235 struct odb_read_stream *nst = attach_stream_filter(st, filter);
236 if (!nst) {
237 odb_read_stream_close(st);
238 return NULL;
239 }
240 st = nst;
241 }
242
243 return st;
244 }
245
246 ssize_t odb_write_stream_read(struct odb_write_stream *st, void *buf, size_t sz)
247 {
248 return st->read(st, buf, sz);
249 }
250
251 void odb_write_stream_release(struct odb_write_stream *st)
252 {
253 free(st->data);
254 }
255
256 int odb_stream_blob_to_fd(struct object_database *odb,
257 int fd,
258 const struct object_id *oid,
259 struct stream_filter *filter,
260 int can_seek)
261 {
262 struct odb_read_stream *st;
263 ssize_t kept = 0;
264 int result = -1;
265
266 st = odb_read_stream_open(odb, oid, filter);
267 if (!st) {
268 if (filter)
269 free_stream_filter(filter);
270 return result;
271 }
272 if (st->type != OBJ_BLOB)
273 goto close_and_exit;
274 for (;;) {
275 char buf[1024 * 16];
276 ssize_t wrote, holeto;
277 ssize_t readlen = odb_read_stream_read(st, buf, sizeof(buf));
278
279 if (readlen < 0)
280 goto close_and_exit;
281 if (!readlen)
282 break;
283 if (can_seek && sizeof(buf) == readlen) {
284 for (holeto = 0; holeto < readlen; holeto++)
285 if (buf[holeto])
286 break;
287 if (readlen == holeto) {
288 kept += holeto;
289 continue;
290 }
291 }
292
293 if (kept && lseek(fd, kept, SEEK_CUR) == (off_t) -1)
294 goto close_and_exit;
295 else
296 kept = 0;
297 wrote = write_in_full(fd, buf, readlen);
298
299 if (wrote < 0)
300 goto close_and_exit;
301 }
302 if (kept && (lseek(fd, kept - 1, SEEK_CUR) == (off_t) -1 ||
303 xwrite(fd, "", 1) != 1))
304 goto close_and_exit;
305 result = 0;
306
307 close_and_exit:
308 odb_read_stream_close(st);
309 return result;
310 }
311
312 struct read_object_fd_data {
313 int fd;
314 size_t remaining;
315 };
316
317 static ssize_t read_object_fd(struct odb_write_stream *stream,
318 unsigned char *buf, size_t len)
319 {
320 struct read_object_fd_data *data = stream->data;
321 ssize_t read_result;
322 size_t count;
323
324 if (stream->is_finished)
325 return 0;
326
327 count = data->remaining < len ? data->remaining : len;
328 read_result = read_in_full(data->fd, buf, count);
329 if (read_result < 0 || (size_t)read_result != count)
330 return -1;
331
332 data->remaining -= count;
333 if (!data->remaining)
334 stream->is_finished = 1;
335
336 return read_result;
337 }
338
339 void odb_write_stream_from_fd(struct odb_write_stream *stream, int fd,
340 size_t size)
341 {
342 struct read_object_fd_data *data;
343
344 CALLOC_ARRAY(data, 1);
345 data->fd = fd;
346 data->remaining = size;
347
348 stream->data = data;
349 stream->read = read_object_fd;
350 stream->is_finished = 0;
351 }