package core import ( "context" "encoding/json" "errors" "fmt" "log/slog" "net/http" "os" "path/filepath" "strings" "time" aptclone "github.com/syonad/clonepack/internal/clone/apt" rpmclone "github.com/syonad/clonepack/internal/clone/rpm" "github.com/syonad/clonepack/internal/store" ) type SyncService struct { repoStore store.RepoStore pendingStore store.PendingPackageStore blockedStore store.BlockedPackageStore cloneSvc *CloneService snapshotSvc *SnapshotService dataDir string httpClient *http.Client } func NewSyncService( repoStore store.RepoStore, pendingStore store.PendingPackageStore, blockedStore store.BlockedPackageStore, cloneSvc *CloneService, snapshotSvc *SnapshotService, dataDir string, ) *SyncService { return &SyncService{ repoStore: repoStore, pendingStore: pendingStore, blockedStore: blockedStore, cloneSvc: cloneSvc, snapshotSvc: snapshotSvc, dataDir: dataDir, httpClient: &http.Client{Timeout: 5 * time.Minute}, } } func (s *SyncService) ValidateRepo(ctx context.Context, repoID int64) error { _, err := s.repoStore.GetRepo(ctx, repoID) return err } func (s *SyncService) ScanRepo(ctx context.Context, repoID int64) error { repo, err := s.repoStore.GetRepo(ctx, repoID) if err != nil { return err } localDir := filepath.Join(s.dataDir, "repos", fmt.Sprintf("%d", repoID), repo.Type) type candidate struct { Name string Version string Arch string Location string Checksum string ChecksumType string Size int64 } var candidates []candidate switch repo.Type { case "rpm": scanner := rpmclone.NewScanner(repo.SourceURL, localDir) pkgs, err := scanner.Scan(ctx) if err != nil { return fmt.Errorf("scan repo %d: %w", repoID, err) } for _, p := range pkgs { candidates = append(candidates, candidate{ Name: p.Name, Version: p.Version, Arch: p.Arch, Location: p.Location, Checksum: p.Checksum, ChecksumType: p.ChecksumType, Size: p.Size, }) } case "apt": var cfg aptclone.Config if err := json.Unmarshal([]byte(repo.Config), &cfg); err != nil { return fmt.Errorf("parse apt config for repo %d: %w", repoID, err) } scanner := aptclone.NewScanner(repo.SourceURL, localDir, cfg) pkgs, err := scanner.Scan(ctx) if err != nil { return fmt.Errorf("scan repo %d: %w", repoID, err) } for _, p := range pkgs { candidates = append(candidates, candidate{ Name: p.Package, Version: p.Version, Arch: p.Architecture, Location: p.Filename, Checksum: p.SHA256, ChecksumType: "sha256", Size: p.Size, }) } default: return fmt.Errorf("sync only supported for rpm and apt repos") } if len(candidates) == 0 { slog.Debug("sync: repo up to date", "repo_id", repoID) return nil } blocked, err := s.blockedStore.ListBlocked(ctx, repoID) if err != nil { return fmt.Errorf("list blocked for repo %d: %w", repoID, err) } blockedLocations := make(map[string]bool, len(blocked)) blockedNames := make(map[string]bool, len(blocked)) for _, b := range blocked { blockedLocations[b.Location] = true if b.Name != "" { blockedNames[b.Name] = true } } filtered := candidates[:0] for _, p := range candidates { if !blockedLocations[p.Location] && !blockedNames[p.Name] { filtered = append(filtered, p) } } candidates = filtered if len(candidates) == 0 { slog.Debug("sync: all new packages are blocked", "repo_id", repoID) return nil } slog.Info("sync: new packages found", "repo_id", repoID, "count", len(candidates)) switch repo.SyncMode { case "auto": _, err = s.cloneSvc.StartClone(ctx, repoID) if err != nil && err != ErrCloneAlreadyRunning { return fmt.Errorf("start clone for repo %d: %w", repoID, err) } case "manual": pending := make([]store.PendingPackage, len(candidates)) for i, p := range candidates { pending[i] = store.PendingPackage{ RepoID: repoID, Name: p.Name, Version: p.Version, Arch: p.Arch, Location: p.Location, Checksum: p.Checksum, ChecksumType: p.ChecksumType, Size: p.Size, } } if err := s.pendingStore.UpsertPending(ctx, pending); err != nil { return fmt.Errorf("upsert pending for repo %d: %w", repoID, err) } default: return fmt.Errorf("unknown sync_mode %q for repo %d", repo.SyncMode, repoID) } return nil } func (s *SyncService) ListPending(ctx context.Context, repoID int64) ([]store.PendingPackage, error) { if _, err := s.repoStore.GetRepo(ctx, repoID); err != nil { return nil, err } return s.pendingStore.ListPending(ctx, repoID) } func (s *SyncService) ApprovePending(ctx context.Context, repoID int64, ids []int64) error { repo, err := s.repoStore.GetRepo(ctx, repoID) if err != nil { return err } all, err := s.pendingStore.ListPending(ctx, repoID) if err != nil { return err } wanted := make(map[int64]bool, len(ids)) for _, id := range ids { wanted[id] = true } localDir := filepath.Join(s.dataDir, "repos", fmt.Sprintf("%d", repoID), repo.Type) var downloadErrors []error var anyApproved bool for _, pkg := range all { if !wanted[pkg.ID] { continue } url := strings.TrimRight(repo.SourceURL, "/") + "/" + pkg.Location dest := filepath.Join(localDir, filepath.FromSlash(pkg.Location)) if err := os.MkdirAll(filepath.Dir(dest), 0o755); err != nil { downloadErrors = append(downloadErrors, fmt.Errorf("mkdirall %s: %w", pkg.Location, err)) continue } slog.Info("downloading package", "repo_id", repoID, "package", pkg.Name, "version", pkg.Version) if _, err := rpmclone.DownloadAndVerify(ctx, s.httpClient, url, dest, pkg.ChecksumType, pkg.Checksum); err != nil { slog.Warn("download failed", "repo_id", repoID, "package", pkg.Name, "error", err) downloadErrors = append(downloadErrors, fmt.Errorf("download %s: %w", pkg.Name, err)) continue } slog.Info("package downloaded", "repo_id", repoID, "package", pkg.Name, "version", pkg.Version) // Supprime immédiatement de la liste — même si la suite échoue, le paquet est acquis. if err := s.pendingStore.DeletePending(ctx, []int64{pkg.ID}); err != nil { slog.Warn("delete pending failed", "pending_id", pkg.ID, "error", err) } anyApproved = true } if anyApproved { switch repo.Type { case "rpm": if err := rpmclone.RegenerateMetadata(localDir); err != nil { slog.Warn("metadata regeneration failed", "repo_id", repoID, "error", err) } label := "auto-" + time.Now().UTC().Format(time.RFC3339) if _, err := s.snapshotSvc.TakeSnapshot(context.Background(), repoID, label); err != nil { slog.Warn("auto-snapshot failed", "repo_id", repoID, "error", err) } case "apt": var cfg aptclone.Config if err := json.Unmarshal([]byte(repo.Config), &cfg); err != nil { slog.Warn("parse apt config failed", "repo_id", repoID, "error", err) } else if err := aptclone.RegenerateMetadata(localDir, cfg); err != nil { slog.Warn("metadata regeneration failed", "repo_id", repoID, "error", err) } } } return errors.Join(downloadErrors...) } func (s *SyncService) BlockPackages(ctx context.Context, repoID int64, pendingIDs []int64) error { if _, err := s.repoStore.GetRepo(ctx, repoID); err != nil { return err } all, err := s.pendingStore.ListPending(ctx, repoID) if err != nil { return err } wanted := make(map[int64]bool, len(pendingIDs)) for _, id := range pendingIDs { wanted[id] = true } var toBlock []store.BlockedPackage var toDelete []int64 for _, p := range all { if wanted[p.ID] { toBlock = append(toBlock, store.BlockedPackage{Name: p.Name, Location: p.Location}) toDelete = append(toDelete, p.ID) } } if len(toBlock) == 0 { return nil } if err := s.blockedStore.BlockPackages(ctx, repoID, toBlock); err != nil { return err } return s.pendingStore.DeletePending(ctx, toDelete) } func (s *SyncService) UnblockPackages(ctx context.Context, repoID int64, ids []int64) error { if _, err := s.repoStore.GetRepo(ctx, repoID); err != nil { return err } return s.blockedStore.UnblockPackages(ctx, repoID, ids) } func (s *SyncService) ListBlocked(ctx context.Context, repoID int64) ([]store.BlockedPackage, error) { if _, err := s.repoStore.GetRepo(ctx, repoID); err != nil { return nil, err } return s.blockedStore.ListBlocked(ctx, repoID) } func (s *SyncService) RejectPending(ctx context.Context, repoID int64, ids []int64) error { if _, err := s.repoStore.GetRepo(ctx, repoID); err != nil { return err } return s.pendingStore.DeletePending(ctx, ids) } func (s *SyncService) StartScheduler(ctx context.Context, interval time.Duration) { go func() { slog.Info("sync scheduler started", "interval", interval) ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ctx.Done(): slog.Info("sync scheduler stopped") return case <-ticker.C: s.runScheduledScan(ctx) } } }() } func (s *SyncService) runScheduledScan(ctx context.Context) { repos, err := s.repoStore.ListRepos(ctx) if err != nil { slog.Error("sync scheduler: list repos failed", "error", err) return } for _, repo := range repos { if err := s.ScanRepo(ctx, repo.ID); err != nil { slog.Warn("sync scheduler: scan failed", "repo_id", repo.ID, "error", err) } } }