1524991adc
ECOSYSTEM-SPEC §2.6 requires list_attention to distinguish "nothing needs attention" from "I cannot currently tell", and to say so when a source is failed or stale. Maven said the first one unconditionally: ListAttention decoded into []map[string]any, the word degraded appeared nowhere, and an empty list answered "ничего не требует внимания". A Praxis with every source dead read as calm. Two halves, because the spec's mechanism does not exist server-side yet. The deployed Praxis answers /api/v1/tools/attention with a bare array and no envelope, so praxisAttention now decodes either shape and believes a degraded array when one arrives. Until one does, an empty list triggers one read of /api/v1/sources, and anything that is not reporting health "ok" is named instead of the all-clear. Zero sources is the same answer: a Praxis that polls nothing knows nothing, which is the state of this box today. A sources read that fails is deliberately not a hedge. The attention call succeeded, and not being able to ask about health is not evidence of a fault. Both hedges also cover the entity-scoped digest, where a per-entity all-clear is the more convincing of the two. New keys attention_degraded and attention_no_sources, in acts_ru_v1.json and the floor. The fake Praxis serves one healthy source by default, so the existing attention tests still assert an all-clear on purpose rather than by omission.
589 lines
21 KiB
Go
589 lines
21 KiB
Go
package main
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/rand"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"log"
|
|
"net/http"
|
|
"net/url"
|
|
"time"
|
|
|
|
hexisclient "github.com/kami/hexis/pkg/client"
|
|
"github.com/kami/maven/internal/config"
|
|
)
|
|
|
|
// ecosystemCorrelationKey carries a per-call correlation ID through context
|
|
// so every ecosystem client (Nexus, Praxis, Hexis) tags its request with the
|
|
// same ID, letting a single Maven-initiated action be traced end to end.
|
|
type ecosystemCorrelationKey struct{}
|
|
|
|
func withCorrelationID(ctx context.Context, id string) context.Context {
|
|
return context.WithValue(ctx, ecosystemCorrelationKey{}, id)
|
|
}
|
|
|
|
func correlationIDFromCtx(ctx context.Context) string {
|
|
id, _ := ctx.Value(ecosystemCorrelationKey{}).(string)
|
|
return id
|
|
}
|
|
|
|
// ecosystemAPIVersion is the contract version Maven speaks to Nexus and
|
|
// Praxis. It is sent on every request so a service that has moved on can
|
|
// refuse or adapt explicitly instead of misreading an older payload.
|
|
const ecosystemAPIVersion = "v1"
|
|
|
|
// mavenRequester identifies the calling system on every ecosystem request, so
|
|
// a trace on the far side can attribute a call to Maven rather than to an
|
|
// anonymous HTTP client.
|
|
const mavenRequester = "maven"
|
|
|
|
// setEcosystemHeaders stamps the version, requester, auth and correlation
|
|
// headers common to every outgoing ecosystem request. token may be empty,
|
|
// which means the transport itself is trusted (loopback or unix socket).
|
|
//
|
|
// The correlation ID is read from the context and never minted here. Minting
|
|
// one per request sent the far side an ID that existed nowhere on this side,
|
|
// and gave a single multi-hop action as many unrelated IDs as it made calls.
|
|
// Callers that start an action assign the ID once (handleHexisAct,
|
|
// handlePraxisAct, resolveEntityReference) and every hop inherits it.
|
|
func setEcosystemHeaders(req *http.Request, ctx context.Context, versionHeader, token string) {
|
|
req.Header.Set("Content-Type", "application/json")
|
|
req.Header.Set(versionHeader, ecosystemAPIVersion)
|
|
req.Header.Set("Accept", "application/json")
|
|
req.Header.Set("X-Requested-By", mavenRequester)
|
|
if token != "" {
|
|
req.Header.Set("Authorization", "Bearer "+token)
|
|
}
|
|
if id := correlationIDFromCtx(ctx); id != "" {
|
|
req.Header.Set("X-Correlation-ID", id)
|
|
}
|
|
}
|
|
|
|
// ecosystemError is the typed failure every ecosystem client returns, so
|
|
// callers can tell a transport failure from a refusal from a contract
|
|
// mismatch without matching on message text. The distinction matters:
|
|
// "the service is down" and "the service rejected my version" degrade the
|
|
// same way to the user but not to whoever reads the trace.
|
|
type ecosystemError struct {
|
|
Service string // "nexus", "praxis", "hexis"
|
|
Op string // logical operation, e.g. "resolve"
|
|
Status int // HTTP status, 0 when the call never got an answer
|
|
Err error
|
|
}
|
|
|
|
func (e *ecosystemError) Error() string {
|
|
if e.Status != 0 {
|
|
return fmt.Sprintf("%s %s: http %d: %v", e.Service, e.Op, e.Status, e.Err)
|
|
}
|
|
return fmt.Sprintf("%s %s: %v", e.Service, e.Op, e.Err)
|
|
}
|
|
|
|
func (e *ecosystemError) Unwrap() error { return e.Err }
|
|
|
|
// Unauthorized reports a rejected or missing credential.
|
|
func (e *ecosystemError) Unauthorized() bool {
|
|
return e.Status == http.StatusUnauthorized || e.Status == http.StatusForbidden
|
|
}
|
|
|
|
// ContractMismatch reports that the far side refused the version Maven speaks.
|
|
func (e *ecosystemError) ContractMismatch() bool {
|
|
return e.Status == http.StatusNotAcceptable || e.Status == http.StatusUpgradeRequired
|
|
}
|
|
|
|
// Unreachable reports a call that never produced an HTTP answer at all
|
|
// (connection refused, timeout, cancelled).
|
|
func (e *ecosystemError) Unreachable() bool { return e.Status == 0 }
|
|
|
|
// httpError builds an ecosystemError from a response status.
|
|
func httpError(service, op string, status int) *ecosystemError {
|
|
return &ecosystemError{
|
|
Service: service, Op: op, Status: status,
|
|
Err: errors.New(http.StatusText(status)),
|
|
}
|
|
}
|
|
|
|
type nexusClient struct {
|
|
baseURL string
|
|
token string
|
|
httpClient *http.Client
|
|
}
|
|
|
|
func newNexusClient(url string) *nexusClient {
|
|
return &nexusClient{
|
|
baseURL: url,
|
|
httpClient: &http.Client{Timeout: 10 * time.Second},
|
|
}
|
|
}
|
|
|
|
// withToken sets the bearer token sent on every request. Returns the client so
|
|
// wiring reads as one expression.
|
|
func (c *nexusClient) withToken(token string) *nexusClient {
|
|
c.token = token
|
|
return c
|
|
}
|
|
|
|
type nexusEntity struct {
|
|
ID string `json:"id"`
|
|
Type string `json:"type"`
|
|
DisplayName string `json:"display_name"`
|
|
Key string `json:"key,omitempty"`
|
|
State string `json:"state,omitempty"`
|
|
}
|
|
|
|
type nexusCandidate struct {
|
|
EntityID string `json:"entity_id"`
|
|
DisplayName string `json:"display_name"`
|
|
Type string `json:"type"`
|
|
Score float64 `json:"score"`
|
|
Evidence string `json:"evidence"`
|
|
}
|
|
|
|
type nexusResolveResult struct {
|
|
Status string `json:"status"`
|
|
Entity *nexusEntity `json:"entity,omitempty"`
|
|
Score float64 `json:"score,omitempty"`
|
|
Candidates []nexusCandidate `json:"candidates,omitempty"`
|
|
|
|
// Flat fields per ECOSYSTEM-SPEC.md §1.5's documented resolve response
|
|
// shape. Nexus emits both this and the nested Entity above; normalize
|
|
// into Entity in UnmarshalJSON so callers only ever look at one place.
|
|
EntityID string `json:"entity_id,omitempty"`
|
|
EntityType string `json:"entity_type,omitempty"`
|
|
DisplayName string `json:"display_name,omitempty"`
|
|
}
|
|
|
|
func (r *nexusResolveResult) UnmarshalJSON(data []byte) error {
|
|
type alias nexusResolveResult
|
|
var a alias
|
|
if err := json.Unmarshal(data, &a); err != nil {
|
|
return err
|
|
}
|
|
*r = nexusResolveResult(a)
|
|
if r.Entity == nil && r.EntityID != "" {
|
|
r.Entity = &nexusEntity{ID: r.EntityID, Type: r.EntityType, DisplayName: r.DisplayName}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (c *nexusClient) Resolve(ctx context.Context, query string, types []string) (*nexusResolveResult, error) {
|
|
body := map[string]any{"query": query}
|
|
if len(types) > 0 {
|
|
body["types"] = types
|
|
}
|
|
|
|
data, _ := json.Marshal(body)
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/api/v1/resolve", bytes.NewReader(data))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create request: %w", err)
|
|
}
|
|
setEcosystemHeaders(req, ctx, "X-Nexus-Version", c.token)
|
|
|
|
resp, err := c.httpClient.Do(req)
|
|
if err != nil {
|
|
return nil, &ecosystemError{Service: "nexus", Op: "resolve", Err: err}
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
bodyBytes, _ := io.ReadAll(resp.Body)
|
|
if resp.StatusCode != 200 {
|
|
return nil, httpError("nexus", "resolve", resp.StatusCode)
|
|
}
|
|
|
|
var result nexusResolveResult
|
|
if err := json.Unmarshal(bodyBytes, &result); err != nil {
|
|
return nil, &ecosystemError{Service: "nexus", Op: "resolve", Status: resp.StatusCode, Err: err}
|
|
}
|
|
return &result, nil
|
|
}
|
|
|
|
func (c *nexusClient) Health(ctx context.Context) error {
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.baseURL+"/health", nil)
|
|
if err != nil {
|
|
return &ecosystemError{Service: "nexus", Op: "health", Err: err}
|
|
}
|
|
setEcosystemHeaders(req, ctx, "X-Nexus-Version", c.token)
|
|
resp, err := c.httpClient.Do(req)
|
|
if err != nil {
|
|
return &ecosystemError{Service: "nexus", Op: "health", Err: err}
|
|
}
|
|
resp.Body.Close()
|
|
if resp.StatusCode != 200 {
|
|
return httpError("nexus", "health", resp.StatusCode)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// praxisClient talks to the Praxis HTTP tools API. Maven must not open Praxis's
|
|
// SQLite store directly (ecosystem invariant: no component reads another's DB),
|
|
// so attention/changes/lifecycle all go over this HTTP contract against praxisd.
|
|
type praxisClient struct {
|
|
baseURL string
|
|
token string
|
|
httpClient *http.Client
|
|
}
|
|
|
|
func newPraxisClient(url string) *praxisClient {
|
|
return &praxisClient{
|
|
baseURL: url,
|
|
httpClient: &http.Client{Timeout: 10 * time.Second},
|
|
}
|
|
}
|
|
|
|
func (c *praxisClient) withToken(token string) *praxisClient {
|
|
c.token = token
|
|
return c
|
|
}
|
|
|
|
// getJSON performs a GET and decodes the JSON body into out. op is the logical
|
|
// operation name for errors and traces: the path carries the query string, and
|
|
// after entity scoping that means an entity id in every log line built from the
|
|
// error, next to a trace that redacts far less than that.
|
|
func (c *praxisClient) getJSON(ctx context.Context, op, path string, out any) error {
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.baseURL+path, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
setEcosystemHeaders(req, ctx, "X-Praxis-Version", c.token)
|
|
resp, err := c.httpClient.Do(req)
|
|
if err != nil {
|
|
return &ecosystemError{Service: "praxis", Op: op, Err: err}
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode != 200 {
|
|
return httpError("praxis", op, resp.StatusCode)
|
|
}
|
|
if err := json.NewDecoder(resp.Body).Decode(out); err != nil {
|
|
return &ecosystemError{Service: "praxis", Op: op, Status: resp.StatusCode, Err: err}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// praxisAttention — an attention response in either of the two shapes Praxis
|
|
// may send (Vikunja #540).
|
|
//
|
|
// ECOSYSTEM-SPEC §2.6 says the response carries `degraded: [source_ids]` when a
|
|
// source is failed or stale, and that Maven is required to say so rather than
|
|
// report all-clear. The deployed Praxis answers with a bare JSON array and no
|
|
// envelope at all, so both are decoded here: an array is the items, an object is
|
|
// the spec envelope. This lands the Maven half without waiting on the server,
|
|
// and the sources read below is what makes the hedge work meanwhile.
|
|
type praxisAttention struct {
|
|
Items []map[string]any
|
|
Degraded []string
|
|
}
|
|
|
|
func (a *praxisAttention) UnmarshalJSON(data []byte) error {
|
|
trimmed := bytes.TrimSpace(data)
|
|
if len(trimmed) > 0 && trimmed[0] == '[' {
|
|
return json.Unmarshal(trimmed, &a.Items)
|
|
}
|
|
var env struct {
|
|
Items []map[string]any `json:"items"`
|
|
Degraded []string `json:"degraded"`
|
|
}
|
|
if err := json.Unmarshal(trimmed, &env); err != nil {
|
|
return err
|
|
}
|
|
a.Items, a.Degraded = env.Items, env.Degraded
|
|
return nil
|
|
}
|
|
|
|
func (c *praxisClient) ListAttention(ctx context.Context, limit int) (praxisAttention, error) {
|
|
var out praxisAttention
|
|
err := c.getJSON(ctx, "attention", fmt.Sprintf("/api/v1/tools/attention?limit=%d", limit), &out)
|
|
return out, err
|
|
}
|
|
|
|
// praxisSource — one polled source, as much of it as the hedge needs. The tools
|
|
// API does not expose sources, so this decodes the plain `/api/v1/sources` rows.
|
|
type praxisSource struct {
|
|
ID string `json:"id"`
|
|
SourceID string `json:"source_id"`
|
|
Health string `json:"health"`
|
|
}
|
|
|
|
func (s praxisSource) name() string {
|
|
if s.SourceID != "" {
|
|
return s.SourceID
|
|
}
|
|
return s.ID
|
|
}
|
|
|
|
// UnhealthySources reports which sources cannot be trusted to have reported,
|
|
// and how many sources Praxis has at all (Vikunja #540).
|
|
//
|
|
// Only read when the attention list came back empty, which is the one turn where
|
|
// an all-clear is at stake. A source whose health field is absent counts as
|
|
// healthy: a Praxis that never reports health would otherwise make every quiet
|
|
// turn a hedge, and an unreported field is not evidence of a fault. Everything it
|
|
// does report other than "ok" — failed, stale, degraded, unknown — counts as
|
|
// cannot-tell, because none of them mean the source has spoken.
|
|
func (c *praxisClient) UnhealthySources(ctx context.Context) (bad []string, total int, err error) {
|
|
var out []praxisSource
|
|
if err := c.getJSON(ctx, "sources", "/api/v1/sources", &out); err != nil {
|
|
return nil, 0, err
|
|
}
|
|
for _, s := range out {
|
|
if s.Health != "" && s.Health != "ok" {
|
|
bad = append(bad, s.name())
|
|
}
|
|
}
|
|
return bad, len(out), nil
|
|
}
|
|
|
|
// ListAttentionForEntity is ListAttention scoped to a single canonical Nexus
|
|
// entity, so callers already holding a resolved entity_id (e.g. after
|
|
// resolveEntityReference) can ask "what needs attention for this entity"
|
|
// instead of filtering the unscoped list client-side.
|
|
func (c *praxisClient) ListAttentionForEntity(ctx context.Context, entityID string, limit int) (praxisAttention, error) {
|
|
var out praxisAttention
|
|
err := c.getJSON(ctx, "attention_for_entity",
|
|
fmt.Sprintf("/api/v1/tools/attention?limit=%d&entity_id=%s", limit, url.QueryEscape(entityID)), &out)
|
|
return out, err
|
|
}
|
|
|
|
func (c *praxisClient) ListChanges(ctx context.Context, limit int) ([]map[string]any, error) {
|
|
var out []map[string]any
|
|
err := c.getJSON(ctx, "changes", fmt.Sprintf("/api/v1/tools/changes?limit=%d", limit), &out)
|
|
return out, err
|
|
}
|
|
|
|
// praxisItem is the typed shape of a Praxis item, decoded from the tools API's
|
|
// itemToMap output (pkg/tools/api.go in the praxis repo). Kept as a distinct
|
|
// type from the raw attention/changes maps above so lifecycle callers get
|
|
// compile-time field checks instead of map[string]any type assertions.
|
|
type praxisItem struct {
|
|
ID string `json:"id"`
|
|
SourceID string `json:"source_id"`
|
|
ExternalID string `json:"external_id"`
|
|
Kind string `json:"kind"`
|
|
Title string `json:"title"`
|
|
Body string `json:"body"`
|
|
State string `json:"state"`
|
|
Importance int `json:"importance"`
|
|
FirstSeenAt string `json:"first_seen_at"`
|
|
LastSeenAt string `json:"last_seen_at"`
|
|
SurfacedAt string `json:"surfaced_at"`
|
|
AckedAt string `json:"acknowledged_at"`
|
|
ResolvedAt string `json:"resolved_at"`
|
|
}
|
|
|
|
// postItemAction posts {"item_id": id} to a Praxis tools lifecycle endpoint
|
|
// and decodes the resulting item. Shared by Surface/Acknowledge/Resolve/Ignore.
|
|
func (c *praxisClient) postItemAction(ctx context.Context, op, path, itemID string) (*praxisItem, error) {
|
|
return c.postJSON(ctx, op, path, map[string]any{"item_id": itemID})
|
|
}
|
|
|
|
// postJSON posts a body to a Praxis lifecycle endpoint and decodes the item.
|
|
// Every failure is a *ecosystemError, including the transport and decode ones:
|
|
// these are the paths that mutate remote state, and the question worth
|
|
// answering afterwards is whether the call never left or was refused.
|
|
func (c *praxisClient) postJSON(ctx context.Context, op, path string, payload map[string]any) (*praxisItem, error) {
|
|
body, _ := json.Marshal(payload)
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+path, bytes.NewReader(body))
|
|
if err != nil {
|
|
return nil, &ecosystemError{Service: "praxis", Op: op, Err: err}
|
|
}
|
|
setEcosystemHeaders(req, ctx, "X-Praxis-Version", c.token)
|
|
resp, err := c.httpClient.Do(req)
|
|
if err != nil {
|
|
return nil, &ecosystemError{Service: "praxis", Op: op, Err: err}
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode != 200 {
|
|
return nil, httpError("praxis", op, resp.StatusCode)
|
|
}
|
|
var out praxisItem
|
|
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
|
|
return nil, &ecosystemError{Service: "praxis", Op: op, Status: resp.StatusCode, Err: err}
|
|
}
|
|
return &out, nil
|
|
}
|
|
|
|
// Surface marks an item read/spoken without acknowledging it (surfaced != acknowledged,
|
|
// ECOSYSTEM-SPEC.md §2.3). Callers that read attention aloud must call this, never
|
|
// Acknowledge, so "I mentioned it" stays distinguishable from "you told me you saw it".
|
|
func (c *praxisClient) Surface(ctx context.Context, itemID string) (*praxisItem, error) {
|
|
return c.postItemAction(ctx, "surface", "/api/v1/tools/surface", itemID)
|
|
}
|
|
|
|
func (c *praxisClient) Acknowledge(ctx context.Context, itemID string) (*praxisItem, error) {
|
|
return c.postItemAction(ctx, "acknowledge", "/api/v1/tools/acknowledge", itemID)
|
|
}
|
|
|
|
func (c *praxisClient) Resolve(ctx context.Context, itemID string) (*praxisItem, error) {
|
|
return c.postItemAction(ctx, "resolve", "/api/v1/tools/resolve", itemID)
|
|
}
|
|
|
|
func (c *praxisClient) Ignore(ctx context.Context, itemID string) (*praxisItem, error) {
|
|
return c.postItemAction(ctx, "ignore", "/api/v1/tools/ignore", itemID)
|
|
}
|
|
|
|
func (c *praxisClient) Pin(ctx context.Context, itemID string, pinned bool) (*praxisItem, error) {
|
|
return c.postJSON(ctx, "pin", "/api/v1/tools/pin", map[string]any{"item_id": itemID, "pinned": pinned})
|
|
}
|
|
|
|
func (c *praxisClient) GetItem(ctx context.Context, itemID string) (*praxisItem, error) {
|
|
var out praxisItem
|
|
err := c.getJSON(ctx, "get_item", "/api/v1/tools/items/"+itemID, &out)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return &out, nil
|
|
}
|
|
|
|
func (c *praxisClient) Search(ctx context.Context, query string, limit int) ([]praxisItem, error) {
|
|
var out []praxisItem
|
|
err := c.getJSON(ctx, "search", fmt.Sprintf("/api/v1/tools/search?q=%s&limit=%d", url.QueryEscape(query), limit), &out)
|
|
return out, err
|
|
}
|
|
|
|
// ecosystemWiring holds the ecosystem service clients.
|
|
type ecosystemWiring struct {
|
|
nexus *nexusClient
|
|
hexis *hexisclient.Client
|
|
praxis *praxisClient
|
|
}
|
|
|
|
func wireEcosystem(cfg *config.Config) *ecosystemWiring {
|
|
w := &ecosystemWiring{}
|
|
|
|
// Nexus identity service
|
|
if cfg.Nexus != nil && cfg.Nexus.URL != "" {
|
|
w.nexus = newNexusClient(cfg.Nexus.URL).withToken(cfg.Nexus.Token)
|
|
log.Printf("ecosystem: nexus at %s", cfg.Nexus.URL)
|
|
} else {
|
|
log.Printf("ecosystem: nexus not configured")
|
|
}
|
|
|
|
// Hexis capability service
|
|
if cfg.Hexis != nil && cfg.Hexis.URL != "" {
|
|
w.hexis = hexisclient.New(cfg.Hexis.URL).WithToken(cfg.Hexis.Token)
|
|
log.Printf("ecosystem: hexis at %s", cfg.Hexis.URL)
|
|
} else {
|
|
log.Printf("ecosystem: hexis not configured")
|
|
}
|
|
|
|
// Praxis attention service (HTTP tools API — never the DB directly)
|
|
if cfg.Praxis != nil && cfg.Praxis.URL != "" {
|
|
w.praxis = newPraxisClient(cfg.Praxis.URL).withToken(cfg.Praxis.Token)
|
|
log.Printf("ecosystem: praxis at %s", cfg.Praxis.URL)
|
|
} else {
|
|
log.Printf("ecosystem: praxis not configured")
|
|
}
|
|
|
|
return w
|
|
}
|
|
|
|
// resolveEntityReference extracts and resolves an entity name from utterance text.
|
|
// Returns the canonical entity ID on a confident resolve. On an ambiguous match it
|
|
// returns candidate display names so the caller can ask for clarification rather
|
|
// than silently guessing (ecosystem invariant: ambiguity blocks mutation).
|
|
func (w *ecosystemWiring) resolveEntityReference(ctx context.Context, text string, entityTypes []string) (entityID string, displayName string, ambiguous []string, err error) {
|
|
if w == nil || w.nexus == nil {
|
|
return "", "", nil, nil
|
|
}
|
|
if correlationIDFromCtx(ctx) == "" {
|
|
ctx = withCorrelationID(ctx, newCorrelationID())
|
|
}
|
|
result, err := w.nexus.Resolve(ctx, text, entityTypes)
|
|
if err != nil {
|
|
log.Printf("ecosystem: nexus resolve error: %v", err)
|
|
return "", "", nil, err
|
|
}
|
|
if result.Status == "resolved" {
|
|
// "resolved" with nothing to resolve to is a contract violation, not a
|
|
// miss. Treating it as "no such entity" let the caller fall straight
|
|
// through to the local executor with his verb intact, which is a
|
|
// dependency failure reaching execution.
|
|
if result.Entity == nil || result.Entity.ID == "" {
|
|
err := &ecosystemError{
|
|
Service: "nexus", Op: "resolve", Status: 200,
|
|
Err: errors.New("resolved status with no entity"),
|
|
}
|
|
log.Printf("ecosystem: %v", err)
|
|
return "", "", nil, err
|
|
}
|
|
return result.Entity.ID, result.Entity.DisplayName, nil, nil
|
|
}
|
|
if result.Status == "ambiguous" {
|
|
names := make([]string, 0, len(result.Candidates))
|
|
for _, c := range result.Candidates {
|
|
names = append(names, c.DisplayName)
|
|
}
|
|
log.Printf("ecosystem: ambiguous entity '%s' — %d candidates", text, len(names))
|
|
return "", "", names, nil
|
|
}
|
|
return "", "", nil, nil
|
|
}
|
|
|
|
// discoverCapabilities returns Hexis capabilities applicable to an entity.
|
|
// A non-nil error means Hexis could not be reached or refused the request —
|
|
// distinct from a nil error with zero capabilities, which means Hexis is
|
|
// healthy and genuinely has nothing registered for this entity. Callers must
|
|
// not conflate the two: a dependency failure must not silently read as "no
|
|
// capabilities" and fall through to unrelated local execution.
|
|
//
|
|
// The correlation header is stamped in the client's do(), so discovery and
|
|
// execution can be joined on the Hexis side as long as both hops carry the
|
|
// same ID through ctx. (This used to say the header went out on Execute only;
|
|
// that was never true of the vendored code and is not true after the 2026-08-01
|
|
// re-vendor.)
|
|
func (w *ecosystemWiring) discoverCapabilities(ctx context.Context, entityID string) ([]hexisclient.Capability, error) {
|
|
if w == nil || w.hexis == nil || entityID == "" {
|
|
return nil, nil
|
|
}
|
|
caps, err := w.hexis.Capabilities(ctx, entityID)
|
|
if err != nil {
|
|
log.Printf("ecosystem: hexis capabilities error: %v", err)
|
|
return nil, err
|
|
}
|
|
return caps, nil
|
|
}
|
|
|
|
// executeCapability runs a Hexis capability, tagging the request with a
|
|
// correlation ID so the call is traceable across services. Returns the
|
|
// correlation ID alongside the outcome.
|
|
func (w *ecosystemWiring) executeCapability(ctx context.Context, capabilityID, targetEntityID string, args map[string]any) (correlationID string, err error) {
|
|
if w == nil || w.hexis == nil {
|
|
return "", fmt.Errorf("hexis not configured")
|
|
}
|
|
correlationID = newCorrelationID()
|
|
ctx = withCorrelationID(ctx, correlationID)
|
|
ctx = hexisclient.WithCorrelationID(ctx, correlationID)
|
|
req := hexisclient.ExecuteRequest{
|
|
CapabilityID: capabilityID,
|
|
TargetEntityID: targetEntityID,
|
|
Arguments: args,
|
|
RequestedBy: map[string]string{"system": "maven", "actor": "user"},
|
|
Origin: map[string]string{"source": "voice"},
|
|
CorrelationID: correlationID,
|
|
}
|
|
|
|
exec, err := w.hexis.Execute(ctx, req)
|
|
if err != nil {
|
|
return correlationID, fmt.Errorf("execute: %w", err)
|
|
}
|
|
if exec.Status == "succeeded" {
|
|
return correlationID, nil
|
|
}
|
|
if exec.Error != "" {
|
|
return correlationID, fmt.Errorf("execution failed: %s", exec.Error)
|
|
}
|
|
return correlationID, fmt.Errorf("execution status: %s", exec.Status)
|
|
}
|
|
|
|
// newCorrelationID returns a short unique ID for cross-service call tracing.
|
|
func newCorrelationID() string {
|
|
var b [8]byte
|
|
if _, err := rand.Read(b[:]); err != nil {
|
|
return fmt.Sprintf("cor-%d", time.Now().UnixNano())
|
|
}
|
|
return "cor-" + hex.EncodeToString(b[:])
|
|
}
|