Files
kami dda4acfbb6 Serve capabilities through one serializer and add GET /api/v1/executions
The two "refactor later" items from REVIEW-2026-07-30.md; they share the wire
types, so they land together.

A capability had four divergent wire shapes — the HTTP handler, the MCP
adapter, pkg/client, and Maven's vendored copy of it. There is now a single
definition in pkg/client, mapped from domain by internal/wire and used by the
HTTP list/create/get paths and all four MCP surfaces. It lives in pkg/client
rather than internal so external consumers need not vendor internal/domain,
and so producer and consumer are literally the same type.

The unified shape is a strict superset of all four predecessors; nothing was
dropped. It adds enabled and requires_confirmation to the list responses
(never omitempty — an absent bool reads as unknown, not false), capability_id
to the MCP and client shapes, and the timing/attribute/version fields
previously only on get-by-ID. target_types and the list itself now serialize
as [] rather than null.

Both `id` and `capability_id` are deliberately kept, carrying the same value.
Maven decodes `id`; the spec and the rest of the API say `capability_id`.
Bearer auth is already a breaking change for that consumer, and stacking a
second silent one is the wrong trade — the redundancy stays until every
consumer is confirmed on capability_id, then `id` goes in an announced
removal. A test pins this and says so.

GET /api/v1/executions?entity_id=&since=&limit= implements spec §4.5, which
the Command Center needs. `since` reuses the changes-feed cursor convention
rather than inventing a second paging idiom. That cursor is the row's implicit
SQLite rowid, which is safe only while nothing deletes executions and nothing
VACUUMs — both would renumber and silently invalidate outstanding cursors. If
retention is ever added, this must become an explicit monotonic column first;
the constraint is documented at the query site.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Uea55zaiWuEByEDC4UBSdd
2026-07-30 23:40:20 +04:00

688 lines
24 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)`,
// Retired: schema_migrations was never read — PRAGMA user_version is the
// real mechanism. Migrations are applied by slice index, so this slot must
// keep its position or every later migration would be renumbered and
// skipped on already-migrated databases. Left as a no-op instead.
`SELECT 1`,
`ALTER TABLE capabilities ADD COLUMN requires_confirmation INTEGER NOT NULL DEFAULT 0`,
`ALTER TABLE capabilities ADD COLUMN enabled INTEGER NOT NULL DEFAULT 1`,
`ALTER TABLE capabilities ADD COLUMN timeout_seconds INTEGER NOT NULL DEFAULT 30`,
`ALTER TABLE executions ADD COLUMN confirmation_id TEXT NOT NULL DEFAULT ''`,
`CREATE TABLE IF NOT EXISTS confirmations (
id TEXT PRIMARY KEY,
capability_id TEXT NOT NULL,
capability_version INTEGER NOT NULL,
target_entity_id TEXT NOT NULL,
args_normalized TEXT NOT NULL DEFAULT '{}',
args_hash TEXT NOT NULL,
requester TEXT NOT NULL DEFAULT '',
created_at TEXT NOT NULL,
expires_at TEXT NOT NULL,
state TEXT NOT NULL DEFAULT 'pending'
)`,
`CREATE INDEX IF NOT EXISTS idx_executions_inflight ON executions(capability_id, target_entity_id, status)`,
// Audit columns missing from executions (spec §4.1) plus a partial UNIQUE
// index enforcing one in-flight execution per (capability, target).
// Pre-existing duplicate in-flight rows would make the index creation fail,
// so they are resolved to 'unknown' first — that is the spec's outcome for
// an execution whose side effect can no longer be determined.
`ALTER TABLE executions ADD COLUMN causation_id TEXT NOT NULL DEFAULT '';
ALTER TABLE executions ADD COLUMN capability_version INTEGER NOT NULL DEFAULT 0;
UPDATE executions SET status = 'unknown'
WHERE status = 'started'
AND id NOT IN (
SELECT MIN(id) FROM executions WHERE status = 'started'
GROUP BY capability_id, target_entity_id
);
DROP INDEX IF EXISTS idx_executions_inflight;
CREATE UNIQUE INDEX IF NOT EXISTS idx_executions_inflight
ON executions(capability_id, target_entity_id) WHERE status = 'started';`,
// Supports the entity_id filter on GET /api/v1/executions (spec §4.5). The
// `since` cursor is the row's implicit rowid, which cannot appear in an
// index (SQLite indexes it as the btree key already), so ordering and
// range-scanning by rowid needs no index of its own.
`CREATE INDEX IF NOT EXISTS idx_executions_entity ON executions(target_entity_id)`,
}
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, requires_confirmation, enabled, timeout_seconds, 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, boolInt(c.RequiresConfirmation), boolInt(c.Enabled), c.TimeoutSeconds, 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,''), requires_confirmation, enabled, timeout_seconds, 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, &c.RequiresConfirmation, &c.Enabled, &c.TimeoutSeconds, &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=?, requires_confirmation=?, enabled=?, timeout_seconds=?, 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, boolInt(c.RequiresConfirmation), boolInt(c.Enabled), c.TimeoutSeconds, 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,''), requires_confirmation, enabled, timeout_seconds, 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, &c.RequiresConfirmation, &c.Enabled, &c.TimeoutSeconds, &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, capability_version, target_entity_id, entity_version, arguments, requested_by, origin, idempotency_key, confirmation_id, status, result, error, resolution_evidence, correlation_id, causation_id, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
e.ID, e.CapabilityID, e.CapabilityVersion, e.TargetEntityID, e.EntityVersion, string(args), string(reqBy), string(origin), nullString(e.IdempotencyKey), e.ConfirmationID, string(e.Status), string(result), e.Error, string(evidence), e.CorrelationID, e.CausationID, formatTime(e.CreatedAt), formatTime(e.UpdatedAt),
)
if err != nil {
if strings.Contains(err.Error(), "UNIQUE") {
// The partial unique index on (capability_id, target_entity_id)
// WHERE status='started' is the in-flight guard, not a replay.
if strings.Contains(err.Error(), "idempotency_key") {
return domain.ErrIdempotencyReplay
}
if strings.Contains(err.Error(), "capability_id") || strings.Contains(err.Error(), "target_entity_id") {
return domain.ErrExecutionInFlight
}
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, capability_version, target_entity_id, entity_version, arguments, requested_by, origin, COALESCE(idempotency_key,''), COALESCE(confirmation_id,''), status, result, COALESCE(error,''), resolution_evidence, COALESCE(correlation_id,''), COALESCE(causation_id,''), created_at, updated_at
FROM executions WHERE id = ?`, id,
)
e := &domain.Execution{}
var args, reqBy, origin, idempKey, confID, status, result, errStr, evidence, corrID, causID, createdAt, updatedAt string
err := row.Scan(&e.ID, &e.CapabilityID, &e.CapabilityVersion, &e.TargetEntityID, &e.EntityVersion, &args, &reqBy, &origin, &idempKey, &confID, &status, &result, &errStr, &evidence, &corrID, &causID, &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.ConfirmationID = confID
e.Status = domain.ExecutionStatus(status)
e.Error = errStr
e.CorrelationID = corrID
e.CausationID = causID
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, capability_version, target_entity_id, entity_version, arguments, requested_by, origin, COALESCE(idempotency_key,''), COALESCE(confirmation_id,''), status, result, COALESCE(error,''), resolution_evidence, COALESCE(correlation_id,''), COALESCE(causation_id,''), created_at, updated_at
FROM executions WHERE idempotency_key = ?`, key,
)
e := &domain.Execution{}
var args, reqBy, origin, idempKey, confID, status, result, errStr, evidence, corrID, causID, createdAt, updatedAt string
err := row.Scan(&e.ID, &e.CapabilityID, &e.CapabilityVersion, &e.TargetEntityID, &e.EntityVersion, &args, &reqBy, &origin, &idempKey, &confID, &status, &result, &errStr, &evidence, &corrID, &causID, &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.ConfirmationID = confID
e.Status = domain.ExecutionStatus(status)
e.Error = errStr
e.CorrelationID = corrID
e.CausationID = causID
e.CreatedAt = parseTime(createdAt)
e.UpdatedAt = parseTime(updatedAt)
return e, nil
}
func (s *Store) GetInFlightExecution(capabilityID, targetEntityID string) (*domain.Execution, error) {
s.mu.RLock()
defer s.mu.RUnlock()
row := s.db.QueryRow(
`SELECT id FROM executions WHERE capability_id = ? AND target_entity_id = ? AND status = ? LIMIT 1`,
capabilityID, targetEntityID, string(domain.ExecutionStarted),
)
var id string
err := row.Scan(&id)
if err == sql.ErrNoRows {
return nil, domain.ErrExecutionNotFound
}
if err != nil {
return nil, err
}
return s.getExecutionLocked(id)
}
// getExecutionLocked reads an execution without acquiring s.mu; callers must
// already hold it (read or write).
func (s *Store) getExecutionLocked(id string) (*domain.Execution, error) {
row := s.db.QueryRow(
`SELECT id, capability_id, capability_version, target_entity_id, entity_version, arguments, requested_by, origin, COALESCE(idempotency_key,''), COALESCE(confirmation_id,''), status, result, COALESCE(error,''), resolution_evidence, COALESCE(correlation_id,''), COALESCE(causation_id,''), created_at, updated_at
FROM executions WHERE id = ?`, id,
)
e := &domain.Execution{}
var args, reqBy, origin, idempKey, confID, status, result, errStr, evidence, corrID, causID, createdAt, updatedAt string
err := row.Scan(&e.ID, &e.CapabilityID, &e.CapabilityVersion, &e.TargetEntityID, &e.EntityVersion, &args, &reqBy, &origin, &idempKey, &confID, &status, &result, &errStr, &evidence, &corrID, &causID, &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.ConfirmationID = confID
e.Status = domain.ExecutionStatus(status)
e.Error = errStr
e.CorrelationID = corrID
e.CausationID = causID
e.CreatedAt = parseTime(createdAt)
e.UpdatedAt = parseTime(updatedAt)
return e, nil
}
// executionColumns is the shared SELECT list for execution reads. `rowid` is
// the insertion sequence and doubles as the pagination cursor for
// ListExecutions; single-row reads ignore it.
//
// Using rowid as the cursor is safe here because executions is an ordinary
// rowid table (its PRIMARY KEY is TEXT, so rowid is a separate hidden counter)
// and nothing in this service deletes execution rows or runs VACUUM — both of
// which could renumber rowids and invalidate outstanding cursors. If either
// ever becomes true, promote this to an explicit monotonic seq column.
const executionColumns = `rowid, id, capability_id, capability_version, target_entity_id, entity_version, arguments, requested_by, origin, COALESCE(idempotency_key,''), COALESCE(confirmation_id,''), status, result, COALESCE(error,''), resolution_evidence, COALESCE(correlation_id,''), COALESCE(causation_id,''), created_at, updated_at`
// scanExecution reads one row selected with executionColumns.
func scanExecution(sc interface{ Scan(...any) error }) (*domain.Execution, error) {
e := &domain.Execution{}
var args, reqBy, origin, idempKey, confID, status, result, errStr, evidence, corrID, causID, createdAt, updatedAt string
err := sc.Scan(&e.Seq, &e.ID, &e.CapabilityID, &e.CapabilityVersion, &e.TargetEntityID, &e.EntityVersion, &args, &reqBy, &origin, &idempKey, &confID, &status, &result, &errStr, &evidence, &corrID, &causID, &createdAt, &updatedAt)
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.ConfirmationID = confID
e.Status = domain.ExecutionStatus(status)
e.Error = errStr
e.CorrelationID = corrID
e.CausationID = causID
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
}
// MaxExecutionPageSize bounds a single ListExecutions page, matching the
// 100-row cap the changes feed uses.
const MaxExecutionPageSize = 100
// ListExecutions returns executions ordered by ascending insertion sequence.
//
// entityID, when non-empty, restricts to that target entity. sinceSeq is an
// exclusive cursor: only rows with Seq > sinceSeq are returned, the same
// convention as EventsAfter on /api/v1/changes. limit is clamped to
// MaxExecutionPageSize.
func (s *Store) ListExecutions(entityID string, sinceSeq int64, limit int) ([]*domain.Execution, error) {
s.mu.RLock()
defer s.mu.RUnlock()
if limit <= 0 || limit > MaxExecutionPageSize {
limit = MaxExecutionPageSize
}
query := `SELECT ` + executionColumns + ` FROM executions WHERE rowid > ?`
args := []any{sinceSeq}
if entityID != "" {
query += ` AND target_entity_id = ?`
args = append(args, entityID)
}
query += ` ORDER BY rowid ASC LIMIT ?`
args = append(args, limit)
rows, err := s.db.Query(query, args...)
if err != nil {
return nil, err
}
defer rows.Close()
out := []*domain.Execution{}
for rows.Next() {
e, err := scanExecution(rows)
if err != nil {
return nil, err
}
out = append(out, e)
}
return out, rows.Err()
}
// Confirmation operations
func (s *Store) CreateConfirmation(c *domain.Confirmation) error {
s.mu.Lock()
defer s.mu.Unlock()
_, err := s.db.Exec(
`INSERT INTO confirmations (id, capability_id, capability_version, target_entity_id, args_normalized, args_hash, requester, created_at, expires_at, state)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
c.ID, c.CapabilityID, c.CapabilityVersion, c.TargetEntityID, c.ArgsNormalized, c.ArgsHash, c.Requester, formatTime(c.CreatedAt), formatTime(c.ExpiresAt), string(c.State),
)
return err
}
func (s *Store) GetConfirmation(id string) (*domain.Confirmation, error) {
s.mu.RLock()
defer s.mu.RUnlock()
row := s.db.QueryRow(
`SELECT id, capability_id, capability_version, target_entity_id, args_normalized, args_hash, requester, created_at, expires_at, state
FROM confirmations WHERE id = ?`, id,
)
c := &domain.Confirmation{}
var createdAt, expiresAt, state string
err := row.Scan(&c.ID, &c.CapabilityID, &c.CapabilityVersion, &c.TargetEntityID, &c.ArgsNormalized, &c.ArgsHash, &c.Requester, &createdAt, &expiresAt, &state)
if err == sql.ErrNoRows {
return nil, domain.ErrConfirmationNotFound
}
if err != nil {
return nil, err
}
c.CreatedAt = parseTime(createdAt)
c.ExpiresAt = parseTime(expiresAt)
c.State = domain.ConfirmationState(state)
return c, nil
}
func (s *Store) UpdateConfirmationState(id string, state domain.ConfirmationState) error {
s.mu.Lock()
defer s.mu.Unlock()
_, err := s.db.Exec(`UPDATE confirmations SET state = ? WHERE id = ?`, string(state), id)
return err
}
// 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, &timestamp, &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)