FFmpeg
Loading...
Searching...
No Matches
shared.c
Go to the documentation of this file.
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
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
78static 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;
88 av_hash_update(ctx, (const uint8_t *) &version, sizeof(version));
89 av_hash_update(ctx, (const uint8_t *) uri, strlen(uri));
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
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
110static 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
123typedef struct Block {
124 atomic_uint state; /* enum BlockState */
125} Block;
126
127typedef struct Spacemap {
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
137} Spacemap;
138
139static_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
156DEF_SET_ONCE(unsigned char, uchar)
157DEF_SET_ONCE(unsigned int, uint)
158DEF_SET_ONCE(unsigned short, ushort)
159DEF_SET_ONCE(unsigned long long, ullong)
160
161typedef struct SharedContext {
162 AVClass *class;
165
166 /* options */
168 int block_shift; ///< requested shift; updated on init if it disagrees
176
177 /* misc state */
178 int64_t pos; ///< current logical position
179 uint8_t *tmp_buf;
181 int write_err; ///< write error occurred
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
189 off_t cache_size; ///< size of mapped memory region (for unmapping)
190 int fd;
191
192 /* space map */
194 char *map_path;
195 off_t map_size;
196 int mapfd;
197
198 /* statistics */
202
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
224static int spacemap_init(URLContext *h, const uint8_t hash[HASH_SIZE]);
226
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
241static 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
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
279static 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
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 */
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
380fail:
381 if (ret < 0)
383 return ret;
384}
385
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
411static 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
451skip_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
471fail:
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
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 */
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);
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
514static 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",
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
569static 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
589static 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
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
619static 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
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];
648 int64_t pending_since = 0;
649 int verify_read = 0, acquired = 0, allocated = 0;
650
651retry:
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 */
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
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 }
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,
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
742 {
743 /* Acquired pending state, proceed to fetch the block */
744 acquired = 1;
745 allocated = (state == BLOCK_NONE || state == BLOCK_FAILED);
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) {
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
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));
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) {
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 */
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) {
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 }
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);
901 if (allocated)
902 atomic_fetch_add_explicit(&s->spacemap->blocks_cached, 1, memory_order_release);
903 }
904 } else {
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
919{
920 SharedContext *s = h->priv_data;
921 int64_t res;
922 if (!s->spacemap)
923 return AVERROR(EIO);
924
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
982{
983 SharedContext *s = h->priv_data;
984 return s->inner ? ffurl_get_file_handle(s->inner) : -1;
985}
986
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
997static 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
1011 .class_name = "shared",
1012 .item_name = av_default_item_name,
1013 .option = options,
1014 .version = LIBAVUTIL_VERSION_INT,
1015};
1016
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};
static int read_block(ALSDecContext *ctx, ALSBlockData *bd)
Read the block data.
Definition alsdec.c:1031
static uint8_t hash[HASH_SIZE]
static AVFormatContext * ctx
static av_cold void close(AVCodecParserContext *s)
Definition apv_parser.c:197
simple assert() macros that are a bit more flexible than ISO C assert().
#define av_assert1(cond)
assert() equivalent, that does not lie in speed critical code.
Definition avassert.h:58
#define av_assert0(cond)
assert() equivalent, that is always enabled.
Definition avassert.h:42
#define D
Definition avdct.c:35
int ff_check_interrupt(AVIOInterruptCB *cb)
Check if the user has requested to interrupt a blocking function associated with cb.
Definition avio.c:929
int ffurl_open_whitelist(URLContext **puc, const char *filename, int flags, const AVIOInterruptCB *int_cb, AVDictionary **options, const char *whitelist, const char *blacklist, URLContext *parent)
Create an URLContext for accessing to the resource indicated by url, and open it.
Definition avio.c:466
int64_t ffurl_size(URLContext *h)
Return the filesize of the resource accessed by h, AVERROR(ENOSYS) if the operation is not supported ...
Definition avio.c:874
int ffurl_closep(URLContext **hh)
Close the resource accessed by the URLContext h, and free the memory used by it.
Definition avio.c:663
int ffurl_close(URLContext *h)
Definition avio.c:686
int ffurl_get_short_seek(void *urlcontext)
Return the current short seek threshold value for this URL.
Definition avio.c:913
int ffurl_get_file_handle(URLContext *h)
Return the file descriptor associated with this URL.
Definition avio.c:889
#define AVSEEK_SIZE
Passing this as the "whence" parameter to a seek function causes it to return the filesize without se...
Definition avio.h:468
#define AVIO_FLAG_NONBLOCK
Use non-blocking mode.
Definition avio.h:636
char * av_asprintf(const char *fmt,...)
Definition avstring.c:115
#define flags(name, subs,...)
Definition cbs_h264.c:74
#define i(width, name, range_min, range_max)
Definition cbs_h264.c:63
#define s(width, name)
Definition cbs_vp9.c:198
#define NULL
Definition coverity.c:32
long long int64_t
Definition coverity.c:34
Public header for CRC hash function implementation.
static int16_t block[64]
Definition dct.c:125
error code definitions
static struct @346255127015250356166251341105367306144006377143 state
static int64_t filesize(AVIOContext *pb)
Definition ffmpeg_mux.c:52
Misc file utilities.
#define fail
Definition test.h:479
@ AV_OPT_TYPE_INT64
Underlying C type is int64_t.
Definition opt.h:262
@ AV_OPT_TYPE_INT
Underlying C type is int.
Definition opt.h:258
@ AV_OPT_TYPE_BOOL
Underlying C type is int.
Definition opt.h:326
@ AV_OPT_TYPE_STRING
Underlying C type is a uint8_t* that is either NULL or points to a C string allocated with the av_mal...
Definition opt.h:275
const AVCRC * av_crc_get_table(AVCRCId crc_id)
Get an initialized standard CRC table.
Definition crc.c:389
uint32_t av_crc(const AVCRC *ctx, uint32_t crc, const uint8_t *buffer, size_t length)
Calculate the CRC of a block.
Definition crc.c:421
@ AV_CRC_32_IEEE
Definition crc.h:52
#define AVERROR_EXIT
Immediate exit was requested; the called function should not be restarted.
Definition error.h:58
#define AVERROR_BUG
Internal bug, also see AVERROR_BUG2.
Definition error.h:52
#define AVERROR_EOF
End of file.
Definition error.h:57
#define av_err2str(errnum)
Convenience macro, the return value should be used only directly in function arguments but never stan...
Definition error.h:122
#define AVERROR(e)
Definition error.h:45
void av_hash_freep(AVHashContext **ctx)
Free hash context and set hash context pointer to NULL.
Definition hash.c:248
void av_hash_init(AVHashContext *ctx)
Initialize or reset a hash context.
Definition hash.c:151
void av_hash_update(AVHashContext *ctx, const uint8_t *src, size_t len)
Update a hash context with additional data.
Definition hash.c:172
int av_hash_alloc(AVHashContext **ctx, const char *name)
Allocate a hash context for the algorithm specified by name.
Definition hash.c:114
void av_hash_final(AVHashContext *ctx, uint8_t *dst)
Finalize a hash context and compute the actual hash value.
Definition hash.c:193
#define AV_LOG_TRACE
Extremely verbose debugging, useful for libav* development.
Definition log.h:236
#define AV_LOG_DEBUG
Stuff which is only useful for libav* developers.
Definition log.h:231
#define AV_LOG_WARNING
Something somehow does not look correct.
Definition log.h:216
#define AV_LOG_VERBOSE
Detailed information.
Definition log.h:226
#define AV_LOG_ERROR
Something went wrong and cannot losslessly be recovered.
Definition log.h:210
const char * av_default_item_name(void *ptr)
Return the context name.
Definition log.c:241
int av_strstart(const char *str, const char *pfx, const char **ptr)
Return non-zero if pfx is a prefix of str.
Definition avstring.c:36
#define LIBAVUTIL_VERSION_INT
Definition version.h:85
int av_hash_get_size(const AVHashContext *ctx)
Definition hash.c:109
Generic hashing API.
const VDPAUPixFmtMap * map
unsigned offset
Definition libaomenc.c:763
static int shift(int a, int b)
Definition bonk.c:261
const char * arg
Definition jacosubdec.c:65
char * ff_data_to_hex(char *buf, const uint8_t *src, int size, int lowercase)
Write hexadecimal string corresponding to given binary data.
Definition utils.c:473
Macro definitions for various function/variable attributes.
#define av_fallthrough
Definition attributes.h:67
int av_file_map_shared(int fd, size_t size, void **bufptr)
Map the beginning of an open file into memory for shared read and write access.
Definition file.c:186
void av_file_unmap_shared(void *bufptr, size_t size)
Unmap the memory mapped by av_file_map_shared().
Definition file.c:198
int avpriv_open(const char *filename, int flags,...)
A wrapper for open() setting O_CLOEXEC.
Definition file_open.c:67
version
Definition libkvazaar.c:313
#define FFMIN(a, b)
Definition macros.h:49
#define FFMAX(a, b)
Definition macros.h:47
#define FFALIGN(x, a)
Definition macros.h:78
Memory handling functions.
#define av_malloc(s)
Definition ops_static.c:52
AVOptions.
miscellaneous OS support macros and functions.
#define O_BINARY
Definition os_support.h:36
const URLProtocol ff_shared_protocol
Definition shared.c:1017
static int shared_read(URLContext *h, unsigned char *buf, int size)
Definition shared.c:619
static int read_cache(SharedContext *s, uint8_t *buf, size_t size, off_t offset)
Definition shared.c:569
static int spacemap_grow(URLContext *h, int64_t block)
Definition shared.c:478
static int is_ignorable_error(int64_t err)
Definition shared.c:264
#define HEADER_VERSION
Definition shared.c:70
#define DEF_SET_ONCE(ctype, atype)
Definition shared.c:142
static int set_filesize(URLContext *h, int64_t new_size)
Definition shared.c:241
static int shared_get_short_seek(URLContext *h)
Definition shared.c:987
#define HASH_METHOD
This hash should be resistant against collision attacks, so that an attacker could not generate e....
Definition shared.c:67
static int shared_open(URLContext *h, const char *arg, int flags, AVDictionary **options)
Definition shared.c:279
static int64_t shared_seek(URLContext *h, int64_t pos, int whence)
Definition shared.c:918
#define RELEASE_PENDING(block, state)
static int shared_close(URLContext *h)
Definition shared.c:203
static int shared_get_file_handle(URLContext *h)
Definition shared.c:981
static int64_t get_filesize(URLContext *h)
Definition shared.c:227
static int hash_uri(uint8_t hash[HASH_SIZE], const char *uri)
Definition shared.c:78
static int cache_map(URLContext *h, int64_t filesize)
Definition shared.c:386
static int write_cache(SharedContext *s, const uint8_t *buf, size_t size, off_t offset)
Definition shared.c:589
static const AVClass shared_context_class
Definition shared.c:1010
static int spacemap_remap(URLContext *h, size_t map_size)
Definition shared.c:411
#define MAX_CORRUPT_BLOCKS
Hard watershed of consecutive failed blocks before we give up on the cache file altogether and assume...
Definition shared.c:76
static uint32_t get_block_crc(const uint8_t *block, size_t block_size)
Definition shared.c:110
#define OFFSET(x)
Definition shared.c:994
static int spacemap_init(URLContext *h, const uint8_t hash[HASH_SIZE])
Definition shared.c:514
#define HASH_SIZE
Definition shared.c:68
BlockState
Definition shared.c:98
@ BLOCK_NONE
block is not cached
Definition shared.c:100
@ BLOCK_PENDING
a thread is currently trying to write this block
Definition shared.c:101
@ BLOCK_FAILED
the underlying I/O source failed to read this block
Definition shared.c:102
#define HEADER_MAGIC
Definition shared.c:69
static int clamp_size(URLContext *h, int size, int64_t pos, int64_t filesize)
Definition shared.c:609
unsigned int pos
Definition spdifenc.c:431
@ memory_order_release
Definition stdatomic.h:32
@ memory_order_relaxed
Definition stdatomic.h:29
@ memory_order_acquire
Definition stdatomic.h:31
#define atomic_fetch_add_explicit(object, operand, order)
Definition stdatomic.h:297
unsigned char atomic_uchar
Definition stdatomic.h:60
FF_ATOMIC_ALIGN64 unsigned long long atomic_ullong
Definition stdatomic.h:68
#define atomic_compare_exchange_strong_explicit(object, expected, desired, success, failure)
Definition stdatomic.h:274
unsigned int atomic_uint
Definition stdatomic.h:64
#define atomic_load_explicit(object, order)
Definition stdatomic.h:247
unsigned short atomic_ushort
Definition stdatomic.h:62
#define atomic_load(object)
Definition stdatomic.h:250
#define atomic_store_explicit(object, desired, order)
Definition stdatomic.h:253
Describe the class of an AVClass context structure.
Definition log.h:76
uint32_t crc
Definition hash.c:70
AVOption.
Definition opt.h:428
atomic_uint state
Definition shared.c:124
char * cache_dir
Definition shared.c:167
int write_err
write error occurred
Definition shared.c:181
int block_size
Definition shared.c:180
int64_t nb_hit
Definition shared.c:199
int block_shift
requested shift; updated on init if it disagrees
Definition shared.c:168
int64_t filesize
once known
Definition shared.c:183
int64_t nb_miss
Definition shared.c:200
Spacemap * spacemap
Definition shared.c:193
uint8_t * tmp_buf
Definition shared.c:179
int64_t timeout
Definition shared.c:170
URLContext * inner
Definition shared.c:163
int retry_errors
Definition shared.c:172
uint8_t * cache_data
optional mapping of the cache file
Definition shared.c:187
int64_t inner_pos
Definition shared.c:164
int num_corrupt
Definition shared.c:182
int ignore_errors
Definition shared.c:171
int64_t blocks_max
maximum number of blocks to cache
Definition shared.c:184
int64_t pos
current logical position
Definition shared.c:178
off_t cache_size
size of mapped memory region (for unmapping)
Definition shared.c:189
int64_t cache_size_max
Definition shared.c:175
int read_only
Definition shared.c:169
off_t map_size
Definition shared.c:195
char * cache_path
Definition shared.c:188
char * map_path
Definition shared.c:194
int retry_corrupt
Definition shared.c:173
atomic_uint header_magic
Definition shared.c:128
atomic_ushort version
Definition shared.c:129
Block blocks[]
Definition shared.c:136
atomic_uchar hash[HASH_SIZE]
Definition shared.c:132
atomic_ushort block_shift
Definition shared.c:130
char reserved[72]
Definition shared.c:134
atomic_ullong blocks_cached
Definition shared.c:133
atomic_ullong filesize
Definition shared.c:131
Definition swscale.c:71
#define av_freep(p)
#define av_log(a,...)
static uint8_t tmp[40]
Definition aes_ctr.c:52
int av_usleep(unsigned usec)
Sleep for a period of time.
Definition time.c:93
int64_t av_gettime_relative(void)
Get the current time in microseconds since some unspecified starting point.
Definition time.c:57
int size
unbuffered private I/O API
static int64_t ffurl_seek(URLContext *h, int64_t pos, int whence)
Change the position that will be used by the next read/write operation on the resource accessed by h.
Definition url.h:225
static int ffurl_read(URLContext *h, uint8_t *buf, int size)
Read up to size bytes from the resource accessed by h, and store the read bytes in buf.
Definition url.h:184