-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathregistrar.go
104 lines (94 loc) · 2.79 KB
/
registrar.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
package main
import (
"fmt"
"os"
"path/filepath"
"github.com/Shopify/sarama"
log "github.com/Sirupsen/logrus"
)
type Registrar struct {
source string
dir string
file *os.File
recordOpt func(int64) error
publishCtrl chan bool
}
const REGISTRAR_DIR string = "record"
func (r *Registrar) SilenceRecordOffset(offset int64) error {
return nil
}
func (r *Registrar) OpenRecord(recordDir string) (*os.File, error) {
if _, err := os.Stat(recordDir); os.IsNotExist(err) {
log.Errorf("%s not exist, create it now", recordDir)
err := os.MkdirAll(recordDir, 0755)
if err != nil {
if err != os.ErrExist {
log.Errorf("MkdirAll %s failed, error: %s", recordDir, err)
return nil, err
}
}
}
path := filepath.Join(recordDir, filepath.Base(r.source))
file, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY, 0666)
if err != nil {
log.Errorf("open %s failed, error:%s", path, err)
return nil, err
}
log.Infof("registrar open record %s ok %v", path, file)
r.file = file
return file, nil
}
func (r *Registrar) RecordOffset(offset int64) error {
if _, err := r.file.WriteAt([]byte(fmt.Sprintf("%d", offset)), 0); err != nil {
log.Errorf("record offset of %s failed, error: %s", r.file.Name(), err)
return err
}
//log.Printf("record offset of %s, offset %d\n", r.file.Name(), offset)
return nil
}
func (r *Registrar) RecordSucceed(offset, rawBytes int64) error {
return r.RecordOffset(offset + rawBytes)
}
func (r *Registrar) doBackup(text string) {
}
func (r *Registrar) RegistrarDo(errorChan <-chan *sarama.ProducerError, succChan <-chan *sarama.ProducerMessage) {
// record root dir: $work_dir/record, path: $record_root/$log_name
recordRoot := filepath.Join("./", "record")
if _, err := r.OpenRecord(recordRoot); err != nil {
log.Debug("set to Silence")
r.recordOpt = r.SilenceRecordOffset
} else {
log.Debug("set to file record")
r.recordOpt = r.RecordOffset
defer r.file.Close()
}
for {
select {
case err, ok := <-errorChan:
if !ok {
return
}
log.Error("Received Error:", err)
fev := err.Msg.Metadata.(*FileEvent)
r.recordOpt(err.Msg.Metadata.(*FileEvent).Offset + fev.RawBytes)
// record to retryer
msg := fmt.Sprint(filepath.Base(*fev.Source)+" ", *fev.Text)
mainRetryer.doBackup(msg)
// set remote serve failure
remoteAvailable = false
r.publishCtrl <- false
case success, ok := <-succChan:
if !ok {
return
}
log.Debug("Received OK:", success.Metadata.(*FileEvent).RawBytes,
success.Metadata.(*FileEvent).Offset,
*success.Metadata.(*FileEvent).Source)
fev := success.Metadata.(*FileEvent)
//log.Printf("registrar(%s), fileEvent(%s)\n", r.file.Name(), *fev.Source)
r.recordOpt(success.Metadata.(*FileEvent).Offset + fev.RawBytes)
// TODO: sync file with a switch flag
//r.file.Sync()
}
}
}