This repository has been archived by the owner on Sep 7, 2022. It is now read-only.
/
batchwriter.go
105 lines (92 loc) · 2.26 KB
/
batchwriter.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
package bsdagwriter
import (
"bytes"
"context"
"fmt"
"io"
"sync"
blocks "github.com/ipfs/go-block-format"
"github.com/ipfs/go-blockservice"
"github.com/ipfs/go-cid"
"github.com/ipld/go-ipld-prime"
cidlink "github.com/ipld/go-ipld-prime/linking/cid"
)
type dagBatchWriter struct {
*ipld.LinkSystem
bs blockservice.BlockService
cache *cachedOperationsStore
}
func (tds *dagBatchWriter) put(lnkCtx ipld.LinkContext) (io.Writer, ipld.BlockWriteCommitter, error) {
buf := bytes.Buffer{}
return &buf, func(lnk ipld.Link) error {
asCidLink, ok := lnk.(cidlink.Link)
if !ok {
return fmt.Errorf("unsupported link type %v", lnk)
}
tds.cache.write(asCidLink.Cid, buf.Bytes())
return nil
}, nil
}
func (tds *dagBatchWriter) Delete(ctx context.Context, lnk ipld.Link) error {
asCidLink, ok := lnk.(cidlink.Link)
if !ok {
return fmt.Errorf("unsupported link type %v", lnk)
}
tds.cache.delete(asCidLink.Cid)
return nil
}
func (tds *dagBatchWriter) Commit() error {
blks, deletes, err := tds.cache.reset()
if err != nil {
return err
}
for _, c := range deletes {
err := tds.bs.DeleteBlock(c)
if err != nil {
return nil
}
}
return tds.bs.AddBlocks(blks)
}
type cacheRecord struct {
data []byte
tombstone bool
}
type cachedOperationsStore struct {
cache map[cid.Cid]cacheRecord
cachelk sync.RWMutex
}
func newCachedOperationsStore() *cachedOperationsStore {
return &cachedOperationsStore{
cache: make(map[cid.Cid]cacheRecord),
}
}
func (cos *cachedOperationsStore) write(c cid.Cid, data []byte) {
cos.cachelk.Lock()
cos.cache[c] = cacheRecord{data, false}
cos.cachelk.Unlock()
}
func (cos *cachedOperationsStore) delete(c cid.Cid) {
cos.cachelk.Lock()
cos.cache[c] = cacheRecord{nil, true}
cos.cachelk.Unlock()
}
func (cos *cachedOperationsStore) reset() ([]blocks.Block, []cid.Cid, error) {
cos.cachelk.Lock()
defer cos.cachelk.Unlock()
blks := make([]blocks.Block, 0, len(cos.cache))
deletes := make([]cid.Cid, 0, len(cos.cache))
for c, record := range cos.cache {
if record.tombstone {
deletes = append(deletes, c)
continue
}
blk, err := blocks.NewBlockWithCid(record.data, c)
if err != nil {
return nil, nil, nil
}
blks = append(blks, blk)
}
cos.cache = make(map[cid.Cid]cacheRecord)
return blks, deletes, nil
}