FFmpeg coverage


Directory: ../../../ffmpeg/
File: src/libavformat/shared.c
Date: 2026-09-28 11:37:26
Exec Total Coverage
Lines: 0 508 0.0%
Functions: 0 22 0.0%
Branches: 0 349 0.0%

Line Branch Exec Source
1 /*
2 * Shared file cache protocol.
3 * Copyright (c) 2026 Niklas Haas
4 *
5 * This file is part of FFmpeg.
6 *
7 * FFmpeg is free software; you can redistribute it and/or
8 * modify it under the terms of the GNU Lesser General Public
9 * License as published by the Free Software Foundation; either
10 * version 2.1 of the License, or (at your option) any later version.
11 *
12 * FFmpeg is distributed in the hope that it will be useful,
13 * but WITHOUT ANY WARRANTY; without even the implied warranty of
14 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
15 * Lesser General Public License for more details.
16 *
17 * You should have received a copy of the GNU Lesser General Public
18 * License along with FFmpeg; if not, write to the Free Software
19 * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
20 *
21 * Based on cache.c by Michael Niedermayer
22 */
23
24 #include "libavutil/attributes.h"
25 #include "libavutil/avassert.h"
26 #include "libavutil/avstring.h"
27 #include "libavutil/crc.h"
28 #include "libavutil/error.h"
29 #include "libavutil/file.h"
30 #include "libavutil/hash.h"
31 #include "libavutil/file_open.h"
32 #include "libavutil/mem.h"
33 #include "libavutil/opt.h"
34 #include "libavutil/time.h"
35
36 #include "internal.h"
37 #include "os_support.h"
38 #include "url.h"
39
40 #include <assert.h>
41 #include <errno.h>
42 #include <fcntl.h>
43 #include <inttypes.h>
44 #include <stdatomic.h>
45 #include <string.h>
46 #include <sys/stat.h>
47 #if HAVE_UNISTD_H
48 #include <unistd.h>
49 #endif
50 #ifndef _WIN32
51 #include <sys/file.h>
52 #endif
53
54 /**
55 * This hash should be resistant against collision attacks, so that an
56 * attacker could not generate e.g. two different URIs that map to the same
57 * cache file. This requires at least 64 bits of collision resistance in
58 * practice (i.e. 128 bits = 16 bytes of hash size). However, we can be
59 * conservative by computing e.g. a 256 bit hash and storing it inside the
60 * file header for verification.
61 *
62 * Note that due to the way we use atomics, we should avoid zero bytes in
63 * the resulting hash; hence we tweak the input slightly to avoid this.
64 * The resulting loss in hash strength is negligible, since 32 bytes is
65 * already much more than needed.
66 */
67 #define HASH_METHOD "SHA512/256"
68 #define HASH_SIZE 32
69 #define HEADER_MAGIC MKTAG(u'\xFF', 'S', 'h', '$')
70 #define HEADER_VERSION 3
71
72 /**
73 * Hard watershed of consecutive failed blocks before we give up on the cache
74 * file altogether and assume it's entirely lost to us.
75 **/
76 #define MAX_CORRUPT_BLOCKS 10
77
78 ✗ static int hash_uri(uint8_t hash[HASH_SIZE], const char *uri)
79 {
80 ✗ struct AVHashContext *ctx = NULL;
81 ✗ int ret = av_hash_alloc(&ctx, HASH_METHOD);
82 ✗ if (ret < 0)
83 ✗ return ret;
84
85 ✗ const int16_t version = HEADER_VERSION;
86 ✗ av_assert0(av_hash_get_size(ctx) == HASH_SIZE);
87 ✗ av_hash_init(ctx);
88 ✗ av_hash_update(ctx, (const uint8_t *) &version, sizeof(version));
89 ✗ av_hash_update(ctx, (const uint8_t *) uri, strlen(uri));
90 ✗ av_hash_final(ctx, hash);
91 ✗ av_hash_freep(&ctx);
92
93 ✗ for (int i = 0; i < HASH_SIZE; i++)
94 ✗ hash[i] = hash[i] ? hash[i] : ~hash[i]; /* prevent zero bytes */
95 ✗ return 0;
96 }
97
98 enum BlockState {
99 /* Reserved block state values */
100 BLOCK_NONE = 0, ///< block is not cached
101 BLOCK_PENDING, ///< a thread is currently trying to write this block
102 BLOCK_FAILED, ///< the underlying I/O source failed to read this block
103
104 /**
105 * All other block states represent valid cached blocks, with the value
106 * being the CRC of the block data.
107 */
108 };
109
110 ✗ static uint32_t get_block_crc(const uint8_t *block, size_t block_size)
111 {
112 ✗ uint32_t crc = av_crc(av_crc_get_table(AV_CRC_32_IEEE), 0, block, block_size);
113 ✗ switch (crc) {
114 ✗ case BLOCK_NONE:
115 case BLOCK_FAILED:
116 case BLOCK_PENDING:
117 ✗ return ~crc; /* avoid reserved block states */
118 ✗ default:
119 ✗ return crc;
120 }
121 }
122
123 typedef struct Block {
124 atomic_uint state; /* enum BlockState */
125 } Block;
126
127 typedef struct Spacemap {
128 atomic_uint header_magic;
129 atomic_ushort version;
130 atomic_ushort block_shift;
131 atomic_ullong filesize; /* byte offset of true EOF, or 0 if unknown */
132 atomic_uchar hash[HASH_SIZE]; /* hash of resource URI / filename */
133 atomic_ullong blocks_cached; /* (lower bound on) the number of blocks cached */
134 char reserved[72];
135
136 Block blocks[];
137 } Spacemap;
138
139 static_assert(offsetof(Spacemap, blocks) == 128, "Spacemap header layout mismatch");
140
141 /* Set to value iff the current value is unset (zero) */
142 #define DEF_SET_ONCE(ctype, atype) \
143 static int set_once_##atype(atomic_##atype *const ptr, const ctype value) \
144 { \
145 ctype prev = 0; \
146 av_assert1(value != 0); \
147 if (atomic_compare_exchange_strong_explicit( \
148 ptr, &prev, value, memory_order_release, memory_order_relaxed)) \
149 return 1; \
150 else if (prev == value) \
151 return 0; \
152 else \
153 return AVERROR(EINVAL); \
154 }
155
156 ✗ DEF_SET_ONCE(unsigned char, uchar)
157 ✗ DEF_SET_ONCE(unsigned int, uint)
158 ✗ DEF_SET_ONCE(unsigned short, ushort)
159 ✗ DEF_SET_ONCE(unsigned long long, ullong)
160
161 typedef struct SharedContext {
162 AVClass *class;
163 URLContext *inner;
164 int64_t inner_pos;
165
166 /* options */
167 char *cache_dir;
168 int block_shift; ///< requested shift; updated on init if it disagrees
169 int read_only;
170 int64_t timeout;
171 int ignore_errors;
172 int retry_errors;
173 int retry_corrupt;
174 int verify;
175 int64_t cache_size_max;
176
177 /* misc state */
178 int64_t pos; ///< current logical position
179 uint8_t *tmp_buf;
180 int block_size;
181 int write_err; ///< write error occurred
182 int num_corrupt;
183 int64_t filesize; ///< once known
184 int64_t blocks_max; ///< maximum number of blocks to cache
185
186 /* cache file */
187 uint8_t *cache_data; ///< optional mapping of the cache file
188 char *cache_path;
189 off_t cache_size; ///< size of mapped memory region (for unmapping)
190 int fd;
191
192 /* space map */
193 Spacemap *spacemap;
194 char *map_path;
195 off_t map_size;
196 int mapfd;
197
198 /* statistics */
199 int64_t nb_hit;
200 int64_t nb_miss;
201 } SharedContext;
202
203 ✗ static int shared_close(URLContext *h)
204 {
205 ✗ SharedContext *s = h->priv_data;
206
207 ✗ ffurl_close(s->inner);
208 ✗ av_file_unmap_shared(s->cache_data, s->cache_size);
209 ✗ av_file_unmap_shared(s->spacemap, s->map_size);
210 ✗ if (s->fd != -1)
211 ✗ close(s->fd);
212 ✗ if (s->mapfd != -1)
213 ✗ close(s->mapfd);
214 ✗ av_freep(&s->cache_path);
215 ✗ av_freep(&s->map_path);
216 ✗ av_freep(&s->tmp_buf);
217
218 ✗ av_log(h, AV_LOG_DEBUG, "Cache statistics: %"PRId64" hits, %"PRId64" misses\n",
219 s->nb_hit, s->nb_miss);
220 ✗ return 0;
221 }
222
223 static int cache_map(URLContext *h, int64_t filesize);
224 static int spacemap_init(URLContext *h, const uint8_t hash[HASH_SIZE]);
225 static int spacemap_grow(URLContext *h, int64_t block);
226
227 ✗ static int64_t get_filesize(URLContext *h)
228 {
229 ✗ SharedContext *s = h->priv_data;
230 ✗ if (!s->filesize) {
231 ✗ uint64_t size = atomic_load_explicit(&s->spacemap->filesize, memory_order_relaxed);
232 ✗ if (size > INT64_MAX)
233 ✗ return AVERROR(EINVAL);
234 ✗ else if (size)
235 ✗ s->filesize = size;
236 }
237
238 ✗ return s->filesize;
239 }
240
241 ✗ static int set_filesize(URLContext *h, int64_t new_size)
242 {
243 ✗ SharedContext *s = h->priv_data;
244 int ret;
245
246 ✗ if (!new_size)
247 ✗ return 0;
248
249 ✗ ret = set_once_ullong(&s->spacemap->filesize, new_size);
250 ✗ if (ret < 0) {
251 ✗ av_log(h, AV_LOG_ERROR, "Cached file size mismatch, expected: "
252 "%"PRId64", got: %"PRIu64"!\n", new_size,
253 ✗ (uint64_t) atomic_load(&s->spacemap->filesize));
254 ✗ return ret;
255 ✗ } else if (ret) {
256 /* Opportunistically map the file; this also sets the correct filesize.
257 * Ignore errors as this is not critical to the cache logic. */
258 ✗ cache_map(h, new_size);
259 }
260
261 ✗ return ret;
262 }
263
264 ✗ static int is_ignorable_error(int64_t err)
265 {
266 ✗ switch (err) {
267 ✗ case AVERROR_EXIT:
268 case AVERROR_EOF:
269 case AVERROR_BUG:
270 case AVERROR(EAGAIN):
271 case AVERROR(ENOSYS):
272 case AVERROR(EINVAL):
273 ✗ return 0;
274 ✗ default:
275 ✗ return err < 0;
276 }
277 }
278
279 ✗ static int shared_open(URLContext *h, const char *arg, int flags, AVDictionary **options)
280 {
281 ✗ SharedContext *s = h->priv_data;
282 int ret;
283
284 ✗ if (!s->cache_dir || !s->cache_dir[0]) {
285 ✗ av_log(h, AV_LOG_ERROR, "Missing path for shared cache! Specify a "
286 "directory using the -cache_dir option.\n");
287 ✗ return AVERROR(EINVAL);
288 }
289
290 ✗ s->fd = s->mapfd = -1; /* Set these early for shared_close() failure path */
291
292 /* Open underlying protocol */
293 ✗ av_strstart(arg, "shared:", &arg);
294 ✗ ret = ffurl_open_whitelist(&s->inner, arg, flags, &h->interrupt_callback,
295 options, h->protocol_whitelist, h->protocol_blacklist, h);
296 ✗ if (is_ignorable_error(ret) && s->ignore_errors) {
297 ✗ av_log(h, AV_LOG_WARNING, "Underlying URL failed to open: %s. "
298 ✗ "Continuing with cache file only.\n", av_err2str(ret));
299 ✗ } else if (ret < 0)
300 ✗ goto fail;
301
302 uint8_t hash[HASH_SIZE];
303 ✗ ret = hash_uri(hash, arg);
304 ✗ if (ret < 0)
305 ✗ goto fail;
306
307 /* 128 bits is enough for collision resistance; we already store the full
308 * hash inside the header for verification */
309 char filename[2 * 16 + 1];
310 ✗ ff_data_to_hex(filename, hash, sizeof(filename)/2, 0);
311 ✗ s->cache_path = av_asprintf("%s/%s.cache", s->cache_dir, filename);
312 ✗ s->map_path = av_asprintf("%s/%s.spacemap", s->cache_dir, filename);
313 ✗ if (!s->cache_path || !s->map_path) {
314 ✗ ret = AVERROR(ENOMEM);
315 ✗ goto fail;
316 }
317
318 ✗ av_log(h, AV_LOG_VERBOSE, "Opening cache file '%s' for URI: '%s'\n",
319 ✗ s->cache_path, s->inner ? s->inner->filename : arg);
320
321 ✗ const int mode = O_RDWR | O_BINARY | (s->inner ? O_CREAT : 0);
322 ✗ s->fd = avpriv_open(s->cache_path, mode, 0660);
323 ✗ s->mapfd = s->fd >= 0 ? avpriv_open(s->map_path, mode, 0660) : -1;
324 ✗ if (s->fd < 0 || s->mapfd < 0) {
325 ✗ ret = AVERROR(errno);
326 ✗ av_log(h, AV_LOG_ERROR, "Failed to open '%s': %s\n",
327 ✗ s->fd < 0 ? s->cache_path : s->map_path, av_err2str(ret));
328 ✗ goto fail;
329 }
330
331 ✗ ret = spacemap_init(h, hash);
332 ✗ if (ret < 0)
333 ✗ goto fail;
334
335 /* s->block_shift is fully settled after spacemap_init() */
336 ✗ s->block_size = 1 << s->block_shift;
337 ✗ s->blocks_max = s->cache_size_max >> s->block_shift;
338
339 ✗ int64_t filesize = get_filesize(h);
340 ✗ if (filesize < 0) {
341 ✗ ret = (int) filesize;
342 ✗ goto fail;
343 ✗ } else if (!filesize) {
344 /* Filesize is not yet known, try to get it from the underlying URL;
345 * go through our own seek function to handle errors and updates */
346 ✗ filesize = ffurl_size(h);
347 ✗ if (filesize < 0 && filesize != AVERROR(ENOSYS)) {
348 ✗ ret = (int) filesize;
349 ✗ goto fail;
350 }
351 }
352
353 ✗ if (filesize > 0) {
354 ✗ int64_t last_pos = filesize - 1;
355 ✗ int64_t last_block = last_pos >> s->block_shift;
356 ✗ ret = spacemap_grow(h, last_block);
357 ✗ if (ret < 0)
358 ✗ goto fail;
359
360 /* If filesize is known, we can directly map the cache file */
361 ✗ ret = cache_map(h, filesize);
362 ✗ if (ret < 0) {
363 ✗ av_log(h, AV_LOG_WARNING, "Failed to map cache file: %s. Falling "
364 ✗ "back to normal read/write\n", av_err2str(ret));
365 ✗ ret = 0;
366 }
367 }
368
369 /* Temporary buffer needed for pread/pwrite() fallback */
370 ✗ s->tmp_buf = av_malloc(s->block_size);
371 ✗ if (!s->tmp_buf) {
372 ✗ ret = AVERROR(ENOMEM);
373 ✗ goto fail;
374 }
375
376 ✗ h->max_packet_size = s->block_size;
377 ✗ h->min_packet_size = s->block_size;
378 ✗ ret = 0;
379
380 ✗ fail:
381 ✗ if (ret < 0)
382 ✗ shared_close(h);
383 ✗ return ret;
384 }
385
386 ✗ static int cache_map(URLContext *h, int64_t filesize)
387 {
388 ✗ SharedContext *s = h->priv_data;
389 ✗ if (s->cache_size >= filesize || filesize > SIZE_MAX)
390 ✗ return 0;
391
392 ✗ if (s->cache_data) {
393 ✗ av_file_unmap_shared(s->cache_data, s->cache_size);
394 ✗ s->cache_data = NULL;
395 ✗ s->cache_size = 0;
396 }
397
398 /* The mapping extends the file to the file size; it can be shorter if
399 * another process wrote the correct filesize to the header but crashed
400 * right before actually successfully resizing the file. */
401 void *map;
402 ✗ int ret = av_file_map_shared(s->fd, filesize, &map);
403 ✗ if (ret < 0)
404 ✗ return ret;
405
406 ✗ s->cache_data = map;
407 ✗ s->cache_size = filesize;
408 ✗ return 0;
409 }
410
411 ✗ static int spacemap_remap(URLContext *h, size_t map_size)
412 {
413 ✗ SharedContext *s = h->priv_data;
414 ✗ int ret, did_grow = 0, locked = 0;
415 ✗ if (map_size <= s->map_size)
416 ✗ return 0;
417
418 /* Opportunistically get current filesize before attempting to lock */
419 struct stat st;
420 ✗ ret = fstat(s->mapfd, &st);
421 ✗ if (ret < 0) {
422 ✗ ret = AVERROR(errno);
423 ✗ goto fail;
424 }
425
426 ✗ if (st.st_size >= map_size)
427 ✗ goto skip_resize;
428
429 /* Lock the spacemap to ensure nobody else is currently resizing it */
430 ✗ ret = flock(s->mapfd, LOCK_EX);
431 ✗ if (ret < 0) {
432 ✗ ret = AVERROR(errno);
433 ✗ goto fail;
434 }
435 ✗ locked = 1;
436
437 /* Refresh filesize after acquiring the lock */
438 ✗ ret = fstat(s->mapfd, &st);
439 ✗ if (ret < 0) {
440 ✗ ret = AVERROR(errno);
441 ✗ goto fail;
442 }
443
444 ✗ if (st.st_size >= map_size)
445 ✗ goto skip_resize;
446
447 /* The new mapping extends the file */
448 ✗ st.st_size = map_size;
449 ✗ did_grow = 1;
450
451 ✗ skip_resize:
452 ✗ av_file_unmap_shared(s->spacemap, s->map_size);
453 ✗ s->spacemap = NULL;
454 ✗ s->map_size = st.st_size;
455
456 void *map;
457 ✗ ret = av_file_map_shared(s->mapfd, s->map_size, &map);
458 ✗ if (ret < 0) {
459 ✗ s->map_size = 0;
460 ✗ goto fail;
461 }
462 ✗ s->spacemap = map;
463
464 ✗ if (locked) {
465 ✗ flock(s->mapfd, LOCK_UN);
466 ✗ locked = 0;
467 }
468
469 ✗ return did_grow;
470
471 ✗ fail:
472 ✗ if (locked)
473 ✗ flock(s->mapfd, LOCK_UN);
474 ✗ av_log(h, AV_LOG_ERROR, "Failed to resize space map: %s\n", av_err2str(ret));
475 ✗ return ret;
476 }
477
478 ✗ static int spacemap_grow(URLContext *h, int64_t block)
479 {
480 ✗ SharedContext *s = h->priv_data;
481 ✗ int64_t num_blocks = block + 1;
482 ✗ size_t map_bytes = sizeof(Spacemap) + num_blocks * sizeof(Block);
483
484 /* When streaming files without known size, round up the number of blocks
485 * to the nearest multiple of the block size to reduce the rate of resizes */
486 ✗ int64_t filesize = get_filesize(h);
487 ✗ if (filesize < 0)
488 ✗ return (int) filesize;
489 ✗ else if (!filesize) {
490 ✗ av_assert0(s->block_size > 0);
491 ✗ map_bytes = FFALIGN(map_bytes, (int64_t) s->block_size);
492 }
493
494 ✗ if (map_bytes < num_blocks)
495 ✗ return AVERROR(EINVAL); /* overflow */
496
497 ✗ const off_t old_size = s->map_size;
498 ✗ int ret = spacemap_remap(h, map_bytes);
499 ✗ if (ret < 0)
500 ✗ return ret;
501
502 /* Report new size after successful grow */
503 ✗ if (s->map_size > old_size) {
504 ✗ num_blocks = (s->map_size - sizeof(Spacemap)) / sizeof(Block);
505 ✗ av_log(h, AV_LOG_DEBUG,
506 "%s %zu bytes, capacity: %"PRId64" blocks = %"PRId64" MB\n",
507 ret ? "Resized spacemap to" : "Mapped spacemap with",
508 ✗ (size_t) s->map_size, num_blocks,
509 ✗ (num_blocks * (int64_t) s->block_size) >> 20);
510 }
511 ✗ return 0;
512 }
513
514 ✗ static int spacemap_init(URLContext *h, const uint8_t hash[HASH_SIZE])
515 {
516 ✗ SharedContext *s = h->priv_data;
517 int ret;
518
519 ✗ ret = spacemap_remap(h, sizeof(Spacemap));
520 ✗ if (ret < 0)
521 ✗ return ret;
522
523 ✗ if ((ret = set_once_uint(&s->spacemap->header_magic, HEADER_MAGIC)) < 0 ||
524 ✗ (ret = set_once_ushort(&s->spacemap->version, HEADER_VERSION)) < 0)
525 {
526 ✗ av_log(h, AV_LOG_ERROR, "Shared cache spacemap header mismatch!\n");
527 ✗ av_log(h, AV_LOG_ERROR, " Expected magic: 0x%X, version: %d\n",
528 HEADER_MAGIC, HEADER_VERSION);
529 ✗ av_log(h, AV_LOG_ERROR, " Got magic: 0x%X, version: %d\n",
530 ✗ atomic_load(&s->spacemap->header_magic),
531 ✗ atomic_load(&s->spacemap->version));
532 ✗ return ret;
533 }
534
535 ✗ ret = set_once_ushort(&s->spacemap->block_shift, s->block_shift);
536 ✗ if (ret < 0) {
537 ✗ const int shift = atomic_load(&s->spacemap->block_shift);
538 ✗ av_log(h, AV_LOG_WARNING, "Shared cache uses block shift %d, "
539 "but requested block shift is %d.\n", shift, s->block_shift);
540 ✗ if (shift < 9 || shift > 30) {
541 ✗ av_log(h, AV_LOG_ERROR, "Invalid block shift %d in cache file!\n", shift);
542 ✗ return AVERROR(EINVAL);
543 }
544 ✗ s->block_shift = shift;
545 }
546
547 ✗ for (int i = 0; i < HASH_SIZE; i++) {
548 ✗ ret = set_once_uchar(&s->spacemap->hash[i], hash[i]);
549 ✗ if (ret < 0) {
550 ✗ av_log(h, AV_LOG_ERROR, "Shared cache spacemap hash mismatch!\n");
551 char hash_hex[2 * HASH_SIZE + 1];
552 ✗ ff_data_to_hex(hash_hex, hash, HASH_SIZE, 0);
553 ✗ av_log(h, AV_LOG_ERROR, " Expected hash: %s\n", hash_hex);
554 uint8_t hash2[HASH_SIZE];
555 ✗ for (int j = 0; j < HASH_SIZE; ++j)
556 ✗ hash2[j] = atomic_load_explicit(&s->spacemap->hash[j], memory_order_relaxed);
557 ✗ ff_data_to_hex(hash_hex, hash2, HASH_SIZE, 0);
558 ✗ av_log(h, AV_LOG_ERROR, " Got hash: %s\n", hash_hex);
559 ✗ return ret;
560 }
561 }
562
563 ✗ if (ret) /* set_once() return 1 if this is the first time setting the value */
564 ✗ av_log(h, AV_LOG_DEBUG, "Initialized new cache spacemap.\n");
565
566 ✗ return ret;
567 }
568
569 ✗ static int read_cache(SharedContext *s, uint8_t *buf, size_t size, off_t offset)
570 {
571 ✗ if (s->cache_data) {
572 av_assert1(offset + size <= s->cache_size);
573 ✗ memcpy(buf, s->cache_data + offset, size);
574 ✗ return 0;
575 }
576
577 ✗ while (size) {
578 ✗ ssize_t ret = pread(s->fd, buf, size, offset);
579 ✗ if (ret <= 0)
580 ✗ return ret ? AVERROR(errno) : AVERROR_EOF;
581 ✗ buf += ret;
582 ✗ offset += ret;
583 ✗ size -= ret;
584 }
585
586 ✗ return 0;
587 }
588
589 ✗ static int write_cache(SharedContext *s, const uint8_t *buf, size_t size, off_t offset)
590 {
591 ✗ if (s->cache_data) {
592 av_assert1(offset + size <= s->cache_size);
593 ✗ memcpy(s->cache_data + offset, buf, size);
594 ✗ return 0;
595 }
596
597 ✗ while (size) {
598 ✗ ssize_t ret = pwrite(s->fd, buf, size, offset);
599 ✗ if (ret <= 0)
600 ✗ return ret ? AVERROR(errno) : AVERROR(EIO);
601 ✗ buf += ret;
602 ✗ offset += ret;
603 ✗ size -= ret;
604 }
605
606 ✗ return 0;
607 }
608
609 ✗ static int clamp_size(URLContext *h, int size, int64_t pos, int64_t filesize)
610 {
611 ✗ if (!filesize)
612 ✗ return size;
613 ✗ else if (pos > filesize)
614 ✗ return 0;
615 else
616 ✗ return FFMIN(filesize - pos, size);
617 }
618
619 ✗ static int shared_read(URLContext *h, unsigned char *buf, int size)
620 {
621 ✗ SharedContext *s = h->priv_data;
622 uint8_t *tmp;
623 int ret;
624 ✗ if (!s->spacemap)
625 ✗ return AVERROR(EIO);
626
627 ✗ if (size <= 0)
628 ✗ return 0;
629
630 ✗ int64_t filesize = get_filesize(h);
631 ✗ if (filesize < 0)
632 ✗ return (int) filesize;
633
634 ✗ size = clamp_size(h, size, s->pos, filesize);
635 ✗ if (size <= 0)
636 ✗ return AVERROR_EOF;
637
638 ✗ const int64_t block_id = s->pos >> s->block_shift;
639 ✗ const int64_t offset = s->pos & (s->block_size - 1);
640 ✗ const int64_t block_pos = block_id * s->block_size;
641 ✗ int block_size = clamp_size(h, s->block_size, block_pos, filesize);
642 ✗ ret = spacemap_grow(h, block_id);
643 ✗ if (ret < 0)
644 ✗ return ret;
645
646 ✗ Block *const block = &s->spacemap->blocks[block_id];
647 ✗ unsigned state = atomic_load_explicit(&block->state, memory_order_acquire);
648 ✗ int64_t pending_since = 0;
649 ✗ int verify_read = 0, acquired = 0, allocated = 0;
650
651 ✗ retry:
652 ✗ switch (state) {
653 ✗ default:
654 ✗ if (s->num_corrupt >= MAX_CORRUPT_BLOCKS)
655 ✗ goto read_block; /* assume broken cache file */
656
657 /* filesize may have become known in the meantime */
658 ✗ filesize = get_filesize(h);
659 ✗ if (filesize < 0)
660 ✗ return (int) filesize;
661
662 /* We always need to read the entire block to verify integrity */
663 ✗ block_size = clamp_size(h, block_size, block_pos, filesize);
664 ✗ if (s->cache_data) {
665 av_assert1(block_pos + block_size <= s->cache_size);
666 ✗ tmp = s->cache_data + block_pos;
667 } else {
668 ✗ tmp = s->tmp_buf;
669 ✗ ret = read_cache(s, tmp, block_size, block_pos);
670 ✗ if (ret < 0) {
671 ✗ av_log(h, AV_LOG_ERROR, "Failed to read from cache file: %s\n", av_err2str(ret));
672 ✗ if (ret == AVERROR_EOF) { /* e.g. cache appears truncated? */
673 ✗ if (s->retry_corrupt) {
674 ✗ s->num_corrupt++;
675 ✗ goto read_block;
676 }
677 ✗ ret = AVERROR(EIO); /* don't propagate EOF to caller */
678 }
679 ✗ return ret;
680 }
681 }
682
683 ✗ uint32_t crc = get_block_crc(tmp, block_size);
684 ✗ if (crc != state) {
685 ✗ av_log(h, AV_LOG_ERROR, "Cache corruption detected for block 0x%"PRIx64" at "
686 "offset 0x%"PRIx64": expected CRC: 0x%08X, got: 0x%08X\n",
687 block_id, block_pos, state, crc);
688 ✗ if (s->retry_corrupt) {
689 ✗ s->num_corrupt++;
690 ✗ goto read_block;
691 }
692 ✗ return AVERROR(EIO);
693 } else
694 ✗ s->num_corrupt = 0; /* reset corrupt block count on success */
695
696 ✗ tmp += (ptrdiff_t) offset;
697 ✗ size = FFMIN(size, block_size - offset);
698 ✗ if (size <= 0)
699 ✗ return AVERROR_EOF;
700 ✗ if (s->verify) {
701 ✗ verify_read = 1;
702 ✗ break; /* fall through to the cache miss logic */
703 }
704
705 ✗ memcpy(buf, tmp, size);
706 ✗ s->nb_hit++;
707 ✗ s->pos += size;
708 ✗ return size;
709
710 ✗ case BLOCK_FAILED:
711 ✗ if (s->retry_errors)
712 ✗ goto read_block;
713 ✗ return AVERROR(EIO);
714
715 ✗ read_block:
716 ✗ if (s->num_corrupt == MAX_CORRUPT_BLOCKS) {
717 ✗ av_log(h, AV_LOG_ERROR, "Too many consecutive corrupt blocks; "
718 "assuming cache file is completely broken.\n");
719 ✗ s->num_corrupt++; /* silence this log on subsequent reads */
720 }
721 av_fallthrough;
722
723 case BLOCK_NONE:
724 ✗ if (s->read_only || s->write_err || !s->inner)
725 break; /* don't mark block as pending */
726 ✗ else if (s->cache_size_max) {
727 ✗ int64_t cached = atomic_load_explicit(&s->spacemap->blocks_cached,
728 memory_order_relaxed);
729 ✗ if (cached >= s->blocks_max) {
730 ✗ av_log(h, AV_LOG_WARNING, "Cache size limit reached (%"PRId64" "
731 "blocks = %"PRId64" bytes), switching to read-only mode.\n",
732 ✗ s->blocks_max, s->blocks_max << s->block_shift);
733 ✗ s->read_only = 1;
734 ✗ break;
735 }
736 }
737
738 ✗ if (atomic_compare_exchange_strong_explicit(&block->state, &state,
739 BLOCK_PENDING,
740 memory_order_acquire,
741 memory_order_acquire))
742 {
743 /* Acquired pending state, proceed to fetch the block */
744 ✗ acquired = 1;
745 ✗ allocated = (state == BLOCK_NONE || state == BLOCK_FAILED);
746 ✗ state = BLOCK_PENDING;
747 ✗ break;
748 }
749 /* CAS failed, another thread changed the state; reload it */
750 ✗ goto retry;
751
752 ✗ case BLOCK_PENDING:
753 /* Another thread is busy fetching this block, wait for it to finish */
754 ✗ if (!s->timeout) {
755 ✗ break; /* no timeout requested, immediately race to fetch block */
756 ✗ } else if (pending_since) {
757 ✗ int64_t new = av_gettime_relative();
758 ✗ if (new - pending_since >= s->timeout)
759 ✗ break; /* timeout expired, try to fetch the block ourselves */
760 } else {
761 ✗ pending_since = av_gettime_relative();
762 }
763
764 ✗ if (h->flags & AVIO_FLAG_NONBLOCK)
765 ✗ return AVERROR(EAGAIN);
766
767 /* Make sure we try a few times before giving up */
768 ✗ av_usleep(FFMIN(s->timeout >> 4, 10000));
769 ✗ if (ff_check_interrupt(&h->interrupt_callback))
770 ✗ return AVERROR_EXIT;
771
772 ✗ state = atomic_load_explicit(&block->state, memory_order_acquire);
773 ✗ goto retry;
774 }
775
776 /* Release pending state on failure to avoid stalling other threads */
777 #define RELEASE_PENDING(block, state) \
778 do { \
779 if (acquired) { \
780 av_assert1(state == BLOCK_PENDING); \
781 atomic_compare_exchange_strong_explicit( \
782 &block->state, &state, BLOCK_NONE, memory_order_relaxed, \
783 memory_order_relaxed); \
784 } \
785 } while (0)
786
787 /* Cache miss, fetch this block from underlying protocol */
788 ✗ s->nb_miss++;
789
790 ✗ if (!s->inner) {
791 ✗ av_log(h, AV_LOG_ERROR, "Cache miss for block 0x%"PRIx64" at offset "
792 "0x%"PRIx64", but underlying protocol is not available!\n",
793 block_id, block_pos);
794 ✗ av_assert0(!acquired);
795 ✗ return AVERROR(EIO);
796 }
797
798 ✗ const int read_only = s->read_only || s->write_err || verify_read;
799 ✗ int64_t inner_pos = read_only ? s->pos : block_pos;
800 ✗ if (s->inner_pos != inner_pos) {
801 ✗ inner_pos = ffurl_seek(s->inner, inner_pos, SEEK_SET);
802 ✗ if (inner_pos < 0) {
803 ✗ av_log(h, AV_LOG_ERROR, "Failed to seek underlying protocol: %s\n",
804 ✗ av_err2str(inner_pos));
805 ✗ RELEASE_PENDING(block, state);
806 ✗ return inner_pos;
807 }
808
809 ✗ av_log(h, AV_LOG_DEBUG, "Inner seek to 0x%"PRIx64"\n", inner_pos);
810 ✗ s->inner_pos = inner_pos;
811 }
812
813 ✗ if (read_only) {
814 /* Directly defer to the underlying protocol */
815 ✗ ret = ffurl_read(s->inner, buf, size);
816 ✗ if (ret < 0) {
817 av_assert1(!acquired);
818 ✗ return ret;
819 } else {
820 ✗ s->inner_pos = inner_pos + ret;
821 }
822
823 /* Verify the read data against the cached data if requested */
824 ✗ if (verify_read && memcmp(buf, tmp, ret)) {
825 ✗ av_log(h, AV_LOG_ERROR, "Cache verification failed for %d bytes "
826 "in block 0x%"PRIx64" at offset 0x%"PRIx64" + %"PRId64"!\n",
827 ret, block_id, block_pos, offset);
828 ✗ return AVERROR(EIO);
829 }
830
831 ✗ s->pos = s->inner_pos;
832 ✗ return ret;
833 }
834
835 ✗ int write_back = 1;
836 ✗ if (s->cache_data && acquired) {
837 /* Read directly into memory mapped cache file */
838 ✗ tmp = s->cache_data + block_pos;
839 ✗ write_back = 0;
840 ✗ } else if (size >= block_size && !offset) {
841 /* Read directly into output buffer if aligned and large enough */
842 ✗ tmp = buf;
843 } else {
844 /* Read into temporary buffer and copy later */
845 ✗ tmp = s->tmp_buf;
846 }
847
848 /* Try and fetch the entire block */
849 ✗ av_assert0(inner_pos == block_pos);
850 ✗ int bytes_read = 0;
851 ✗ while (bytes_read < block_size) {
852 ✗ ret = ffurl_read(s->inner, &tmp[bytes_read], block_size - bytes_read);
853 ✗ if (!ret || ret == AVERROR_EOF)
854 break;
855 ✗ else if (ret < 0) {
856 ✗ av_log(h, AV_LOG_ERROR, "Failed to read block 0x%"PRIx64": %s\n",
857 ✗ block_id, av_err2str(ret));
858 ✗ if (ret == AVERROR(EAGAIN) || ret == AVERROR_EXIT) {
859 ✗ RELEASE_PENDING(block, state);
860 ✗ return ret; /* transient error, allow retries */
861 }
862
863 /* Try to mark block as failed; ignore errors - any mismatch
864 * here will mean that either another thread already marked it
865 * as failed, or successfully cached it in the meantime */
866 ✗ atomic_compare_exchange_strong_explicit(&block->state, &state,
867 BLOCK_FAILED,
868 memory_order_relaxed,
869 memory_order_relaxed);
870 ✗ return ret;
871 }
872
873 ✗ bytes_read += ret;
874 ✗ s->inner_pos += ret;
875 }
876
877 ✗ if (bytes_read < block_size) {
878 /* Learned location of true EOF, update filesize */
879 ✗ ret = set_filesize(h, inner_pos + bytes_read);
880 ✗ if (ret < 0) {
881 ✗ RELEASE_PENDING(block, state);
882 ✗ return ret;
883 }
884 }
885
886 ✗ if (bytes_read > 0) {
887 ✗ ret = write_back ? write_cache(s, tmp, bytes_read, block_pos) : 0;
888 ✗ if (ret < 0) {
889 ✗ if (ret != AVERROR(EINTR)) {
890 ✗ av_log(h, AV_LOG_ERROR, "Failed to write to cache file: %s\n",
891 ✗ av_err2str(ret));
892 ✗ s->write_err = 1;
893 }
894 ✗ RELEASE_PENDING(block, state);
895 } else {
896 ✗ uint32_t crc = get_block_crc(tmp, bytes_read);
897 ✗ av_log(h, AV_LOG_TRACE, "Cached %d bytes to block 0x%"PRIx64" at "
898 "offset 0x%"PRIx64", CRC 0x%08X\n", bytes_read, block_id,
899 block_pos, crc);
900 ✗ atomic_store_explicit(&block->state, crc, memory_order_release);
901 ✗ if (allocated)
902 ✗ atomic_fetch_add_explicit(&s->spacemap->blocks_cached, 1, memory_order_release);
903 }
904 } else {
905 ✗ RELEASE_PENDING(block, state);
906 ✗ return AVERROR_EOF;
907 }
908
909 ✗ size = FFMIN(bytes_read - offset, size);
910 ✗ if (size <= 0)
911 ✗ return AVERROR_EOF;
912 ✗ if (tmp != buf)
913 ✗ memcpy(buf, &tmp[offset], size);
914 ✗ s->pos += size;
915 ✗ return size;
916 }
917
918 ✗ static int64_t shared_seek(URLContext *h, int64_t pos, int whence)
919 {
920 ✗ SharedContext *s = h->priv_data;
921 int64_t res;
922 ✗ if (!s->spacemap)
923 ✗ return AVERROR(EIO);
924
925 ✗ const int64_t filesize = get_filesize(h);
926 ✗ if (filesize < 0)
927 ✗ return filesize;
928
929 ✗ switch (whence) {
930 ✗ case AVSEEK_SIZE:
931 ✗ if (filesize)
932 ✗ return filesize;
933 ✗ res = s->inner ? ffurl_seek(s->inner, pos, whence) : AVERROR(ENOSYS);
934 ✗ if (res > 0) {
935 ✗ if (set_filesize(h, res) < 0)
936 ✗ return AVERROR(EINVAL);
937 ✗ } else if (is_ignorable_error(res) && s->ignore_errors) {
938 ✗ av_log(h, AV_LOG_WARNING, "Underlying URL failed to get size: %s. "
939 ✗ "Continuing with cache file only.\n", av_err2str(res));
940 ✗ ffurl_closep(&s->inner);
941 ✗ res = AVERROR(ENOSYS);
942 }
943 ✗ return res;
944 ✗ case SEEK_SET:
945 ✗ break;
946 ✗ case SEEK_CUR:
947 ✗ pos += s->pos;
948 ✗ break;
949 ✗ case SEEK_END:
950 ✗ if (filesize) {
951 ✗ pos += filesize;
952 ✗ break;
953 }
954
955 /* Defer to underlying protocol if filesize is unknown */
956 ✗ res = s->inner ? ffurl_seek(s->inner, pos, whence) : AVERROR(ENOSYS);
957 ✗ if (is_ignorable_error(res) && s->ignore_errors) {
958 ✗ av_log(h, AV_LOG_WARNING, "Underlying URL failed to seek: %s. "
959 ✗ "Continuing with cache file only.\n", av_err2str(res));
960 ✗ ffurl_closep(&s->inner);
961 ✗ return AVERROR(ENOSYS);
962 ✗ } else if (res < 0)
963 ✗ return res;
964
965 /* Opportunistically update known filesize */
966 ✗ if (set_filesize(h, res - pos) < 0)
967 ✗ return AVERROR(EINVAL);
968 ✗ av_log(h, AV_LOG_DEBUG, "Inner seek to 0x%"PRIx64"\n", res);
969 ✗ return s->pos = s->inner_pos = res;
970 ✗ default:
971 ✗ return AVERROR(EINVAL);
972 }
973
974 ✗ if (pos < 0)
975 ✗ return AVERROR(EINVAL);
976
977 ✗ av_log(h, AV_LOG_DEBUG, "Virtual seek to 0x%"PRIx64"\n", pos);
978 ✗ return s->pos = pos;
979 }
980
981 ✗ static int shared_get_file_handle(URLContext *h)
982 {
983 ✗ SharedContext *s = h->priv_data;
984 ✗ return s->inner ? ffurl_get_file_handle(s->inner) : -1;
985 }
986
987 ✗ static int shared_get_short_seek(URLContext *h)
988 {
989 ✗ SharedContext *s = h->priv_data;
990 ✗ int ret = s->inner ? ffurl_get_short_seek(s->inner) : 0;
991 ✗ return ret > 0 ? FFMAX(ret, s->block_size) : s->block_size;
992 }
993
994 #define OFFSET(x) offsetof(SharedContext, x)
995 #define D AV_OPT_FLAG_DECODING_PARAM
996
997 static const AVOption options[] = {
998 { "cache_dir", "Directory path for shared file cache", OFFSET(cache_dir), AV_OPT_TYPE_STRING, {.str = NULL}, .flags = D },
999 { "block_shift", "Set the base 2 logarithm of the block size", OFFSET(block_shift), AV_OPT_TYPE_INT, {.i64 = 15}, 9, 30, .flags = D },
1000 { "read_only", "Don't write data to the cache, only read from it", OFFSET(read_only), AV_OPT_TYPE_BOOL, {.i64 = 0}, 0, 1, .flags = D },
1001 { "cache_verify", "Verify correctness of the cache against the source", OFFSET(verify), AV_OPT_TYPE_BOOL, {.i64 = 0}, 0, 1, .flags = D },
1002 { "cache_timeout", "Time in us to wait before re-fetching pending blocks", OFFSET(timeout), AV_OPT_TYPE_INT64, {.i64 = 10000}, 0, INT64_MAX, .flags = D },
1003 { "ignore_errors", "Continue even if the inner URL failed", OFFSET(ignore_errors), AV_OPT_TYPE_BOOL, {.i64 = 0}, 0, 1, .flags = D },
1004 { "retry_errors", "Re-request blocks even if they previously failed", OFFSET(retry_errors), AV_OPT_TYPE_BOOL, {.i64 = 1}, 0, 1, .flags = D },
1005 { "retry_corrupt", "Re-request blocks that fail the CRC check", OFFSET(retry_corrupt), AV_OPT_TYPE_BOOL, {.i64 = 1}, 0, 1, .flags = D },
1006 { "cache_size_max", "Limit the maximum amount of data cached", OFFSET(cache_size_max), AV_OPT_TYPE_INT64, {.i64 = 0}, 0, INT64_MAX, .flags = D },
1007 {0},
1008 };
1009
1010 static const AVClass shared_context_class = {
1011 .class_name = "shared",
1012 .item_name = av_default_item_name,
1013 .option = options,
1014 .version = LIBAVUTIL_VERSION_INT,
1015 };
1016
1017 const URLProtocol ff_shared_protocol = {
1018 .name = "shared",
1019 .url_open2 = shared_open,
1020 .url_read = shared_read,
1021 .url_seek = shared_seek,
1022 .url_close = shared_close,
1023 .url_get_file_handle = shared_get_file_handle,
1024 .url_get_short_seek = shared_get_short_seek,
1025 .priv_data_size = sizeof(SharedContext),
1026 .priv_data_class = &shared_context_class,
1027 };
1028