Changed-lines coverage: PR changed C/C++ lines covered by tests: 85.74% (511/596) Uncovered changed code (with context): ================================================================================ src/Processors/Executors/PipelineExecutor.cpp ================================================================================ --- uncovered block 343-343 --- 341 | all_processors_finished = false; 342 | if (!is_cancelled) >> 343 | break; 344 | } 345 | } --- uncovered block 387-387 --- 385 | 386 | if (read_progress->counters.total_bytes) >> 387 | read_progress_callback->addTotalBytes(read_progress->counters.total_bytes); 388 | 389 | /// We are finalizing the execution, so no need to call onProgress if there is nothing to report. ================================================================================ src/Processors/Formats/IOutputFormat.h ================================================================================ --- uncovered block 178-178 --- 176 | virtual void consume(Chunk) = 0; 177 | virtual void consumeTotals(Chunk) {} >> 178 | virtual void consumeExtremes(Chunk) {} 179 | virtual void finalizeImpl() {} 180 | virtual void finalizeBuffers() {} --- uncovered block 188-188 --- 186 | /// Override to write the deferred statistics section and close the document. 187 | /// Called in phase 2 after all progress from connection draining has been collected. >> 188 | virtual void writeDeferredStatisticsAndFinalize() {} 189 | 190 | /// Copy the current values of the shared rows-before-* counters into the format state ================================================================================ src/Processors/Formats/Impl/JSONColumnsWithMetadataBlockOutputFormat.cpp ================================================================================ --- uncovered block 104-106 --- 102 | return; 103 | >> 104 | JSONUtils::writeObjectEnd(*ostr); >> 105 | writeChar('\n', *ostr); >> 106 | ostr->next(); 107 | } 108 | ================================================================================ src/Processors/Formats/Impl/JSONCompactEachRowWithProgressRowOutputFormat.cpp ================================================================================ --- uncovered block 97-98 --- 95 | { 96 | writeCString("{\"rows_before_aggregation\":", *ostr); >> 97 | writeIntText(statistics.rows_before_aggregation, *ostr); >> 98 | writeCString("}\n", *ostr); 99 | } 100 | } --- uncovered block 136-136 --- 134 | if (!exception_message.empty()) 135 | { >> 136 | writeCString("{\"exception\":", *ostr); 137 | writeJSONString(exception_message, *ostr, settings); 138 | writeCString("}\n", *ostr); ================================================================================ src/Processors/Formats/Impl/JSONEachRowWithProgressRowOutputFormat.cpp ================================================================================ --- uncovered block 99-100 --- 97 | { 98 | writeCString("{\"rows_before_aggregation\":", *ostr); >> 99 | writeIntText(statistics.rows_before_aggregation, *ostr); >> 100 | writeCString("}\n", *ostr); 101 | } 102 | } --- uncovered block 138-138 --- 136 | if (!exception_message.empty()) 137 | { >> 138 | writeCString("{\"exception\":", *ostr); 139 | writeJSONString(exception_message, *ostr, settings); 140 | writeCString("}\n", *ostr); ================================================================================ src/Processors/Formats/Impl/JSONRowOutputFormat.cpp ================================================================================ --- uncovered block 173-178 --- 171 | if (!exception_message.empty()) 172 | { >> 173 | writeCString(",\n\n", *ostr); >> 174 | JSONUtils::writeException(exception_message, *ostr, settings, 1); >> 175 | JSONUtils::writeObjectEnd(*ostr); >> 176 | writeChar('\n', *ostr); >> 177 | ostr->next(); >> 178 | return; 179 | } 180 | ================================================================================ src/Processors/Formats/Impl/ParallelFormattingOutputFormat.cpp ================================================================================ --- uncovered block 58-58 --- 56 | formatter->finalizeImpl(); 57 | if (formatter->hasDeferredStatistics()) >> 58 | formatter->writeDeferredStatisticsAndFinalize(); 59 | } 60 | --- uncovered block 71-71 --- 69 | /// the exception trailer instead of a success one. 70 | if (!exception_message.empty()) >> 71 | formatter->setException(exception_message); 72 | formatter->writeDeferredStatisticsAndFinalize(); 73 | /// Flush and finalize the formatter's internal write buffers ================================================================================ src/Processors/Formats/Impl/XMLRowOutputFormat.cpp ================================================================================ --- uncovered block 227-227 --- 225 | /// otherwise the document would look like a complete successful result. 226 | if (!exception_message.empty()) >> 227 | writeException(); 228 | /// hasDeferredStatistics() only guarantees there was no exception at phase 1; the statistics 229 | /// section is still gated on write_statistics. ================================================================================ src/Processors/Formats/PullingOutputFormat.h ================================================================================ --- uncovered block 25-26 --- 23 | { 24 | /// Refresh rows-before-* from the shared counters; see LazyOutputFormat::getProfileInfo. >> 25 | snapshotRowsBeforeCounters(); >> 26 | return info; 27 | } 28 | ================================================================================ src/Processors/Sources/RemoteSource.cpp ================================================================================ --- uncovered block 283-283 --- 281 | 282 | if (!cancel_reason.compare_exchange_strong(expected, reason, std::memory_order_acq_rel)) >> 283 | return; 284 | 285 | /// The `PartialResult` reason we are upgrading could have been published by `onUpdatePorts`, --- uncovered block 299-299 --- 297 | catch (...) 298 | { >> 299 | tryLogCurrentException(getLogger("RemoteSource"), "Error occurs on cancellation upgrade."); 300 | } 301 | } ================================================================================ src/Processors/tests/gtest_finalize_progress_replay_break_overflow_mode.cpp ================================================================================ --- uncovered block 43-44 --- 41 | void work() override 42 | { >> 43 | ++work_calls; >> 44 | ISource::work(); 45 | } 46 | --- uncovered block 55-58 --- 53 | ++progress_polls; 54 | if (progress_polls <= work_calls) >> 55 | return std::nullopt; 56 | 57 | if (late_progress_reported) >> 58 | return std::nullopt; 59 | late_progress_reported = true; 60 | --- uncovered block 70-75 --- 68 | Chunk generate() override 69 | { >> 70 | if (produced) >> 71 | return {}; >> 72 | produced = true; >> 73 | auto column = ColumnUInt64::create(); >> 74 | column->insertValue(42); >> 75 | return Chunk(Columns{std::move(column)}, 1); 76 | } 77 | ================================================================================ src/QueryPipeline/RemoteQueryExecutor.cpp ================================================================================ --- uncovered block 719-719 --- 717 | { 718 | if (was_cancelled) >> 719 | return ReadResult(Block()); 720 | 721 | auto anything = processPacket(std::move(packet)); --- uncovered block 751-758 --- 749 | /// hard-cancellation path into an exception, and it may still have unread 750 | /// packets, so prevent the connections from returning to the pool. >> 751 | tryLogCurrentException(log, "Error while cancelling remote query."); >> 752 | try 753 | { >> 754 | connections->disconnect(); 755 | } >> 756 | catch (...) 757 | { >> 758 | tryLogCurrentException(log, "Error while disconnecting cancelled remote query."); 759 | } 760 | } --- uncovered block 782-789 --- 780 | /// into an exception. The connection may have unread packets, so do not return it 781 | /// to the pool. >> 782 | tryLogCurrentException(log, "Error while draining cancelled remote query."); >> 783 | try 784 | { >> 785 | connections->disconnect(); 786 | } >> 787 | catch (...) 788 | { >> 789 | tryLogCurrentException(log, "Error while disconnecting cancelled remote query."); 790 | } 791 | } --- uncovered block 926-926 --- 924 | return ReadResult(Block{}); 925 | } >> 926 | break; 927 | } 928 | --- uncovered block 985-991 --- 983 | 984 | case Protocol::Server::TimezoneUpdate: >> 985 | break; 986 | >> 987 | default: >> 988 | got_unknown_packet_from_replica.store(true, std::memory_order_release); >> 989 | throw Exception( >> 990 | ErrorCodes::UNKNOWN_PACKET_FROM_SERVER, >> 991 | "Unknown packet {} from one of the following replicas: {}", 992 | packet.type, 993 | connections->dumpAddresses()); --- uncovered block 1118-1118 --- 1116 | /// reasoning as in the branch below. 1117 | if (connections && sent_query) >> 1118 | connections->disconnect(); 1119 | 1120 | finished = true; --- uncovered block 1200-1200 --- 1198 | /// If connections weren't created yet, query wasn't sent or was already finished, nothing to do. 1199 | if (!connections || !sent_query || finished) >> 1200 | return; 1201 | 1202 | /// Take ownership of the drain loop before releasing the mutex. A concurrent `finish` that --- uncovered block 1208-1208 --- 1206 | } 1207 | >> 1208 | drainConnections(std::nullopt); 1209 | } 1210 | --- uncovered block 1309-1309 --- 1307 | break; 1308 | >> 1309 | case Protocol::Server::Exception: 1310 | /// The drain owner called `tryCancel` before entering the drain (see `finish`), which 1311 | /// sends `Cancel` to the replicas. A replica that was actively processing responds with --- uncovered block 1325-1340 --- 1323 | /// still active, and would cause subsequent `finish` / `cancelUnlocked` calls to 1324 | /// return early as if a real replica error had occurred. >> 1325 | if (packet.exception->code() == ErrorCodes::QUERY_WAS_CANCELLED_BY_CLIENT >> 1326 | || packet.exception->code() == ErrorCodes::QUERY_WAS_CANCELLED) 1327 | { >> 1328 | if (log) >> 1329 | LOG_TRACE(log, "Replica reported expected cancellation during drain: {}", packet.exception->displayText()); >> 1330 | break; 1331 | } >> 1332 | if (shouldIgnoreShardException(packet.exception->code())) 1333 | { >> 1334 | if (log) >> 1335 | LOG_ERROR(log, >> 1336 | "Ignoring exception from connection(s) {} due to `skip_unavailable_shards_mode` setting: {}", >> 1337 | connections->dumpAddresses(), >> 1338 | packet.exception->displayText()); 1339 | >> 1340 | reportShardSkipped(); 1341 | 1342 | /// Treat an ignored shard exception like `EndOfStream` for that one connection. --- uncovered block 1353-1358 --- 1351 | /// flag is published by the `SCOPE_EXIT` in `drainConnections` once every replica 1352 | /// is drained. >> 1353 | break; 1354 | } 1355 | >> 1356 | got_exception_from_replica.store(true, std::memory_order_release); >> 1357 | packet.exception->rethrow(); >> 1358 | break; 1359 | 1360 | case Protocol::Server::Log: --- uncovered block 1370-1370 --- 1368 | if (auto profile_queue = CurrentThread::getInternalProfileEventsQueue()) 1369 | if (!profile_queue->emplace(std::move(packet.block))) >> 1370 | throw Exception(ErrorCodes::SYSTEM_ERROR, "Could not push into profile queue"); 1371 | break; 1372 | --- uncovered block 1384-1401 --- 1382 | break; 1383 | >> 1384 | case Protocol::Server::Totals: 1385 | /// A replica may deliver its `Totals` block among the trailing packets drained here 1386 | /// rather than during the normal read (the same window the 1387 | /// `tcp_handler_sleep_before_secondary_query_trailing_packets` failpoint forces). 1388 | /// Store it exactly as `processPacket` does; otherwise `RemoteTotalsSource` would later 1389 | /// read an empty block and a `WITH TOTALS` query would silently lose its totals section. >> 1390 | totals = packet.block; >> 1391 | if (!totals.empty()) >> 1392 | totals = adaptBlockStructure(totals, *header); >> 1393 | break; 1394 | >> 1395 | case Protocol::Server::Extremes: 1396 | /// Likewise for `Extremes` (`extremes = 1`), which would otherwise be dropped here and 1397 | /// leave `RemoteExtremesSource` with an empty block. >> 1398 | extremes = packet.block; >> 1399 | if (!extremes.empty()) >> 1400 | extremes = adaptBlockStructure(packet.block, *header); >> 1401 | break; 1402 | 1403 | case Protocol::Server::MergeTreeReadTaskRequest: WARNING: Failed to get start time for [Print Uncovered Code] - start time and duration won't be set --- Coverage counts --- Lines : baseline 1,041,507/1,175,182 -> current 1,042,040/1,175,665 (delta +533 / +483) Functions : baseline 836,365/908,530 -> current 836,496/908,586 (delta +131 / +56) Branches : baseline 345,873/427,418 -> current 346,030/427,618 (delta +157 / +200)