Files
jbergner 440423c5b6
All checks were successful
release-tag / release-image (push) Successful in 2m43s
RC-3
2026-08-09 11:29:13 +02:00

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
}