268 lines
7.3 KiB
Go
268 lines
7.3 KiB
Go
package core
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"net/http"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"time"
|
|
|
|
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
|
|
}
|
|
if repo.Type != "rpm" {
|
|
return fmt.Errorf("sync only supported for rpm repos")
|
|
}
|
|
|
|
localDir := filepath.Join(s.dataDir, "repos", fmt.Sprintf("%d", repoID), "rpm")
|
|
scanner := rpmclone.NewScanner(repo.SourceURL, localDir)
|
|
newPkgs, err := scanner.Scan(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("scan repo %d: %w", repoID, err)
|
|
}
|
|
|
|
if len(newPkgs) == 0 {
|
|
log.Printf("sync: repo %d is up to date", 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 := newPkgs[:0]
|
|
for _, p := range newPkgs {
|
|
if !blockedLocations[p.Location] && !blockedNames[p.Name] {
|
|
filtered = append(filtered, p)
|
|
}
|
|
}
|
|
newPkgs = filtered
|
|
|
|
if len(newPkgs) == 0 {
|
|
log.Printf("sync: repo %d is up to date (all new packages are blocked)", repoID)
|
|
return nil
|
|
}
|
|
log.Printf("sync: repo %d has %d new package(s)", repoID, len(newPkgs))
|
|
|
|
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(newPkgs))
|
|
for i, p := range newPkgs {
|
|
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), "rpm")
|
|
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
|
|
}
|
|
if _, err := rpmclone.DownloadAndVerify(ctx, s.httpClient, url, dest, pkg.ChecksumType, pkg.Checksum); err != nil {
|
|
downloadErrors = append(downloadErrors, fmt.Errorf("download %s: %w", pkg.Name, err))
|
|
continue
|
|
}
|
|
// 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 {
|
|
log.Printf("delete pending %d: %v", pkg.ID, err)
|
|
}
|
|
anyApproved = true
|
|
}
|
|
|
|
if anyApproved {
|
|
if repo.Type == "rpm" {
|
|
if err := rpmclone.RegenerateMetadata(localDir); err != nil {
|
|
log.Printf("metadata regeneration for repo %d failed: %v", repoID, err)
|
|
}
|
|
}
|
|
label := "auto-" + time.Now().UTC().Format(time.RFC3339)
|
|
if _, err := s.snapshotSvc.TakeSnapshot(context.Background(), repoID, label); err != nil {
|
|
log.Printf("auto-snapshot for repo %d failed: %v", repoID, 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() {
|
|
log.Printf("sync scheduler started (interval=%s)", interval)
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
log.Printf("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 {
|
|
log.Printf("sync scheduler: list repos error: %v", err)
|
|
return
|
|
}
|
|
for _, repo := range repos {
|
|
if err := s.ScanRepo(ctx, repo.ID); err != nil {
|
|
log.Printf("sync scheduler: scan repo %d error: %v", repo.ID, err)
|
|
}
|
|
}
|
|
}
|