clonepack/internal/core/sync.go
GnomeZworc e1aff8d16b
add logging
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-04-25 20:00:00 +02:00

317 lines
9.1 KiB
Go

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)
}
}
}