Files
dolthub--dolt/go/store/nbs/journal_chunk_source.go
wehub-resource-sync 5357c39144
Fuzzer / Run Fuzzer (push) Has been cancelled
Race tests / Go race tests (ubuntu-22.04) (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 13:01:40 +08:00

345 lines
9.1 KiB
Go

// Copyright 2022 Dolthub, Inc.
//
// Licensed 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.
package nbs
import (
"context"
"errors"
"fmt"
"io"
"runtime/trace"
"sort"
"sync"
"golang.org/x/sync/errgroup"
dherrors "github.com/dolthub/dolt/go/libraries/utils/errors"
"github.com/dolthub/dolt/go/store/chunks"
"github.com/dolthub/dolt/go/store/hash"
)
// journalChunkSource is a chunkSource that reads chunks
// from a ChunkJournal. Unlike other NBS chunkSources,
// it is not immutable and its set of chunks grows as
// more commits are made to the ChunkJournal.
type journalChunkSource struct {
journal *journalWriter
}
var _ chunkSource = journalChunkSource{}
func (s journalChunkSource) has(h hash.Hash, keeper keeperF) (bool, gcBehavior, error) {
res := s.journal.hasAddr(h)
if res && keeper != nil && keeper(h) {
return false, gcBehavior_Block, nil
}
return res, gcBehavior_Continue, nil
}
func (s journalChunkSource) hasMany(addrs []hasRecord, keeper keeperF) (bool, gcBehavior, error) {
missing := false
for i := range addrs {
h := *addrs[i].a
ok := s.journal.hasAddr(h)
if ok {
if keeper != nil && keeper(h) {
return true, gcBehavior_Block, nil
}
addrs[i].has = true
} else {
missing = true
}
}
return missing, gcBehavior_Continue, nil
}
func (s journalChunkSource) getCompressed(ctx context.Context, h hash.Hash, _ *Stats) (CompressedChunk, error) {
defer trace.StartRegion(ctx, "journalChunkSource.getCompressed").End()
return s.journal.getCompressedChunk(h)
}
func (s journalChunkSource) get(ctx context.Context, h hash.Hash, keeper keeperF, _ *Stats) ([]byte, gcBehavior, error) {
defer trace.StartRegion(ctx, "journalChunkSource.get").End()
cc, err := s.journal.getCompressedChunk(h)
if err != nil {
return nil, gcBehavior_Continue, err
} else if cc.IsEmpty() {
return nil, gcBehavior_Continue, nil
}
if keeper != nil && keeper(h) {
return nil, gcBehavior_Block, nil
}
ch, err := cc.ToChunk()
if err != nil {
return nil, gcBehavior_Continue, err
}
return ch.Data(), gcBehavior_Continue, nil
}
type journalRecord struct {
// r is the journal range for this chunk
r Range
// idx is the array offset into the shared |reqs|
idx int
}
func (s journalChunkSource) getMany(ctx context.Context, eg *errgroup.Group, reqs []getRecord, found func(context.Context, *chunks.Chunk), keeper keeperF, stats *Stats) (bool, gcBehavior, error) {
return s.getManyCompressed(ctx, eg, reqs, func(ctx context.Context, cc ToChunker) {
ch, err := cc.ToChunk()
if err != nil {
eg.Go(func() error {
return err
})
return
}
chWHash := chunks.NewChunkWithHash(cc.Hash(), ch.Data())
found(ctx, &chWHash)
}, keeper, stats)
}
// getManyCompressed implements chunkReader. Here we (1) synchronously check
// the journal index for read ranges, (2) record if the source misses any
// needed remaining chunks, (3) sort the lookups for efficient disk access,
// and then (4) asynchronously perform reads. We release the journal read
// lock after returning when all reads are completed, which can be after the
// function returns.
func (s journalChunkSource) getManyCompressed(ctx context.Context, eg *errgroup.Group, reqs []getRecord, found func(context.Context, ToChunker), keeper keeperF, stats *Stats) (bool, gcBehavior, error) {
defer trace.StartRegion(ctx, "journalChunkSource.getManyCompressed").End()
var remaining bool
var jReqs []journalRecord
var wg sync.WaitGroup
s.journal.lock.RLock()
for i, r := range reqs {
if r.found {
continue
}
h := *r.a
rang, ok := s.journal.ranges.get(h)
if !ok {
remaining = true
continue
}
if keeper != nil && keeper(h) {
s.journal.lock.RUnlock()
return true, gcBehavior_Block, nil
}
jReqs = append(jReqs, journalRecord{r: rang, idx: i})
reqs[i].found = true
}
// sort chunks by journal locality
sort.Slice(jReqs, func(i, j int) bool {
return jReqs[i].r.Offset < jReqs[j].r.Offset
})
wg.Add(len(jReqs))
go func() {
wg.Wait()
s.journal.lock.RUnlock()
}()
for i := range jReqs {
// workers populate the parent error group
// record local workers for releasing lock
eg.Go(func() error {
defer wg.Done()
if ctx.Err() != nil {
return ctx.Err()
}
rec := jReqs[i]
a := reqs[rec.idx].a
if cc, err := s.journal.getCompressedChunkAtRange(rec.r, *a); err != nil {
return err
} else if cc.IsEmpty() {
return errors.New("chunk in journal index was empty.")
} else {
found(ctx, cc)
return nil
}
})
}
return remaining, gcBehavior_Continue, nil
}
func (s journalChunkSource) count() uint32 {
return s.journal.recordCount()
}
func (s journalChunkSource) uncompressedLen() (uint64, error) {
return s.journal.uncompressedSize(), nil
}
func (s journalChunkSource) hash() hash.Hash {
return journalAddr
}
func (s journalChunkSource) suffix() string {
return ""
}
// reader implements chunkSource.
func (s journalChunkSource) reader(ctx context.Context, behavior dherrors.FatalBehavior) (io.ReadCloser, uint64, error) {
rdr, sz, err := s.journal.snapshot(ctx, behavior)
return rdr, uint64(sz), err
}
func (s journalChunkSource) getRecordRanges(ctx context.Context, behavior dherrors.FatalBehavior, requests []getRecord, keeper keeperF) (map[hash.Hash]Range, gcBehavior, error) {
ranges := make(map[hash.Hash]Range, len(requests))
for _, req := range requests {
if req.found {
continue
}
h := *req.a
rng, ok, err := s.journal.getRange(ctx, behavior, h)
if err != nil {
return nil, gcBehavior_Continue, err
} else if !ok {
continue
}
if keeper != nil && keeper(h) {
return nil, gcBehavior_Block, nil
}
req.found = true // update |requests|
ranges[h] = rng
}
return ranges, gcBehavior_Continue, nil
}
// size implements chunkSource.
// size returns the total size of the chunkSource: chunks, index, and footer
func (s journalChunkSource) currentSize() uint64 {
return uint64(s.journal.currentSize())
}
// index implements chunkSource.
func (s journalChunkSource) index() (tableIndex, error) {
return nil, fmt.Errorf("journalChunkSource cannot be conjoined")
}
func (s journalChunkSource) clone() (chunkSource, error) {
return s, nil
}
func (s journalChunkSource) close() error {
// |s.journal| closed via ChunkJournal
return nil
}
func (s journalChunkSource) iterateAllChunks(ctx context.Context, cb func(chunks.Chunk), _ *Stats) error {
// TODO - a less time consuming lock is possible here. Using s.journal.snapshot and processJournalRecords()
// would allow for no locking. Need to filter out the journal records which are actually chunks, then convert
// those to chunks and pass them to cb. When we support online FSCK this will allow the server to keep running uninterrupted.
s.journal.lock.RLock()
defer s.journal.lock.RUnlock()
for h, r := range s.journal.ranges.novel {
if ctx.Err() != nil {
return context.Cause(ctx)
}
cchk, err := s.journal.getCompressedChunkAtRange(r, h)
if err != nil {
return err
}
chunk, err := cchk.ToChunk()
if err != nil {
return err
}
cb(chunk)
}
for a16, r := range s.journal.ranges.cached {
if ctx.Err() != nil {
return context.Cause(ctx)
}
// We only have 16 bytes of the hash. The value returned here will have 4 0x00 bytes at the end.
var h hash.Hash
copy(h[:], a16[:])
cchk, err := s.journal.getCompressedChunkAtRange(r, h)
if err != nil {
return err
}
chunk, err := cchk.ToChunk()
if err != nil {
return err
}
cb(chunk)
}
return nil
}
func (s journalChunkSource) tolerantIterateAllChunks(ctx context.Context, cb func(chunks.Chunk), errCb func(error), _ *Stats) {
s.journal.lock.RLock()
defer s.journal.lock.RUnlock()
for h, r := range s.journal.ranges.novel {
if ctx.Err() != nil {
return
}
cchk, err := s.journal.getCompressedChunkAtRange(r, h)
if err != nil {
errCb(fmt.Errorf("chunk %s: %w", h.String(), err))
continue
}
chunk, err := cchk.ToChunk()
if err != nil {
errCb(fmt.Errorf("chunk %s: decompress error: %w", h.String(), err))
continue
}
cb(chunk)
}
for a16, r := range s.journal.ranges.cached {
if ctx.Err() != nil {
return
}
var h hash.Hash
copy(h[:], a16[:])
cchk, err := s.journal.getCompressedChunkAtRange(r, h)
if err != nil {
errCb(fmt.Errorf("chunk %s: %w", h.String(), err))
continue
}
chunk, err := cchk.ToChunk()
if err != nil {
errCb(fmt.Errorf("chunk %s: decompress error: %w", h.String(), err))
continue
}
cb(chunk)
}
}
func equalSpecs(left, right []tableSpec) bool {
if len(left) != len(right) {
return false
}
l := make(map[hash.Hash]struct{}, len(left))
for _, s := range left {
l[s.name] = struct{}{}
}
for _, s := range right {
if _, ok := l[s.name]; !ok {
return false
}
}
return true
}