cca63269e1
Go daemon (hexisd/hexisctl) implementing capability registry, guarded execution (confirmations, blessed-entity checks), systemd/workspace-mcp providers per ECOSYSTEM-SPEC.md. Snapshotting existing working state before further development.
450 lines
13 KiB
Go
450 lines
13 KiB
Go
package storage
|
|
|
|
import (
|
|
"database/sql"
|
|
"encoding/json"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/kami/hexis/internal/domain"
|
|
_ "modernc.org/sqlite"
|
|
)
|
|
|
|
type Store struct {
|
|
mu sync.RWMutex
|
|
db *sql.DB
|
|
path string
|
|
}
|
|
|
|
func Open(path string) (*Store, error) {
|
|
dir := filepath.Dir(path)
|
|
if err := os.MkdirAll(dir, 0755); err != nil {
|
|
return nil, fmt.Errorf("create directory: %w", err)
|
|
}
|
|
|
|
db, err := sql.Open("sqlite", path+"?_pragma=journal_mode(WAL)&_pragma=foreign_keys(1)")
|
|
if err != nil {
|
|
return nil, fmt.Errorf("open database: %w", err)
|
|
}
|
|
|
|
db.SetMaxOpenConns(1)
|
|
|
|
store := &Store{db: db, path: path}
|
|
if err := store.migrate(); err != nil {
|
|
return nil, fmt.Errorf("migrate: %w", err)
|
|
}
|
|
return store, nil
|
|
}
|
|
|
|
func (s *Store) Close() error {
|
|
return s.db.Close()
|
|
}
|
|
|
|
func (s *Store) DB() *sql.DB {
|
|
return s.db
|
|
}
|
|
|
|
func (s *Store) migrate() error {
|
|
tx, err := s.db.Begin()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
var v int
|
|
err = tx.QueryRow("PRAGMA user_version").Scan(&v)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if v < len(migrations) {
|
|
for i, m := range migrations[v:] {
|
|
if _, err := tx.Exec(m); err != nil {
|
|
return fmt.Errorf("migration %d: %w", v+i+1, err)
|
|
}
|
|
if _, err := tx.Exec(fmt.Sprintf("PRAGMA user_version = %d", v+i+1)); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
|
|
return tx.Commit()
|
|
}
|
|
|
|
var migrations = []string{
|
|
`CREATE TABLE IF NOT EXISTS capabilities (
|
|
id TEXT PRIMARY KEY,
|
|
name TEXT NOT NULL,
|
|
description TEXT NOT NULL DEFAULT '',
|
|
target_types TEXT NOT NULL DEFAULT '[]',
|
|
target_entity_id TEXT NOT NULL DEFAULT '',
|
|
provider TEXT NOT NULL,
|
|
operation TEXT NOT NULL,
|
|
risk TEXT NOT NULL DEFAULT '',
|
|
read_only INTEGER NOT NULL DEFAULT 1,
|
|
expected_side_effects TEXT NOT NULL DEFAULT '',
|
|
attributes TEXT NOT NULL DEFAULT '{}',
|
|
created_at TEXT NOT NULL,
|
|
updated_at TEXT NOT NULL,
|
|
version INTEGER NOT NULL DEFAULT 0
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS executions (
|
|
id TEXT PRIMARY KEY,
|
|
capability_id TEXT NOT NULL,
|
|
target_entity_id TEXT NOT NULL,
|
|
entity_version INTEGER NOT NULL DEFAULT 0,
|
|
arguments TEXT NOT NULL DEFAULT '{}',
|
|
requested_by TEXT NOT NULL DEFAULT '{}',
|
|
origin TEXT NOT NULL DEFAULT '{}',
|
|
idempotency_key TEXT UNIQUE,
|
|
status TEXT NOT NULL DEFAULT 'started',
|
|
result TEXT NOT NULL DEFAULT '{}',
|
|
error TEXT NOT NULL DEFAULT '',
|
|
resolution_evidence TEXT NOT NULL DEFAULT '[]',
|
|
correlation_id TEXT NOT NULL DEFAULT '',
|
|
created_at TEXT NOT NULL,
|
|
updated_at TEXT NOT NULL
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS hexis_events (
|
|
id TEXT PRIMARY KEY,
|
|
sequence INTEGER NOT NULL,
|
|
type TEXT NOT NULL,
|
|
timestamp TEXT NOT NULL,
|
|
actor TEXT NOT NULL DEFAULT '',
|
|
correlation_id TEXT NOT NULL DEFAULT '',
|
|
causation_id TEXT NOT NULL DEFAULT '',
|
|
payload TEXT NOT NULL DEFAULT '{}'
|
|
)`,
|
|
`CREATE INDEX IF NOT EXISTS idx_hevents_sequence ON hexis_events(sequence)`,
|
|
`CREATE TABLE IF NOT EXISTS schema_migrations (
|
|
version INTEGER PRIMARY KEY,
|
|
applied_at TEXT NOT NULL
|
|
)`,
|
|
}
|
|
|
|
const timeFmt = "2006-01-02T15:04:05.999999999Z07:00"
|
|
|
|
func formatTime(t time.Time) string {
|
|
return t.UTC().Format(timeFmt)
|
|
}
|
|
|
|
func parseTime(s string) time.Time {
|
|
t, err := time.Parse(timeFmt, s)
|
|
if err != nil {
|
|
t, err = time.Parse(time.RFC3339, s)
|
|
if err != nil {
|
|
return time.Time{}
|
|
}
|
|
}
|
|
return t
|
|
}
|
|
|
|
// Capability operations
|
|
|
|
func (s *Store) CreateCapability(c *domain.Capability) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
targetTypes, _ := json.Marshal(c.TargetTypes)
|
|
attrs, _ := json.Marshal(c.Attributes)
|
|
|
|
_, err := s.db.Exec(
|
|
`INSERT INTO capabilities (id, name, description, target_types, target_entity_id, provider, operation, risk, read_only, expected_side_effects, attributes, created_at, updated_at, version)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
|
c.ID, c.Name, c.Description, string(targetTypes), c.TargetEntityID, c.Provider, c.Operation, c.Risk, boolInt(c.ReadOnly), c.ExpectedSideEffects, string(attrs), formatTime(c.CreatedAt), formatTime(c.UpdatedAt), c.Version,
|
|
)
|
|
return err
|
|
}
|
|
|
|
func (s *Store) GetCapability(id string) (*domain.Capability, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
|
|
row := s.db.QueryRow(
|
|
`SELECT id, name, description, target_types, COALESCE(target_entity_id,''), provider, operation, risk, read_only, COALESCE(expected_side_effects,''), attributes, created_at, updated_at, version
|
|
FROM capabilities WHERE id = ?`, id,
|
|
)
|
|
c := &domain.Capability{}
|
|
var targetTypes, attrs, createdAt, updatedAt string
|
|
err := row.Scan(&c.ID, &c.Name, &c.Description, &targetTypes, &c.TargetEntityID, &c.Provider, &c.Operation, &c.Risk, &c.ReadOnly, &c.ExpectedSideEffects, &attrs, &createdAt, &updatedAt, &c.Version)
|
|
if err == sql.ErrNoRows {
|
|
return nil, domain.ErrCapabilityNotFound
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
json.Unmarshal([]byte(targetTypes), &c.TargetTypes)
|
|
json.Unmarshal([]byte(attrs), &c.Attributes)
|
|
c.CreatedAt = parseTime(createdAt)
|
|
c.UpdatedAt = parseTime(updatedAt)
|
|
if c.TargetTypes == nil {
|
|
c.TargetTypes = []string{}
|
|
}
|
|
if c.Attributes == nil {
|
|
c.Attributes = map[string]any{}
|
|
}
|
|
return c, nil
|
|
}
|
|
|
|
func (s *Store) UpdateCapability(c *domain.Capability) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
targetTypes, _ := json.Marshal(c.TargetTypes)
|
|
attrs, _ := json.Marshal(c.Attributes)
|
|
|
|
res, err := s.db.Exec(
|
|
`UPDATE capabilities SET name=?, description=?, target_types=?, target_entity_id=?, provider=?, operation=?, risk=?, read_only=?, expected_side_effects=?, attributes=?, updated_at=?, version=version+1
|
|
WHERE id=? AND version=?`,
|
|
c.Name, c.Description, string(targetTypes), c.TargetEntityID, c.Provider, c.Operation, c.Risk, boolInt(c.ReadOnly), c.ExpectedSideEffects, string(attrs), formatTime(c.UpdatedAt), c.ID, c.Version,
|
|
)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
n, _ := res.RowsAffected()
|
|
if n == 0 {
|
|
return domain.ErrConflict
|
|
}
|
|
c.Version++
|
|
return nil
|
|
}
|
|
|
|
func (s *Store) ListCapabilities(entityID string) ([]*domain.Capability, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
|
|
query := `SELECT id, name, description, target_types, COALESCE(target_entity_id,''), provider, operation, risk, read_only, COALESCE(expected_side_effects,''), attributes, created_at, updated_at, version
|
|
FROM capabilities`
|
|
args := []any{}
|
|
if entityID != "" {
|
|
query += " WHERE target_entity_id = ?"
|
|
args = append(args, entityID)
|
|
}
|
|
query += " ORDER BY name ASC"
|
|
|
|
rows, err := s.db.Query(query, args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var result []*domain.Capability
|
|
for rows.Next() {
|
|
c := &domain.Capability{}
|
|
var targetTypes, attrs, createdAt, updatedAt string
|
|
if err := rows.Scan(&c.ID, &c.Name, &c.Description, &targetTypes, &c.TargetEntityID, &c.Provider, &c.Operation, &c.Risk, &c.ReadOnly, &c.ExpectedSideEffects, &attrs, &createdAt, &updatedAt, &c.Version); err != nil {
|
|
return nil, err
|
|
}
|
|
json.Unmarshal([]byte(targetTypes), &c.TargetTypes)
|
|
json.Unmarshal([]byte(attrs), &c.Attributes)
|
|
c.CreatedAt = parseTime(createdAt)
|
|
c.UpdatedAt = parseTime(updatedAt)
|
|
if c.TargetTypes == nil {
|
|
c.TargetTypes = []string{}
|
|
}
|
|
if c.Attributes == nil {
|
|
c.Attributes = map[string]any{}
|
|
}
|
|
result = append(result, c)
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func (s *Store) DeleteCapability(id string) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
_, err := s.db.Exec(`DELETE FROM capabilities WHERE id = ?`, id)
|
|
return err
|
|
}
|
|
|
|
// Execution operations
|
|
|
|
func (s *Store) CreateExecution(e *domain.Execution) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
args, _ := json.Marshal(e.Arguments)
|
|
reqBy, _ := json.Marshal(e.RequestedBy)
|
|
origin, _ := json.Marshal(e.Origin)
|
|
result, _ := json.Marshal(e.Result)
|
|
evidence, _ := json.Marshal(e.ResolutionEvidence)
|
|
|
|
_, err := s.db.Exec(
|
|
`INSERT INTO executions (id, capability_id, target_entity_id, entity_version, arguments, requested_by, origin, idempotency_key, status, result, error, resolution_evidence, correlation_id, created_at, updated_at)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
|
e.ID, e.CapabilityID, e.TargetEntityID, e.EntityVersion, string(args), string(reqBy), string(origin), nullString(e.IdempotencyKey), string(e.Status), string(result), e.Error, string(evidence), e.CorrelationID, formatTime(e.CreatedAt), formatTime(e.UpdatedAt),
|
|
)
|
|
if err != nil {
|
|
if strings.Contains(err.Error(), "UNIQUE") {
|
|
return domain.ErrIdempotencyReplay
|
|
}
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Store) GetExecution(id string) (*domain.Execution, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
|
|
row := s.db.QueryRow(
|
|
`SELECT id, capability_id, target_entity_id, entity_version, arguments, requested_by, origin, COALESCE(idempotency_key,''), status, result, COALESCE(error,''), resolution_evidence, COALESCE(correlation_id,''), created_at, updated_at
|
|
FROM executions WHERE id = ?`, id,
|
|
)
|
|
e := &domain.Execution{}
|
|
var args, reqBy, origin, idempKey, status, result, errStr, evidence, corrID, createdAt, updatedAt string
|
|
err := row.Scan(&e.ID, &e.CapabilityID, &e.TargetEntityID, &e.EntityVersion, &args, &reqBy, &origin, &idempKey, &status, &result, &errStr, &evidence, &corrID, &createdAt, &updatedAt)
|
|
if err == sql.ErrNoRows {
|
|
return nil, domain.ErrExecutionNotFound
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
json.Unmarshal([]byte(args), &e.Arguments)
|
|
json.Unmarshal([]byte(reqBy), &e.RequestedBy)
|
|
json.Unmarshal([]byte(origin), &e.Origin)
|
|
json.Unmarshal([]byte(result), &e.Result)
|
|
json.Unmarshal([]byte(evidence), &e.ResolutionEvidence)
|
|
e.IdempotencyKey = idempKey
|
|
e.Status = domain.ExecutionStatus(status)
|
|
e.Error = errStr
|
|
e.CorrelationID = corrID
|
|
e.CreatedAt = parseTime(createdAt)
|
|
e.UpdatedAt = parseTime(updatedAt)
|
|
if e.Arguments == nil {
|
|
e.Arguments = map[string]any{}
|
|
}
|
|
if e.RequestedBy == nil {
|
|
e.RequestedBy = map[string]string{}
|
|
}
|
|
if e.Origin == nil {
|
|
e.Origin = map[string]string{}
|
|
}
|
|
if e.Result == nil {
|
|
e.Result = map[string]any{}
|
|
}
|
|
return e, nil
|
|
}
|
|
|
|
func (s *Store) UpdateExecution(e *domain.Execution) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
result, _ := json.Marshal(e.Result)
|
|
|
|
_, err := s.db.Exec(
|
|
`UPDATE executions SET status=?, result=?, error=?, updated_at=? WHERE id=?`,
|
|
string(e.Status), string(result), e.Error, formatTime(e.UpdatedAt), e.ID,
|
|
)
|
|
return err
|
|
}
|
|
|
|
func (s *Store) GetExecutionByIdempotencyKey(key string) (*domain.Execution, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
|
|
row := s.db.QueryRow(
|
|
`SELECT id, capability_id, target_entity_id, entity_version, arguments, requested_by, origin, idempotency_key, status, result, error, resolution_evidence, correlation_id, created_at, updated_at
|
|
FROM executions WHERE idempotency_key = ?`, key,
|
|
)
|
|
e := &domain.Execution{}
|
|
var args, reqBy, origin, idempKey, status, result, errStr, evidence, corrID, createdAt, updatedAt string
|
|
err := row.Scan(&e.ID, &e.CapabilityID, &e.TargetEntityID, &e.EntityVersion, &args, &reqBy, &origin, &idempKey, &status, &result, &errStr, &evidence, &corrID, &createdAt, &updatedAt)
|
|
if err == sql.ErrNoRows {
|
|
return nil, domain.ErrExecutionNotFound
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
json.Unmarshal([]byte(args), &e.Arguments)
|
|
json.Unmarshal([]byte(reqBy), &e.RequestedBy)
|
|
json.Unmarshal([]byte(origin), &e.Origin)
|
|
json.Unmarshal([]byte(result), &e.Result)
|
|
json.Unmarshal([]byte(evidence), &e.ResolutionEvidence)
|
|
e.IdempotencyKey = idempKey
|
|
e.Status = domain.ExecutionStatus(status)
|
|
e.Error = errStr
|
|
e.CorrelationID = corrID
|
|
e.CreatedAt = parseTime(createdAt)
|
|
e.UpdatedAt = parseTime(updatedAt)
|
|
return e, nil
|
|
}
|
|
|
|
// Event operations
|
|
|
|
func (s *Store) AppendEvent(evt *domain.Event) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
|
|
var maxSeq sql.NullInt64
|
|
s.db.QueryRow(`SELECT MAX(sequence) FROM hexis_events`).Scan(&maxSeq)
|
|
evt.Sequence = maxSeq.Int64 + 1
|
|
|
|
payload, _ := json.Marshal(evt.Payload)
|
|
|
|
_, err := s.db.Exec(
|
|
`INSERT INTO hexis_events (id, sequence, type, timestamp, actor, correlation_id, causation_id, payload)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`,
|
|
evt.ID, evt.Sequence, string(evt.Type), formatTime(evt.Timestamp), evt.Actor, evt.CorrelationID, evt.CausationID, string(payload),
|
|
)
|
|
return err
|
|
}
|
|
|
|
func (s *Store) EventsAfter(seq int64, limit int) ([]*domain.Event, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
|
|
rows, err := s.db.Query(
|
|
`SELECT id, sequence, type, timestamp, COALESCE(actor,''), COALESCE(correlation_id,''), COALESCE(causation_id,''), payload
|
|
FROM hexis_events WHERE sequence > ? ORDER BY sequence LIMIT ?`, seq, limit,
|
|
)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var result []*domain.Event
|
|
for rows.Next() {
|
|
e := &domain.Event{}
|
|
var typ, payload, timestamp string
|
|
if err := rows.Scan(&e.ID, &e.Sequence, &typ, ×tamp, &e.Actor, &e.CorrelationID, &e.CausationID, &payload); err != nil {
|
|
return nil, err
|
|
}
|
|
e.Type = domain.HexisEventType(typ)
|
|
e.Timestamp = parseTime(timestamp)
|
|
json.Unmarshal([]byte(payload), &e.Payload)
|
|
result = append(result, e)
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func (s *Store) LatestSequence() (int64, error) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
|
|
var seq sql.NullInt64
|
|
s.db.QueryRow(`SELECT MAX(sequence) FROM hexis_events`).Scan(&seq)
|
|
return seq.Int64, nil
|
|
}
|
|
|
|
func nullString(s string) interface{} {
|
|
if s == "" {
|
|
return nil
|
|
}
|
|
return s
|
|
}
|
|
|
|
func boolInt(b bool) int {
|
|
if b {
|
|
return 1
|
|
}
|
|
return 0
|
|
}
|
|
|
|
var _ Interface = (*Store)(nil)
|