-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathfile_writer.go
More file actions
174 lines (160 loc) · 3.79 KB
/
Copy pathfile_writer.go
File metadata and controls
174 lines (160 loc) · 3.79 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
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
package loggo
import (
"bufio"
"fmt"
"os"
"sync"
"sync/atomic"
"time"
)
var (
DEF_CHAN_LEN = 1024
DEF_FILE_PEM os.FileMode = 0644
DEF_BUF_SIZE = 8 * 1024
)
type FileWriter struct {
filename string
bw *bufio.Writer
file *os.File
closed uint32
ch chan Buffer
quit chan chan error
reopen chan chan error
putpool *sync.Pool
}
func (fw *FileWriter) WriteBuffer(b Buffer) (int, error) {
if atomic.LoadUint32(&fw.closed) == 1 {
return 0, fmt.Errorf("closed")
}
fw.ch <- b
return b.Len(), nil
}
func (fw *FileWriter) Write(b []byte) (int, error) {
if atomic.LoadUint32(&fw.closed) == 1 {
return 0, fmt.Errorf("closed")
}
fw.ch <- ByteWarp(b)
return len(b), nil
}
func (fw *FileWriter) WriteString(b string) (int, error) {
if atomic.LoadUint32(&fw.closed) == 1 {
//fmt.Printf("use closed chan.\n")
return 0, fmt.Errorf("closed")
}
fw.ch <- StringWarp(b)
return len(b), nil
}
func (fw *FileWriter) Close() {
//close(fw.ch)
//<-fw.quit
//fw.file.Close()
atomic.StoreUint32(&fw.closed, 1)
quited := make(chan error)
fw.quit <- quited
err := <-quited
if err != nil {
fmt.Printf("LOGGO %s close ret: %v\n", fw.filename, err)
}
}
func (fw *FileWriter) Reopen() error {
finish := make(chan error)
fw.reopen <- finish
err := <-finish
return err
}
func NewFileWriter(filename string) *FileWriter {
return newFileWriter(filename, nil)
}
func newFileWriter(filename string, put *sync.Pool) *FileWriter {
bw, file, err := newBufWriter(filename)
if err != nil {
panic(fmt.Sprintf("newBufWriter err. %v", err))
}
fw := &FileWriter{
filename: filename,
file: file,
bw: bw,
closed: 0,
ch: make(chan Buffer, DEF_CHAN_LEN),
//quit: make(chan struct{}),
quit: make(chan chan error),
reopen: make(chan chan error),
putpool: put,
}
syn := make(chan struct{})
go fw.loop(syn)
<-syn
return fw
}
func newBufWriter(filename string) (*bufio.Writer, *os.File, error) {
file, err := os.OpenFile(filename, os.O_WRONLY|os.O_CREATE|os.O_APPEND, DEF_FILE_PEM)
if err != nil {
return nil, nil, err
}
w := bufio.NewWriterSize(file, DEF_BUF_SIZE)
return w, file, nil
}
func (fw *FileWriter) loop(syn chan struct{}) {
flushTimer := time.NewTicker(time.Millisecond * 500)
close(syn)
for {
//bw := fw.bw
select {
case d := <-fw.ch:
//case d, ok := <-fw.ch:
//if !ok {
// goto END_FOR
//}
if _, err := fw.bw.Write(d.Bytes()); err != nil {
fmt.Printf("LOGGO ERROR log file write err. %v\n", err)
}
if fw.putpool != nil {
fw.putpool.Put(d)
}
case quited := <-fw.quit:
var reterr error
lost := len(fw.ch)
//fmt.Printf("quit lost %d\n", lost)
for i := 0; i < lost; i++ {
d := <-fw.ch
//fw.bw.Write(d)
fw.bw.Write(d.Bytes())
}
if err := fw.bw.Flush(); err != nil {
reterr = fmt.Errorf("LOGGO ERROR log file flush err. %v", err)
}
quited <- reterr
case <-flushTimer.C:
if err := fw.bw.Flush(); err != nil {
fmt.Printf("LOGGO ERROR log file flush err. %v\n", err)
}
case finish := <-fw.reopen:
bw, file, reterr := fw.doReopen()
if bw != nil && file != nil {
fw.bw = bw
fw.file = file
}
finish <- reterr
}
}
//END_FOR:
//if err := fw.bw.Flush(); err != nil {
// fmt.Printf("log quit flush err. %v\n", err)
//}
//close(fw.quit)
}
func (fw *FileWriter) doReopen() (*bufio.Writer, *os.File, error) {
var reterr error
if err := fw.bw.Flush(); err != nil {
reterr = fmt.Errorf("LOGGO ERROR reopen flush err: %v.", err)
}
bw, file, err := newBufWriter(fw.filename)
if err != nil {
reterr = fmt.Errorf("LOGGO ERROR log reopen newbuf err: %v. %v", err, reterr)
return nil, nil, reterr
}
if err := fw.file.Close(); err != nil {
reterr = fmt.Errorf("LOGGO ERROR log reopen close err: %v. %v", err, reterr)
}
return bw, file, reterr
}