All checks were successful
release-tag / release-image (push) Successful in 2m43s
393 lines
12 KiB
Go
393 lines
12 KiB
Go
package sourceagent
|
|
|
|
import (
|
|
"bufio"
|
|
"context"
|
|
"encoding/binary"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"math"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/local/glpi-neural-brain/internal/vectorgraph"
|
|
)
|
|
|
|
const (
|
|
ComputeKindVectorGraph = "vector_graph"
|
|
vectorJobMagic = "NBVJOB01"
|
|
)
|
|
|
|
type VectorGraphComputeHeader struct {
|
|
SchemaVersion int `json:"schema_version"`
|
|
JobID string `json:"job_id"`
|
|
Kind string `json:"kind"`
|
|
GraphVersion uint64 `json:"graph_version"`
|
|
CreatedAt time.Time `json:"created_at"`
|
|
Primary vectorgraph.Config `json:"primary"`
|
|
OrphanPass bool `json:"orphan_pass"`
|
|
Orphan vectorgraph.Config `json:"orphan"`
|
|
OrphanFocusIDs []string `json:"orphan_focus_ids,omitempty"`
|
|
}
|
|
|
|
type VectorGraphComputeRequest struct {
|
|
Header VectorGraphComputeHeader
|
|
Entries []vectorgraph.Entry
|
|
}
|
|
|
|
type VectorGraphComputeResult struct {
|
|
SchemaVersion int `json:"schema_version"`
|
|
JobID string `json:"job_id"`
|
|
Kind string `json:"kind"`
|
|
GraphVersion uint64 `json:"graph_version"`
|
|
AgentID string `json:"agent_id,omitempty"`
|
|
DurationMS int64 `json:"duration_ms"`
|
|
Primary vectorgraph.Result `json:"primary"`
|
|
Orphan vectorgraph.Result `json:"orphan"`
|
|
Error string `json:"error,omitempty"`
|
|
}
|
|
|
|
type computeJobState struct {
|
|
request VectorGraphComputeRequest
|
|
status string
|
|
claimedBy string
|
|
leaseUntil time.Time
|
|
done chan VectorGraphComputeResult
|
|
}
|
|
|
|
func (s *Store) SubmitVectorGraphJob(ctx context.Context, request VectorGraphComputeRequest) (VectorGraphComputeResult, error) {
|
|
if s == nil {
|
|
return VectorGraphComputeResult{}, errors.New("source agent store unavailable")
|
|
}
|
|
if len(request.Entries) < 2 {
|
|
return VectorGraphComputeResult{}, errors.New("vector graph job requires at least two entries")
|
|
}
|
|
if strings.TrimSpace(request.Header.JobID) == "" {
|
|
request.Header.JobID = randomID("compute")
|
|
}
|
|
request.Header.SchemaVersion = SchemaVersion
|
|
request.Header.Kind = ComputeKindVectorGraph
|
|
if request.Header.CreatedAt.IsZero() {
|
|
request.Header.CreatedAt = time.Now().UTC()
|
|
}
|
|
request.Header.OrphanFocusIDs = uniqueStrings(request.Header.OrphanFocusIDs)
|
|
state := &computeJobState{request: request, status: "queued", done: make(chan VectorGraphComputeResult, 1)}
|
|
|
|
s.computeMu.Lock()
|
|
if s.computeJobs == nil {
|
|
s.computeJobs = map[string]*computeJobState{}
|
|
}
|
|
if _, exists := s.computeJobs[request.Header.JobID]; exists {
|
|
s.computeMu.Unlock()
|
|
return VectorGraphComputeResult{}, fmt.Errorf("compute job %s already exists", request.Header.JobID)
|
|
}
|
|
s.computeJobs[request.Header.JobID] = state
|
|
s.computeOrder = append(s.computeOrder, request.Header.JobID)
|
|
s.computeMu.Unlock()
|
|
|
|
select {
|
|
case result := <-state.done:
|
|
if result.Error != "" {
|
|
return result, errors.New(result.Error)
|
|
}
|
|
return result, nil
|
|
case <-ctx.Done():
|
|
s.computeMu.Lock()
|
|
delete(s.computeJobs, request.Header.JobID)
|
|
s.removeComputeOrderLocked(request.Header.JobID)
|
|
s.computeMu.Unlock()
|
|
return VectorGraphComputeResult{}, ctx.Err()
|
|
}
|
|
}
|
|
|
|
func (s *Store) ClaimVectorGraphJob(agentID string, lease time.Duration) (VectorGraphComputeRequest, bool) {
|
|
if s == nil {
|
|
return VectorGraphComputeRequest{}, false
|
|
}
|
|
if lease < 30*time.Second {
|
|
lease = 30 * time.Second
|
|
}
|
|
now := time.Now().UTC()
|
|
s.computeMu.Lock()
|
|
defer s.computeMu.Unlock()
|
|
for _, id := range s.computeOrder {
|
|
state := s.computeJobs[id]
|
|
if state == nil {
|
|
continue
|
|
}
|
|
if state.status == "claimed" && now.Before(state.leaseUntil) {
|
|
continue
|
|
}
|
|
if state.status != "queued" && state.status != "claimed" {
|
|
continue
|
|
}
|
|
state.status = "claimed"
|
|
state.claimedBy = strings.TrimSpace(agentID)
|
|
state.leaseUntil = now.Add(lease)
|
|
return state.request, true
|
|
}
|
|
return VectorGraphComputeRequest{}, false
|
|
}
|
|
|
|
func (s *Store) CompleteVectorGraphJob(agentID string, result VectorGraphComputeResult) error {
|
|
if s == nil {
|
|
return errors.New("source agent store unavailable")
|
|
}
|
|
if strings.TrimSpace(result.JobID) == "" {
|
|
return errors.New("compute result job_id required")
|
|
}
|
|
s.computeMu.Lock()
|
|
state := s.computeJobs[result.JobID]
|
|
if state == nil {
|
|
s.computeMu.Unlock()
|
|
return errors.New("compute job not found or no longer active")
|
|
}
|
|
if state.status != "claimed" || strings.TrimSpace(state.claimedBy) == "" {
|
|
s.computeMu.Unlock()
|
|
return errors.New("compute job has not been claimed")
|
|
}
|
|
if strings.TrimSpace(agentID) != state.claimedBy {
|
|
s.computeMu.Unlock()
|
|
return errors.New("compute job is claimed by another agent")
|
|
}
|
|
if result.Kind == "" {
|
|
result.Kind = ComputeKindVectorGraph
|
|
}
|
|
if result.Kind != ComputeKindVectorGraph {
|
|
s.computeMu.Unlock()
|
|
return errors.New("unsupported compute result kind")
|
|
}
|
|
result.SchemaVersion = SchemaVersion
|
|
result.AgentID = strings.TrimSpace(agentID)
|
|
result.GraphVersion = state.request.Header.GraphVersion
|
|
delete(s.computeJobs, result.JobID)
|
|
s.removeComputeOrderLocked(result.JobID)
|
|
s.computeMu.Unlock()
|
|
|
|
select {
|
|
case state.done <- result:
|
|
default:
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Store) removeComputeOrderLocked(id string) {
|
|
for i, value := range s.computeOrder {
|
|
if value != id {
|
|
continue
|
|
}
|
|
s.computeOrder = append(s.computeOrder[:i], s.computeOrder[i+1:]...)
|
|
return
|
|
}
|
|
}
|
|
|
|
func (s *Store) ComputeStats() map[string]any {
|
|
if s == nil {
|
|
return map[string]any{"queued": 0, "claimed": 0, "expired_claims": 0, "article_quality_queued": 0, "article_quality_claimed": 0, "article_quality_expired_claims": 0}
|
|
}
|
|
now := time.Now().UTC()
|
|
queued, claimed, expired := 0, 0, 0
|
|
s.computeMu.Lock()
|
|
for _, state := range s.computeJobs {
|
|
switch state.status {
|
|
case "queued":
|
|
queued++
|
|
case "claimed":
|
|
if now.After(state.leaseUntil) {
|
|
expired++
|
|
} else {
|
|
claimed++
|
|
}
|
|
}
|
|
}
|
|
s.computeMu.Unlock()
|
|
articleQueued, articleClaimed, articleExpired := 0, 0, 0
|
|
s.articleQualityMu.Lock()
|
|
for _, state := range s.articleQualityJobs {
|
|
switch state.status {
|
|
case "queued":
|
|
articleQueued++
|
|
case "claimed":
|
|
if now.After(state.leaseUntil) {
|
|
articleExpired++
|
|
} else {
|
|
articleClaimed++
|
|
}
|
|
}
|
|
}
|
|
s.articleQualityMu.Unlock()
|
|
return map[string]any{"queued": queued, "claimed": claimed, "expired_claims": expired, "article_quality_queued": articleQueued, "article_quality_claimed": articleClaimed, "article_quality_expired_claims": articleExpired}
|
|
}
|
|
|
|
// WriteVectorGraphJob encodes a compute job as a compact binary stream. Vectors
|
|
// remain float32 instead of expanding to decimal JSON, which keeps a 21k x 768
|
|
// snapshot near its natural ~62 MiB size on the wire.
|
|
func WriteVectorGraphJob(w io.Writer, request VectorGraphComputeRequest) error {
|
|
bw := bufio.NewWriterSize(w, 256<<10)
|
|
if _, err := bw.WriteString(vectorJobMagic); err != nil {
|
|
return err
|
|
}
|
|
header, err := json.Marshal(request.Header)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(header) > 16<<20 {
|
|
return errors.New("vector job header too large")
|
|
}
|
|
if err := binary.Write(bw, binary.LittleEndian, uint32(len(header))); err != nil {
|
|
return err
|
|
}
|
|
if _, err := bw.Write(header); err != nil {
|
|
return err
|
|
}
|
|
if err := binary.Write(bw, binary.LittleEndian, uint32(len(request.Entries))); err != nil {
|
|
return err
|
|
}
|
|
for _, entry := range request.Entries {
|
|
if len(entry.ID) == 0 || len(entry.ID) > math.MaxUint16 {
|
|
return errors.New("invalid vector job entry id")
|
|
}
|
|
if len(entry.Vector) == 0 || len(entry.Vector) > 4096 {
|
|
return errors.New("invalid vector job dimensions")
|
|
}
|
|
if err := binary.Write(bw, binary.LittleEndian, uint16(len(entry.ID))); err != nil {
|
|
return err
|
|
}
|
|
if _, err := bw.WriteString(entry.ID); err != nil {
|
|
return err
|
|
}
|
|
if err := binary.Write(bw, binary.LittleEndian, uint16(len(entry.Vector))); err != nil {
|
|
return err
|
|
}
|
|
for _, value := range entry.Vector {
|
|
if err := binary.Write(bw, binary.LittleEndian, value); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
return bw.Flush()
|
|
}
|
|
|
|
func ReadVectorGraphJob(r io.Reader, maxBytes int64) (VectorGraphComputeRequest, error) {
|
|
if maxBytes <= 0 {
|
|
maxBytes = 128 << 20
|
|
}
|
|
limited := &countingReader{r: r, remaining: maxBytes}
|
|
br := bufio.NewReaderSize(limited, 256<<10)
|
|
magic := make([]byte, len(vectorJobMagic))
|
|
if _, err := io.ReadFull(br, magic); err != nil {
|
|
return VectorGraphComputeRequest{}, err
|
|
}
|
|
if string(magic) != vectorJobMagic {
|
|
return VectorGraphComputeRequest{}, errors.New("invalid vector compute payload magic")
|
|
}
|
|
var headerLen uint32
|
|
if err := binary.Read(br, binary.LittleEndian, &headerLen); err != nil {
|
|
return VectorGraphComputeRequest{}, err
|
|
}
|
|
if headerLen == 0 || headerLen > 16<<20 {
|
|
return VectorGraphComputeRequest{}, errors.New("invalid vector compute header length")
|
|
}
|
|
headerRaw := make([]byte, headerLen)
|
|
if _, err := io.ReadFull(br, headerRaw); err != nil {
|
|
return VectorGraphComputeRequest{}, err
|
|
}
|
|
var request VectorGraphComputeRequest
|
|
if err := json.Unmarshal(headerRaw, &request.Header); err != nil {
|
|
return VectorGraphComputeRequest{}, err
|
|
}
|
|
if request.Header.SchemaVersion != SchemaVersion {
|
|
return VectorGraphComputeRequest{}, fmt.Errorf("unsupported vector compute schema_version %d", request.Header.SchemaVersion)
|
|
}
|
|
if request.Header.Kind != ComputeKindVectorGraph || strings.TrimSpace(request.Header.JobID) == "" {
|
|
return VectorGraphComputeRequest{}, errors.New("invalid vector compute header")
|
|
}
|
|
var count uint32
|
|
if err := binary.Read(br, binary.LittleEndian, &count); err != nil {
|
|
return VectorGraphComputeRequest{}, err
|
|
}
|
|
if count < 2 || count > 100000 {
|
|
return VectorGraphComputeRequest{}, errors.New("invalid vector compute entry count")
|
|
}
|
|
request.Entries = make([]vectorgraph.Entry, 0, int(count))
|
|
for i := uint32(0); i < count; i++ {
|
|
var idLen uint16
|
|
if err := binary.Read(br, binary.LittleEndian, &idLen); err != nil {
|
|
return VectorGraphComputeRequest{}, err
|
|
}
|
|
if idLen == 0 || idLen > 1024 {
|
|
return VectorGraphComputeRequest{}, errors.New("invalid vector compute id length")
|
|
}
|
|
idRaw := make([]byte, int(idLen))
|
|
if _, err := io.ReadFull(br, idRaw); err != nil {
|
|
return VectorGraphComputeRequest{}, err
|
|
}
|
|
var dimensions uint16
|
|
if err := binary.Read(br, binary.LittleEndian, &dimensions); err != nil {
|
|
return VectorGraphComputeRequest{}, err
|
|
}
|
|
if dimensions == 0 || dimensions > 4096 {
|
|
return VectorGraphComputeRequest{}, errors.New("invalid vector compute dimensions")
|
|
}
|
|
vector := make([]float32, int(dimensions))
|
|
for j := range vector {
|
|
if err := binary.Read(br, binary.LittleEndian, &vector[j]); err != nil {
|
|
return VectorGraphComputeRequest{}, err
|
|
}
|
|
}
|
|
request.Entries = append(request.Entries, vectorgraph.Entry{ID: string(idRaw), Vector: vector})
|
|
}
|
|
return request, nil
|
|
}
|
|
|
|
type countingReader struct {
|
|
r io.Reader
|
|
remaining int64
|
|
}
|
|
|
|
func (r *countingReader) Read(p []byte) (int, error) {
|
|
if r.remaining <= 0 {
|
|
return 0, errors.New("vector compute payload exceeds configured limit")
|
|
}
|
|
if int64(len(p)) > r.remaining {
|
|
p = p[:r.remaining]
|
|
}
|
|
n, err := r.r.Read(p)
|
|
r.remaining -= int64(n)
|
|
if r.remaining <= 0 && err == nil {
|
|
return n, errors.New("vector compute payload exceeds configured limit")
|
|
}
|
|
return n, err
|
|
}
|
|
|
|
func ExecuteVectorGraphJob(request VectorGraphComputeRequest) VectorGraphComputeResult {
|
|
started := time.Now()
|
|
result := VectorGraphComputeResult{SchemaVersion: SchemaVersion, JobID: request.Header.JobID, Kind: ComputeKindVectorGraph, GraphVersion: request.Header.GraphVersion}
|
|
result.Primary = vectorgraph.Build(request.Entries, request.Header.Primary)
|
|
if request.Header.OrphanPass && len(request.Header.OrphanFocusIDs) > 0 {
|
|
focus := make(map[string]bool, len(request.Header.OrphanFocusIDs))
|
|
for _, id := range request.Header.OrphanFocusIDs {
|
|
focus[id] = true
|
|
}
|
|
for _, link := range result.Primary.Links {
|
|
delete(focus, link.Source)
|
|
delete(focus, link.Target)
|
|
}
|
|
result.Orphan = vectorgraph.BuildFocused(request.Entries, focus, request.Header.Orphan)
|
|
}
|
|
result.DurationMS = time.Since(started).Milliseconds()
|
|
return result
|
|
}
|
|
|
|
func SortedFocusIDs(values map[string]bool) []string {
|
|
out := make([]string, 0, len(values))
|
|
for id := range values {
|
|
out = append(out, id)
|
|
}
|
|
sort.Strings(out)
|
|
return out
|
|
}
|