package sourceagent import ( "context" "errors" "fmt" "strings" "time" "github.com/local/glpi-neural-brain/internal/articlequality" ) const ComputeKindArticleQuality = "article_quality" type ArticleQualityComputeRequest struct { SchemaVersion int `json:"schema_version"` JobID string `json:"job_id"` Kind string `json:"kind"` CreatedAt time.Time `json:"created_at"` Payload articlequality.Request `json:"payload"` } type ArticleQualityComputeResult struct { SchemaVersion int `json:"schema_version"` JobID string `json:"job_id"` Kind string `json:"kind"` AgentID string `json:"agent_id,omitempty"` DurationMS int64 `json:"duration_ms"` Quality articlequality.Result `json:"quality"` Error string `json:"error,omitempty"` } type articleQualityJobState struct { request ArticleQualityComputeRequest status, claimedBy string leaseUntil time.Time done chan ArticleQualityComputeResult } func (s *Store) SubmitArticleQualityJob(ctx context.Context, request ArticleQualityComputeRequest) (ArticleQualityComputeResult, error) { if s == nil { return ArticleQualityComputeResult{}, errors.New("source agent store unavailable") } if strings.TrimSpace(request.JobID) == "" { request.JobID = randomID("article-quality") } request.SchemaVersion = SchemaVersion request.Kind = ComputeKindArticleQuality if request.CreatedAt.IsZero() { request.CreatedAt = time.Now().UTC() } request.Payload.JobID = request.JobID st := &articleQualityJobState{request: request, status: "queued", done: make(chan ArticleQualityComputeResult, 1)} s.articleQualityMu.Lock() if s.articleQualityJobs == nil { s.articleQualityJobs = map[string]*articleQualityJobState{} } if _, ok := s.articleQualityJobs[request.JobID]; ok { s.articleQualityMu.Unlock() return ArticleQualityComputeResult{}, fmt.Errorf("article quality compute job %s already exists", request.JobID) } s.articleQualityJobs[request.JobID] = st s.articleQualityOrder = append(s.articleQualityOrder, request.JobID) s.articleQualityMu.Unlock() select { case r := <-st.done: if r.Error != "" { return r, errors.New(r.Error) } return r, nil case <-ctx.Done(): s.articleQualityMu.Lock() delete(s.articleQualityJobs, request.JobID) s.removeArticleQualityOrderLocked(request.JobID) s.articleQualityMu.Unlock() return ArticleQualityComputeResult{}, ctx.Err() } } func (s *Store) ClaimArticleQualityJob(agentID string, lease time.Duration) (ArticleQualityComputeRequest, bool) { if s == nil { return ArticleQualityComputeRequest{}, false } if lease < 30*time.Second { lease = 30 * time.Second } now := time.Now().UTC() s.articleQualityMu.Lock() defer s.articleQualityMu.Unlock() for _, id := range s.articleQualityOrder { st := s.articleQualityJobs[id] if st == nil { continue } if st.status == "claimed" && now.Before(st.leaseUntil) { continue } if st.status != "queued" && st.status != "claimed" { continue } st.status = "claimed" st.claimedBy = strings.TrimSpace(agentID) st.leaseUntil = now.Add(lease) return st.request, true } return ArticleQualityComputeRequest{}, false } func (s *Store) CompleteArticleQualityJob(agentID string, result ArticleQualityComputeResult) error { if s == nil { return errors.New("source agent store unavailable") } if strings.TrimSpace(result.JobID) == "" { return errors.New("article quality result job_id required") } s.articleQualityMu.Lock() st := s.articleQualityJobs[result.JobID] if st == nil { s.articleQualityMu.Unlock() return errors.New("article quality compute job not found or no longer active") } if st.status != "claimed" || strings.TrimSpace(st.claimedBy) == "" { s.articleQualityMu.Unlock() return errors.New("article quality compute job has not been claimed") } if strings.TrimSpace(agentID) != st.claimedBy { s.articleQualityMu.Unlock() return errors.New("article quality compute job is claimed by another agent") } if result.Kind == "" { result.Kind = ComputeKindArticleQuality } if result.Kind != ComputeKindArticleQuality { s.articleQualityMu.Unlock() return errors.New("unsupported article quality result kind") } result.SchemaVersion = SchemaVersion result.AgentID = strings.TrimSpace(agentID) delete(s.articleQualityJobs, result.JobID) s.removeArticleQualityOrderLocked(result.JobID) s.articleQualityMu.Unlock() select { case st.done <- result: default: } return nil } func (s *Store) removeArticleQualityOrderLocked(id string) { for i, v := range s.articleQualityOrder { if v == id { s.articleQualityOrder = append(s.articleQualityOrder[:i], s.articleQualityOrder[i+1:]...) return } } } func ExecuteArticleQualityJob(request ArticleQualityComputeRequest) ArticleQualityComputeResult { started := time.Now() r := ArticleQualityComputeResult{SchemaVersion: SchemaVersion, JobID: request.JobID, Kind: ComputeKindArticleQuality} r.Quality = articlequality.Evaluate(request.Payload) r.DurationMS = time.Since(started).Milliseconds() return r }