2018-03-30 20:43:18 +00:00
|
|
|
package archiver
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
|
2018-05-08 20:28:37 +00:00
|
|
|
"github.com/restic/restic/internal/debug"
|
2018-03-30 20:43:18 +00:00
|
|
|
"github.com/restic/restic/internal/restic"
|
2022-05-27 17:08:50 +00:00
|
|
|
"golang.org/x/sync/errgroup"
|
2018-03-30 20:43:18 +00:00
|
|
|
)
|
|
|
|
|
|
|
|
// Saver allows saving a blob.
|
|
|
|
type Saver interface {
|
2022-05-01 12:26:57 +00:00
|
|
|
SaveBlob(ctx context.Context, t restic.BlobType, data []byte, id restic.ID, storeDuplicate bool) (restic.ID, bool, int, error)
|
2018-03-30 20:43:18 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// BlobSaver concurrently saves incoming blobs to the repo.
|
|
|
|
type BlobSaver struct {
|
|
|
|
repo Saver
|
2020-06-06 20:20:44 +00:00
|
|
|
ch chan<- saveBlobJob
|
2018-03-30 20:43:18 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// NewBlobSaver returns a new blob. A worker pool is started, it is stopped
|
|
|
|
// when ctx is cancelled.
|
2022-05-27 17:08:50 +00:00
|
|
|
func NewBlobSaver(ctx context.Context, wg *errgroup.Group, repo Saver, workers uint) *BlobSaver {
|
2018-04-30 13:13:03 +00:00
|
|
|
ch := make(chan saveBlobJob)
|
2018-03-30 20:43:18 +00:00
|
|
|
s := &BlobSaver{
|
2020-06-06 20:20:44 +00:00
|
|
|
repo: repo,
|
|
|
|
ch: ch,
|
2018-03-30 20:43:18 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
for i := uint(0); i < workers; i++ {
|
2022-05-27 17:08:50 +00:00
|
|
|
wg.Go(func() error {
|
|
|
|
return s.worker(ctx, ch)
|
2018-05-08 20:28:37 +00:00
|
|
|
})
|
2018-03-30 20:43:18 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
return s
|
|
|
|
}
|
|
|
|
|
2022-05-27 17:08:50 +00:00
|
|
|
func (s *BlobSaver) TriggerShutdown() {
|
|
|
|
close(s.ch)
|
|
|
|
}
|
|
|
|
|
2018-03-30 20:43:18 +00:00
|
|
|
// Save stores a blob in the repo. It checks the index and the known blobs
|
2020-06-11 11:34:05 +00:00
|
|
|
// before saving anything. It takes ownership of the buffer passed in.
|
2022-10-07 18:23:38 +00:00
|
|
|
func (s *BlobSaver) Save(ctx context.Context, t restic.BlobType, buf *Buffer, cb func(res SaveBlobResponse)) {
|
2018-05-08 20:28:37 +00:00
|
|
|
select {
|
2022-10-07 18:23:38 +00:00
|
|
|
case s.ch <- saveBlobJob{BlobType: t, buf: buf, cb: cb}:
|
2018-05-08 20:28:37 +00:00
|
|
|
case <-ctx.Done():
|
|
|
|
debug.Log("not sending job, context is cancelled")
|
|
|
|
}
|
2022-05-01 12:41:36 +00:00
|
|
|
}
|
|
|
|
|
2018-03-30 20:43:18 +00:00
|
|
|
type saveBlobJob struct {
|
|
|
|
restic.BlobType
|
2018-04-29 13:34:41 +00:00
|
|
|
buf *Buffer
|
2022-10-07 18:23:38 +00:00
|
|
|
cb func(res SaveBlobResponse)
|
2018-03-30 20:43:18 +00:00
|
|
|
}
|
|
|
|
|
2022-05-22 13:14:25 +00:00
|
|
|
type SaveBlobResponse struct {
|
|
|
|
id restic.ID
|
|
|
|
length int
|
|
|
|
sizeInRepo int
|
|
|
|
known bool
|
2018-03-30 20:43:18 +00:00
|
|
|
}
|
|
|
|
|
2022-05-22 13:14:25 +00:00
|
|
|
func (s *BlobSaver) saveBlob(ctx context.Context, t restic.BlobType, buf []byte) (SaveBlobResponse, error) {
|
|
|
|
id, known, sizeInRepo, err := s.repo.SaveBlob(ctx, t, buf, restic.ID{}, false)
|
2018-03-30 20:43:18 +00:00
|
|
|
|
2018-05-08 20:28:37 +00:00
|
|
|
if err != nil {
|
2022-05-22 13:14:25 +00:00
|
|
|
return SaveBlobResponse{}, err
|
2018-05-08 20:28:37 +00:00
|
|
|
}
|
|
|
|
|
2022-05-22 13:14:25 +00:00
|
|
|
return SaveBlobResponse{
|
|
|
|
id: id,
|
|
|
|
length: len(buf),
|
|
|
|
sizeInRepo: sizeInRepo,
|
|
|
|
known: known,
|
2018-05-08 20:28:37 +00:00
|
|
|
}, nil
|
2018-03-30 20:43:18 +00:00
|
|
|
}
|
|
|
|
|
2018-05-08 20:28:37 +00:00
|
|
|
func (s *BlobSaver) worker(ctx context.Context, jobs <-chan saveBlobJob) error {
|
2018-03-30 20:43:18 +00:00
|
|
|
for {
|
|
|
|
var job saveBlobJob
|
2022-05-27 17:08:50 +00:00
|
|
|
var ok bool
|
2018-03-30 20:43:18 +00:00
|
|
|
select {
|
|
|
|
case <-ctx.Done():
|
2018-05-08 20:28:37 +00:00
|
|
|
return nil
|
2022-05-27 17:08:50 +00:00
|
|
|
case job, ok = <-jobs:
|
|
|
|
if !ok {
|
|
|
|
return nil
|
|
|
|
}
|
2018-03-30 20:43:18 +00:00
|
|
|
}
|
|
|
|
|
2018-05-08 20:28:37 +00:00
|
|
|
res, err := s.saveBlob(ctx, job.BlobType, job.buf.Data)
|
|
|
|
if err != nil {
|
2018-05-12 19:40:31 +00:00
|
|
|
debug.Log("saveBlob returned error, exiting: %v", err)
|
2018-05-08 20:28:37 +00:00
|
|
|
return err
|
|
|
|
}
|
2022-10-07 18:23:38 +00:00
|
|
|
job.cb(res)
|
2018-03-30 20:43:18 +00:00
|
|
|
job.buf.Release()
|
|
|
|
}
|
|
|
|
}
|