6a9d8a4dd5
Two caller-side halves of the same review.
«экосистема недоступна» named nothing. Nexus, Praxis and Hexis fail
independently, and every one of the six call sites already knew which one it was
talking to — it writes that name into the trace on the line above. So eco_down
and eco_denied now take {name}, and he hears which service refused him.
The list entries are single-variant and placeholder-only, so an empty list has
no shorter wording to fall back on: attention_list would render as its own label
and a colon. Both Praxis readers checked the response length and neither checked
what survived formatting, so an item with no title counted toward a list it
could not appear in. They skip the untitled item and fall to the _none entry
when nothing is left.
The ecosystem tests asserted the substring "выполнена", which was a literal out
of the act file that review has now reworded. Seventeen sites go through actRan,
which asks the file.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01XGTGCWX33aX8SMBSRz9VmS
670 lines
25 KiB
Go
670 lines
25 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"strings"
|
|
"time"
|
|
|
|
hexisclient "github.com/kami/hexis/pkg/client"
|
|
"github.com/kami/maven/internal/phraser"
|
|
"github.com/kami/maven/internal/router"
|
|
"github.com/kami/maven/internal/store"
|
|
)
|
|
|
|
// The three services, spelled the way she says them out loud. A service that is
|
|
// down or refusing has to be named: they degrade independently, so "не
|
|
// отвечает" on its own tells him nothing he can act on, and each call site
|
|
// already knows which one it was talking to — it records the same name in the
|
|
// trace (Vikunja #521).
|
|
const (
|
|
serviceNexus = "Nexus"
|
|
serviceHexis = "Hexis"
|
|
)
|
|
|
|
// serviceVars — the one-key map the eco_down and eco_denied lines take.
|
|
func serviceVars(name string) map[string]string { return map[string]string{"name": name} }
|
|
|
|
// praxisCapability is one arm of the Praxis act dispatch. This is an interface
|
|
// rather than a map[string]func because each arm carries its own state: the
|
|
// verb aliases it answers to, the trace name it records, and its own reply
|
|
// formatting. The dispatch grows an arm per Praxis capability, so a new one is
|
|
// added to praxisCapabilities below and nothing else changes.
|
|
type praxisCapability interface {
|
|
// aliases are the verbs (router fn slots, EN and RU) this capability answers to.
|
|
aliases() []string
|
|
// handle runs the capability and returns the user-facing reply.
|
|
handle(ctx context.Context, h *reactiveHandler, px *praxisClient, dec router.Decision) string
|
|
}
|
|
|
|
// praxisCapabilities is the registry handlePraxisAct consults, in order.
|
|
var praxisCapabilities = []praxisCapability{
|
|
listAttentionCapability{},
|
|
praxisItemAction{
|
|
verbs: []string{"acknowledge_item", "принято", "понял", "поняла"},
|
|
ask: "какой пункт отметить принятым?",
|
|
op: "acknowledge",
|
|
failure: "не получилось отметить принятым.",
|
|
success: "принято.",
|
|
call: func(ctx context.Context, px *praxisClient, id string) error {
|
|
_, err := px.Acknowledge(ctx, id)
|
|
return err
|
|
},
|
|
},
|
|
praxisItemAction{
|
|
verbs: []string{"resolve_item", "сделано", "готово", "решено"},
|
|
ask: "какой пункт отметить сделанным?",
|
|
op: "resolve",
|
|
failure: "не получилось отметить сделанным.",
|
|
success: "отмечено как сделано.",
|
|
call: func(ctx context.Context, px *praxisClient, id string) error {
|
|
_, err := px.Resolve(ctx, id)
|
|
return err
|
|
},
|
|
},
|
|
praxisItemAction{
|
|
verbs: []string{"ignore_item", "игнорировать", "неважно"},
|
|
ask: "какой пункт игнорировать?",
|
|
op: "ignore",
|
|
failure: "не получилось проигнорировать.",
|
|
success: "проигнорировано.",
|
|
call: func(ctx context.Context, px *praxisClient, id string) error {
|
|
_, err := px.Ignore(ctx, id)
|
|
return err
|
|
},
|
|
},
|
|
praxisItemAction{
|
|
verbs: []string{"pin_item", "закрепить"},
|
|
ask: "какой пункт закрепить?",
|
|
op: "pin",
|
|
failure: "не получилось закрепить.",
|
|
success: "закреплено.",
|
|
call: func(ctx context.Context, px *praxisClient, id string) error {
|
|
_, err := px.Pin(ctx, id, true)
|
|
return err
|
|
},
|
|
},
|
|
listChangesCapability{},
|
|
entityAttentionCapability{},
|
|
}
|
|
|
|
// handlePraxisAct — dispatches ecosystem tool acts through the Praxis tools API.
|
|
// Returns "" when the act is not a Praxis verb (the caller falls through to the
|
|
// system command executor). Returns a reply string otherwise.
|
|
func (h *reactiveHandler) handlePraxisAct(ctx context.Context, dec router.Decision) string {
|
|
if h.ecosystem == nil || h.ecosystem.praxis == nil {
|
|
return ""
|
|
}
|
|
// Every hop of this action shares one correlation ID, assigned here, so a
|
|
// digest that calls attention once and surface N times reads as one turn
|
|
// on the Praxis side instead of N+1 unrelated request ids.
|
|
if correlationIDFromCtx(ctx) == "" {
|
|
ctx = withCorrelationID(ctx, newCorrelationID())
|
|
}
|
|
px := h.ecosystem.praxis
|
|
for _, capability := range praxisCapabilities {
|
|
for _, alias := range capability.aliases() {
|
|
if alias == dec.Slots.Fn {
|
|
return capability.handle(ctx, h, px, dec)
|
|
}
|
|
}
|
|
}
|
|
// Not a Praxis verb — let the caller fall through.
|
|
return ""
|
|
}
|
|
|
|
// praxisItemAction is the shared shape of the item-lifecycle capabilities: take
|
|
// an item id from the value slot, call one Praxis endpoint, trace the result.
|
|
type praxisItemAction struct {
|
|
verbs []string
|
|
ask string // reply when no item id was given
|
|
op string // trace + log name of the operation
|
|
failure string // reply when the Praxis call errors
|
|
success string
|
|
call func(ctx context.Context, px *praxisClient, id string) error
|
|
}
|
|
|
|
func (a praxisItemAction) aliases() []string { return a.verbs }
|
|
|
|
func (a praxisItemAction) handle(ctx context.Context, h *reactiveHandler, px *praxisClient, dec router.Decision) string {
|
|
id := dec.Slots.Value
|
|
if id == "" {
|
|
return a.ask
|
|
}
|
|
started := h.now()
|
|
if err := a.call(ctx, px, id); err != nil {
|
|
log.Printf("ecosystem: praxis %s %s: %v", a.op, id, err)
|
|
h.recordEcosystemTrace(ctx, "praxis", a.op, traceStatusForError(err), started,
|
|
mergeFields(traceErrorFields(err), map[string]any{"item_id": id}))
|
|
return a.failure
|
|
}
|
|
h.recordPraxisTrace(ctx, a.op, started, map[string]any{"item_id": id})
|
|
return a.success
|
|
}
|
|
|
|
// listAttentionCapability reads the attention digest and surfaces every item it speaks.
|
|
type listAttentionCapability struct{}
|
|
|
|
func (listAttentionCapability) aliases() []string {
|
|
return []string{"list_attention", "attention", "внимание", "что требует внимания", "что нового"}
|
|
}
|
|
|
|
func (listAttentionCapability) handle(ctx context.Context, h *reactiveHandler, px *praxisClient, _ router.Decision) string {
|
|
started := h.now()
|
|
items, err := px.ListAttention(ctx, 20)
|
|
if err != nil {
|
|
log.Printf("ecosystem: praxis attention: %v", err)
|
|
h.recordEcosystemTrace(ctx, "praxis", "list_attention", traceStatusForError(err),
|
|
started, traceErrorFields(err))
|
|
return phraser.A(phraser.AttentionFail, nil)
|
|
}
|
|
if len(items) == 0 {
|
|
return phraser.A(phraser.AttentionNone, nil)
|
|
}
|
|
h.recordPraxisTrace(ctx, "list_attention", started, map[string]any{"count": len(items)})
|
|
var parts []string
|
|
for _, item := range items {
|
|
title, _ := item["title"].(string)
|
|
// importance arrives as JSON number ⇒ float64 over the HTTP contract.
|
|
importance, _ := item["importance"].(float64)
|
|
rule, _ := item["rule"].(string)
|
|
s := title
|
|
if s == "" {
|
|
// An item Praxis returned without a title is not an item she can
|
|
// read out. Counting it would put an empty slot in the list.
|
|
continue
|
|
}
|
|
if importance > 0 {
|
|
s += fmt.Sprintf(" (важность %d", int(importance))
|
|
if rule != "" {
|
|
s += ": " + rule
|
|
}
|
|
s += ")"
|
|
}
|
|
parts = append(parts, s)
|
|
|
|
// Speaking an item surfaces it, it does not acknowledge it
|
|
// (ECOSYSTEM-SPEC.md §2.3: surfaced != acknowledged). Best-effort:
|
|
// a failed surface call must not block delivering the digest.
|
|
if id, ok := item["id"].(string); ok && id != "" {
|
|
if _, err := px.Surface(ctx, id); err != nil {
|
|
log.Printf("ecosystem: praxis surface %s: %v", id, err)
|
|
}
|
|
}
|
|
}
|
|
if len(parts) == 0 {
|
|
// Praxis returned items and not one of them could be said. "ничего не
|
|
// требует внимания" is the honest answer; the list line would render as
|
|
// its own label and a colon (Vikunja #521).
|
|
return phraser.A(phraser.AttentionNone, nil)
|
|
}
|
|
return phraser.A(phraser.AttentionList, map[string]string{"items": strings.Join(parts, "; ")})
|
|
}
|
|
|
|
// listChangesCapability reads the recent-changes feed.
|
|
type listChangesCapability struct{}
|
|
|
|
func (listChangesCapability) aliases() []string {
|
|
return []string{"list_changes", "changes", "изменения", "что изменилось"}
|
|
}
|
|
|
|
func (listChangesCapability) handle(ctx context.Context, h *reactiveHandler, px *praxisClient, _ router.Decision) string {
|
|
started := h.now()
|
|
changes, err := px.ListChanges(ctx, 20)
|
|
if err != nil {
|
|
log.Printf("ecosystem: praxis changes: %v", err)
|
|
h.recordEcosystemTrace(ctx, "praxis", "list_changes", traceStatusForError(err),
|
|
started, traceErrorFields(err))
|
|
return phraser.A(phraser.ChangesFail, nil)
|
|
}
|
|
if len(changes) == 0 {
|
|
return phraser.A(phraser.ChangesNone, nil)
|
|
}
|
|
h.recordPraxisTrace(ctx, "list_changes", started, map[string]any{"count": len(changes)})
|
|
var parts []string
|
|
for _, c := range changes {
|
|
title, _ := c["title"].(string)
|
|
if title == "" {
|
|
continue
|
|
}
|
|
typ, _ := c["change_type"].(string)
|
|
if typ == "" {
|
|
parts = append(parts, title)
|
|
continue
|
|
}
|
|
parts = append(parts, fmt.Sprintf("%s (%s)", title, typ))
|
|
}
|
|
if len(parts) == 0 {
|
|
return phraser.A(phraser.ChangesNone, nil)
|
|
}
|
|
return phraser.A(phraser.ChangesList, map[string]string{"items": strings.Join(parts, "; ")})
|
|
}
|
|
|
|
// entityAttentionCapability answers "what's going on with X" by resolving X to
|
|
// a canonical Nexus entity and asking Praxis for that entity's attention items
|
|
// (Vikunja #272). Unlike listAttentionCapability it is scoped: the entity_id
|
|
// travels to Praxis as a query parameter instead of Maven filtering an unscoped
|
|
// list client-side, which is what makes the ref canonical end to end.
|
|
//
|
|
// It also folds in what Maven herself knows about the same entity — facts the
|
|
// enrichment worker has already resolved to that entity_id — so one question
|
|
// gets one answer across both stores.
|
|
type entityAttentionCapability struct{}
|
|
|
|
// aliases are matched against Slots.Fn, which carries a function slot from the
|
|
// act grammar and never free Russian, so only grammar names belong here.
|
|
func (entityAttentionCapability) aliases() []string {
|
|
return []string{"entity_attention", "entity_status"}
|
|
}
|
|
|
|
func (entityAttentionCapability) handle(ctx context.Context, h *reactiveHandler, px *praxisClient, dec router.Decision) string {
|
|
subject := dec.Slots.Value
|
|
if subject == "" {
|
|
subject = dec.Slots.Text
|
|
}
|
|
if subject == "" {
|
|
return phraser.A(phraser.EcoAboutWhat, nil)
|
|
}
|
|
if h.ecosystem == nil || h.ecosystem.nexus == nil {
|
|
// Without Nexus there is no canonical ref to scope by. Say so rather
|
|
// than quietly answering about something else.
|
|
return phraser.A(phraser.EcoNoNexus, nil)
|
|
}
|
|
|
|
started := h.now()
|
|
entityID, displayName, ambiguous, err := h.ecosystem.resolveEntityReference(ctx, subject, nil)
|
|
if err != nil {
|
|
// The subject is his words, so the log gets the same redaction the
|
|
// trace gets. A trace that stores a rune count next to a log line
|
|
// storing the runes is not redacted at all.
|
|
log.Printf("ecosystem: entity attention resolve %s: %v", redactSubject(subject), err)
|
|
h.recordEcosystemTrace(ctx, "nexus", "resolve", traceStatusForError(err), started,
|
|
mergeFields(traceErrorFields(err), map[string]any{"subject": redactSubject(subject)}))
|
|
if unauthorizedEcosystemError(err) {
|
|
return phraser.A(phraser.EcoDenied, serviceVars(serviceNexus))
|
|
}
|
|
return phraser.A(phraser.EcoDown, serviceVars(serviceNexus))
|
|
}
|
|
if len(ambiguous) > 0 {
|
|
return phraser.A(phraser.EcoAmbiguous, map[string]string{"items": strings.Join(ambiguous, ", ")})
|
|
}
|
|
if entityID == "" {
|
|
return phraser.A(phraser.EcoUnknownEntity, nil)
|
|
}
|
|
if displayName == "" {
|
|
displayName = subject
|
|
}
|
|
|
|
queried := h.now()
|
|
items, err := px.ListAttentionForEntity(ctx, entityID, 20)
|
|
if err != nil {
|
|
log.Printf("ecosystem: praxis attention for %s: %v", entityID, err)
|
|
h.recordEcosystemTrace(ctx, "praxis", "entity_attention", traceStatusForError(err),
|
|
queried, mergeFields(traceErrorFields(err), map[string]any{"entity_id": entityID}))
|
|
return phraser.A(phraser.AttentionFailEntity, map[string]string{"name": displayName})
|
|
}
|
|
items, scoped := scopedToEntity(items, entityID)
|
|
if !scoped {
|
|
// A Praxis old enough to ignore an unknown query parameter answers the
|
|
// scoped question with the unscoped list. Reading that back as "по
|
|
// «X»: ..." is the exact fabrication the entity ref exists to prevent,
|
|
// so refuse the answer instead of relabelling someone else's items.
|
|
log.Printf("ecosystem: praxis returned unscoped items for %s, refusing to answer", entityID)
|
|
h.recordEcosystemTrace(ctx, "praxis", "entity_attention", traceFailed, queried,
|
|
map[string]any{"entity_id": entityID, "class": "unscoped_response"})
|
|
return phraser.A(phraser.AttentionFailEntity, map[string]string{"name": displayName})
|
|
}
|
|
h.recordPraxisTrace(ctx, "entity_attention", queried, map[string]any{
|
|
"entity_id": entityID, "count": len(items),
|
|
})
|
|
|
|
var parts []string
|
|
for _, item := range items {
|
|
title, _ := item["title"].(string)
|
|
if title == "" {
|
|
continue
|
|
}
|
|
parts = append(parts, title)
|
|
// Same surfaced != acknowledged rule as the unscoped digest.
|
|
if id, ok := item["id"].(string); ok && id != "" {
|
|
if _, err := px.Surface(ctx, id); err != nil {
|
|
log.Printf("ecosystem: praxis surface %s: %v", id, err)
|
|
}
|
|
}
|
|
}
|
|
if known := h.localFactsForEntity(ctx, entityID); known != "" {
|
|
parts = append(parts, known)
|
|
}
|
|
if len(parts) == 0 {
|
|
return phraser.A(phraser.AttentionNoneEntity, map[string]string{"name": displayName})
|
|
}
|
|
return phraser.A(phraser.AttentionListEntity, map[string]string{"name": displayName, "items": strings.Join(parts, "; ")})
|
|
}
|
|
|
|
// scopedToEntity drops items that carry an entity_id other than the one asked
|
|
// about, and reports whether the response can be trusted as scoped at all. An
|
|
// item without an entity_id is kept only when at least one sibling carries the
|
|
// matching id: a whole page with no entity_id is a Praxis that ignored the
|
|
// scope, not a page of untagged items.
|
|
func scopedToEntity(items []map[string]any, entityID string) ([]map[string]any, bool) {
|
|
if len(items) == 0 {
|
|
return items, true
|
|
}
|
|
var kept []map[string]any
|
|
var sawMatch, sawMismatch bool
|
|
for _, item := range items {
|
|
id, _ := item["entity_id"].(string)
|
|
switch {
|
|
case id == entityID:
|
|
sawMatch = true
|
|
kept = append(kept, item)
|
|
case id != "":
|
|
sawMismatch = true
|
|
default:
|
|
kept = append(kept, item)
|
|
}
|
|
}
|
|
if sawMatch {
|
|
return kept, true
|
|
}
|
|
if sawMismatch {
|
|
// Some items were tagged and none matched: the far side answered about
|
|
// other entities, so nothing here belongs to this one.
|
|
return nil, true
|
|
}
|
|
return nil, false
|
|
}
|
|
|
|
// localFactsForEntity summarises Maven's own facts already resolved to this
|
|
// canonical entity. Empty when the store is unavailable or nothing matched —
|
|
// entity-scoped memory is an enrichment of the answer, never a precondition.
|
|
func (h *reactiveHandler) localFactsForEntity(ctx context.Context, entityID string) string {
|
|
if h.dataStore == nil || entityID == "" {
|
|
return ""
|
|
}
|
|
const spoken = 3
|
|
// One over the spoken limit, so a truncation can be named rather than
|
|
// passed off as everything she knows.
|
|
facts, err := h.dataStore.FactsByEntity(ctx, entityID, spoken+1)
|
|
if err != nil {
|
|
log.Printf("ecosystem: facts by entity %s: %v", entityID, err)
|
|
return ""
|
|
}
|
|
more := false
|
|
if len(facts) > spoken {
|
|
facts, more = facts[:spoken], true
|
|
}
|
|
var parts []string
|
|
for _, f := range facts {
|
|
if f.Value != "" {
|
|
parts = append(parts, f.Value)
|
|
}
|
|
}
|
|
if len(parts) == 0 {
|
|
return ""
|
|
}
|
|
out := phraser.A(phraser.EcoRecall, map[string]string{"items": strings.Join(parts, ", ")})
|
|
if more {
|
|
out += ", и это не всё"
|
|
}
|
|
return out
|
|
}
|
|
|
|
// mergeFields overlays b onto a and returns a.
|
|
func mergeFields(a, b map[string]any) map[string]any {
|
|
for k, v := range b {
|
|
a[k] = v
|
|
}
|
|
return a
|
|
}
|
|
|
|
// recordPraxisTrace — records a completed Praxis call. Thin wrapper over
|
|
// recordEcosystemTrace so every ecosystem hop lands in one table with one
|
|
// shape.
|
|
func (h *reactiveHandler) recordPraxisTrace(ctx context.Context, operation string, started time.Time, details map[string]any) {
|
|
h.recordEcosystemTrace(ctx, "praxis", operation, traceOK, started, details)
|
|
}
|
|
|
|
// traceStatus classifies an ecosystem call for the trace record. Kept coarse
|
|
// on purpose: a trace is read to answer "did this hop work, and how long did
|
|
// it take", not to re-derive the error.
|
|
const (
|
|
traceOK = "ok"
|
|
traceFailed = "failed" // the call never got an answer
|
|
traceRefused = "refused" // the far side answered, and said no
|
|
traceAmbig = "ambiguous"
|
|
traceNotFound = "not_found"
|
|
tracePending = "pending" // deliberately not done yet, awaiting a confirm
|
|
)
|
|
|
|
// traceStatusForError distinguishes "I could not reach it" from "it answered
|
|
// and refused". Both degrade the same way for him and not at all the same way
|
|
// for whoever reads the trace: one is a network or a dead service, the other
|
|
// is a token, a version or a rejected argument.
|
|
func traceStatusForError(err error) string {
|
|
var ee *ecosystemError
|
|
if errors.As(err, &ee) && !ee.Unreachable() {
|
|
return traceRefused
|
|
}
|
|
return traceFailed
|
|
}
|
|
|
|
// redactSubject reduces a user utterance to something safe to persist in a
|
|
// trace: its length only. Traces are diagnostics, and his words are not
|
|
// diagnostics — the correlation ID is what ties a trace to the turn.
|
|
func redactSubject(s string) string {
|
|
return fmt.Sprintf("<%d chars>", len([]rune(s)))
|
|
}
|
|
|
|
// recordEcosystemTrace writes one hop of a cross-service call: which service,
|
|
// which operation, the outcome, how long it took, and the correlation ID that
|
|
// stitches the hops together. It is written for every outcome, not only
|
|
// success — an unrecorded failure is exactly the hop you need when something
|
|
// went wrong at 3am.
|
|
//
|
|
// Traces go to their own store table, never to facts. One act turn produces
|
|
// three or four of them, at machine rate, while facts arrive at human rate:
|
|
// sharing the table meant the habit profile's 2000-row window, memeval's
|
|
// prompt snapshot and the /dash and /history pages all filled with traces and
|
|
// stopped seeing his actual facts.
|
|
func (h *reactiveHandler) recordEcosystemTrace(ctx context.Context, service, op, status string, started time.Time, fields map[string]any) {
|
|
if h.dataStore == nil {
|
|
return
|
|
}
|
|
tr := store.EcosystemTrace{
|
|
Ts: h.now(),
|
|
Service: service,
|
|
Operation: op,
|
|
Status: status,
|
|
DurationMs: h.now().Sub(started).Milliseconds(),
|
|
CorrelationID: correlationIDFromCtx(ctx),
|
|
Fields: map[string]any{},
|
|
}
|
|
for k, v := range fields {
|
|
switch k {
|
|
case "causation_id":
|
|
tr.CausationID, _ = v.(string)
|
|
case "http_status":
|
|
if n, ok := v.(int); ok {
|
|
tr.HTTPStatus = n
|
|
continue
|
|
}
|
|
tr.Fields[k] = v
|
|
default:
|
|
tr.Fields[k] = v
|
|
}
|
|
}
|
|
if _, err := h.dataStore.WriteEcosystemTrace(ctx, tr); err != nil {
|
|
log.Printf("ecosystem: record trace %s:%s: %v", service, op, err)
|
|
}
|
|
}
|
|
|
|
// unauthorizedEcosystemError reports a credential the far side rejected. It
|
|
// gets its own reply: a missing or wrong token looks exactly like an outage to
|
|
// him, and "try again" is advice that will never work.
|
|
func unauthorizedEcosystemError(err error) bool {
|
|
var ee *ecosystemError
|
|
return errors.As(err, &ee) && ee.Unauthorized()
|
|
}
|
|
|
|
// traceErrorFields describes an ecosystemError for a trace without leaking the
|
|
// payload: the HTTP status and the failure class, nothing else.
|
|
func traceErrorFields(err error) map[string]any {
|
|
fields := map[string]any{}
|
|
var ee *ecosystemError
|
|
if errors.As(err, &ee) {
|
|
fields["http_status"] = ee.Status
|
|
switch {
|
|
case ee.Unauthorized():
|
|
fields["class"] = "unauthorized"
|
|
case ee.ContractMismatch():
|
|
fields["class"] = "contract_mismatch"
|
|
case ee.Unreachable():
|
|
fields["class"] = "unreachable"
|
|
default:
|
|
fields["class"] = "error"
|
|
}
|
|
return fields
|
|
}
|
|
fields["class"] = "error"
|
|
return fields
|
|
}
|
|
|
|
// handleHexisAct — resolves entity references through Nexus and executes
|
|
// matching capabilities through Hexis. Returns a reply string when handled,
|
|
// or "" to fall through to the system command executor.
|
|
func (h *reactiveHandler) handleHexisAct(ctx context.Context, dec router.Decision) string {
|
|
if h.ecosystem == nil {
|
|
return ""
|
|
}
|
|
|
|
// Every hop of this action shares one correlation ID, assigned here so
|
|
// resolution and discovery are traceable even when execution never
|
|
// happens.
|
|
if correlationIDFromCtx(ctx) == "" {
|
|
ctx = withCorrelationID(ctx, newCorrelationID())
|
|
}
|
|
|
|
// Resolve the utterance text as an entity reference through Nexus. An
|
|
// ambiguous match must stop and clarify — never guess a mutation target.
|
|
started := h.now()
|
|
entityID, displayName, ambiguous, err := h.ecosystem.resolveEntityReference(ctx, dec.Slots.Text, nil)
|
|
if err != nil {
|
|
h.recordEcosystemTrace(ctx, "nexus", "resolve", traceStatusForError(err), started,
|
|
mergeFields(traceErrorFields(err), map[string]any{"subject": redactSubject(dec.Slots.Text)}))
|
|
if unauthorizedEcosystemError(err) {
|
|
return phraser.A(phraser.EcoDenied, serviceVars(serviceNexus))
|
|
}
|
|
// A genuine Nexus dependency failure, not "no such entity" — stop here
|
|
// and report degradation rather than silently falling through to the
|
|
// local command executor (ECOSYSTEM-SPEC.md: services degrade
|
|
// independently, never a silent all-clear).
|
|
return phraser.A(phraser.EcoDown, serviceVars(serviceNexus))
|
|
}
|
|
if len(ambiguous) > 0 {
|
|
h.recordEcosystemTrace(ctx, "nexus", "resolve", traceAmbig, started,
|
|
map[string]any{"candidates": len(ambiguous)})
|
|
return phraser.A(phraser.EcoAmbiguous, map[string]string{"items": strings.Join(ambiguous, ", ")})
|
|
}
|
|
if entityID == "" {
|
|
h.recordEcosystemTrace(ctx, "nexus", "resolve", traceNotFound, started,
|
|
map[string]any{"subject": redactSubject(dec.Slots.Text)})
|
|
return ""
|
|
}
|
|
h.recordEcosystemTrace(ctx, "nexus", "resolve", traceOK, started,
|
|
map[string]any{"entity_id": entityID})
|
|
|
|
// Discover Hexis capabilities for this entity. A resolved entity with a
|
|
// genuine Hexis failure must not be treated as "no capabilities" and
|
|
// fall through to unrelated local execution.
|
|
discovered := h.now()
|
|
caps, err := h.ecosystem.discoverCapabilities(ctx, entityID)
|
|
if err != nil {
|
|
h.recordEcosystemTrace(ctx, "hexis", "capabilities", traceStatusForError(err), discovered,
|
|
mergeFields(traceErrorFields(err), map[string]any{"entity_id": entityID}))
|
|
if unauthorizedEcosystemError(err) {
|
|
return phraser.A(phraser.EcoDenied, serviceVars(serviceHexis))
|
|
}
|
|
return phraser.A(phraser.EcoDown, serviceVars(serviceHexis))
|
|
}
|
|
h.recordEcosystemTrace(ctx, "hexis", "capabilities", traceOK, discovered,
|
|
map[string]any{"entity_id": entityID, "count": len(caps)})
|
|
if len(caps) == 0 {
|
|
return ""
|
|
}
|
|
|
|
// Match the user's verb to a capability by name/description. Collect all
|
|
// matches: more than one is itself ambiguous, so we ask rather than pick
|
|
// the first (ecosystem invariant: no arbitrary target for mutation).
|
|
verb := dec.Slots.Fn
|
|
if verb == "" {
|
|
verb = dec.Slots.Text
|
|
}
|
|
verbLower := strings.ToLower(verb)
|
|
|
|
var matches []*hexisclient.Capability
|
|
for i, c := range caps {
|
|
if strings.Contains(strings.ToLower(c.Name), verbLower) ||
|
|
(c.Description != "" && strings.Contains(strings.ToLower(c.Description), verbLower)) {
|
|
matches = append(matches, &caps[i])
|
|
}
|
|
}
|
|
if len(matches) == 0 {
|
|
return ""
|
|
}
|
|
if len(matches) > 1 {
|
|
var names []string
|
|
for _, m := range matches {
|
|
names = append(names, m.Name)
|
|
}
|
|
return phraser.A(phraser.ActWhich, map[string]string{"name": displayName, "items": strings.Join(names, ", ")})
|
|
}
|
|
matched := matches[0]
|
|
|
|
// Read-only capabilities run immediately; mutating ones are parked for an
|
|
// explicit spoken confirm bound to this capability + target.
|
|
if !matched.ReadOnly {
|
|
h.mu.Lock()
|
|
h.pendingHexis = &pendingHexisExec{
|
|
capabilityID: matched.ID,
|
|
capName: matched.Name,
|
|
entityID: entityID,
|
|
displayName: displayName,
|
|
expiry: h.now().Add(confirmTTL),
|
|
}
|
|
h.mu.Unlock()
|
|
h.recordEcosystemTrace(ctx, "hexis", "confirmation", tracePending, started,
|
|
map[string]any{"entity_id": entityID, "capability": matched.Name})
|
|
return phraser.A(phraser.ActConfirmEntity, map[string]string{"name": matched.Name, "name_entity": displayName})
|
|
}
|
|
|
|
return h.execHexis(ctx, matched.ID, matched.Name, entityID, displayName)
|
|
}
|
|
|
|
// execHexis runs a resolved capability and records a cross-service trace with
|
|
// the correlation ID. It reports command success, never operational recovery
|
|
// (Praxis observes recovery independently).
|
|
func (h *reactiveHandler) execHexis(ctx context.Context, capID, capName, entityID, displayName string) string {
|
|
started := h.now()
|
|
causationID := correlationIDFromCtx(ctx)
|
|
correlationID, err := h.ecosystem.executeCapability(ctx, capID, entityID, nil)
|
|
traced := withCorrelationID(ctx, correlationID)
|
|
if err != nil {
|
|
log.Printf("ecosystem: hexis execute error (cor=%s): %v", correlationID, err)
|
|
h.recordEcosystemTrace(traced, "hexis", "execute", traceStatusForError(err), started,
|
|
mergeFields(traceErrorFields(err), map[string]any{
|
|
"entity_id": entityID, "capability": capName, "causation_id": causationID,
|
|
}))
|
|
return phraser.A(phraser.ActFailEntity, map[string]string{"name": displayName})
|
|
}
|
|
// One record per hop: the second write this used to make said the same
|
|
// thing under a different key, in a different shape.
|
|
h.recordEcosystemTrace(traced, "hexis", "execute", traceOK, started, map[string]any{
|
|
"entity_id": entityID, "entity_name": displayName,
|
|
"capability": capName, "causation_id": causationID,
|
|
})
|
|
return phraser.A(phraser.ActDoneEntity, map[string]string{"name": displayName})
|
|
}
|