airborne12 commented on code in PR #67657:
URL: https://github.com/apache/doris/pull/67657#discussion_r4128932417


##########
be/src/storage/index/snii/writer/bounded_run_merge.cpp:
##########
@@ -0,0 +1,497 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include <sys/resource.h>
+#include <unistd.h>
+
+#include <algorithm>
+#include <array>
+#include <atomic>
+#include <climits>
+#include <cstdio>
+#include <cstring>
+#include <limits>
+#include <memory>
+#include <queue>
+#include <utility>
+
+#include "storage/index/snii/encoding/crc32c.h"
+#include "storage/index/snii/writer/encoded_spill_run.h"
+#include "storage/index/snii/writer/spill_run_codec.h"
+#include "storage/index/snii/writer/temp_dir.h"
+
+namespace doris::snii::writer {
+namespace {
+
+Status invalid_run(const char* reason) {
+    return Status::Error<ErrorCode::INVERTED_INDEX_FILE_CORRUPTED, 
false>("spill merge: {}",
+                                                                          
reason);
+}
+
+std::atomic<uint64_t> g_run_compactions {0};
+
+struct RunRange {
+    std::string path;
+    uint64_t begin = 0;
+    uint64_t end = UINT64_MAX;
+};
+
+// Caller-owned fixture files, or a bounded directory of sealed ranges in one
+// ingestion spool. The latter has no resident path/offset array proportional 
to
+// the number of spills.
+struct RunInputs {
+    const std::vector<std::string>* files = nullptr;
+    const std::string* spool = nullptr;
+    PostingByteBuffer* ends = nullptr;
+    size_t count = 0;
+    size_t fan_in_limit = 0;
+};
+
+using Readers = std::vector<std::unique_ptr<StreamingRunReader>>;
+
+std::string merge_run_path() {
+    static std::atomic<uint64_t> sequence {0};
+    return resolve_temp_dir() + "/snii_merge_" + std::to_string(::getpid()) + 
"_" +
+           std::to_string(sequence.fetch_add(1, std::memory_order_relaxed)) + 
".run";
+}
+
+// One-token lookahead survives a fill boundary, including a document split
+// across any number of input fragments. Positions are appended incrementally
+// into the consumer's replayable buffer, never accumulated per document here.
+class RunTermSource final : public TermPostingSource {
+public:
+    RunTermSource(Readers* readers, const std::vector<size_t>* matching, bool 
positions)
+            : readers_(readers), matching_(matching), positions_(positions) {}
+
+    Status fill(uint32_t target_docs, TermPostingBuffer* out, bool* exhausted) 
override {
+        if (target_docs == 0 || out == nullptr || exhausted == nullptr || 
!out->empty()) {
+            return Status::Error<ErrorCode::INVALID_ARGUMENT, false>(
+                    "spill merge: invalid source fill arguments");
+        }
+        if (!pending_) {
+            RETURN_IF_ERROR(next());
+        }
+        while (pending_ && out->document_count() < target_docs) {
+            const uint32_t docid = docid_;
+            MutableTermPostingSpan document;
+            RETURN_IF_ERROR(out->grow_uninitialized(1, true, 0, &document));
+            document.docids[0] = docid;
+            uint64_t frequency = 0;
+            do {
+                if (positions_) {
+                    RETURN_IF_ERROR(out->append_position(position_));
+                }
+                if (++frequency > std::numeric_limits<uint32_t>::max()) {
+                    return invalid_run("document frequency overflows uint32");
+                }
+                RETURN_IF_ERROR(next());
+            } while (pending_ && docid_ == docid);
+            // Appending positions can replace the buffer's vectors. Acquire 
the
+            // frequency slot again rather than retaining an invalidated span.
+            out->set_last_frequency(static_cast<uint32_t>(frequency));
+            if (pending_ && docid_ < docid) {
+                return invalid_run("docids go backwards across run fragments");
+            }
+        }
+        *exhausted = !pending_ && run_ == matching_->size();
+        exhausted_ = *exhausted;
+        return Status::OK();
+    }
+
+    bool exhausted() const { return exhausted_; }
+
+private:
+    Status next() {
+        pending_ = false;
+        while (run_ < matching_->size()) {
+            bool end = false;
+            
RETURN_IF_ERROR((*readers_)[(*matching_)[run_]]->next_token(&docid_, 
&position_, &end));
+            if (!end) {
+                pending_ = true;
+                break;
+            }
+            ++run_;
+        }
+        return Status::OK();
+    }
+
+    Readers* readers_;
+    const std::vector<size_t>* matching_;
+    const bool positions_;
+    size_t run_ = 0;
+    uint32_t docid_ = 0;
+    uint32_t position_ = 0;
+    bool pending_ = false;
+    bool exhausted_ = false;
+};
+
+struct HeapItem {
+    uint32_t id;
+    size_t run;
+};
+struct Greater {
+    const std::vector<uint32_t>* rank;
+    bool operator()(const HeapItem& left, const HeapItem& right) const {
+        const uint32_t a = (*rank)[left.id];
+        const uint32_t b = (*rank)[right.id];
+        return a == b ? left.run > right.run : a > b;
+    }
+};
+
+using RunHeap = std::priority_queue<HeapItem, std::vector<HeapItem>, Greater>;
+
+Status take_matching_runs(RunHeap* heap, const Readers& readers, 
EncodedRunTerm* term,
+                          std::vector<size_t>* matching) {
+    matching->clear();
+    while (!heap->empty() && heap->top().id == term->term_id) {
+        const size_t run = heap->top().run;
+        heap->pop();
+        const auto& input = readers[run]->current();
+        if (input.has_positions != term->has_positions) {
+            return invalid_run("posting shape differs");
+        }
+        if (input.document_groups > UINT64_MAX - term->document_groups ||
+            input.tokens > UINT64_MAX - term->tokens) {
+            return invalid_run("term counts overflow");
+        }
+        term->document_groups += input.document_groups;
+        term->tokens += input.tokens;
+        matching->push_back(run);
+    }
+    return Status::OK();
+}
+
+template <typename Consumer>
+Status merge_group(const std::vector<RunRange>& paths, const 
std::vector<uint32_t>& rank,
+                   bool positions, MemoryReporter* reporter, bool allow_legacy,
+                   const Consumer& consume) {
+    auto metadata = reporter == nullptr ? MemoryReporter::Reservation()
+                                        : 
reporter->make_postings_reservation();
+    if (reporter != nullptr) {
+        // Reader objects own fixed, separately charged payload buffers. 
Reserve
+        // their object storage plus both heap and matching-index capacities.
+        RETURN_IF_ERROR(metadata.set_bytes(paths.size() * 
(StreamingRunReader::object_bytes() +
+                                                           sizeof(HeapItem) + 
2 * sizeof(size_t))));
+    }
+    Readers readers;
+    readers.reserve(paths.size());
+    std::vector<HeapItem> heap_storage;
+    heap_storage.reserve(paths.size());
+    RunHeap heap(Greater {&rank}, std::move(heap_storage));
+    for (size_t run = 0; run < paths.size(); ++run) {
+        auto reader = std::make_unique<StreamingRunReader>(reporter);
+        RETURN_IF_ERROR(reader->open(paths[run].path, positions, allow_legacy, 
paths[run].begin,
+                                     paths[run].end));
+        if (!reader->exhausted()) {
+            if (reader->current().term_id >= rank.size()) {
+                return invalid_run("term id out of range");
+            }
+            heap.push({reader->current().term_id, run});
+        }
+        readers.push_back(std::move(reader));
+    }
+    std::vector<size_t> matching;
+    matching.reserve(paths.size());
+    while (!heap.empty()) {
+        const uint32_t id = heap.top().id;
+        EncodedRunTerm term {.term_id = id,
+                             .has_positions = 
readers[heap.top().run]->current().has_positions};
+        RETURN_IF_ERROR(take_matching_runs(&heap, readers, &term, &matching));
+        RETURN_IF_ERROR(consume(term, &readers, matching));
+        for (size_t run : matching) {
+            RETURN_IF_ERROR(readers[run]->advance());
+            if (!readers[run]->exhausted()) {
+                const uint32_t next_id = readers[run]->current().term_id;
+                if (next_id >= rank.size()) {
+                    return invalid_run("term id out of range");
+                }
+                if (rank[next_id] <= rank[id]) {
+                    return invalid_run("terms are not strictly ordered");
+                }
+                heap.push({next_id, run});
+            }
+        }
+    }
+    return Status::OK();
+}
+
+Status compact_group(const std::vector<RunRange>& paths, const 
std::vector<uint32_t>& rank,
+                     bool positions, const std::string& output, 
MemoryReporter* reporter,
+                     bool allow_legacy) {
+    EncodedRunWriter writer(reporter);
+    RETURN_IF_ERROR(writer.open(output));
+    RETURN_IF_ERROR(merge_group(
+            paths, rank, positions, reporter, allow_legacy,
+            [&](const EncodedRunTerm& term, Readers* readers, const 
std::vector<size_t>& matching) {
+                RETURN_IF_ERROR(writer.begin_term(term));
+                for (size_t run : matching) {
+                    
RETURN_IF_ERROR((*readers)[run]->copy_fragments_to(&writer));
+                }
+                return writer.end_term();
+            }));
+    RETURN_IF_ERROR(writer.close());
+    g_run_compactions.fetch_add(1, std::memory_order_relaxed);
+    return Status::OK();
+}
+
+class ReducedRuns {
+public:
+    ~ReducedRuns() {
+        // Manifests also retain names from failed/partially completed passes.
+        // Cleanup reads directly into stack storage and cannot need a budget
+        // reservation while unwinding a memory-limit failure.
+        std::array<char, PATH_MAX> path {};
+        for (auto& manifest : manifests_) {
+            uint64_t offset = 0;
+            while (offset < manifest->size()) {
+                uint32_t length = 0;
+                if (!manifest->read_at(offset, 
{reinterpret_cast<uint8_t*>(&length), 4}).ok()) {

Review Comment:
   Confirmed and fixed in 9f44f73e3f5.
   
   **Verification.** Traced the failure path independently: 
`PostingByteBuffer::flush()` leaves `buffered_` set when `write()` fails, and 
`read_at()` called `flush()` unconditionally whenever the buffer had spilled. 
So once a manifest flush hit ENOSPC, the `ReducedRuns` destructor's very first 
`read_at()` retried that write, failed again, and `break`ed before removing 
anything, including outputs whose manifest records were already on disk. 
`TmpFileDirs::init()` only sweeps the temp directory on BE restart, so the 
leaked `.run` files stayed until then.
   
   **Reproduction (TDD, red first).** Added an ENOSPC debug point in 
`PostingByteBuffer::write_all` and two tests in 
`bounded_posting_codec_test.cpp`:
   - `ByteBufferReadsPendingBytesWithoutFlushing`: with the point armed after a 
spill with 4 pending bytes, `read_at(0, 20)` failed with `No space left on 
device` on the old code.
   - `FailedManifestFlushStillRemovesIntermediateRuns`: 55 runs, 1 MiB budget 
(fan-in 2), point armed for the whole `compact_runs` call. Old code left all 28 
first-pass `snii_merge_*.run` outputs in the temp root; original inputs 
untouched.
   
   **Fix.** `read_at()` now serves the flushed prefix with `pread` and the 
pending tail straight from the resident buffer, so a read never writes. The 
destructor logs a `WARNING` when manifest cleanup has to stop. Both tests pass 
after the change; the rest of the suite is unchanged.
   
   On the current partial output: its name is appended to the manifest before 
`compact_group` runs, so it is already covered by the same cleanup once the 
manifest is readable. I did not add a separate unlink for it. Please point out 
if you see a path where that name is not yet recorded.
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to