-
Notifications
You must be signed in to change notification settings - Fork 3
/
ingester.go
175 lines (148 loc) · 4.01 KB
/
ingester.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
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
package ipcs
import (
"context"
"io"
"time"
"github.com/containerd/containerd/content"
"github.com/containerd/containerd/errdefs"
"github.com/hinshun/ipcs/digestconv"
files "github.com/ipfs/go-ipfs-files"
iface "github.com/ipfs/interface-go-ipfs-core"
"github.com/ipfs/interface-go-ipfs-core/options"
"github.com/ipfs/interface-go-ipfs-core/path"
digest "github.com/opencontainers/go-digest"
"github.com/pkg/errors"
)
func (s *store) Writer(ctx context.Context, opts ...content.WriterOpt) (content.Writer, error) {
var wOpts content.WriterOpts
for _, opt := range opts {
if err := opt(&wOpts); err != nil {
return nil, err
}
}
if wOpts.Desc.Digest != "" {
c, err := digestconv.DigestToCid(wOpts.Desc.Digest)
if err != nil {
return nil, errors.Wrapf(err, "failed to convert digest '%s' to cid", wOpts.Desc.Digest)
}
_, err = s.cln.Unixfs().Get(ctx, path.IpfsPath(c))
if err == nil {
return nil, errors.Wrapf(errdefs.ErrAlreadyExists, "content %v", wOpts.Desc.Digest)
}
}
w := &writer{
ctx: ctx,
cln: s.cln,
ref: wOpts.Ref,
total: wOpts.Desc.Size,
}
err := w.Truncate(0)
if err != nil {
return nil, errors.Wrap(err, "failed to truncate writer")
}
return w, nil
}
type writer struct {
ctx context.Context
cln iface.CoreAPI
ref string
offset int64
total int64
dgst digest.Digest
startedAt time.Time
updatedAt time.Time
pw io.Writer
ipfsErr error
cancel func() error
}
// Write writes len(p) bytes from p to the underlying data stream.
// It returns the number of bytes written from p (0 <= n <= len(p))
// and any error encountered that caused the write to stop early.
// Write must return a non-nil error if it returns n < len(p).
// Write must not modify the slice data, even temporarily.
//
// Implementations must not retain p.
func (w *writer) Write(p []byte) (n int, err error) {
if w.ipfsErr != nil {
return 0, w.ipfsErr
}
n, err = w.pw.Write(p)
w.offset += int64(n)
w.updatedAt = time.Now()
return n, err
}
// Close closes the writer, if the writer has not been
// committed this allows resuming or aborting.
// Calling Close on a closed writer will not error.
func (w *writer) Close() error {
return w.cancel()
}
// Digest may return empty digest or panics until committed.
func (w *writer) Digest() digest.Digest {
return w.dgst
}
// Commit commits the blob (but no roll-back is guaranteed on an error).
// size and expected can be zero-value when unknown.
// Commit always closes the writer, even on error.
// ErrAlreadyExists aborts the writer.
func (w *writer) Commit(ctx context.Context, size int64, expected digest.Digest, opts ...content.Opt) error {
if size > 0 && size != w.offset {
return errors.Wrapf(errdefs.ErrFailedPrecondition, "unexpected commit size %d, expected %d", w.offset, size)
}
if expected != "" && expected != w.dgst {
return errors.Wrapf(errdefs.ErrFailedPrecondition, "unexpected commit digest %s, expected %s", w.dgst, expected)
}
return w.Close()
}
// Status returns the current state of write
func (w *writer) Status() (content.Status, error) {
return content.Status{
Ref: w.ref,
Offset: w.offset,
Total: w.total,
StartedAt: w.startedAt,
UpdatedAt: w.updatedAt,
}, nil
}
// Truncate updates the size of the target blob
func (w *writer) Truncate(size int64) error {
if size != 0 {
return errors.New("Truncate: unsupported size")
}
if w.cancel != nil {
err := w.cancel()
if err != nil {
return err
}
}
var r io.ReadCloser
r, w.pw = io.Pipe()
ctx, cancel := context.WithCancel(w.ctx)
go func() {
p, err := w.cln.Unixfs().Add(ctx, files.NewReaderFile(r), options.Unixfs.Pin(true))
if err != nil {
w.ipfsErr = err
return
}
dgst, err := digestconv.CidToDigest(p.Cid())
if err != nil {
w.ipfsErr = err
return
}
w.dgst = dgst
}()
w.cancel = func() error {
cancel()
w.ipfsErr = nil
err := w.Close()
if err != nil {
return err
}
return r.Close()
}
now := time.Now()
w.startedAt = now
w.updatedAt = now
w.offset = 0
return nil
}