6 Commits

3 changed files with 30 additions and 3 deletions
+1
View File
@@ -1 +1,2 @@
media/ media/
goencode
+27 -2
View File
@@ -138,7 +138,7 @@ func (m *Manager) processWorker() {
case <-m.stopChan: case <-m.stopChan:
return return
case path := <-m.processChan: case path := <-m.processChan:
m.processFile(path) go m.processFile(path)
} }
} }
} }
@@ -151,7 +151,7 @@ func (m *Manager) handleEvent(filePath string) {
t.Stop() t.Stop()
} }
m.timers[filePath] = time.AfterFunc(5*time.Second, func() { m.timers[filePath] = time.AfterFunc(2*time.Minute, func() {
m.timersMu.Lock() m.timersMu.Lock()
delete(m.timers, filePath) delete(m.timers, filePath)
m.timersMu.Unlock() m.timersMu.Unlock()
@@ -176,6 +176,31 @@ func (m *Manager) processFile(filePath string) {
return return
} }
// Wait until the file is fully copied and stable
var lastSize int64 = info.Size()
var lastModTime time.Time = info.ModTime()
var stableCount int
for {
time.Sleep(5 * time.Second)
currentInfo, err := os.Stat(filePath)
if err != nil {
return // File was likely deleted during copying
}
if currentInfo.Size() == lastSize && currentInfo.ModTime().Equal(lastModTime) {
stableCount++
if stableCount >= 3 { // Stable for 15 seconds
break
}
} else {
stableCount = 0
lastSize = currentInfo.Size()
lastModTime = currentInfo.ModTime()
}
}
// Re-check after waiting
alreadyInQueue, err := db.IsFileAlreadyProcessedOrQueued(filePath) alreadyInQueue, err := db.IsFileAlreadyProcessedOrQueued(filePath)
if err != nil { if err != nil {
log.Printf("Error checking DB for %s: %v", filePath, err) log.Printf("Error checking DB for %s: %v", filePath, err)
+1
View File
@@ -165,5 +165,6 @@ func (s *Server) handleRequeueJob(w http.ResponseWriter, r *http.Request) {
s.qm.NotifySSE("queue_updated", nil) s.qm.NotifySSE("queue_updated", nil)
s.qm.NotifySSE("job_added", nil) s.qm.NotifySSE("job_added", nil)
s.qm.Trigger()
w.WriteHeader(http.StatusOK) w.WriteHeader(http.StatusOK)
} }