| Directory: | cvmfs/ |
|---|---|
| File: | cvmfs/ingestion/task_chunk.cc |
| Date: | 2025-11-30 02:35:17 |
| Exec | Total | Coverage | |
|---|---|---|---|
| Lines: | 89 | 91 | 97.8% |
| Branches: | 90 | 149 | 60.4% |
| Line | Branch | Exec | Source |
|---|---|---|---|
| 1 | /** | ||
| 2 | * This file is part of the CernVM File System. | ||
| 3 | */ | ||
| 4 | |||
| 5 | #include "ingestion/task_chunk.h" | ||
| 6 | |||
| 7 | #include <unistd.h> | ||
| 8 | |||
| 9 | #include <cassert> | ||
| 10 | |||
| 11 | #include "util/exception.h" | ||
| 12 | |||
| 13 | /** | ||
| 14 | * The tags from the read stage in the pipeline and the tags given in the | ||
| 15 | * chunking stage can safely overlap. Nevertheless, debugging might be easier | ||
| 16 | * if they don't. So let's start with a high number. | ||
| 17 | */ | ||
| 18 | atomic_int64 TaskChunk::tag_seq_ = 2 << 28; | ||
| 19 | |||
| 20 | /** | ||
| 21 | * Consumes the stream of input blocks and produces new output blocks according | ||
| 22 | * to cut marks. The output blocks correspond to chunks. | ||
| 23 | */ | ||
| 24 | 25815300 | void TaskChunk::Process(BlockItem *input_block) { | |
| 25 | 25815300 | FileItem *file_item = input_block->file_item(); | |
| 26 | 25814242 | const int64_t input_tag = input_block->tag(); | |
| 27 |
3/4✓ Branch 0 taken 25836874 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 25825098 times.
✓ Branch 3 taken 11776 times.
|
25835724 | assert((file_item != NULL) && (input_tag >= 0)); |
| 28 | |||
| 29 | 25825098 | ChunkInfo chunk_info; | |
| 30 | // Do we see blocks of the file for the first time? | ||
| 31 |
3/4✓ Branch 1 taken 25870684 times.
✗ Branch 2 not taken.
✓ Branch 3 taken 11493862 times.
✓ Branch 4 taken 14376822 times.
|
25825650 | if (!tag_map_.Lookup(input_tag, &chunk_info)) { |
| 32 | // We may have only regular chunks, only a bulk chunk, or both. We may | ||
| 33 | // end up in a situation where we produced only a single non-bulk chunk. | ||
| 34 | // This needs to be fixed up later in the pipeline by the write task. | ||
| 35 |
2/2✓ Branch 1 taken 690 times.
✓ Branch 2 taken 11491056 times.
|
11493862 | if (file_item->may_have_chunks()) { |
| 36 |
2/4✓ Branch 1 taken 690 times.
✗ Branch 2 not taken.
✓ Branch 4 taken 690 times.
✗ Branch 5 not taken.
|
690 | chunk_info.next_chunk = new ChunkItem(file_item, 0); |
| 37 | 690 | chunk_info.output_tag_chunk = atomic_xadd64(&tag_seq_, 1); | |
| 38 |
2/2✓ Branch 1 taken 460 times.
✓ Branch 2 taken 184 times.
|
644 | if (file_item->has_legacy_bulk_chunk()) { |
| 39 |
2/4✓ Branch 1 taken 506 times.
✗ Branch 2 not taken.
✓ Branch 4 taken 460 times.
✗ Branch 5 not taken.
|
460 | chunk_info.bulk_chunk = new ChunkItem(file_item, 0); |
| 40 | } | ||
| 41 | } else { | ||
| 42 |
2/4✓ Branch 1 taken 11494000 times.
✗ Branch 2 not taken.
✓ Branch 4 taken 11488342 times.
✗ Branch 5 not taken.
|
11491056 | chunk_info.bulk_chunk = new ChunkItem(file_item, 0); |
| 43 | } | ||
| 44 | |||
| 45 |
1/2✓ Branch 0 taken 11488986 times.
✗ Branch 1 not taken.
|
11488986 | if (chunk_info.bulk_chunk != NULL) { |
| 46 |
1/2✓ Branch 1 taken 11483834 times.
✗ Branch 2 not taken.
|
11488986 | chunk_info.bulk_chunk->MakeBulkChunk(); |
| 47 | 11483834 | chunk_info.bulk_chunk->set_size(file_item->size()); | |
| 48 | 11482546 | chunk_info.output_tag_bulk = atomic_xadd64(&tag_seq_, 1); | |
| 49 | } | ||
| 50 |
1/2✓ Branch 1 taken 11485812 times.
✗ Branch 2 not taken.
|
11500256 | tag_map_.Insert(input_tag, chunk_info); |
| 51 | } | ||
| 52 |
3/4✓ Branch 0 taken 472696 times.
✓ Branch 1 taken 25389938 times.
✗ Branch 2 not taken.
✓ Branch 3 taken 472696 times.
|
25862634 | assert((chunk_info.bulk_chunk != NULL) || (chunk_info.next_chunk != NULL)); |
| 53 | |||
| 54 | 25862634 | BlockItem *output_block_bulk = NULL; | |
| 55 |
2/2✓ Branch 0 taken 25412938 times.
✓ Branch 1 taken 449696 times.
|
25862634 | if (chunk_info.bulk_chunk != NULL) { |
| 56 |
2/4✓ Branch 1 taken 25416664 times.
✗ Branch 2 not taken.
✓ Branch 4 taken 25393756 times.
✗ Branch 5 not taken.
|
25412938 | output_block_bulk = new BlockItem(chunk_info.output_tag_bulk, allocator_); |
| 57 |
1/2✓ Branch 1 taken 25383958 times.
✗ Branch 2 not taken.
|
25393756 | output_block_bulk->SetFileItem(file_item); |
| 58 |
1/2✓ Branch 1 taken 25388926 times.
✗ Branch 2 not taken.
|
25383958 | output_block_bulk->SetChunkItem(chunk_info.bulk_chunk); |
| 59 | } | ||
| 60 | |||
| 61 | 25838622 | ChunkDetector *chunk_detector = file_item->chunk_detector(); | |
| 62 |
2/3✓ Branch 1 taken 11477348 times.
✓ Branch 2 taken 14421028 times.
✗ Branch 3 not taken.
|
25861760 | switch (input_block->type()) { |
| 63 | 11477348 | case BlockItem::kBlockStop: | |
| 64 | // End of the file, no more new chunks | ||
| 65 | 11477348 | file_item->set_is_fully_chunked(); | |
| 66 |
1/2✓ Branch 0 taken 11493356 times.
✗ Branch 1 not taken.
|
11493310 | if (output_block_bulk) |
| 67 |
1/2✓ Branch 1 taken 11477256 times.
✗ Branch 2 not taken.
|
11493356 | output_block_bulk->MakeStop(); |
| 68 |
2/2✓ Branch 0 taken 690 times.
✓ Branch 1 taken 11476520 times.
|
11477210 | if (chunk_info.next_chunk != NULL) { |
| 69 |
1/2✗ Branch 2 not taken.
✓ Branch 3 taken 690 times.
|
690 | assert(file_item->size() >= chunk_info.next_chunk->offset()); |
| 70 | 1380 | chunk_info.next_chunk->set_size(file_item->size() | |
| 71 | 690 | - chunk_info.next_chunk->offset()); | |
| 72 | BlockItem *block_stop = new BlockItem(chunk_info.output_tag_chunk, | ||
| 73 |
2/4✓ Branch 1 taken 690 times.
✗ Branch 2 not taken.
✓ Branch 4 taken 690 times.
✗ Branch 5 not taken.
|
690 | allocator_); |
| 74 |
1/2✓ Branch 1 taken 690 times.
✗ Branch 2 not taken.
|
690 | block_stop->SetFileItem(file_item); |
| 75 |
1/2✓ Branch 1 taken 690 times.
✗ Branch 2 not taken.
|
690 | block_stop->SetChunkItem(chunk_info.next_chunk); |
| 76 |
1/2✓ Branch 1 taken 690 times.
✗ Branch 2 not taken.
|
690 | block_stop->MakeStop(); |
| 77 |
1/2✓ Branch 1 taken 690 times.
✗ Branch 2 not taken.
|
690 | tubes_out_->Dispatch(block_stop); |
| 78 | } | ||
| 79 |
1/2✓ Branch 1 taken 11479004 times.
✗ Branch 2 not taken.
|
11477210 | tag_map_.Erase(input_tag); |
| 80 | 11479004 | break; | |
| 81 | |||
| 82 | 14421028 | case BlockItem::kBlockData: | |
| 83 |
2/2✓ Branch 0 taken 13943686 times.
✓ Branch 1 taken 477342 times.
|
14421028 | if (output_block_bulk) { |
| 84 |
2/2✓ Branch 0 taken 1246048 times.
✓ Branch 1 taken 12697638 times.
|
13943686 | if (chunk_info.next_chunk != NULL) { |
| 85 | // Reserve zero-copy for the regular chunk | ||
| 86 |
1/2✓ Branch 3 taken 1243518 times.
✗ Branch 4 not taken.
|
1246048 | output_block_bulk->MakeDataCopy(input_block->data(), |
| 87 | input_block->size()); | ||
| 88 | } else { | ||
| 89 | // There is only the bulk chunk, zero copy | ||
| 90 |
1/2✓ Branch 1 taken 12694418 times.
✗ Branch 2 not taken.
|
12697638 | output_block_bulk->MakeDataMove(input_block); |
| 91 | } | ||
| 92 | } | ||
| 93 | |||
| 94 |
2/2✓ Branch 0 taken 1720492 times.
✓ Branch 1 taken 12694786 times.
|
14415278 | if (chunk_info.next_chunk != NULL) { |
| 95 | 1720492 | unsigned offset_in_block = 0; | |
| 96 | 1720492 | uint64_t cut_mark = 0; | |
| 97 |
3/4✓ Branch 1 taken 1968432 times.
✗ Branch 2 not taken.
✓ Branch 3 taken 248354 times.
✓ Branch 4 taken 1720078 times.
|
1968846 | while ((cut_mark = chunk_detector->FindNextCutMark(input_block)) != 0) { |
| 98 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 248354 times.
|
248354 | assert(cut_mark >= chunk_info.offset + offset_in_block); |
| 99 | 248354 | const uint64_t cut_mark_in_block = cut_mark - chunk_info.offset; | |
| 100 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 248354 times.
|
248354 | assert(cut_mark_in_block >= offset_in_block); |
| 101 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 248354 times.
|
248354 | assert(cut_mark_in_block <= input_block->size()); |
| 102 | 248354 | const unsigned tail_size = cut_mark_in_block - offset_in_block; | |
| 103 | |||
| 104 |
1/2✓ Branch 0 taken 248354 times.
✗ Branch 1 not taken.
|
248354 | if (tail_size > 0) { |
| 105 | BlockItem *block_tail = new BlockItem(chunk_info.output_tag_chunk, | ||
| 106 |
2/4✓ Branch 1 taken 248354 times.
✗ Branch 2 not taken.
✓ Branch 4 taken 248354 times.
✗ Branch 5 not taken.
|
248354 | allocator_); |
| 107 |
1/2✓ Branch 1 taken 248354 times.
✗ Branch 2 not taken.
|
248354 | block_tail->SetFileItem(file_item); |
| 108 |
1/2✓ Branch 1 taken 248354 times.
✗ Branch 2 not taken.
|
248354 | block_tail->SetChunkItem(chunk_info.next_chunk); |
| 109 |
1/2✓ Branch 2 taken 248354 times.
✗ Branch 3 not taken.
|
248354 | block_tail->MakeDataCopy(input_block->data() + offset_in_block, |
| 110 | tail_size); | ||
| 111 |
1/2✓ Branch 1 taken 248354 times.
✗ Branch 2 not taken.
|
248354 | tubes_out_->Dispatch(block_tail); |
| 112 | } | ||
| 113 | |||
| 114 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 248354 times.
|
248354 | assert(cut_mark >= chunk_info.next_chunk->offset()); |
| 115 | // If the cut mark happens to at the end of file, let the final | ||
| 116 | // incoming stop block schedule dispatch of the chunk stop block | ||
| 117 |
2/2✓ Branch 1 taken 248308 times.
✓ Branch 2 taken 46 times.
|
248354 | if (cut_mark < file_item->size()) { |
| 118 | 248308 | chunk_info.next_chunk->set_size(cut_mark | |
| 119 | 248308 | - chunk_info.next_chunk->offset()); | |
| 120 | BlockItem *block_stop = new BlockItem(chunk_info.output_tag_chunk, | ||
| 121 |
2/4✓ Branch 1 taken 248308 times.
✗ Branch 2 not taken.
✓ Branch 4 taken 248308 times.
✗ Branch 5 not taken.
|
248308 | allocator_); |
| 122 |
1/2✓ Branch 1 taken 248308 times.
✗ Branch 2 not taken.
|
248308 | block_stop->SetFileItem(file_item); |
| 123 |
1/2✓ Branch 1 taken 248308 times.
✗ Branch 2 not taken.
|
248308 | block_stop->SetChunkItem(chunk_info.next_chunk); |
| 124 |
1/2✓ Branch 1 taken 248308 times.
✗ Branch 2 not taken.
|
248308 | block_stop->MakeStop(); |
| 125 |
1/2✓ Branch 1 taken 248308 times.
✗ Branch 2 not taken.
|
248308 | tubes_out_->Dispatch(block_stop); |
| 126 | |||
| 127 |
2/4✓ Branch 1 taken 248308 times.
✗ Branch 2 not taken.
✓ Branch 4 taken 248308 times.
✗ Branch 5 not taken.
|
248308 | chunk_info.next_chunk = new ChunkItem(file_item, cut_mark); |
| 128 | 248308 | chunk_info.output_tag_chunk = atomic_xadd64(&tag_seq_, 1); | |
| 129 | } | ||
| 130 | 248354 | offset_in_block = cut_mark_in_block; | |
| 131 | } | ||
| 132 | 1720078 | chunk_info.offset += offset_in_block; | |
| 133 | |||
| 134 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 1720262 times.
|
1720078 | assert(input_block->size() >= offset_in_block); |
| 135 | 1720262 | const unsigned tail_size = input_block->size() - offset_in_block; | |
| 136 |
2/2✓ Branch 0 taken 1717686 times.
✓ Branch 1 taken 2530 times.
|
1720216 | if (tail_size > 0) { |
| 137 | BlockItem *block_tail = new BlockItem(chunk_info.output_tag_chunk, | ||
| 138 |
2/4✓ Branch 1 taken 1717410 times.
✗ Branch 2 not taken.
✓ Branch 4 taken 1717042 times.
✗ Branch 5 not taken.
|
1717686 | allocator_); |
| 139 |
1/2✓ Branch 1 taken 1716950 times.
✗ Branch 2 not taken.
|
1717042 | block_tail->SetFileItem(file_item); |
| 140 |
1/2✓ Branch 1 taken 1716950 times.
✗ Branch 2 not taken.
|
1716950 | block_tail->SetChunkItem(chunk_info.next_chunk); |
| 141 |
1/2✓ Branch 2 taken 1718054 times.
✗ Branch 3 not taken.
|
1716950 | block_tail->MakeDataCopy(input_block->data() + offset_in_block, |
| 142 | tail_size); | ||
| 143 |
1/2✓ Branch 1 taken 1717548 times.
✗ Branch 2 not taken.
|
1718054 | tubes_out_->Dispatch(block_tail); |
| 144 | 1717548 | chunk_info.offset += tail_size; | |
| 145 | } | ||
| 146 | |||
| 147 | // Delete data from incoming block | ||
| 148 |
1/2✓ Branch 1 taken 1720308 times.
✗ Branch 2 not taken.
|
1720078 | input_block->Reset(); |
| 149 | } | ||
| 150 | |||
| 151 |
1/2✓ Branch 1 taken 14423972 times.
✗ Branch 2 not taken.
|
14415094 | tag_map_.Insert(input_tag, chunk_info); |
| 152 | 14423972 | break; | |
| 153 | |||
| 154 | ✗ | default: | |
| 155 | ✗ | PANIC(NULL); | |
| 156 | } | ||
| 157 | |||
| 158 |
2/2✓ Branch 0 taken 25892166 times.
✓ Branch 1 taken 10810 times.
|
25902976 | delete input_block; |
| 159 |
2/2✓ Branch 0 taken 25447070 times.
✓ Branch 1 taken 472696 times.
|
25919766 | if (output_block_bulk) |
| 160 |
1/2✓ Branch 1 taken 25372090 times.
✗ Branch 2 not taken.
|
25447070 | tubes_out_->Dispatch(output_block_bulk); |
| 161 | 25844786 | } | |
| 162 |