Skip to content

Commit 36cb53c

Browse files
authored
fix(search): reload replaced shards with unchanged mtimes (#1116)
1 parent 44b77f9 commit 36cb53c

2 files changed

Lines changed: 108 additions & 22 deletions

File tree

search/watcher.go

Lines changed: 42 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -37,9 +37,9 @@ type shardLoader interface {
3737
}
3838

3939
type DirectoryWatcher struct {
40-
dir string
41-
timestamps map[string]time.Time
42-
loader shardLoader
40+
dir string
41+
files map[string]watchedFile
42+
loader shardLoader
4343

4444
// closed once ready
4545
ready chan struct{}
@@ -52,6 +52,22 @@ type DirectoryWatcher struct {
5252
stopped chan struct{}
5353
}
5454

55+
type watchedFile struct {
56+
shard os.FileInfo
57+
meta os.FileInfo
58+
modTime time.Time
59+
}
60+
61+
func (f watchedFile) equal(other watchedFile) bool {
62+
if !f.modTime.Equal(other.modTime) || !os.SameFile(f.shard, other.shard) {
63+
return false
64+
}
65+
if (f.meta == nil) != (other.meta == nil) {
66+
return false
67+
}
68+
return f.meta == nil || os.SameFile(f.meta, other.meta)
69+
}
70+
5571
func (sw *DirectoryWatcher) Stop() {
5672
sw.closeOnce.Do(func() {
5773
close(sw.quit)
@@ -61,12 +77,12 @@ func (sw *DirectoryWatcher) Stop() {
6177

6278
func newDirectoryWatcher(dir string, loader shardLoader) (*DirectoryWatcher, error) {
6379
sw := &DirectoryWatcher{
64-
dir: dir,
65-
timestamps: map[string]time.Time{},
66-
loader: loader,
67-
ready: make(chan struct{}),
68-
quit: make(chan struct{}),
69-
stopped: make(chan struct{}),
80+
dir: dir,
81+
files: map[string]watchedFile{},
82+
loader: loader,
83+
ready: make(chan struct{}),
84+
quit: make(chan struct{}),
85+
stopped: make(chan struct{}),
7086
}
7187

7288
go func() {
@@ -140,7 +156,7 @@ func (s *DirectoryWatcher) scan() error {
140156
}
141157
}
142158

143-
ts := map[string]time.Time{}
159+
files := map[string]watchedFile{}
144160
for _, fn := range fs {
145161
if name, version := versionFromPath(fn); latest[name] != version {
146162
continue
@@ -151,31 +167,35 @@ func (s *DirectoryWatcher) scan() error {
151167
continue
152168
}
153169

154-
ts[fn] = fi.ModTime()
170+
current := watchedFile{
171+
shard: fi,
172+
modTime: fi.ModTime(),
173+
}
155174

156175
fiMeta, err := os.Lstat(fn + ".meta")
157-
if err != nil {
158-
continue
159-
}
160-
if fiMeta.ModTime().After(fi.ModTime()) {
161-
ts[fn] = fiMeta.ModTime()
176+
if err == nil {
177+
current.meta = fiMeta
178+
if fiMeta.ModTime().After(current.modTime) {
179+
current.modTime = fiMeta.ModTime()
180+
}
162181
}
182+
files[fn] = current
163183
}
164184

165185
var toLoad []string
166-
for k, mtime := range ts {
167-
if t, ok := s.timestamps[k]; !ok || t != mtime {
186+
for k, current := range files {
187+
if previous, ok := s.files[k]; !ok || !previous.equal(current) {
168188
toLoad = append(toLoad, k)
169-
s.timestamps[k] = mtime
189+
s.files[k] = current
170190
}
171191
}
172192

173193
var toDrop []string
174194
// Unload deleted shards.
175-
for k := range s.timestamps {
176-
if _, ok := ts[k]; !ok {
195+
for k := range s.files {
196+
if _, ok := files[k]; !ok {
177197
toDrop = append(toDrop, k)
178-
delete(s.timestamps, k)
198+
delete(s.files, k)
179199
}
180200
}
181201

search/watcher_test.go

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,72 @@ func advanceFS() {
4545
time.Sleep(10 * time.Millisecond)
4646
}
4747

48+
func TestDirWatcherReloadsReplacementWithSameModTime(t *testing.T) {
49+
for _, replaced := range []string{"shard", "meta"} {
50+
t.Run(replaced, func(t *testing.T) {
51+
dir := t.TempDir()
52+
shard := filepath.Join(dir, "foo.zoekt")
53+
meta := shard + ".meta"
54+
fixedTime := time.Unix(1_700_000_000, 0)
55+
56+
writeFileWithModTime(t, shard, "first shard", fixedTime)
57+
writeFileWithModTime(t, meta, "first meta", fixedTime)
58+
59+
logger := &loggingLoader{
60+
loads: make(chan string, 2),
61+
drops: make(chan string, 1),
62+
}
63+
dw := &DirectoryWatcher{
64+
dir: dir,
65+
files: map[string]watchedFile{},
66+
loader: logger,
67+
}
68+
if err := dw.scan(); err != nil {
69+
t.Fatal(err)
70+
}
71+
if got := <-logger.loads; got != shard {
72+
t.Fatalf("got initial load %q, want %q", got, shard)
73+
}
74+
75+
path := shard
76+
if replaced == "meta" {
77+
path = meta
78+
}
79+
replacement := path + ".replacement"
80+
writeFileWithModTime(t, replacement, "replacement", fixedTime)
81+
if err := os.Rename(replacement, path); err != nil {
82+
t.Fatal(err)
83+
}
84+
85+
if err := dw.scan(); err != nil {
86+
t.Fatal(err)
87+
}
88+
if got := <-logger.loads; got != shard {
89+
t.Fatalf("got replacement load %q, want %q", got, shard)
90+
}
91+
92+
if err := dw.scan(); err != nil {
93+
t.Fatal(err)
94+
}
95+
select {
96+
case got := <-logger.loads:
97+
t.Fatalf("unchanged replacement triggered another load of %q", got)
98+
default:
99+
}
100+
})
101+
}
102+
}
103+
104+
func writeFileWithModTime(t *testing.T, path, content string, modTime time.Time) {
105+
t.Helper()
106+
if err := os.WriteFile(path, []byte(content), 0o644); err != nil {
107+
t.Fatal(err)
108+
}
109+
if err := os.Chtimes(path, modTime, modTime); err != nil {
110+
t.Fatal(err)
111+
}
112+
}
113+
48114
func TestDirWatcherUnloadOnce(t *testing.T) {
49115
dir := t.TempDir()
50116

0 commit comments

Comments
 (0)