thread-master/apps/backend/internal/module/scout/usecase/run_publish.go

113 lines
3.4 KiB
Go
Raw Permalink Normal View History

2026-08-13 02:22:24 +00:00
package usecase
import (
"context"
"fmt"
"apps/backend/internal/module/scout/domain"
)
// StageRunPosts writes a run's candidate set behind the repository visibility
// barrier. Repositories expose these posts only after the run is succeeded.
func (s *Service) StageRunPosts(ctx context.Context, ownerUID int64, runID string, posts []*domain.Post) error {
if s == nil || s.Repo == nil || runID == "" {
return domain.ErrValidation
}
r, err := s.Repo.GetRun(ctx, ownerUID, runID)
if err != nil {
return err
}
if r.Status != domain.RunRunning {
return domain.ErrIllegalRunStatus
}
for _, post := range posts {
if post == nil || post.ID == "" || post.OwnerUID != ownerUID || post.RunID != runID {
return fmt.Errorf("%w: staged post must belong to run", domain.ErrValidation)
}
}
return s.Repo.PublishRunPosts(ctx, ownerUID, runID, posts)
}
// PublishRun atomically changes visibility by transitioning running →
// succeeded only after all staged posts have been written. If writing fails,
// the run remains non-visible and the caller can fail it safely.
func (s *Service) PublishRun(ctx context.Context, ownerUID int64, runID string, posts []*domain.Post) error {
if err := s.StageRunPosts(ctx, ownerUID, runID, posts); err != nil {
return err
}
r, err := s.Repo.GetRun(ctx, ownerUID, runID)
if err != nil {
return err
}
r.EligibleCount = len(posts)
r.PendingCount = countPendingRunPosts(posts)
if r.TargetCount > 0 && len(posts) < r.TargetCount {
r.ShortfallCount = r.TargetCount - len(posts)
if len(r.ShortfallReasons) == 0 {
r.ShortfallReasons = []string{domain.ShortfallSourceExhausted}
}
} else {
r.ShortfallCount = 0
r.ShortfallReasons = []string{}
}
if err := r.Transition(domain.RunSucceeded, domain.NowNano()); err != nil {
return err
}
if err := s.Repo.ReplaceRunGuarded(ctx, ownerUID, runID, []string{domain.RunRunning}, r); err != nil {
return err
}
// Seen identities are recorded only after the visibility barrier opens.
// Marking is idempotent; a failure here must not turn a visible successful
// run into a misleading failed response.
for _, post := range posts {
identity := canonicalPostIdentity(post.Permalink)
if identity == "" {
identity = canonicalPostIdentity(post.ExternalID)
}
if identity == "" && post.ID != "" {
identity = "id:" + post.ID
}
if identity != "" {
_ = s.Repo.MarkSeenIdentity(ctx, ownerUID, identity, post.ID, domain.NowNano())
}
}
return nil
}
// FailRun is idempotent for terminal runs and guarded for queued/running
// states, so a retry cannot move a completed run backwards.
func (s *Service) FailRun(ctx context.Context, ownerUID int64, runID, reason string) error {
if s == nil || s.Repo == nil || runID == "" {
return domain.ErrValidation
}
r, err := s.Repo.GetRun(ctx, ownerUID, runID)
if err != nil {
return err
}
if domain.IsRunTerminal(r.Status) {
return nil
}
if err := r.Transition(domain.RunFailed, domain.NowNano()); err != nil {
return err
}
r.Error = safeRunError(reason)
return s.Repo.ReplaceRunGuarded(ctx, ownerUID, runID, []string{domain.RunQueued, domain.RunRunning}, r)
}
func countPendingRunPosts(posts []*domain.Post) int {
n := 0
for _, post := range posts {
if post != nil && (post.OutreachStatus == domain.OutreachNew || post.OutreachStatus == domain.OutreachDrafted) {
n++
}
}
return n
}
func safeRunError(reason string) string {
if len(reason) > 500 {
return reason[:500]
}
return reason
}