Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ require (
github.com/skratchdot/open-golang v0.0.0-20200116055534-eef842397966
github.com/spf13/cobra v1.10.2
github.com/stretchr/testify v1.11.1
go.etcd.io/etcd/client/v3 v3.6.5
go.opentelemetry.io/contrib/instrumentation/github.com/gin-gonic/gin/otelgin v0.67.0
go.opentelemetry.io/otel v1.42.0
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.42.0
Expand Down Expand Up @@ -69,6 +70,8 @@ require (
github.com/clipperhouse/uax29/v2 v2.7.0 // indirect
github.com/cloudwego/base64x v0.1.6 // indirect
github.com/containerd/console v1.0.5 // indirect
github.com/coreos/go-semver v0.3.1 // indirect
github.com/coreos/go-systemd/v22 v22.5.0 // indirect
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect
github.com/dgraph-io/ristretto v0.1.1 // indirect
github.com/ebitengine/purego v0.10.0 // indirect
Expand Down Expand Up @@ -134,6 +137,8 @@ require (
github.com/ugorji/go/codec v1.3.1 // indirect
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect
github.com/yusufpapurcu/wmi v1.2.4 // indirect
go.etcd.io/etcd/api/v3 v3.6.5 // indirect
go.etcd.io/etcd/client/pkg/v3 v3.6.5 // indirect
go.mongodb.org/mongo-driver/v2 v2.5.0 // indirect
go.opencensus.io v0.22.5 // indirect
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
Expand Down
11 changes: 11 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,10 @@ github.com/containerd/console v1.0.5/go.mod h1:YynlIjWYF8myEu6sdkwKIvGQq+cOckRm6
github.com/coreos/etcd v3.3.10+incompatible/go.mod h1:uF7uidLiAD3TWHmW31ZFd/JWoc32PjwdhPthX9715RE=
github.com/coreos/go-etcd v2.0.0+incompatible/go.mod h1:Jez6KQU2B/sWsbdaef3ED8NzMklzPG4d5KIOhIy30Tk=
github.com/coreos/go-semver v0.2.0/go.mod h1:nnelYz7RCh+5ahJtPPxZlU+153eP4D4r3EedlOD2RNk=
github.com/coreos/go-semver v0.3.1 h1:yi21YpKnrx1gt5R+la8n5WgS0kCrsPp33dmEyHReZr4=
github.com/coreos/go-semver v0.3.1/go.mod h1:irMmmIw/7yzSRPWryHsK7EYSg09caPQL03VsM8rvUec=
github.com/coreos/go-systemd/v22 v22.5.0 h1:RrqgGjYQKalulkV8NGVIfkXQf6YYmOyiJKk8iXXhfZs=
github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc=
github.com/cpuguy83/go-md2man v1.0.10/go.mod h1:SmD6nW6nTyfqj6ABTjUi3V3JVMnlJmwcJI5acqYI6dE=
github.com/cpuguy83/go-md2man/v2 v2.0.6/go.mod h1:oOW0eioCTA6cOiMLiUPZOpcVxMig6NIQQ7OS05n1F4g=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
Expand Down Expand Up @@ -157,6 +161,7 @@ github.com/goccy/go-json v0.10.5 h1:Fq85nIqj+gXn/S5ahsiTlK3TmC85qgirsdTP/+DeaC4=
github.com/goccy/go-json v0.10.5/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M=
github.com/goccy/go-yaml v1.19.2 h1:PmFC1S6h8ljIz6gMRBopkjP1TVT7xuwrButHID66PoM=
github.com/goccy/go-yaml v1.19.2/go.mod h1:XBurs7gK8ATbW4ZPGKgcbrY1Br56PdM69F7LkFRi1kA=
github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA=
github.com/gofrs/flock v0.13.0 h1:95JolYOvGMqeH31+FC7D2+uULf6mG61mEZ/A8dRYMzw=
github.com/gofrs/flock v0.13.0/go.mod h1:jxeyy9R1auM5S6JYDBhDt+E2TCo7DkratH4Pgi8P+Z0=
github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q=
Expand Down Expand Up @@ -347,6 +352,12 @@ github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9dec
github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY=
github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0=
github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0=
go.etcd.io/etcd/api/v3 v3.6.5 h1:pMMc42276sgR1j1raO/Qv3QI9Af/AuyQUW6CBAWuntA=
go.etcd.io/etcd/api/v3 v3.6.5/go.mod h1:ob0/oWA/UQQlT1BmaEkWQzI0sJ1M0Et0mMpaABxguOQ=
go.etcd.io/etcd/client/pkg/v3 v3.6.5 h1:Duz9fAzIZFhYWgRjp/FgNq2gO1jId9Yae/rLn3RrBP8=
go.etcd.io/etcd/client/pkg/v3 v3.6.5/go.mod h1:8Wx3eGRPiy0qOFMZT/hfvdos+DjEaPxdIDiCDUv/FQk=
go.etcd.io/etcd/client/v3 v3.6.5 h1:yRwZNFBx/35VKHTcLDeO7XVLbCBFbPi+XV4OC3QJf2U=
go.etcd.io/etcd/client/v3 v3.6.5/go.mod h1:ZqwG/7TAFZ0BJ0jXRPoJjKQJtbFo/9NIY8uoFFKcCyo=
go.mongodb.org/mongo-driver/v2 v2.5.0 h1:yXUhImUjjAInNcpTcAlPHiT7bIXhshCTL3jVBkF3xaE=
go.mongodb.org/mongo-driver/v2 v2.5.0/go.mod h1:yOI9kBsufol30iFsl1slpdq1I0eHPzybRWdyYUs8K/0=
go.opencensus.io v0.22.5 h1:dntmOdLpSpHlVqbW5Eay97DelsZHe+55D+xC6i0dDS0=
Expand Down
35 changes: 35 additions & 0 deletions internal/command/controller/run.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"os"
"path/filepath"
"strconv"
"strings"
"time"

configpkg "github.com/cirruslabs/orchard/internal/config"
Expand All @@ -34,6 +35,9 @@ var experimentalRPCV2 bool
var noExperimentalRPCV2 bool
var experimentalPingInterval time.Duration
var experimentalDisableDBCompression bool
var experimentalStoreBackend string
var experimentalEtcdEndpoints string
var experimentalEtcdKeyPrefix string
var workerOfflineTimeout time.Duration
var execSessionRetentionTTL time.Duration
var execSSHConnectionKeepaliveInterval time.Duration
Expand Down Expand Up @@ -85,6 +89,12 @@ func newRunCommand() *cobra.Command {
"smaller than the controller's default 30 second interval")
cmd.Flags().BoolVar(&experimentalDisableDBCompression, "experimental-disable-db-compression", false,
"disable database compression, which might reduce RAM usage in some scenarios")
cmd.Flags().StringVar(&experimentalStoreBackend, "experimental-store", string(controller.StoreBackendBadger),
"store backend to use (badger or etcd)")
cmd.Flags().StringVar(&experimentalEtcdEndpoints, "experimental-etcd-endpoints", "localhost:2379",
"comma-separated etcd endpoints used when --experimental-store=etcd")
cmd.Flags().StringVar(&experimentalEtcdKeyPrefix, "experimental-etcd-key-prefix", "/orchard",
"etcd key prefix used when --experimental-store=etcd")
cmd.Flags().DurationVar(&workerOfflineTimeout, "worker-offline-timeout", 3*time.Minute,
"duration (e.g. 60s or 5m30s) after which a worker is considered offline for the purposes "+
"of scheduling (no new VMs will be scheduled on such worker and already assigned VMs will be "+
Expand Down Expand Up @@ -219,6 +229,17 @@ func runController(cmd *cobra.Command, args []string) (err error) {
controllerOpts = append(controllerOpts, controller.WithDisableDBCompression())
}

switch controller.StoreBackend(experimentalStoreBackend) {
case controller.StoreBackendBadger:
// Default store backend.
case controller.StoreBackendEtcd:
controllerOpts = append(controllerOpts, controller.WithEtcdStore(
splitEtcdEndpoints(experimentalEtcdEndpoints), experimentalEtcdKeyPrefix,
))
default:
return fmt.Errorf("unsupported --experimental-store value %q", experimentalStoreBackend)
}

if execSSHConnectionKeepaliveInterval < 5*time.Second {
return fmt.Errorf("--exec-ssh-connection-keepalive-interval's value cannot be less than 5 seconds")
}
Expand Down Expand Up @@ -246,6 +267,20 @@ func runController(cmd *cobra.Command, args []string) (err error) {
return controllerInstance.Run(cmd.Context())
}

func splitEtcdEndpoints(rawEndpoints string) []string {
var endpoints []string
for _, endpoint := range strings.Split(rawEndpoints, ",") {
endpoint = strings.TrimSpace(endpoint)
if endpoint == "" {
continue
}

endpoints = append(endpoints, endpoint)
}

return endpoints
}

func createBootstrapContext(
controllerAddress string,
controllerCert tls.Certificate,
Expand Down
19 changes: 17 additions & 2 deletions internal/controller/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import (
"github.com/cirruslabs/orchard/internal/controller/sshserver"
storepkg "github.com/cirruslabs/orchard/internal/controller/store"
"github.com/cirruslabs/orchard/internal/controller/store/badger"
etcdstore "github.com/cirruslabs/orchard/internal/controller/store/etcd"
"github.com/cirruslabs/orchard/internal/netconstants"
"github.com/cirruslabs/orchard/internal/opentelemetry"
v1 "github.com/cirruslabs/orchard/pkg/resource/v1"
Expand Down Expand Up @@ -59,6 +60,9 @@ type Controller struct {
execSSHConnectionKeepaliveInterval time.Duration
experimentalRPCV2 bool
disableDBCompression bool
storeBackend StoreBackend
etcdEndpoints []string
etcdKeyPrefix string
pingInterval time.Duration
synthetic bool

Expand Down Expand Up @@ -108,8 +112,7 @@ func New(opts ...Option) (*Controller, error) {
)

// Instantiate the database
store, err := badger.NewBadgerStore(controller.dataDir.DBPath(), controller.disableDBCompression,
controller.logger)
store, err := controller.initStore()
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -198,6 +201,18 @@ func New(opts ...Option) (*Controller, error) {
return controller, nil
}

func (controller *Controller) initStore() (storepkg.Store, error) {
switch controller.storeBackend {
case "", StoreBackendBadger:
return badger.NewBadgerStore(controller.dataDir.DBPath(), controller.disableDBCompression,
controller.logger)
case StoreBackendEtcd:
return etcdstore.NewEtcdStore(controller.etcdEndpoints, controller.etcdKeyPrefix, controller.logger)
default:
return nil, fmt.Errorf("%w: unsupported store backend %q", ErrInitFailed, controller.storeBackend)
}
}

func (controller *Controller) vmsEnsurePlatformDefaults() error {
return controller.store.Update(func(txn storepkg.Transaction) error {
vms, err := txn.ListVMs()
Expand Down
15 changes: 15 additions & 0 deletions internal/controller/option.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,13 @@ import (

type Option func(*Controller)

type StoreBackend string

const (
StoreBackendBadger StoreBackend = "badger"
StoreBackendEtcd StoreBackend = "etcd"
)

func WithDataDir(dataDir *DataDir) Option {
return func(controller *Controller) {
controller.dataDir = dataDir
Expand Down Expand Up @@ -84,6 +91,14 @@ func WithDisableDBCompression() Option {
}
}

func WithEtcdStore(endpoints []string, keyPrefix string) Option {
return func(controller *Controller) {
controller.storeBackend = StoreBackendEtcd
controller.etcdEndpoints = endpoints
controller.etcdKeyPrefix = keyPrefix
}
}

func WithPingInterval(pingInterval time.Duration) Option {
return func(controller *Controller) {
controller.pingInterval = pingInterval
Expand Down
156 changes: 156 additions & 0 deletions internal/controller/store/etcd/etcd_events.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,156 @@
package etcd

import (
"bytes"
"encoding/json"
"fmt"
"path"
"time"

storepkg "github.com/cirruslabs/orchard/internal/controller/store"
"github.com/cirruslabs/orchard/pkg/resource/v1"
clientv3 "go.etcd.io/etcd/client/v3"
)

const SpaceEvents = "/events"

func scopePrefix(scope []string) string {
keyParts := []string{SpaceEvents}
keyParts = append(keyParts, scope...)

return path.Join(keyParts...)
}

func (txn *Transaction) AppendEvents(events []v1.Event, scope ...string) error {
injectionTime := time.Now().UnixNano()

for index, event := range events {
valueBytes, err := json.Marshal(event)
if err != nil {
return err
}

eventUID := fmt.Sprintf("/%d-%d-%06d",
event.Timestamp,
injectionTime,
index,
)
eventKey := txn.store.key(scopePrefix(scope) + eventUID)

txn.puts[eventKey] = string(valueBytes)
delete(txn.deletes, eventKey)
}

return nil
}

func (txn *Transaction) ListEvents(scope ...string) ([]v1.Event, error) {
page, err := txn.ListEventsPage(storepkg.ListOptions{}, scope...)
if err != nil {
return nil, err
}

return page.Items, nil
}

func (txn *Transaction) ListEventsPage(options storepkg.ListOptions, scope ...string) (
storepkg.Page[v1.Event],
error,
) {
var result storepkg.Page[v1.Event]
result.Items = []v1.Event{}

logicalPrefix := scopePrefix(scope)
physicalPrefix := txn.store.keyPrefix(logicalPrefix)
getKey, getOptions := txn.listEventsPageQueryOptions(physicalPrefix, logicalPrefix, options)
response, err := txn.store.client.Get(txn.ctx, getKey, getOptions...)
if err != nil {
return result, mapErr(err)
}
txn.prefixReadRevisions[physicalPrefix] = response.Header.Revision

limit := options.Limit
for _, kv := range response.Kvs {
key := string(kv.Key)
if txn.isDeleted(key) {
continue
}

var event v1.Event
if err := json.Unmarshal(kv.Value, &event); err != nil {
return result, err
}
if limit > 0 && len(result.Items) >= limit {
break
}

result.Items = append(result.Items, event)

if limit > 0 && len(result.Items) == limit && len(response.Kvs) > limit {
result.NextCursor = bytes.TrimPrefix([]byte(key), []byte(physicalPrefix))
}
}

return result, nil
}

func (txn *Transaction) DeleteEvents(scope ...string) error {
physicalPrefix := txn.store.keyPrefix(scopePrefix(scope))
response, err := txn.store.client.Get(txn.ctx, physicalPrefix, clientv3.WithPrefix(), clientv3.WithKeysOnly(),
clientv3.WithLimit(1))
if err != nil {
return mapErr(err)
}
txn.prefixReadRevisions[physicalPrefix] = response.Header.Revision

txn.prefixDeletes[physicalPrefix] = struct{}{}
for key := range txn.puts {
if hasPrefix(key, physicalPrefix) {
delete(txn.puts, key)
}
}

return nil
}

func (txn *Transaction) listEventsPageQueryOptions(
physicalPrefix string,
logicalPrefix string,
options storepkg.ListOptions,
) (string, []clientv3.OpOption) {
rangeEnd := clientv3.GetPrefixRangeEnd(physicalPrefix)
getKey := physicalPrefix
getRangeEnd := rangeEnd
if len(options.Cursor) > 0 {
cursor := eventCursor(physicalPrefix, logicalPrefix, options.Cursor)
if options.Order == storepkg.ListOrderDesc {
getRangeEnd = cursor
} else {
getKey = cursor + "\x00"
}
}

getOptions := []clientv3.OpOption{
clientv3.WithRange(getRangeEnd),
clientv3.WithSort(clientv3.SortByKey, clientv3.SortAscend),
}
if options.Order == storepkg.ListOrderDesc {
getOptions[1] = clientv3.WithSort(clientv3.SortByKey, clientv3.SortDescend)
}
if options.Limit > 0 {
getOptions = append(getOptions, clientv3.WithLimit(int64(options.Limit+1)))
}

return getKey, getOptions
}

func eventCursor(physicalPrefix, logicalPrefix string, cursor []byte) string {
if bytes.HasPrefix(cursor, []byte(physicalPrefix)) {
return string(cursor)
}
if bytes.HasPrefix(cursor, []byte(logicalPrefix)) {
return path.Join(physicalPrefix, string(bytes.TrimPrefix(cursor, []byte(logicalPrefix))))
}

return physicalPrefix + string(cursor)
}
Loading