/
qm_sync.go
118 lines (96 loc) · 2.38 KB
/
qm_sync.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
package server
import (
"bufio"
"context"
"fmt"
"net/http"
"time"
"github.com/labstack/echo/v4"
)
const _qmFileKey = "_data/qm_cids.csv"
func (ss *MediorumServer) writeQmFile() error {
ctx := context.Background()
bail := func(err error) error {
if err != nil {
ss.bucket.Delete(ctx, _qmFileKey)
}
return err
}
// if exists do nothing
if exists, _ := ss.bucket.Exists(ctx, _qmFileKey); exists {
return nil
}
// blob writer
blobWriter, err := ss.bucket.NewWriter(ctx, _qmFileKey, nil)
if err != nil {
return bail(err)
}
// db conn
conn, err := ss.pgPool.Acquire(ctx)
if err != nil {
return bail(err)
}
defer conn.Release()
// doit
_, err = conn.Conn().PgConn().CopyTo(ctx, blobWriter, "COPY qm_cids TO STDOUT")
if err != nil {
return bail(err)
}
return bail(blobWriter.Close())
}
func (ss *MediorumServer) serveInternalQmCsv(c echo.Context) error {
r, err := ss.bucket.NewReader(c.Request().Context(), _qmFileKey, nil)
if err != nil {
return err
}
return c.Stream(200, "text/plain", r)
}
func (ss *MediorumServer) pullQmFromPeer(host string) error {
ctx := context.Background()
done := false
ss.pgPool.QueryRow(ctx, "select count(*) = 1 from qm_sync where host = $1", host).Scan(&done)
if done {
return nil
}
req, err := http.Get(apiPath(host, "internal/qm.csv"))
if err != nil {
return err
}
defer req.Body.Close()
if req.StatusCode != 200 {
return fmt.Errorf("bad status %d", req.StatusCode)
}
tx, err := ss.pgPool.Begin(context.Background())
if err != nil {
return err
}
defer tx.Rollback(context.Background())
scanner := bufio.NewScanner(req.Body)
for scanner.Scan() {
_, err = tx.Exec(ctx, "insert into qm_cids values ($1) on conflict do nothing", scanner.Text())
if err != nil {
return err
}
}
err = tx.Commit(ctx)
if err != nil {
return err
}
_, err = ss.pgPool.Exec(ctx, "insert into qm_sync values($1)", host)
return err
}
func (ss *MediorumServer) startQmSyncer() {
time.Sleep(time.Minute * 1)
err := ss.writeQmFile()
if err != nil {
ss.logger.Error("qmSync: failed to write qm.csv file", "err", err)
}
time.Sleep(time.Minute * 1)
for _, peer := range ss.findHealthyPeers(time.Hour) {
if err = ss.pullQmFromPeer(peer); err != nil {
ss.logger.Error("qmSync: failed to pull qm.csv from peer", "peer", peer, "err", err)
} else {
ss.logger.Info("qmSync: pulled qm.csv from peer", "peer", peer)
}
}
}