ship(P02): node mgmt merged into milestone

This commit is contained in:
Jon Chery
2026-06-03 12:39:13 +00:00
11 changed files with 618 additions and 12 deletions
+20 -1
View File
@@ -1,10 +1,29 @@
module git.cloudinit.dev/coreci/orca
go 1.25
go 1.25.0
require github.com/spf13/cobra v1.8.1
require (
github.com/agext/levenshtein v1.2.1 // indirect
github.com/apparentlymart/go-textseg/v15 v15.0.0 // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/hashicorp/hcl/v2 v2.24.0 // indirect
github.com/inconshreveable/mousetrap v1.1.0 // indirect
github.com/mattn/go-isatty v0.0.20 // indirect
github.com/mitchellh/go-wordwrap v1.0.1 // indirect
github.com/ncruces/go-strftime v1.0.0 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
github.com/spf13/pflag v1.0.5 // indirect
github.com/zclconf/go-cty v1.16.3 // indirect
golang.org/x/mod v0.33.0 // indirect
golang.org/x/sync v0.20.0 // indirect
golang.org/x/sys v0.42.0 // indirect
golang.org/x/text v0.25.0 // indirect
golang.org/x/tools v0.42.0 // indirect
modernc.org/libc v1.72.3 // indirect
modernc.org/mathutil v1.7.1 // indirect
modernc.org/memory v1.11.0 // indirect
modernc.org/sqlite v1.51.0 // indirect
)
+39
View File
@@ -1,10 +1,49 @@
github.com/agext/levenshtein v1.2.1 h1:QmvMAjj2aEICytGiWzmxoE0x2KZvE0fvmqMOfy2tjT8=
github.com/agext/levenshtein v1.2.1/go.mod h1:JEDfjyjHDjOF/1e4FlBE/PkbqA9OfWu2ki2W0IB5558=
github.com/apparentlymart/go-textseg/v15 v15.0.0 h1:uYvfpb3DyLSCGWnctWKGj857c6ew1u1fNQOlOtuGxQY=
github.com/apparentlymart/go-textseg/v15 v15.0.0/go.mod h1:K8XmNZdhEBkdlyDdvbmmsvpAG721bKi0joRfFdHIWJ4=
github.com/cpuguy83/go-md2man/v2 v2.0.4/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46tRHOmNcaadrF8o=
github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/hashicorp/hcl/v2 v2.24.0 h1:2QJdZ454DSsYGoaE6QheQZjtKZSUs9Nh2izTWiwQxvE=
github.com/hashicorp/hcl/v2 v2.24.0/go.mod h1:oGoO1FIQYfn/AgyOhlg9qLC6/nOJPX3qGbkZpYAcqfM=
github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw=
github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY=
github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y=
github.com/mitchellh/go-wordwrap v1.0.1 h1:TLuKupo69TCn6TQSyGxwI1EblZZEsQ0vMlAFQflz0v0=
github.com/mitchellh/go-wordwrap v1.0.1/go.mod h1:R62XHJLzvMFRBbcrT7m7WgmE1eOyTSsCt+hzestvNj0=
github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w=
github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls=
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE=
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo=
github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
github.com/spf13/cobra v1.8.1 h1:e5/vxKd/rZsfSJMUX1agtjeTDf+qv1/JdBF8gg5k9ZM=
github.com/spf13/cobra v1.8.1/go.mod h1:wHxEcudfqmLYa8iTfL+OuZPbBZkmvliBWKIezN3kD9Y=
github.com/spf13/pflag v1.0.5 h1:iy+VFUOCP1a+8yFto/drg2CJ5u0yRoB7fZw3DKv/JXA=
github.com/spf13/pflag v1.0.5/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg=
github.com/zclconf/go-cty v1.16.3 h1:osr++gw2T61A8KVYHoQiFbFd1Lh3JOCXc/jFLJXKTxk=
github.com/zclconf/go-cty v1.16.3/go.mod h1:VvMs5i0vgZdhYawQNq5kePSpLAoz8u1xvZgrPIxfnZE=
golang.org/x/mod v0.33.0 h1:tHFzIWbBifEmbwtGz65eaWyGiGZatSrT9prnU8DbVL8=
golang.org/x/mod v0.33.0/go.mod h1:swjeQEj+6r7fODbD2cqrnje9PnziFuw4bmLbBZFrQ5w=
golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4=
golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo=
golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/text v0.25.0 h1:qVyWApTSYLk/drJRO5mDlNYskwQznZmkpV2c8q9zls4=
golang.org/x/text v0.25.0/go.mod h1:WEdwpYrmk1qmdHvhkSTNPm3app7v4rsT8F2UD6+VHIA=
golang.org/x/tools v0.42.0 h1:uNgphsn75Tdz5Ji2q36v/nsFSfR/9BRFvqhGBaJGd5k=
golang.org/x/tools v0.42.0/go.mod h1:Ma6lCIwGZvHK6XtgbswSoWroEkhugApmsXyrUmBhfr0=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
modernc.org/libc v1.72.3 h1:ZnDF4tXn4NBXFutMMQC4vtbTFSXhhKzR73fv0beZEAU=
modernc.org/libc v1.72.3/go.mod h1:dn0dZNnnn1clLyvRxLxYExxiKRZIRENOfqQ8XEeg4Qs=
modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU=
modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg=
modernc.org/memory v1.11.0 h1:o4QC8aMQzmcwCK3t3Ux/ZHmwFPzE6hf2Y5LbkRs+hbI=
modernc.org/memory v1.11.0/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw=
modernc.org/sqlite v1.51.0 h1:aH/MMSoayAIhozZ7uJbVTT9QO/VhzBf0J9tymmmuC/U=
modernc.org/sqlite v1.51.0/go.mod h1:tcNzv5p84E0skkmJn038y+hWJbLQXQqEnQfeh5r2JLM=
+141 -11
View File
@@ -1,43 +1,173 @@
package cli
import (
"context"
"database/sql"
"fmt"
"log/slog"
"os"
"path/filepath"
"time"
"github.com/google/uuid"
"github.com/spf13/cobra"
"git.cloudinit.dev/coreci/orca/internal/engine"
"git.cloudinit.dev/coreci/orca/internal/model"
"git.cloudinit.dev/coreci/orca/internal/store"
)
func dbPath() string {
if p := os.Getenv("ORCA_DB"); p != "" {
return p
}
home, _ := os.UserHomeDir()
return filepath.Join(home, ".orca", "orca.db")
}
func openDB() (*sql.DB, func() error, error) {
db, err := store.Open(dbPath())
if err != nil {
return nil, nil, err
}
return db, db.Close, nil
}
func newLogger() *slog.Logger {
return slog.New(slog.NewJSONHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelInfo}))
}
func nodeRegistry() (*engine.NodeRegistry, func() error, error) {
db, closer, err := openDB()
if err != nil {
return nil, nil, err
}
repo := store.NewNodeRepo(db)
return engine.NewNodeRegistry(repo, newLogger()), closer, nil
}
var (
joinName string
joinAddr string
leaveID string
)
var nodeCmd = &cobra.Command{
Use: "node",
Short: "Manage orca nodes",
Long: "Join, leave, or list orca nodes in the cluster.",
Long: "Join, leave, or list orca nodes in the registry.",
}
var nodeJoinCmd = &cobra.Command{
Use: "join",
Short: "Join a node to the orca cluster",
Long: "Register the local node with the orca cluster. Implemented in Phase 2.",
Short: "Join a node to the orca registry",
Long: "Register a node in the local orca registry. Persisted to SQLite.",
RunE: func(cmd *cobra.Command, args []string) error {
return notImplemented("orca node join")
if joinName == "" {
return fmt.Errorf("--name is required")
}
if joinAddr == "" {
joinAddr = "localhost:8443"
}
ctx, cancel := context.WithTimeout(cmd.Context(), 5*time.Second)
defer cancel()
registry, closer, err := nodeRegistry()
if err != nil {
return err
}
defer closer()
node := &model.Node{
ID: uuid.NewString(),
Name: joinName,
Address: joinAddr,
State: model.NodeStateReady,
JoinedAt: time.Now().UTC(),
LastSeen: time.Now().UTC(),
}
if err := registry.Join(ctx, node); err != nil {
return err
}
if jsonOutput {
return printJSON(node)
}
fmt.Fprintf(cmd.OutOrStdout(), "✓ Node joined: %s (%s) at %s\n", node.ID, node.Name, node.Address)
return nil
},
}
var nodeLeaveCmd = &cobra.Command{
Use: "leave",
Short: "Remove a node from the orca cluster",
Long: "Deregister a node from the orca cluster. Implemented in Phase 2.",
Use: "leave [node-id]",
Short: "Remove a node from the orca registry",
Long: "Mark a node as left. Use --id to specify, or pass as argument.",
Args: cobra.MaximumNArgs(1),
RunE: func(cmd *cobra.Command, args []string) error {
return notImplemented("orca node leave")
id := leaveID
if id == "" && len(args) > 0 {
id = args[0]
}
if id == "" {
return fmt.Errorf("node id required (use --id or pass as argument)")
}
ctx, cancel := context.WithTimeout(cmd.Context(), 5*time.Second)
defer cancel()
registry, closer, err := nodeRegistry()
if err != nil {
return err
}
defer closer()
if err := registry.Leave(ctx, id); err != nil {
return err
}
if jsonOutput {
return printJSON(map[string]string{"id": id, "state": "left"})
}
fmt.Fprintf(cmd.OutOrStdout(), "✓ Node left: %s\n", id)
return nil
},
}
var nodeListCmd = &cobra.Command{
Use: "list",
Short: "List all nodes in the orca cluster",
Long: "Display all registered nodes. Implemented in Phase 2.",
Short: "List all nodes in the orca registry",
Long: "Display all registered nodes and their state.",
RunE: func(cmd *cobra.Command, args []string) error {
return notImplemented("orca node list")
ctx, cancel := context.WithTimeout(cmd.Context(), 5*time.Second)
defer cancel()
registry, closer, err := nodeRegistry()
if err != nil {
return err
}
defer closer()
nodes, err := registry.List(ctx)
if err != nil {
return err
}
if jsonOutput {
return printJSON(nodes)
}
if len(nodes) == 0 {
fmt.Fprintln(cmd.OutOrStdout(), "No nodes registered. Use 'orca node join' to add one.")
return nil
}
fmt.Fprintf(cmd.OutOrStdout(), "%-36s %-20s %-22s %-10s\n", "ID", "NAME", "ADDRESS", "STATE")
for _, n := range nodes {
fmt.Fprintf(cmd.OutOrStdout(), "%-36s %-20s %-22s %-10s\n", n.ID, n.Name, n.Address, n.State)
}
return nil
},
}
func init() {
nodeJoinCmd.Flags().StringVar(&joinName, "name", "", "node name (required)")
nodeJoinCmd.Flags().StringVar(&joinAddr, "addr", "", "node address (default localhost:8443)")
nodeLeaveCmd.Flags().StringVar(&leaveID, "id", "", "node id")
nodeCmd.AddCommand(nodeJoinCmd)
nodeCmd.AddCommand(nodeLeaveCmd)
nodeCmd.AddCommand(nodeListCmd)
+57
View File
@@ -0,0 +1,57 @@
package engine
import (
"context"
"fmt"
"log/slog"
"git.cloudinit.dev/coreci/orca/internal/model"
"git.cloudinit.dev/coreci/orca/internal/store"
)
type NodeRegistry struct {
repo *store.NodeRepo
log *slog.Logger
}
func NewNodeRegistry(repo *store.NodeRepo, log *slog.Logger) *NodeRegistry {
if log == nil {
log = slog.Default()
}
return &NodeRegistry{repo: repo, log: log}
}
func (r *NodeRegistry) Join(ctx context.Context, n *model.Node) error {
if err := r.repo.Insert(ctx, n); err != nil {
return fmt.Errorf("join node: %w", err)
}
r.log.Info("node joined",
slog.String("node_id", n.ID),
slog.String("name", n.Name),
slog.String("address", n.Address))
return nil
}
func (r *NodeRegistry) Leave(ctx context.Context, id string) error {
if err := r.repo.UpdateState(ctx, id, model.NodeStateLeft); err != nil {
return fmt.Errorf("leave node: %w", err)
}
r.log.Info("node left", slog.String("node_id", id))
return nil
}
func (r *NodeRegistry) Forget(ctx context.Context, id string) error {
if err := r.repo.Delete(ctx, id); err != nil {
return fmt.Errorf("forget node: %w", err)
}
r.log.Info("node removed from registry", slog.String("node_id", id))
return nil
}
func (r *NodeRegistry) List(ctx context.Context) ([]*model.Node, error) {
return r.repo.List(ctx)
}
func (r *NodeRegistry) Get(ctx context.Context, id string) (*model.Node, error) {
return r.repo.Get(ctx, id)
}
+21
View File
@@ -0,0 +1,21 @@
package model
import "time"
type NodeState string
const (
NodeStatePending NodeState = "pending"
NodeStateReady NodeState = "ready"
NodeStateLeft NodeState = "left"
)
type Node struct {
ID string `json:"id"`
Name string `json:"name"`
Address string `json:"address"`
State NodeState `json:"state"`
JoinedAt time.Time `json:"joined_at"`
LastSeen time.Time `json:"last_seen"`
Metadata map[string]string `json:"metadata,omitempty"`
}
+53
View File
@@ -0,0 +1,53 @@
package store
import (
"context"
"database/sql"
"embed"
"fmt"
"sort"
"strings"
)
//go:embed migrations/*.sql
var migrationsFS embed.FS
func migrate(db *sql.DB) error {
entries, err := migrationsFS.ReadDir("migrations")
if err != nil {
return fmt.Errorf("read migrations dir: %w", err)
}
names := make([]string, 0, len(entries))
for _, e := range entries {
if !e.IsDir() && strings.HasSuffix(e.Name(), ".sql") {
names = append(names, e.Name())
}
}
sort.Strings(names)
if _, err := db.ExecContext(context.Background(), `CREATE TABLE IF NOT EXISTS schema_migrations (name TEXT PRIMARY KEY, applied_at DATETIME NOT NULL)`); err != nil {
return fmt.Errorf("create schema_migrations: %w", err)
}
for _, name := range names {
var existing string
err := db.QueryRowContext(context.Background(), `SELECT name FROM schema_migrations WHERE name = ?`, name).Scan(&existing)
if err == nil {
continue
}
if err != sql.ErrNoRows {
return fmt.Errorf("check migration %s: %w", name, err)
}
sqlBytes, err := migrationsFS.ReadFile("migrations/" + name)
if err != nil {
return fmt.Errorf("read migration %s: %w", name, err)
}
if _, err := db.ExecContext(context.Background(), string(sqlBytes)); err != nil {
return fmt.Errorf("apply migration %s: %w", name, err)
}
if _, err := db.ExecContext(context.Background(), `INSERT INTO schema_migrations (name, applied_at) VALUES (?, datetime('now'))`, name); err != nil {
return fmt.Errorf("record migration %s: %w", name, err)
}
}
return nil
}
+13
View File
@@ -0,0 +1,13 @@
-- Node registry
CREATE TABLE IF NOT EXISTS nodes (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
address TEXT NOT NULL,
state TEXT NOT NULL DEFAULT 'pending',
joined_at DATETIME NOT NULL,
last_seen DATETIME NOT NULL,
metadata TEXT
);
CREATE INDEX IF NOT EXISTS idx_nodes_state ON nodes(state);
CREATE INDEX IF NOT EXISTS idx_nodes_name ON nodes(name);
+122
View File
@@ -0,0 +1,122 @@
package store
import (
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"time"
"git.cloudinit.dev/coreci/orca/internal/model"
)
var ErrNotFound = errors.New("not found")
type NodeRepo struct {
db *sql.DB
}
func NewNodeRepo(db *sql.DB) *NodeRepo {
return &NodeRepo{db: db}
}
func (r *NodeRepo) Insert(ctx context.Context, n *model.Node) error {
if n.JoinedAt.IsZero() {
n.JoinedAt = time.Now().UTC()
}
if n.LastSeen.IsZero() {
n.LastSeen = n.JoinedAt
}
if n.State == "" {
n.State = model.NodeStateReady
}
metaJSON, err := json.Marshal(n.Metadata)
if err != nil {
return fmt.Errorf("marshal metadata: %w", err)
}
_, err = r.db.ExecContext(ctx,
`INSERT INTO nodes (id, name, address, state, joined_at, last_seen, metadata) VALUES (?, ?, ?, ?, ?, ?, ?)`,
n.ID, n.Name, n.Address, string(n.State), n.JoinedAt, n.LastSeen, string(metaJSON))
if err != nil {
return fmt.Errorf("insert node: %w", err)
}
return nil
}
func (r *NodeRepo) Get(ctx context.Context, id string) (*model.Node, error) {
row := r.db.QueryRowContext(ctx,
`SELECT id, name, address, state, joined_at, last_seen, metadata FROM nodes WHERE id = ?`, id)
return scanNode(row)
}
func (r *NodeRepo) List(ctx context.Context) ([]*model.Node, error) {
rows, err := r.db.QueryContext(ctx,
`SELECT id, name, address, state, joined_at, last_seen, metadata FROM nodes ORDER BY joined_at ASC`)
if err != nil {
return nil, fmt.Errorf("list nodes: %w", err)
}
defer rows.Close()
var nodes []*model.Node
for rows.Next() {
n, err := scanNode(rows)
if err != nil {
return nil, err
}
nodes = append(nodes, n)
}
return nodes, rows.Err()
}
func (r *NodeRepo) UpdateState(ctx context.Context, id string, state model.NodeState) error {
res, err := r.db.ExecContext(ctx,
`UPDATE nodes SET state = ?, last_seen = ? WHERE id = ?`,
string(state), time.Now().UTC(), id)
if err != nil {
return fmt.Errorf("update node state: %w", err)
}
rows, _ := res.RowsAffected()
if rows == 0 {
return ErrNotFound
}
return nil
}
func (r *NodeRepo) Delete(ctx context.Context, id string) error {
res, err := r.db.ExecContext(ctx, `DELETE FROM nodes WHERE id = ?`, id)
if err != nil {
return fmt.Errorf("delete node: %w", err)
}
rows, _ := res.RowsAffected()
if rows == 0 {
return ErrNotFound
}
return nil
}
type scanner interface {
Scan(dest ...any) error
}
func scanNode(s scanner) (*model.Node, error) {
var (
n model.Node
state string
metaJSON sql.NullString
)
err := s.Scan(&n.ID, &n.Name, &n.Address, &state, &n.JoinedAt, &n.LastSeen, &metaJSON)
if err == sql.ErrNoRows {
return nil, ErrNotFound
}
if err != nil {
return nil, fmt.Errorf("scan node: %w", err)
}
n.State = model.NodeState(state)
if metaJSON.Valid && metaJSON.String != "" {
if err := json.Unmarshal([]byte(metaJSON.String), &n.Metadata); err != nil {
return nil, fmt.Errorf("unmarshal metadata: %w", err)
}
}
return &n, nil
}
+102
View File
@@ -0,0 +1,102 @@
package store
import (
"context"
"path/filepath"
"testing"
"time"
"git.cloudinit.dev/coreci/orca/internal/model"
)
func openTestDB(t *testing.T) (*NodeRepo, func()) {
t.Helper()
path := filepath.Join(t.TempDir(), "test.db")
db, err := Open(path)
if err != nil {
t.Fatalf("open db: %v", err)
}
return NewNodeRepo(db), func() { _ = db.Close() }
}
func TestNodeRepo_InsertAndGet(t *testing.T) {
repo, cleanup := openTestDB(t)
defer cleanup()
ctx := context.Background()
n := &model.Node{
ID: "test-id-1",
Name: "alpha",
Address: "localhost:8443",
State: model.NodeStateReady,
JoinedAt: time.Now().UTC(),
LastSeen: time.Now().UTC(),
}
if err := repo.Insert(ctx, n); err != nil {
t.Fatalf("insert: %v", err)
}
got, err := repo.Get(ctx, "test-id-1")
if err != nil {
t.Fatalf("get: %v", err)
}
if got.Name != "alpha" || got.Address != "localhost:8443" {
t.Errorf("unexpected node: %+v", got)
}
if got.State != model.NodeStateReady {
t.Errorf("expected state ready, got %s", got.State)
}
}
func TestNodeRepo_List(t *testing.T) {
repo, cleanup := openTestDB(t)
defer cleanup()
ctx := context.Background()
for _, name := range []string{"a", "b", "c"} {
_ = repo.Insert(ctx, &model.Node{
ID: name, Name: name, Address: "addr",
JoinedAt: time.Now().UTC(), LastSeen: time.Now().UTC(),
})
}
nodes, err := repo.List(ctx)
if err != nil {
t.Fatalf("list: %v", err)
}
if len(nodes) != 3 {
t.Errorf("expected 3 nodes, got %d", len(nodes))
}
}
func TestNodeRepo_UpdateState(t *testing.T) {
repo, cleanup := openTestDB(t)
defer cleanup()
ctx := context.Background()
_ = repo.Insert(ctx, &model.Node{
ID: "x", Name: "x", Address: "a", JoinedAt: time.Now().UTC(), LastSeen: time.Now().UTC(),
})
if err := repo.UpdateState(ctx, "x", model.NodeStateLeft); err != nil {
t.Fatalf("update: %v", err)
}
got, _ := repo.Get(ctx, "x")
if got.State != model.NodeStateLeft {
t.Errorf("expected left, got %s", got.State)
}
}
func TestNodeRepo_Delete(t *testing.T) {
repo, cleanup := openTestDB(t)
defer cleanup()
ctx := context.Background()
_ = repo.Insert(ctx, &model.Node{
ID: "y", Name: "y", Address: "a", JoinedAt: time.Now().UTC(), LastSeen: time.Now().UTC(),
})
if err := repo.Delete(ctx, "y"); err != nil {
t.Fatalf("delete: %v", err)
}
_, err := repo.Get(ctx, "y")
if err != ErrNotFound {
t.Errorf("expected ErrNotFound, got %v", err)
}
}
+36
View File
@@ -0,0 +1,36 @@
package store
import (
"database/sql"
"fmt"
"os"
"path/filepath"
_ "modernc.org/sqlite"
)
func Open(path string) (*sql.DB, error) {
if path == "" {
home, err := os.UserHomeDir()
if err != nil {
return nil, fmt.Errorf("get home dir: %w", err)
}
path = filepath.Join(home, ".orca", "orca.db")
}
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
return nil, fmt.Errorf("create db dir: %w", err)
}
db, err := sql.Open("sqlite", path+"?_pragma=journal_mode(WAL)&_pragma=foreign_keys(ON)")
if err != nil {
return nil, fmt.Errorf("open sqlite: %w", err)
}
if err := db.Ping(); err != nil {
_ = db.Close()
return nil, fmt.Errorf("ping sqlite: %w", err)
}
if err := migrate(db); err != nil {
_ = db.Close()
return nil, fmt.Errorf("migrate: %w", err)
}
return db, nil
}
+14
View File
@@ -5,6 +5,20 @@
set -uo pipefail
# Source .env if present (provides GITEA_TOKEN)
# Look in repo root, then cwd, then script dir
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
REPO_ROOT="$(cd "$SCRIPT_DIR/.." && pwd)"
for env_file in "$REPO_ROOT/.env" "$PWD/.env" "./.env"; do
if [ -f "$env_file" ]; then
set -a
# shellcheck disable=SC1090
. "$env_file"
set +a
break
fi
done
CORECI_URL="${CORECI_URL:-https://git.cloudinit.dev/coreci/coreci}"
GITEA_TOKEN="${GITEA_TOKEN:-}"