diff --git a/go.mod b/go.mod index f159fdb..8ea842a 100644 --- a/go.mod +++ b/go.mod @@ -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 ) diff --git a/go.sum b/go.sum index 912390a..6078fd6 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/internal/cli/node.go b/internal/cli/node.go index 518efe4..1dd63e0 100644 --- a/internal/cli/node.go +++ b/internal/cli/node.go @@ -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) diff --git a/internal/engine/registry.go b/internal/engine/registry.go new file mode 100644 index 0000000..63dd52b --- /dev/null +++ b/internal/engine/registry.go @@ -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) +} diff --git a/internal/model/node.go b/internal/model/node.go new file mode 100644 index 0000000..dd10dce --- /dev/null +++ b/internal/model/node.go @@ -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"` +} diff --git a/internal/store/migrate.go b/internal/store/migrate.go new file mode 100644 index 0000000..affe47e --- /dev/null +++ b/internal/store/migrate.go @@ -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 +} diff --git a/internal/store/migrations/0001_nodes.sql b/internal/store/migrations/0001_nodes.sql new file mode 100644 index 0000000..3bdc071 --- /dev/null +++ b/internal/store/migrations/0001_nodes.sql @@ -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); diff --git a/internal/store/node_repo.go b/internal/store/node_repo.go new file mode 100644 index 0000000..82cc4d1 --- /dev/null +++ b/internal/store/node_repo.go @@ -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 +} diff --git a/internal/store/node_repo_test.go b/internal/store/node_repo_test.go new file mode 100644 index 0000000..711f10a --- /dev/null +++ b/internal/store/node_repo_test.go @@ -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) + } +} diff --git a/internal/store/store.go b/internal/store/store.go new file mode 100644 index 0000000..ad9e170 --- /dev/null +++ b/internal/store/store.go @@ -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 +} diff --git a/scripts/trigger_coreci.sh b/scripts/trigger_coreci.sh index 4eb7094..758cf25 100755 --- a/scripts/trigger_coreci.sh +++ b/scripts/trigger_coreci.sh @@ -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:-}"