refactor: remove dataset cron records

This commit is contained in:
2026-07-23 23:08:25 +08:00
parent 3705d75e59
commit fbd142803a
14 changed files with 125 additions and 380 deletions

View File

@@ -27,7 +27,6 @@ var (
const (
defaultItemPageSize = 50
maxItemPageSize = 200
datasetRunRetention = 30 * 24 * time.Hour
maxConcurrentFetches = 4
)
@@ -91,9 +90,10 @@ type ItemPage struct {
Offset int
}
type queuedSource struct {
Source models.SaDatasetSource
Run models.SaDatasetCron
type SyncResult struct {
SourceIdentity string
Status string
Result string
}
func NewService(database *gorm.DB) *Service {
@@ -228,9 +228,6 @@ func (s *Service) DeleteSource(userID uint, identity string) error {
if source.SeedKey != nil {
return ErrBuiltInSource
}
if err := tx.Where("source_id = ?", source.ID).Delete(&models.SaDatasetCron{}).Error; err != nil {
return err
}
if err := tx.Where("source_id = ?", source.ID).Delete(&models.SaDatasetItem{}).Error; err != nil {
return err
}
@@ -363,93 +360,32 @@ func (s *Service) DepositItem(userID uint, itemIdentity, projectIdentity string)
return DepositResult{NoteIdentity: note.Identity, ProjectIdentity: project.Identity}, nil
}
func (s *Service) ListCrons(userID uint) ([]models.SaDatasetCron, error) {
var crons []models.SaDatasetCron
if err := s.db.Model(&models.SaDatasetCron{}).
Select("sa_dataset_crons.*").
Joins("JOIN sa_dataset_sources ON sa_dataset_sources.id = sa_dataset_crons.source_id").
Where("sa_dataset_sources.owner_id = ?", userID).
Order("sa_dataset_crons.created_at desc, sa_dataset_crons.id desc").
Limit(100).Find(&crons).Error; err != nil {
return nil, err
}
if crons == nil {
crons = []models.SaDatasetCron{}
}
return crons, nil
}
func (s *Service) SyncSources(ctx context.Context, userID uint) ([]models.SaDatasetCron, error) {
s.syncMu.Lock()
defer s.syncMu.Unlock()
return s.syncSources(ctx, userID, "scheduled")
}
func (s *Service) QueueSources(userID uint) ([]models.SaDatasetCron, error) {
func (s *Service) SyncSources(ctx context.Context, userID uint) ([]SyncResult, error) {
if !s.syncMu.TryLock() {
return nil, ErrSyncInProgress
}
sources, err := s.enabledRSSSources(userID)
if err != nil {
s.syncMu.Unlock()
return nil, err
}
queued, err := s.createDatasetRuns(userID, sources, "manual")
if err != nil {
s.syncMu.Unlock()
return nil, err
}
runs := make([]models.SaDatasetCron, 0, len(queued))
for _, entry := range queued {
runs = append(runs, entry.Run)
}
if len(queued) == 0 {
s.syncMu.Unlock()
return runs, nil
}
go func() {
defer s.syncMu.Unlock()
_, _ = s.executeDatasetRuns(context.Background(), queued)
}()
return runs, nil
}
defer s.syncMu.Unlock()
func (s *Service) SyncAllSources(ctx context.Context) ([]models.SaDatasetCron, error) {
var userIDs []uint
if err := s.db.Model(&models.SaDatasetSource{}).
Distinct("owner_id").
Where("enabled = ? AND kind = ?", true, "rss").
Order("owner_id asc").
Pluck("owner_id", &userIDs).Error; err != nil {
return nil, err
}
result := make([]models.SaDatasetCron, 0)
var syncErrors []error
for _, userID := range userIDs {
crons, err := s.SyncSources(ctx, userID)
result = append(result, crons...)
if err != nil {
syncErrors = append(syncErrors, fmt.Errorf("sync dataset sources for user %d: %w", userID, err))
}
if ctx.Err() != nil {
break
}
}
return result, errors.Join(syncErrors...)
}
func (s *Service) syncSources(ctx context.Context, userID uint, trigger string) ([]models.SaDatasetCron, error) {
sources, err := s.enabledRSSSources(userID)
if err != nil {
return nil, err
}
queued, err := s.createDatasetRuns(userID, sources, trigger)
if err != nil {
return s.executeDatasetSources(ctx, sources)
}
func (s *Service) SyncAllSources(ctx context.Context) ([]SyncResult, error) {
if !s.syncMu.TryLock() {
return nil, ErrSyncInProgress
}
defer s.syncMu.Unlock()
var sources []models.SaDatasetSource
if err := s.db.Where("enabled = ? AND kind = ?", true, "rss").
Order("owner_id asc, id asc").
Find(&sources).Error; err != nil {
return nil, err
}
return s.executeDatasetRuns(ctx, queued)
return s.executeDatasetSources(ctx, sources)
}
func (s *Service) enabledRSSSources(userID uint) ([]models.SaDatasetSource, error) {
@@ -461,78 +397,41 @@ func (s *Service) enabledRSSSources(userID uint) ([]models.SaDatasetSource, erro
return sources, nil
}
func (s *Service) createDatasetRuns(userID uint, sources []models.SaDatasetSource, trigger string) ([]queuedSource, error) {
queued := make([]queuedSource, 0, len(sources))
now := time.Now().UTC()
err := s.db.Transaction(func(tx *gorm.DB) error {
if err := tx.Where("created_at < ?", now.Add(-datasetRunRetention)).
Delete(&models.SaDatasetCron{}).Error; err != nil {
return err
}
for _, source := range sources {
run := models.SaDatasetCron{
SourceID: source.ID, CreatedBy: userID, Schedule: trigger,
Status: "pending", Enabled: true, NextRunAt: &now,
}
if err := tx.Create(&run).Error; err != nil {
return err
}
queued = append(queued, queuedSource{Source: source, Run: run})
}
return nil
})
return queued, err
}
func (s *Service) executeDatasetRuns(ctx context.Context, queued []queuedSource) ([]models.SaDatasetCron, error) {
result := make([]models.SaDatasetCron, len(queued))
runErrors := make([]error, len(queued))
func (s *Service) executeDatasetSources(ctx context.Context, sources []models.SaDatasetSource) ([]SyncResult, error) {
result := make([]SyncResult, len(sources))
runErrors := make([]error, len(sources))
concurrency := maxConcurrentFetches
if s.db.Dialector.Name() == "sqlite" {
concurrency = 1
}
semaphore := make(chan struct{}, concurrency)
var workers sync.WaitGroup
for index, entry := range queued {
for index, source := range sources {
workers.Add(1)
go func() {
defer workers.Done()
semaphore <- struct{}{}
defer func() { <-semaphore }()
result[index], runErrors[index] = s.executeDatasetRun(ctx, entry)
result[index], runErrors[index] = s.executeDatasetSource(ctx, source)
}()
}
workers.Wait()
return result, errors.Join(runErrors...)
}
func (s *Service) executeDatasetRun(ctx context.Context, entry queuedSource) (models.SaDatasetCron, error) {
source := entry.Source
cron := entry.Run
func (s *Service) executeDatasetSource(ctx context.Context, source models.SaDatasetSource) (SyncResult, error) {
result := SyncResult{SourceIdentity: source.Identity}
if err := ctx.Err(); err != nil {
return cron, err
}
now := time.Now().UTC()
cron.Status = "running"
cron.LastRunAt = &now
if err := s.db.Model(&cron).Updates(map[string]any{
"status": "running", "last_run_at": now,
}).Error; err != nil {
return cron, err
result.Status = "failed"
result.Result = truncateResult(err.Error())
return result, err
}
feed, fetchErr := s.fetcher.Fetch(ctx, source.URL)
if fetchErr != nil {
cron.Status = "failed"
cron.Enabled = false
cron.NextRunAt = nil
cron.LastResult = truncateResult(fetchErr.Error())
if err := s.db.Model(&cron).Updates(map[string]any{
"status": cron.Status, "last_result": cron.LastResult, "enabled": false, "next_run_at": nil,
}).Error; err != nil {
return cron, err
}
return cron, ctx.Err()
result.Status = "failed"
result.Result = truncateResult(fetchErr.Error())
return result, nil
}
inserted := 0
@@ -543,30 +442,16 @@ func (s *Service) executeDatasetRun(ctx context.Context, entry queuedSource) (mo
if storeErr != nil {
return storeErr
}
cron.Status = "completed"
cron.Enabled = false
cron.NextRunAt = nil
cron.LastResult = fmt.Sprintf("format=%s fetched=%d inserted=%d", feed.Format, len(feed.Items), inserted)
if err := tx.Model(&cron).Updates(map[string]any{
"status": cron.Status, "last_result": cron.LastResult, "enabled": false, "next_run_at": nil,
}).Error; err != nil {
return err
}
return tx.Model(&source).Update("last_synced_at", completedAt).Error
})
if err == nil {
return cron, nil
if err != nil {
result.Status = "failed"
result.Result = truncateResult("store feed items: " + err.Error())
return result, nil
}
cron.Status = "failed"
cron.Enabled = false
cron.NextRunAt = nil
cron.LastResult = truncateResult("store feed items: " + err.Error())
if updateErr := s.db.Model(&cron).Updates(map[string]any{
"status": cron.Status, "last_result": cron.LastResult, "enabled": false, "next_run_at": nil,
}).Error; updateErr != nil {
return cron, updateErr
}
return cron, nil
result.Status = "completed"
result.Result = fmt.Sprintf("format=%s fetched=%d inserted=%d", feed.Format, len(feed.Items), inserted)
return result, nil
}
func storeFeedItems(tx *gorm.DB, source models.SaDatasetSource, items []FeedItem) (int, error) {