This repository has been archived by the owner on Feb 12, 2019. It is now read-only.
/
journal_server_util.go
157 lines (146 loc) · 4.21 KB
/
journal_server_util.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
// Copyright 2016 Keybase Inc. All rights reserved.
// Use of this source code is governed by a BSD
// license that can be found in the LICENSE file.
package libkbfs
import (
"errors"
"time"
"github.com/keybase/client/go/logger"
"github.com/keybase/kbfs/tlf"
"golang.org/x/net/context"
"golang.org/x/sync/errgroup"
)
// GetJournalServer returns the JournalServer tied to a particular
// config.
func GetJournalServer(config Config) (*JournalServer, error) {
bserver := config.BlockServer()
jbserver, ok := bserver.(journalBlockServer)
if !ok {
return nil, errors.New("Write journal not enabled")
}
return jbserver.jServer, nil
}
// TLFJournalEnabled returns true if journaling is enabled for the
// given TLF.
func TLFJournalEnabled(config Config, tlfID tlf.ID) bool {
if jServer, err := GetJournalServer(config); err == nil {
_, err := jServer.JournalStatus(tlfID)
return err == nil
}
return false
}
// WaitForTLFJournal waits for the corresponding journal to flush, if
// one exists.
func WaitForTLFJournal(ctx context.Context, config Config, tlfID tlf.ID,
log logger.Logger) error {
if jServer, err := GetJournalServer(config); err == nil {
log.CDebugf(ctx, "Waiting for journal to flush")
if err := jServer.Wait(ctx, tlfID); err != nil {
return err
}
}
return nil
}
// FillInJournalStatusUnflushedPaths adds the unflushed paths to the
// given journal status.
func FillInJournalStatusUnflushedPaths(ctx context.Context, config Config,
jStatus *JournalServerStatus, tlfIDs []tlf.ID) error {
if len(tlfIDs) == 0 {
// Nothing to do.
return nil
}
// Get the folder statuses in parallel.
eg, groupCtx := errgroup.WithContext(ctx)
statusesToFetch := make(chan tlf.ID, len(tlfIDs))
unflushedPaths := make(chan []string, len(tlfIDs))
storedBytes := make(chan int64, len(tlfIDs))
unflushedBytes := make(chan int64, len(tlfIDs))
endEstimates := make(chan *time.Time, len(tlfIDs))
errIncomplete := errors.New("Incomplete status")
statusFn := func() error {
for tlfID := range statusesToFetch {
select {
case <-groupCtx.Done():
return groupCtx.Err()
default:
}
status, _, err := config.KBFSOps().FolderStatus(
groupCtx, FolderBranch{Tlf: tlfID, Branch: MasterBranch})
if err != nil {
return err
}
if status.Journal == nil {
continue
}
up := status.Journal.UnflushedPaths
unflushedPaths <- up
if len(up) > 0 && up[len(up)-1] == incompleteUnflushedPathsMarker {
// There were too many paths to process. Return an
// error to stop the other statuses since we have
// enough to return now.
return errIncomplete
}
storedBytes <- status.Journal.StoredBytes
unflushedBytes <- status.Journal.UnflushedBytes
endEstimates <- status.Journal.EndEstimate
}
return nil
}
// Do up to 10 statuses at a time.
numWorkers := len(tlfIDs)
if numWorkers > 10 {
numWorkers = 10
}
for i := 0; i < numWorkers; i++ {
eg.Go(statusFn)
}
for _, tlfID := range tlfIDs {
statusesToFetch <- tlfID
}
close(statusesToFetch)
if err := eg.Wait(); err != nil && err != errIncomplete {
return err
}
close(unflushedPaths)
close(storedBytes)
close(unflushedBytes)
close(endEstimates)
// Aggregate all the paths together, but only allow one incomplete
// marker, at the very end.
incomplete := false
for up := range unflushedPaths {
for _, p := range up {
if p == incompleteUnflushedPathsMarker {
incomplete = true
continue
}
jStatus.UnflushedPaths = append(jStatus.UnflushedPaths, p)
}
}
if incomplete {
jStatus.UnflushedPaths = append(jStatus.UnflushedPaths,
incompleteUnflushedPathsMarker)
} else {
// Replace the existing unflushed byte count with one
// that's guaranteed consistent with the unflushed
// paths, and also replace the existing stored byte
// count with one that's guaranteed consistent with
// the new unflushed byte count.
jStatus.StoredBytes = 0
for sb := range storedBytes {
jStatus.StoredBytes += sb
}
jStatus.UnflushedBytes = 0
for ub := range unflushedBytes {
jStatus.UnflushedBytes += ub
}
// Pick the latest end estimate.
for e := range endEstimates {
if e != nil &&
(jStatus.EndEstimate == nil || jStatus.EndEstimate.Before(*e)) {
jStatus.EndEstimate = e
}
}
}
return nil
}