Implements the foundational layer that all SynapBus features depend on: - internal/storage: SQLite connection manager (WAL mode, busy_timeout, foreign_keys) and embedded migration runner using modernc.org/sqlite - internal/messaging: MessagingService with send, read inbox, claim, mark done/failed, and FTS5 search. SQLite-backed MessageStore with conversation auto-creation and read/unread tracking via inbox_state. - internal/agents: AgentService with register, authenticate (bcrypt), update, deregister, discover by capability. HTTP auth middleware. - internal/mcp: MCP server using mark3labs/mcp-go with 9 registered tools (send_message, read_inbox, claim_messages, mark_done, search_messages, register_agent, discover_agents, update_agent, deregister_agent). SSE transport, health endpoint, connection manager. - internal/trace: Async trace recorder with buffered channel for recording agent actions to SQLite traces table. - cmd/synapbus: Updated main.go wiring storage, migrations, services, MCP server, chi router, and graceful shutdown. All code compiles with CGO_ENABLED=0. Full test suite passes. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
72 lines
1.5 KiB
Go
72 lines
1.5 KiB
Go
// Package storage provides SQLite storage layer for SynapBus.
|
|
package storage
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"fmt"
|
|
"log/slog"
|
|
"os"
|
|
"path/filepath"
|
|
|
|
_ "modernc.org/sqlite"
|
|
)
|
|
|
|
// DB wraps a *sql.DB with SynapBus-specific configuration.
|
|
type DB struct {
|
|
*sql.DB
|
|
}
|
|
|
|
// New opens a SQLite database with WAL mode, busy_timeout, and foreign keys enabled.
|
|
// If dataDir is empty or ":memory:", an in-memory database is used.
|
|
func New(ctx context.Context, dataDir string) (*DB, error) {
|
|
var dsn string
|
|
|
|
if dataDir == "" || dataDir == ":memory:" {
|
|
dsn = ":memory:"
|
|
} else {
|
|
if err := os.MkdirAll(dataDir, 0o755); err != nil {
|
|
return nil, fmt.Errorf("create data directory: %w", err)
|
|
}
|
|
dsn = filepath.Join(dataDir, "synapbus.db")
|
|
}
|
|
|
|
db, err := sql.Open("sqlite", dsn)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("open database: %w", err)
|
|
}
|
|
|
|
// Configure SQLite pragmas
|
|
pragmas := []string{
|
|
"PRAGMA journal_mode=WAL",
|
|
"PRAGMA busy_timeout=5000",
|
|
"PRAGMA foreign_keys=ON",
|
|
}
|
|
|
|
for _, pragma := range pragmas {
|
|
if _, err := db.ExecContext(ctx, pragma); err != nil {
|
|
db.Close()
|
|
return nil, fmt.Errorf("execute %s: %w", pragma, err)
|
|
}
|
|
}
|
|
|
|
// Verify settings
|
|
var journalMode string
|
|
if err := db.QueryRowContext(ctx, "PRAGMA journal_mode").Scan(&journalMode); err != nil {
|
|
db.Close()
|
|
return nil, fmt.Errorf("verify journal_mode: %w", err)
|
|
}
|
|
|
|
slog.Info("database opened",
|
|
"dsn", dsn,
|
|
"journal_mode", journalMode,
|
|
)
|
|
|
|
return &DB{DB: db}, nil
|
|
}
|
|
|
|
// Close closes the database connection.
|
|
func (db *DB) Close() error {
|
|
return db.DB.Close()
|
|
}
|