date range
This commit is contained in:
+72
-51
@@ -1,6 +1,7 @@
|
||||
package bsdaily
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
@@ -8,27 +9,20 @@ import (
|
||||
"time"
|
||||
)
|
||||
|
||||
func Run(targetDate *time.Time) error {
|
||||
func Run(targetDates []time.Time) error {
|
||||
snapshotDir, snapshotDate, err := FindLatestDailySnapshot()
|
||||
if err != nil {
|
||||
return fmt.Errorf("finding latest snapshot: %w", err)
|
||||
}
|
||||
slog.Info("found latest daily snapshot", "dir", snapshotDir, "snapshot_date", snapshotDate.Format("2006-01-02"))
|
||||
|
||||
var targetDay time.Time
|
||||
if targetDate != nil {
|
||||
targetDay = *targetDate
|
||||
} else {
|
||||
targetDay = snapshotDate.AddDate(0, 0, -1)
|
||||
if len(targetDates) == 0 {
|
||||
targetDates = []time.Time{snapshotDate.AddDate(0, 0, -1)}
|
||||
}
|
||||
slog.Info("target day for extraction", "date", targetDay.Format("2006-01-02"))
|
||||
|
||||
// Check if output already exists
|
||||
outputDir := filepath.Join(DailiesBase, targetDay.Format("2006-01"))
|
||||
outputFinal := filepath.Join(outputDir, targetDay.Format("2006-01-02")+".sql.zst")
|
||||
if _, err := os.Stat(outputFinal); err == nil {
|
||||
return fmt.Errorf("output file already exists: %s", outputFinal)
|
||||
}
|
||||
slog.Info("target days for extraction", "count", len(targetDates),
|
||||
"first", targetDates[0].Format("2006-01-02"),
|
||||
"last", targetDates[len(targetDates)-1].Format("2006-01-02"))
|
||||
|
||||
// Check disk space
|
||||
if err := CheckFreeSpace(TmpBase, MinTmpFreeBytes, "tmpBase"); err != nil {
|
||||
@@ -82,46 +76,73 @@ func Run(targetDate *time.Time) error {
|
||||
}
|
||||
}
|
||||
|
||||
if err := CheckFreeSpace(TmpBase, PostCopyTmpMinFree, "tmpBase (post-copy)"); err != nil {
|
||||
return err
|
||||
// Process each day
|
||||
processed := 0
|
||||
skipped := 0
|
||||
|
||||
for _, targetDay := range targetDates {
|
||||
dayStr := targetDay.Format("2006-01-02")
|
||||
slog.Info("processing day", "date", dayStr)
|
||||
|
||||
// Check if output already exists
|
||||
outputDir := filepath.Join(DailiesBase, targetDay.Format("2006-01"))
|
||||
outputFinal := filepath.Join(outputDir, dayStr+".sql.zst")
|
||||
if _, err := os.Stat(outputFinal); err == nil {
|
||||
slog.Info("output already exists, skipping", "path", outputFinal)
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
|
||||
// Extract target day into a per-day database
|
||||
extractedDB := filepath.Join(tmpDir, "extracted-"+dayStr+".db")
|
||||
slog.Info("extracting target day", "src", dstDB, "dst", extractedDB)
|
||||
if err := ExtractDay(dstDB, extractedDB, targetDay); err != nil {
|
||||
if errors.Is(err, ErrNoPosts) {
|
||||
slog.Warn("no posts found, skipping day", "date", dayStr)
|
||||
os.Remove(extractedDB)
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
return fmt.Errorf("extracting day %s: %w", dayStr, err)
|
||||
}
|
||||
|
||||
// Dump to SQL and compress
|
||||
if err := os.MkdirAll(outputDir, 0755); err != nil {
|
||||
return fmt.Errorf("creating output directory %s: %w", outputDir, err)
|
||||
}
|
||||
|
||||
outputTmp := filepath.Join(outputDir, "."+dayStr+".sql.zst.tmp")
|
||||
|
||||
slog.Info("dumping and compressing", "tmp_output", outputTmp)
|
||||
if err := DumpAndCompress(extractedDB, outputTmp); err != nil {
|
||||
os.Remove(outputTmp)
|
||||
return fmt.Errorf("dump and compress for %s: %w", dayStr, err)
|
||||
}
|
||||
|
||||
slog.Info("verifying compressed output")
|
||||
if err := VerifyOutput(outputTmp); err != nil {
|
||||
os.Remove(outputTmp)
|
||||
return fmt.Errorf("verification failed for %s: %w", dayStr, err)
|
||||
}
|
||||
|
||||
// Atomic rename to final path
|
||||
slog.Info("renaming to final output", "from", outputTmp, "to", outputFinal)
|
||||
if err := os.Rename(outputTmp, outputFinal); err != nil {
|
||||
return fmt.Errorf("atomic rename for %s: %w", dayStr, err)
|
||||
}
|
||||
|
||||
info, err := os.Stat(outputFinal)
|
||||
if err != nil {
|
||||
return fmt.Errorf("stat final output: %w", err)
|
||||
}
|
||||
slog.Info("day completed", "date", dayStr, "path", outputFinal, "size_bytes", info.Size())
|
||||
|
||||
// Remove extracted DB to reclaim space
|
||||
os.Remove(extractedDB)
|
||||
processed++
|
||||
}
|
||||
|
||||
// Prune database to target day only
|
||||
slog.Info("opening database for pruning", "path", dstDB)
|
||||
if err := PruneDatabase(dstDB, targetDay); err != nil {
|
||||
return fmt.Errorf("pruning database: %w", err)
|
||||
}
|
||||
|
||||
// Dump to SQL and compress
|
||||
if err := os.MkdirAll(outputDir, 0755); err != nil {
|
||||
return fmt.Errorf("creating output directory %s: %w", outputDir, err)
|
||||
}
|
||||
|
||||
outputTmp := filepath.Join(outputDir, "."+targetDay.Format("2006-01-02")+".sql.zst.tmp")
|
||||
|
||||
slog.Info("dumping and compressing", "tmp_output", outputTmp)
|
||||
if err := DumpAndCompress(dstDB, outputTmp); err != nil {
|
||||
os.Remove(outputTmp)
|
||||
return fmt.Errorf("dump and compress: %w", err)
|
||||
}
|
||||
|
||||
slog.Info("verifying compressed output")
|
||||
if err := VerifyOutput(outputTmp); err != nil {
|
||||
os.Remove(outputTmp)
|
||||
return fmt.Errorf("verification failed: %w", err)
|
||||
}
|
||||
|
||||
// Atomic rename to final path
|
||||
slog.Info("renaming to final output", "from", outputTmp, "to", outputFinal)
|
||||
if err := os.Rename(outputTmp, outputFinal); err != nil {
|
||||
return fmt.Errorf("atomic rename: %w", err)
|
||||
}
|
||||
|
||||
info, err := os.Stat(outputFinal)
|
||||
if err != nil {
|
||||
return fmt.Errorf("stat final output: %w", err)
|
||||
}
|
||||
slog.Info("final output written", "path", outputFinal, "size_bytes", info.Size())
|
||||
slog.Info("run summary", "processed", processed, "skipped", skipped, "total", len(targetDates))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user