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 }