Files
hexis/internal/api/handler.go
T
kami 74b19e091e Wire hexis.resolve_target to real Nexus, fix changes cursor, pass full execute fields over MCP
resolve_target previously returned a hardcoded "requires_nexus_resolution"
placeholder; it now calls Nexus's /api/v1/resolve via a new minimal
internal/nexusclient, configurable with -nexus (default localhost:8987).

/api/v1/changes ignored the since query param and always returned from
sequence 0 (`since = 0` regardless of what was parsed) — fixed to actually
parse and use it, so change-cursor polling works.

The MCP hexis.execute tool only forwarded capability_id/target_entity_id/
arguments/idempotency_key, silently dropping entity_version, requested_by,
origin, correlation_id, causation_id, resolution_evidence, and
confirmation_id even though the native HTTP API and domain.ExecuteRequest
already supported all of them — MCP callers now get full parity.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018ghELqYhZNLub2TXGMazqA
2026-07-20 11:28:13 +04:00

366 lines
11 KiB
Go

package api
import (
"encoding/json"
"errors"
"net/http"
"strconv"
"strings"
"time"
"github.com/kami/hexis/internal/domain"
"github.com/kami/hexis/internal/execution"
"github.com/kami/hexis/internal/storage"
)
type Handler struct {
store storage.Interface
engine *execution.Engine
}
func NewHandler(store storage.Interface, engine *execution.Engine) *Handler {
return &Handler{store: store, engine: engine}
}
// SupportedAPIVersion is the version this server implements. A request
// carrying X-Hexis-Version set to anything else is rejected — clients that
// don't send the header at all are allowed through unversioned, to avoid
// breaking callers mid-rollout.
const SupportedAPIVersion = "v1"
func (h *Handler) Register(mux *http.ServeMux) {
mux.HandleFunc("/health", h.health)
mux.HandleFunc("/ready", h.ready)
api := http.NewServeMux()
api.HandleFunc("/api/v1/capabilities", h.handleCapabilities)
api.HandleFunc("/api/v1/capabilities/", h.handleCapabilityByID)
api.HandleFunc("/api/v1/execute", h.handleExecute)
api.HandleFunc("/api/v1/confirmations", h.handleConfirmations)
api.HandleFunc("/api/v1/executions/", h.handleExecutionByID)
api.HandleFunc("/api/v1/changes", h.handleChanges)
mux.Handle("/api/v1/", versionCheck(api))
}
func versionCheck(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if v := r.Header.Get("X-Hexis-Version"); v != "" && v != SupportedAPIVersion {
writeJSON(w, http.StatusPreconditionFailed, map[string]string{
"error": "unsupported API version",
"requested_version": v,
"supported_version": SupportedAPIVersion,
})
return
}
next.ServeHTTP(w, r)
})
}
func (h *Handler) health(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, map[string]string{"status": "ok"})
}
func (h *Handler) ready(w http.ResponseWriter, r *http.Request) {
_, err := h.store.LatestSequence()
if err != nil {
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"status": "not_ready"})
return
}
writeJSON(w, http.StatusOK, map[string]string{"status": "ready"})
}
func (h *Handler) handleCapabilities(w http.ResponseWriter, r *http.Request) {
switch r.Method {
case http.MethodGet:
h.listCapabilities(w, r)
case http.MethodPost:
h.createCapability(w, r)
default:
writeJSON(w, http.StatusMethodNotAllowed, errorResponse("method not allowed"))
}
}
func (h *Handler) listCapabilities(w http.ResponseWriter, r *http.Request) {
entityID := r.URL.Query().Get("entity_id")
caps, err := h.store.ListCapabilities(entityID)
if err != nil {
writeJSON(w, http.StatusInternalServerError, errorResponse(err.Error()))
return
}
if caps == nil {
caps = []*domain.Capability{}
}
var apiCaps []map[string]any
for _, c := range caps {
apiCaps = append(apiCaps, map[string]any{
"capability_id": c.ID,
"id": c.ID,
"name": c.Name,
"description": c.Description,
"target_types": c.TargetTypes,
"target_entity_id": c.TargetEntityID,
"provider": c.Provider,
"operation": c.Operation,
"risk": c.Risk,
"read_only": c.ReadOnly,
"expected_side_effects": c.ExpectedSideEffects,
})
}
writeJSON(w, http.StatusOK, apiCaps)
}
func (h *Handler) createCapability(w http.ResponseWriter, r *http.Request) {
var req struct {
Name string `json:"name"`
Description string `json:"description,omitempty"`
TargetTypes []string `json:"target_types"`
TargetEntityID string `json:"target_entity_id,omitempty"`
Provider string `json:"provider"`
Operation string `json:"operation"`
Risk string `json:"risk,omitempty"`
ReadOnly bool `json:"read_only"`
ExpectedSideEffects string `json:"expected_side_effects,omitempty"`
RequiresConfirmation bool `json:"requires_confirmation,omitempty"`
Enabled *bool `json:"enabled,omitempty"`
TimeoutSeconds int `json:"timeout_seconds,omitempty"`
Attributes map[string]any `json:"attributes,omitempty"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, errorResponse("invalid JSON"))
return
}
if req.Name == "" || req.Provider == "" || req.Operation == "" {
writeJSON(w, http.StatusBadRequest, errorResponse("name, provider, and operation are required"))
return
}
// Destructive capabilities are disabled by default and must be turned on
// explicitly (ECOSYSTEM-SPEC.md §4.3).
enabled := req.Risk != "destructive"
if req.Enabled != nil {
enabled = *req.Enabled
}
now := time.Now().UTC()
cap := &domain.Capability{
ID: domain.NewCapabilityID(),
Name: req.Name,
Description: req.Description,
TargetTypes: req.TargetTypes,
TargetEntityID: req.TargetEntityID,
Provider: req.Provider,
Operation: req.Operation,
Risk: req.Risk,
ReadOnly: req.ReadOnly,
ExpectedSideEffects: req.ExpectedSideEffects,
RequiresConfirmation: req.RequiresConfirmation,
Enabled: enabled,
TimeoutSeconds: req.TimeoutSeconds,
Attributes: req.Attributes,
CreatedAt: now,
UpdatedAt: now,
Version: 1,
}
if cap.TargetTypes == nil {
cap.TargetTypes = []string{}
}
if cap.Attributes == nil {
cap.Attributes = map[string]any{}
}
if err := h.store.CreateCapability(cap); err != nil {
writeJSON(w, http.StatusInternalServerError, errorResponse(err.Error()))
return
}
h.store.AppendEvent(&domain.Event{
ID: domain.NewEventID(),
Type: domain.EventCapabilityRegistered,
Timestamp: now,
Payload: map[string]any{"capability_id": cap.ID, "name": cap.Name},
})
writeJSON(w, http.StatusCreated, cap)
}
func (h *Handler) handleCapabilityByID(w http.ResponseWriter, r *http.Request) {
id := strings.TrimPrefix(r.URL.Path, "/api/v1/capabilities/")
if id == "" {
writeJSON(w, http.StatusBadRequest, errorResponse("id required"))
return
}
switch r.Method {
case http.MethodGet:
h.getCapability(w, r, id)
case http.MethodDelete:
h.deleteCapability(w, r, id)
default:
writeJSON(w, http.StatusMethodNotAllowed, errorResponse("method not allowed"))
}
}
func (h *Handler) getCapability(w http.ResponseWriter, r *http.Request, id string) {
cap, err := h.store.GetCapability(id)
if err != nil {
writeJSON(w, http.StatusNotFound, errorResponse(err.Error()))
return
}
writeJSON(w, http.StatusOK, cap)
}
func (h *Handler) deleteCapability(w http.ResponseWriter, r *http.Request, id string) {
if err := h.store.DeleteCapability(id); err != nil {
writeJSON(w, http.StatusNotFound, errorResponse(err.Error()))
return
}
writeJSON(w, http.StatusOK, map[string]string{"status": "deleted"})
}
func (h *Handler) handleExecute(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, errorResponse("method not allowed"))
return
}
var req domain.ExecuteRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, errorResponse("invalid JSON"))
return
}
if req.CapabilityID == "" || req.TargetEntityID == "" {
writeJSON(w, http.StatusBadRequest, errorResponse("capability_id and target_entity_id are required"))
return
}
if req.CorrelationID == "" {
req.CorrelationID = r.Header.Get("X-Correlation-ID")
}
if req.CausationID == "" {
req.CausationID = r.Header.Get("X-Causation-ID")
}
result, err := h.engine.Execute(&req)
if err != nil {
writeJSON(w, executeErrorStatus(err), errorResponse(err.Error()))
return
}
writeJSON(w, http.StatusOK, result.Execution)
}
func executeErrorStatus(err error) int {
switch {
case errors.Is(err, domain.ErrExecutionInFlight):
return http.StatusConflict
case errors.Is(err, domain.ErrCapabilityNotFound):
return http.StatusNotFound
case errors.Is(err, domain.ErrConfirmationRequired),
errors.Is(err, domain.ErrConfirmationInvalid),
errors.Is(err, domain.ErrConfirmationExpired),
errors.Is(err, domain.ErrConfirmationConsumed),
errors.Is(err, domain.ErrConfirmationNotFound),
errors.Is(err, domain.ErrCapabilityDisabled),
errors.Is(err, domain.ErrCapabilityNotBound):
return http.StatusForbidden
default:
return http.StatusBadRequest
}
}
func (h *Handler) handleConfirmations(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSON(w, http.StatusMethodNotAllowed, errorResponse("method not allowed"))
return
}
var req struct {
CapabilityID string `json:"capability_id"`
TargetEntityID string `json:"target_entity_id"`
Arguments map[string]any `json:"arguments,omitempty"`
Requester string `json:"requester,omitempty"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, errorResponse("invalid JSON"))
return
}
if req.CapabilityID == "" || req.TargetEntityID == "" {
writeJSON(w, http.StatusBadRequest, errorResponse("capability_id and target_entity_id are required"))
return
}
conf, err := h.engine.CreateConfirmation(req.CapabilityID, req.TargetEntityID, req.Requester, req.Arguments)
if err != nil {
status := http.StatusBadRequest
if errors.Is(err, domain.ErrCapabilityNotFound) {
status = http.StatusNotFound
} else if errors.Is(err, domain.ErrCapabilityNotBound) {
status = http.StatusForbidden
}
writeJSON(w, status, errorResponse(err.Error()))
return
}
writeJSON(w, http.StatusCreated, conf)
}
func (h *Handler) handleExecutionByID(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, errorResponse("method not allowed"))
return
}
id := strings.TrimPrefix(r.URL.Path, "/api/v1/executions/")
if id == "" {
writeJSON(w, http.StatusBadRequest, errorResponse("id required"))
return
}
exec, err := h.store.GetExecution(id)
if err != nil {
writeJSON(w, http.StatusNotFound, errorResponse(err.Error()))
return
}
writeJSON(w, http.StatusOK, exec)
}
func (h *Handler) handleChanges(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
writeJSON(w, http.StatusMethodNotAllowed, errorResponse("method not allowed"))
return
}
seqStr := r.URL.Query().Get("since")
var since int64
if seqStr != "" {
parsed, err := strconv.ParseInt(seqStr, 10, 64)
if err != nil {
writeJSON(w, http.StatusBadRequest, errorResponse("since must be an integer sequence"))
return
}
since = parsed
}
events, err := h.store.EventsAfter(since, 100)
if err != nil {
writeJSON(w, http.StatusInternalServerError, errorResponse(err.Error()))
return
}
if events == nil {
events = []*domain.Event{}
}
writeJSON(w, http.StatusOK, events)
}
func writeJSON(w http.ResponseWriter, status int, v any) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
json.NewEncoder(w).Encode(v)
}
func errorResponse(msg string) map[string]string {
return map[string]string{"error": msg}
}