All checks were successful
release-tag / release-image (push) Successful in 2m43s
159 lines
5.1 KiB
Go
159 lines
5.1 KiB
Go
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
|
|
}
|