-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathproducer.go
More file actions
81 lines (65 loc) · 1.7 KB
/
Copy pathproducer.go
File metadata and controls
81 lines (65 loc) · 1.7 KB
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
package feedx
import (
"context"
"github.com/bsm/bfs"
)
// ProduceFunc is a callback which is run by the producer on every iteration.
type ProduceFunc func(*Writer) error
// Producer instances push data feeds to remote locations.
type Producer struct {
remote *bfs.Object
ownRemote bool
}
// NewProducer inits a new feed producer.
func NewProducer(ctx context.Context, remoteURL string) (*Producer, error) {
remote, err := bfs.NewObject(ctx, remoteURL)
if err != nil {
return nil, err
}
pcr := NewProducerForRemote(remote)
pcr.ownRemote = true
return pcr, nil
}
// NewProducerForRemote starts a new feed producer with a remote.
func NewProducerForRemote(remote *bfs.Object) *Producer {
return &Producer{remote: remote}
}
// Close stops the producer.
func (p *Producer) Close() error {
if p.ownRemote && p.remote != nil {
err := p.remote.Close()
p.remote = nil
return err
}
return nil
}
func (p *Producer) Produce(ctx context.Context, version int64, opt *WriterOptions, pfn ProduceFunc) (*Status, error) {
status := Status{LocalVersion: version}
// retrieve previous remote version
remoteVersion, err := fetchRemoteVersion(ctx, p.remote)
if err != nil {
return nil, err
}
status.RemoteVersion = remoteVersion
// skip if not modified
if skipSync(version, remoteVersion) {
status.Skipped = true
return &status, nil
}
// set version for writer
if opt == nil {
opt = new(WriterOptions)
}
opt.Version = version
// init writer and perform
writer := NewWriter(ctx, p.remote, opt)
defer writer.Discard()
if err := pfn(writer); err != nil {
return nil, err
}
if err := writer.Commit(); err != nil {
return nil, err
}
status.NumItems = writer.NumWritten()
return &status, nil
}