feat: scanner for directory reconciliation
This commit is contained in:
@@ -0,0 +1,109 @@
|
||||
package scanner
|
||||
|
||||
import (
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"time"
|
||||
"yoresee_dropbox/internal/store"
|
||||
)
|
||||
|
||||
type Scanner struct {
|
||||
store *store.Store
|
||||
workspace string
|
||||
interval time.Duration
|
||||
stopCh chan struct{}
|
||||
}
|
||||
|
||||
func New(s *store.Store, workspace string, interval time.Duration) *Scanner {
|
||||
return &Scanner{
|
||||
store: s,
|
||||
workspace: workspace,
|
||||
interval: interval,
|
||||
stopCh: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
func (sc *Scanner) Start() {
|
||||
go sc.run()
|
||||
}
|
||||
|
||||
func (sc *Scanner) Stop() {
|
||||
close(sc.stopCh)
|
||||
}
|
||||
|
||||
func (sc *Scanner) run() {
|
||||
ticker := time.NewTicker(sc.interval)
|
||||
defer ticker.Stop()
|
||||
|
||||
sc.scan()
|
||||
for {
|
||||
select {
|
||||
case <-sc.stopCh:
|
||||
return
|
||||
case <-ticker.C:
|
||||
sc.scan()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (sc *Scanner) scan() {
|
||||
sc.scanDir("inbox")
|
||||
sc.scanDir("outbox")
|
||||
}
|
||||
|
||||
func (sc *Scanner) scanDir(dir string) {
|
||||
dirPath := filepath.Join(sc.workspace, dir)
|
||||
entries, err := os.ReadDir(dirPath)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
dbFiles, err := sc.store.FileList(dir)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
dbMap := make(map[string]*store.File)
|
||||
for _, f := range dbFiles {
|
||||
dbMap[f.StorageName] = f
|
||||
}
|
||||
|
||||
diskMap := make(map[string]bool)
|
||||
for _, entry := range entries {
|
||||
if entry.IsDir() {
|
||||
continue
|
||||
}
|
||||
info, err := entry.Info()
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
diskMap[entry.Name()] = true
|
||||
|
||||
if _, exists := dbMap[entry.Name()]; !exists {
|
||||
id := generateUUID()
|
||||
f := &store.File{
|
||||
ID: id,
|
||||
OriginalName: entry.Name(),
|
||||
StorageName: entry.Name(),
|
||||
Dir: dir,
|
||||
Size: info.Size(),
|
||||
CreatedAt: time.Now().Unix(),
|
||||
}
|
||||
sc.store.FileCreate(f)
|
||||
}
|
||||
}
|
||||
|
||||
for name, f := range dbMap {
|
||||
if !diskMap[name] {
|
||||
sc.store.FileDelete(f.ID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func generateUUID() string {
|
||||
b := make([]byte, 16)
|
||||
rand.Read(b)
|
||||
return hex.EncodeToString(b)
|
||||
}
|
||||
@@ -0,0 +1,44 @@
|
||||
package scanner
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
"yoresee_dropbox/internal/store"
|
||||
)
|
||||
|
||||
func TestScannerReconciles(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
inbox := filepath.Join(dir, "inbox")
|
||||
outbox := filepath.Join(dir, "outbox")
|
||||
os.MkdirAll(inbox, 0755)
|
||||
os.MkdirAll(outbox, 0755)
|
||||
|
||||
dbPath := filepath.Join(dir, "test.db")
|
||||
s, err := store.Open(dbPath)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer s.Close()
|
||||
|
||||
// Create a file on disk
|
||||
testFile := filepath.Join(inbox, "test.txt")
|
||||
os.WriteFile(testFile, []byte("hello"), 0644)
|
||||
|
||||
sc := New(s, dir, 100*time.Millisecond)
|
||||
sc.Start()
|
||||
time.Sleep(300 * time.Millisecond)
|
||||
sc.Stop()
|
||||
|
||||
files, err := s.FileList("inbox")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(files) != 1 {
|
||||
t.Errorf("Expected 1 file, got %d", len(files))
|
||||
}
|
||||
if files[0].OriginalName != "test.txt" {
|
||||
t.Errorf("OriginalName = %q, want %q", files[0].OriginalName, "test.txt")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user