feat(ecosystem): compliant Praxis/Hexis integration + vendored build
Bring the Nexus/Praxis/Hexis integration in line with MAVEN_ECOSYSTEM_ARCHITECTURE.md: - Praxis over HTTP: drop the in-process praxis.db open (praxisstore/ praxistools) and call praxisd's /api/v1/tools/* API via a new praxisClient. Honors the "no component reads another's DB" invariant (AC#12). PraxisConfig.DBPath -> URL. - Hexis confirmation gate: mutating capabilities (ReadOnly=false) now park a bound pendingHexis confirmation and require a spoken "да" before executing; read-only run immediately (AC#7, no auto attention->action). - Capability safety: >1 verb match is ambiguous -> ask instead of firing the first; ambiguous Nexus resolution asks for clarification (AC#2). - Correlation IDs on Hexis execute, recorded in the cross-service trace. - Bug: importance arrives as JSON float64 over HTTP, not int. - Tests: confirm-gate, decline, read-only, and ambiguity paths. Build: vendor/ bakes in the hexis client (replace-directed at a sibling repo outside the Docker context); Dockerfile builds from vendor and no longer `go mod download`s the unreachable replace paths. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,259 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/config"
|
||||
hexisclient "github.com/kami/hexis/pkg/client"
|
||||
)
|
||||
|
||||
type nexusClient struct {
|
||||
baseURL string
|
||||
httpClient *http.Client
|
||||
}
|
||||
|
||||
func newNexusClient(url string) *nexusClient {
|
||||
return &nexusClient{
|
||||
baseURL: url,
|
||||
httpClient: &http.Client{Timeout: 10 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
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"`
|
||||
}
|
||||
|
||||
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)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("do request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
bodyBytes, _ := io.ReadAll(resp.Body)
|
||||
if resp.StatusCode != 200 {
|
||||
return nil, fmt.Errorf("nexus: %s", http.StatusText(resp.StatusCode))
|
||||
}
|
||||
|
||||
var result nexusResolveResult
|
||||
if err := json.Unmarshal(bodyBytes, &result); err != nil {
|
||||
return nil, fmt.Errorf("decode: %w", 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 err
|
||||
}
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
resp.Body.Close()
|
||||
if resp.StatusCode != 200 {
|
||||
return fmt.Errorf("nexus health: %s", http.StatusText(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
|
||||
httpClient *http.Client
|
||||
}
|
||||
|
||||
func newPraxisClient(url string) *praxisClient {
|
||||
return &praxisClient{
|
||||
baseURL: url,
|
||||
httpClient: &http.Client{Timeout: 10 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
// getJSON performs a GET and decodes the JSON body into out.
|
||||
func (c *praxisClient) getJSON(ctx context.Context, path string, out any) error {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.baseURL+path, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
resp, err := c.httpClient.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != 200 {
|
||||
return fmt.Errorf("praxis: %s", http.StatusText(resp.StatusCode))
|
||||
}
|
||||
return json.NewDecoder(resp.Body).Decode(out)
|
||||
}
|
||||
|
||||
func (c *praxisClient) ListAttention(ctx context.Context, limit int) ([]map[string]any, error) {
|
||||
var out []map[string]any
|
||||
err := c.getJSON(ctx, fmt.Sprintf("/api/v1/tools/attention?limit=%d", limit), &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, fmt.Sprintf("/api/v1/tools/changes?limit=%d", 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)
|
||||
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)
|
||||
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)
|
||||
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
|
||||
}
|
||||
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" && result.Entity != nil {
|
||||
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.
|
||||
func (w *ecosystemWiring) discoverCapabilities(ctx context.Context, entityID string) []hexisclient.Capability {
|
||||
if w == nil || w.hexis == nil || entityID == "" {
|
||||
return nil
|
||||
}
|
||||
caps, err := w.hexis.Capabilities(ctx, entityID)
|
||||
if err != nil {
|
||||
log.Printf("ecosystem: hexis capabilities error: %v", err)
|
||||
return nil
|
||||
}
|
||||
return caps
|
||||
}
|
||||
|
||||
// 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()
|
||||
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[:])
|
||||
}
|
||||
@@ -0,0 +1,142 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
hexisclient "github.com/kami/hexis/pkg/client"
|
||||
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/router"
|
||||
)
|
||||
|
||||
// stubEcosystem wires nexus+hexis clients at the given base URLs.
|
||||
func stubEcosystem(nexusURL, hexisURL string) *ecosystemWiring {
|
||||
return &ecosystemWiring{
|
||||
nexus: newNexusClient(nexusURL),
|
||||
hexis: hexisclient.New(hexisURL),
|
||||
}
|
||||
}
|
||||
|
||||
// newHexisTestHandler builds a reactiveHandler backed by fake nexus+hexis
|
||||
// servers. resolveBody is returned verbatim from nexus /resolve; caps is the
|
||||
// capability list; executed records whether Hexis /execute was called.
|
||||
func newHexisTestHandler(t *testing.T, resolveBody string, caps string) (*reactiveHandler, *bool) {
|
||||
t.Helper()
|
||||
executed := new(bool)
|
||||
|
||||
nexus := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.Write([]byte(resolveBody))
|
||||
}))
|
||||
t.Cleanup(nexus.Close)
|
||||
|
||||
hexis := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
switch {
|
||||
case strings.HasPrefix(r.URL.Path, "/api/v1/capabilities"):
|
||||
w.Write([]byte(caps))
|
||||
case r.URL.Path == "/api/v1/execute":
|
||||
*executed = true
|
||||
w.Write([]byte(`{"id":"exec_1","status":"succeeded"}`))
|
||||
default:
|
||||
http.NotFound(w, r)
|
||||
}
|
||||
}))
|
||||
t.Cleanup(hexis.Close)
|
||||
|
||||
st := newTestStore(t)
|
||||
now := time.Now()
|
||||
return &reactiveHandler{
|
||||
api: ipc.NewStoreAPI(st),
|
||||
dataStore: st,
|
||||
now: func() time.Time { return now },
|
||||
ecosystem: stubEcosystem(nexus.URL, hexis.URL),
|
||||
}, executed
|
||||
}
|
||||
|
||||
func actDec(text string) router.Decision {
|
||||
return router.Decision{Intent: router.IntentAct, Slots: router.Slots{Text: text, Fn: "restart", HasFn: true}}
|
||||
}
|
||||
|
||||
func TestHexisMutatingRequiresConfirm(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
resolved := `{"status":"resolved","entity":{"id":"ent_muzick","display_name":"Muzick indexer","type":"service"}}`
|
||||
caps := `[{"id":"cap_restart","name":"restart","read_only":false,"risk":"high"}]`
|
||||
h, executed := newHexisTestHandler(t, resolved, caps)
|
||||
|
||||
reply := h.handleHexisAct(ctx, actDec("muzick indexer"))
|
||||
if !strings.Contains(reply, "да") {
|
||||
t.Fatalf("mutating cap should ask to confirm, got %q", reply)
|
||||
}
|
||||
if *executed {
|
||||
t.Fatal("mutating cap must NOT execute before confirmation")
|
||||
}
|
||||
if h.pendingHexis == nil || h.pendingHexis.capabilityID != "cap_restart" {
|
||||
t.Fatalf("expected pending hexis bound to cap_restart, got %+v", h.pendingHexis)
|
||||
}
|
||||
|
||||
// The follow-up "да" turn executes exactly the parked capability.
|
||||
confirmReply, handled := h.resolveConfirm(ctx, "да")
|
||||
if !handled || !strings.Contains(confirmReply, "выполнена") {
|
||||
t.Fatalf("confirm should execute, got handled=%v reply=%q", handled, confirmReply)
|
||||
}
|
||||
if !*executed {
|
||||
t.Fatal("confirmed mutating cap should have executed")
|
||||
}
|
||||
if h.pendingHexis != nil {
|
||||
t.Fatal("pending should be cleared after confirm")
|
||||
}
|
||||
}
|
||||
|
||||
func TestHexisConfirmNoDoesNotExecute(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
resolved := `{"status":"resolved","entity":{"id":"ent_muzick","display_name":"Muzick indexer","type":"service"}}`
|
||||
caps := `[{"id":"cap_restart","name":"restart","read_only":false}]`
|
||||
h, executed := newHexisTestHandler(t, resolved, caps)
|
||||
|
||||
_ = h.handleHexisAct(ctx, actDec("muzick indexer"))
|
||||
reply, handled := h.resolveConfirm(ctx, "нет")
|
||||
if !handled || !strings.Contains(reply, "отменила") {
|
||||
t.Fatalf("no should cancel, got handled=%v reply=%q", handled, reply)
|
||||
}
|
||||
if *executed {
|
||||
t.Fatal("declined cap must not execute")
|
||||
}
|
||||
}
|
||||
|
||||
func TestHexisReadOnlyExecutesImmediately(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
resolved := `{"status":"resolved","entity":{"id":"ent_muzick","display_name":"Muzick indexer","type":"service"}}`
|
||||
caps := `[{"id":"cap_status","name":"restart","read_only":true}]`
|
||||
h, executed := newHexisTestHandler(t, resolved, caps)
|
||||
|
||||
reply := h.handleHexisAct(ctx, actDec("muzick indexer"))
|
||||
if !*executed {
|
||||
t.Fatal("read-only cap should execute without confirmation")
|
||||
}
|
||||
if h.pendingHexis != nil {
|
||||
t.Fatal("read-only cap should not park a confirmation")
|
||||
}
|
||||
if !strings.Contains(reply, "выполнена") {
|
||||
t.Fatalf("unexpected reply %q", reply)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHexisAmbiguousAsksClarification(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
ambiguous := `{"status":"ambiguous","candidates":[{"entity_id":"ent_muzick","display_name":"Muzick indexer"},{"entity_id":"ent_manga","display_name":"Manga indexer"}]}`
|
||||
h, executed := newHexisTestHandler(t, ambiguous, `[]`)
|
||||
|
||||
reply := h.handleHexisAct(ctx, actDec("the indexer"))
|
||||
if !strings.Contains(reply, "Muzick indexer") || !strings.Contains(reply, "Manga indexer") {
|
||||
t.Fatalf("ambiguous should list candidates, got %q", reply)
|
||||
}
|
||||
if *executed {
|
||||
t.Fatal("ambiguous target must never execute")
|
||||
}
|
||||
}
|
||||
+16
-10
@@ -57,10 +57,10 @@ import (
|
||||
"github.com/kami/maven/internal/delivery/ntfysink"
|
||||
"github.com/kami/maven/internal/delivery/telegramsink"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
"github.com/kami/maven/internal/loop"
|
||||
"github.com/kami/maven/internal/phraser"
|
||||
"github.com/kami/maven/internal/store"
|
||||
"github.com/kami/maven/internal/webauthn"
|
||||
"github.com/kami/maven/internal/loop"
|
||||
)
|
||||
|
||||
var errLocked = errors.New("mavend: daemon locked — complete passkey assertion first")
|
||||
@@ -241,13 +241,14 @@ func run(args []string) error {
|
||||
// ----- daemon components (only wired when unlocked) -----
|
||||
// Pre-declare so the unlock path can wire them later.
|
||||
var (
|
||||
gatherer *loop.Gatherer
|
||||
rules []loop.Rule
|
||||
phr phraser.Phraser
|
||||
voiceW *voiceWiring
|
||||
dispatcher *delivery.Dispatcher
|
||||
tl *tickLoop
|
||||
coreAPI ipc.CoreAPI
|
||||
gatherer *loop.Gatherer
|
||||
rules []loop.Rule
|
||||
phr phraser.Phraser
|
||||
voiceW *voiceWiring
|
||||
dispatcher *delivery.Dispatcher
|
||||
tl *tickLoop
|
||||
coreAPI ipc.CoreAPI
|
||||
eco *ecosystemWiring
|
||||
)
|
||||
|
||||
if !locked {
|
||||
@@ -288,8 +289,11 @@ func run(args []string) error {
|
||||
}
|
||||
}
|
||||
|
||||
// ecosystem — nexus + hexis + praxis (all over HTTP; no direct DB access)
|
||||
eco = wireEcosystem(cfg)
|
||||
|
||||
// voice
|
||||
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st)
|
||||
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st, eco)
|
||||
if err != nil {
|
||||
return fmt.Errorf("wire voice: %w", err)
|
||||
}
|
||||
@@ -444,7 +448,9 @@ func run(args []string) error {
|
||||
}
|
||||
}
|
||||
|
||||
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st)
|
||||
eco = wireEcosystem(cfg)
|
||||
|
||||
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st, eco)
|
||||
if err != nil {
|
||||
return fmt.Errorf("wire voice: %w", err)
|
||||
}
|
||||
|
||||
+232
-1
@@ -45,6 +45,7 @@ package main
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
@@ -74,6 +75,7 @@ import (
|
||||
"github.com/kami/maven/internal/voice"
|
||||
"github.com/kami/maven/internal/weather"
|
||||
"github.com/kami/maven/internal/worker"
|
||||
hexisclient "github.com/kami/hexis/pkg/client"
|
||||
)
|
||||
|
||||
// voiceWiring — everything the daemon needs to run the audio path. Held by
|
||||
@@ -116,7 +118,7 @@ func (w *voiceWiring) close() {
|
||||
//
|
||||
// When voice is enabled, MUST wire a voicesink into the dispatcher's Voice
|
||||
// slot using w.sessions (the caller does that — see main.go).
|
||||
func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, memStore memory.Store, dataStore *store.Store) (*voiceWiring, error) {
|
||||
func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, memStore memory.Store, dataStore *store.Store, eco *ecosystemWiring) (*voiceWiring, error) {
|
||||
if cfg.Voice == nil || !cfg.Voice.Enabled {
|
||||
return nil, nil
|
||||
}
|
||||
@@ -250,6 +252,7 @@ func wireVoice(cfg *config.Config, coreAPI ipc.CoreAPI, phr phraser.Phraser, mem
|
||||
dialogueSessions: dialogueSessions,
|
||||
queryMinScore: cfg.Voice.QueryMinScore,
|
||||
timeParser: router.NewPythonDateParser(),
|
||||
ecosystem: eco,
|
||||
}
|
||||
|
||||
// ----- the server (TCP listener) -----
|
||||
@@ -308,6 +311,21 @@ type reactiveHandler struct {
|
||||
mu sync.Mutex
|
||||
pending *pendingAct
|
||||
pendingRoutine *pendingRoutineConfirm // routine proposal awaiting y/n
|
||||
pendingHexis *pendingHexisExec // mutating Hexis capability awaiting y/n
|
||||
|
||||
ecosystem *ecosystemWiring // nexus + hexis + praxis clients
|
||||
}
|
||||
|
||||
// pendingHexisExec — a mutating Hexis capability parked awaiting a spoken
|
||||
// confirm. The confirmation is bound to the resolved capability + canonical
|
||||
// target entity so a later "да" can only execute exactly what was proposed
|
||||
// (ecosystem invariant: protected actions require bound confirmation).
|
||||
type pendingHexisExec struct {
|
||||
capabilityID string
|
||||
capName string
|
||||
entityID string
|
||||
displayName string
|
||||
expiry time.Time
|
||||
}
|
||||
|
||||
// pendingRoutineConfirm — a proposed routine awaiting a spoken y/n to become
|
||||
@@ -595,6 +613,22 @@ func (h *reactiveHandler) applyAction(ctx context.Context, dec router.Decision)
|
||||
dec.Slots.Fn, dec.Slots.Args, dec.Slots.HasFn = fn, args, true
|
||||
}
|
||||
}
|
||||
|
||||
// Praxis ecosystem tools: intercept before the system command executor.
|
||||
if h.ecosystem != nil && h.ecosystem.praxis != nil && dec.Slots.HasFn {
|
||||
if reply := h.handlePraxisAct(ctx, dec); reply != "" {
|
||||
return reply
|
||||
}
|
||||
}
|
||||
|
||||
// Hexis ecosystem action: if ecosystem is configured and we have a verb
|
||||
// + entity text, try to resolve the entity and execute via Hexis.
|
||||
if h.ecosystem != nil && h.ecosystem.hexis != nil && dec.Slots.Text != "" {
|
||||
if reply := h.handleHexisAct(ctx, dec); reply != "" {
|
||||
return reply
|
||||
}
|
||||
}
|
||||
|
||||
// HasFn still false ⇒ no allowlist match: scaffold a 'proposed' tool
|
||||
// the user can enable on the authed surface ("earn the right to ask").
|
||||
if !dec.Slots.HasFn {
|
||||
@@ -1097,6 +1131,183 @@ func loadSeedFile(c *router.Classifier, intent router.Intent) (int, error) {
|
||||
return count, nil
|
||||
}
|
||||
|
||||
// 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 ""
|
||||
}
|
||||
px := h.ecosystem.praxis
|
||||
fn := dec.Slots.Fn
|
||||
|
||||
// Map verbs and Russian aliases to Praxis tool calls.
|
||||
// Each case: if the verb matches, call the tool and return a user-facing reply.
|
||||
switch fn {
|
||||
case "list_attention", "attention", "внимание", "что требует внимания", "что нового":
|
||||
items, err := px.ListAttention(ctx, 20)
|
||||
if err != nil {
|
||||
log.Printf("ecosystem: praxis attention: %v", err)
|
||||
return "не могу сейчас узнать, что требует внимания."
|
||||
}
|
||||
if len(items) == 0 {
|
||||
return "ничего не требует внимания."
|
||||
}
|
||||
h.recordPraxisTrace(ctx, "list_attention", 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 importance > 0 {
|
||||
s += fmt.Sprintf(" (важность %d", int(importance))
|
||||
if rule != "" {
|
||||
s += ": " + rule
|
||||
}
|
||||
s += ")"
|
||||
}
|
||||
parts = append(parts, s)
|
||||
}
|
||||
return "требует внимания: " + strings.Join(parts, "; ")
|
||||
|
||||
case "list_changes", "changes", "изменения", "что изменилось":
|
||||
changes, err := px.ListChanges(ctx, 20)
|
||||
if err != nil {
|
||||
log.Printf("ecosystem: praxis changes: %v", err)
|
||||
return "не могу сейчас узнать об изменениях."
|
||||
}
|
||||
if len(changes) == 0 {
|
||||
return "нет изменений."
|
||||
}
|
||||
h.recordPraxisTrace(ctx, "list_changes", map[string]any{"count": len(changes)})
|
||||
var parts []string
|
||||
for _, c := range changes {
|
||||
title, _ := c["title"].(string)
|
||||
typ, _ := c["change_type"].(string)
|
||||
parts = append(parts, fmt.Sprintf("%s (%s)", title, typ))
|
||||
}
|
||||
return "изменения: " + strings.Join(parts, "; ")
|
||||
|
||||
default:
|
||||
// Not a Praxis verb — let the caller fall through.
|
||||
return ""
|
||||
}
|
||||
}
|
||||
|
||||
// recordPraxisTrace — writes a fact recording a cross-service ecosystem call.
|
||||
// The fact is stored with source "praxis:trace" so the proactive loop can
|
||||
// reference it and the dashboard can display recent ecosystem activity.
|
||||
func (h *reactiveHandler) recordPraxisTrace(ctx context.Context, operation string, details map[string]any) {
|
||||
now := h.now()
|
||||
value := operation
|
||||
if len(details) > 0 {
|
||||
if b, err := json.Marshal(details); err == nil {
|
||||
value = operation + " " + string(b)
|
||||
}
|
||||
}
|
||||
_, _ = h.api.WriteFact(ctx, ipc.WriteFactReq{
|
||||
Ts: now,
|
||||
Kind: "system",
|
||||
Key: "praxis:" + operation,
|
||||
Value: value,
|
||||
Source: "praxis:trace",
|
||||
Confidence: 1.0,
|
||||
})
|
||||
}
|
||||
|
||||
// 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 ""
|
||||
}
|
||||
|
||||
// Resolve the utterance text as an entity reference through Nexus. An
|
||||
// ambiguous match must stop and clarify — never guess a mutation target.
|
||||
entityID, displayName, ambiguous, err := h.ecosystem.resolveEntityReference(ctx, dec.Slots.Text, nil)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
if len(ambiguous) > 0 {
|
||||
return "уточни, что именно: " + strings.Join(ambiguous, ", ") + "?"
|
||||
}
|
||||
if entityID == "" {
|
||||
return ""
|
||||
}
|
||||
|
||||
// Discover Hexis capabilities for this entity.
|
||||
caps := h.ecosystem.discoverCapabilities(ctx, entityID)
|
||||
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 "какую команду для " + displayName + ": " + 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()
|
||||
return "выполнить «" + matched.Name + "» для " + 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 {
|
||||
correlationID, err := h.ecosystem.executeCapability(ctx, capID, entityID, nil)
|
||||
if err != nil {
|
||||
log.Printf("ecosystem: hexis execute error (cor=%s): %v", correlationID, err)
|
||||
return "не получилось выполнить команду для " + displayName + "."
|
||||
}
|
||||
h.recordPraxisTrace(ctx, "hexis:"+capName, map[string]any{
|
||||
"entity_id": entityID,
|
||||
"entity_name": displayName,
|
||||
"capability": capName,
|
||||
"correlation_id": correlationID,
|
||||
})
|
||||
return "команда выполнена для " + displayName + "."
|
||||
}
|
||||
|
||||
// jsonString — a one-line JSON string encoder without dragging encoding/json
|
||||
// into the top of this file. Used to wrap a reminder payload's text field;
|
||||
// the router's reminder Slots are already absolute (DateTimeParser resolved
|
||||
@@ -1166,6 +1377,26 @@ func (h *reactiveHandler) resolveConfirm(ctx context.Context, text string) (stri
|
||||
h.pendingRoutine = nil
|
||||
}
|
||||
|
||||
// Check pending Hexis execution confirm. Bound to the exact capability +
|
||||
// target that was proposed; a stray "да" can only run that, nothing else.
|
||||
if hx := h.pendingHexis; hx != nil {
|
||||
if h.now().After(hx.expiry) {
|
||||
h.pendingHexis = nil
|
||||
} else {
|
||||
switch classifyConfirm(text) {
|
||||
case confirmYes:
|
||||
h.pendingHexis = nil
|
||||
return h.execHexis(ctx, hx.capabilityID, hx.capName, hx.entityID, hx.displayName), true
|
||||
case confirmNo:
|
||||
h.pendingHexis = nil
|
||||
return "отменила.", true
|
||||
default:
|
||||
h.pendingHexis = nil
|
||||
return "", false
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Check tool confirm (existing behavior).
|
||||
p := h.pending
|
||||
if p == nil {
|
||||
|
||||
Reference in New Issue
Block a user