Skip to content
Merged
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
10 changes: 10 additions & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@ require (
github.com/aws/aws-sdk-go-v2 v1.42.0
github.com/aws/aws-sdk-go-v2/credentials v1.19.25
github.com/aws/aws-sdk-go-v2/service/s3 v1.104.1
github.com/aws/aws-sdk-go-v2/service/sns v1.40.2
github.com/aws/aws-sdk-go-v2/service/sqs v1.44.1
github.com/aws/smithy-go v1.27.1
github.com/compose-spec/compose-go/v2 v2.12.1
github.com/go-playground/validator/v10 v10.30.3
Expand All @@ -21,7 +23,10 @@ require (
github.com/jackc/pgx/v5 v5.10.0
github.com/moby/moby/api v1.54.2
github.com/moby/moby/client v0.4.1
github.com/nats-io/nats.go v1.52.0
github.com/spf13/cobra v1.10.2
github.com/twmb/franz-go v1.21.4
github.com/twmb/franz-go/pkg/kadm v1.18.0
github.com/zalando/go-keyring v0.2.8
golang.org/x/mod v0.37.0
golang.org/x/sync v0.20.0
Expand Down Expand Up @@ -74,6 +79,7 @@ require (
github.com/inconshreveable/mousetrap v1.1.0 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/klauspost/compress v1.18.6 // indirect
github.com/leodido/go-urn v1.4.0 // indirect
github.com/lucasb-eyer/go-colorful v1.3.0 // indirect
github.com/mattn/go-isatty v0.0.20 // indirect
Expand All @@ -86,14 +92,18 @@ require (
github.com/muesli/mango-cobra v1.2.0 // indirect
github.com/muesli/mango-pflag v0.1.0 // indirect
github.com/muesli/roff v0.1.0 // indirect
github.com/nats-io/nkeys v0.4.15 // indirect
github.com/nats-io/nuid v1.0.1 // indirect
github.com/ncruces/go-strftime v1.0.0 // indirect
github.com/opencontainers/go-digest v1.0.0 // indirect
github.com/opencontainers/image-spec v1.1.1 // indirect
github.com/pierrec/lz4/v4 v4.1.26 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
github.com/rivo/uniseg v0.4.7 // indirect
github.com/santhosh-tekuri/jsonschema/v6 v6.0.1 // indirect
github.com/sirupsen/logrus v1.9.3 // indirect
github.com/spf13/pflag v1.0.9 // indirect
github.com/twmb/franz-go/pkg/kmsg v1.13.1 // indirect
github.com/xhit/go-str2duration/v2 v2.1.0 // indirect
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect
go.opentelemetry.io/auto/sdk v1.1.0 // indirect
Expand Down
20 changes: 20 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,10 @@ github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.30 h1:4HbXxyipSYxex
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.30/go.mod h1:G7RP+uhagpKtKhd1BM9N6JQqjCcGEU47K5lBVZQyRQw=
github.com/aws/aws-sdk-go-v2/service/s3 v1.104.1 h1:yb03KevaOAG5e8suo79Af74vjIQvoeKmjl79WQchLrs=
github.com/aws/aws-sdk-go-v2/service/s3 v1.104.1/go.mod h1:mreYODw0Y4yv7xeczvqC6vciwFao8lPE9k1l1ulfY6E=
github.com/aws/aws-sdk-go-v2/service/sns v1.40.2 h1:00dZG/qsR/Uwn5SSF6DKnV2uazRaI8JA7kS+nV4sg30=
github.com/aws/aws-sdk-go-v2/service/sns v1.40.2/go.mod h1:V9szvM64GdG5VJUeDRstvLmt/ozgWiSNg3gYnp3mSkk=
github.com/aws/aws-sdk-go-v2/service/sqs v1.44.1 h1:Iu9jGoXkvMGZJ7kJLruiPkNQ964pvuqYdoClTMIkj8M=
github.com/aws/aws-sdk-go-v2/service/sqs v1.44.1/go.mod h1:d7eKytKiwDFJeNAP4VWo47VCiNM9z539tkoCGJ6PjXw=
github.com/aws/smithy-go v1.27.1 h1:4T340VFndXtADGF52gYa1POyL7s9E4Z1OeZ1hCscIw8=
github.com/aws/smithy-go v1.27.1/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc=
github.com/aymanbagabas/go-udiff v0.4.1 h1:OEIrQ8maEeDBXQDoGCbbTTXYJMYRCRO1fnodZ12Gv5o=
Expand Down Expand Up @@ -147,6 +151,8 @@ github.com/jackc/pgx/v5 v5.10.0 h1:VhSvgU2jSli8o3AqIEOTJr7rZwAEUVo4E4XhR94Zfr0=
github.com/jackc/pgx/v5 v5.10.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4=
github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo=
github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/klauspost/compress v1.18.6 h1:2jupLlAwFm95+YDR+NwD2MEfFO9d4z4Prjl1XXDjuao=
github.com/klauspost/compress v1.18.6/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
github.com/leodido/go-urn v1.4.0 h1:WT9HwE9SGECu3lg4d/dIA+jxlljEa1/ffXKmRjqdmIQ=
github.com/leodido/go-urn v1.4.0/go.mod h1:bvxc+MVxLKB4z00jd1z+Dvzr47oO32F/QSNjSBOlFxI=
github.com/lucasb-eyer/go-colorful v1.3.0 h1:2/yBRLdWBZKrf7gB40FoiKfAWYQ0lqNcbuQwVHXptag=
Expand Down Expand Up @@ -175,12 +181,20 @@ github.com/muesli/mango-pflag v0.1.0 h1:UADqbYgpUyRoBja3g6LUL+3LErjpsOwaC9ywvBWe
github.com/muesli/mango-pflag v0.1.0/go.mod h1:YEQomTxaCUp8PrbhFh10UfbhbQrM/xJ4i2PB8VTLLW0=
github.com/muesli/roff v0.1.0 h1:YD0lalCotmYuF5HhZliKWlIx7IEhiXeSfq7hNjFqGF8=
github.com/muesli/roff v0.1.0/go.mod h1:pjAHQM9hdUUwm/krAfrLGgJkXJ+YuhtsfZ42kieB2Ig=
github.com/nats-io/nats.go v1.52.0 h1:n3avV4VBsCgsdwh71TppsTwtv+QdPs7ntSKM8qJLGsc=
github.com/nats-io/nats.go v1.52.0/go.mod h1:26HypzazeOkyO3/mqd1zZd53STJN0EjCYF9Uy2ZOBno=
github.com/nats-io/nkeys v0.4.15 h1:JACV5jRVO9V856KOapQ7x+EY8Jo3qw1vJt/9Jpwzkk4=
github.com/nats-io/nkeys v0.4.15/go.mod h1:CpMchTXC9fxA5zrMo4KpySxNjiDVvr8ANOSZdiNfUrs=
github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c=
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/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U=
github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM=
github.com/opencontainers/image-spec v1.1.1 h1:y0fUlFfIZhPF1W537XOLg0/fcx6zcHCJwooC2xJA040=
github.com/opencontainers/image-spec v1.1.1/go.mod h1:qpqAh3Dmcf36wStyyWU+kCeDgrGnAve2nCC8+7h8Q0M=
github.com/pierrec/lz4/v4 v4.1.26 h1:GrpZw1gZttORinvzBdXPUXATeqlJjqUG/D87TKMnhjY=
github.com/pierrec/lz4/v4 v4.1.26/go.mod h1:EoQMVJgeeEOMsCqCzqFm2O0cJvljX2nGZjcRIPL34O4=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE=
Expand All @@ -203,6 +217,12 @@ github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UV
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/twmb/franz-go v1.21.4 h1:skglTjGHOHHKxVdUG3A563gynBDhvSWFBBHXKOOMS8M=
github.com/twmb/franz-go v1.21.4/go.mod h1:rfoMTnVk7107fhTGxfEKIHP/e7tPe6oyij/ywzO0czk=
github.com/twmb/franz-go/pkg/kadm v1.18.0 h1:WRf/LZmDdcDXwX7WMbtDU++v+b3NzYh2bCGoPMmzirw=
github.com/twmb/franz-go/pkg/kadm v1.18.0/go.mod h1:XeLhGoLXLFzK8/ryv5FfpxPxGwj4oFEGpPJMB/x6KDE=
github.com/twmb/franz-go/pkg/kmsg v1.13.1 h1:fG5kItwysTk5UXqVwb64EpQEy3TydF3vYYK21nUQ+bI=
github.com/twmb/franz-go/pkg/kmsg v1.13.1/go.mod h1:+DPt4NC8RmI6hqb8G09+3giKObE6uD2Eya6CfqBpeJY=
github.com/xhit/go-str2duration/v2 v2.1.0 h1:lxklc02Drh6ynqX+DdPyp5pCKLUQpRT8bp8Ydu2Bstc=
github.com/xhit/go-str2duration/v2 v2.1.0/go.mod h1:ohY8p+0f07DiV6Em5LKB0s2YpLtXVyJfNt1+BlmyAsU=
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e h1:JVG44RsyaB9T2KIHavMF/ppJZNG9ZpyihvCd0w101no=
Expand Down
129 changes: 129 additions & 0 deletions internal/cli/messaging.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
package cli

import (
"fmt"

"github.com/spf13/cobra"

"github.com/open-source-cloud/devstack/internal/config"
"github.com/open-source-cloud/devstack/internal/orchestrate"
"github.com/open-source-cloud/devstack/internal/state"
)

// This file holds the shared plumbing for the messaging resource groups (spec 29
// §messaging): queue / topic / stream. Each verb group mirrors the db/s3 flow —
// create/rm go through the up-saga's lock → overlay → provisioner → ledger → event
// path via internal/orchestrate; list is a lock-free ledger read. The DECISION is
// NATIVE-by-default: NATS backs queues+streams, Redpanda/Kafka backs topics+kafka
// streams, and LocalStack (SQS/SNS) is opt-in via --engine sqs|sns. A bare verb
// with no --engine infers from the ACTIVE shared engines (never auto-starting one)
// and errors clearly when unsatisfiable.

// msgPrefixed computes the DNS-safe tenant-scoped resource name <project>-<name>
// (spec 29 §tenant naming) — hyphen-separated, valid for NATS streams, Kafka
// topics, SQS/SNS names and Redis keys alike — unless --no-prefix keeps the literal.
func msgPrefixed(project, name string, noPrefix bool) string {
if noPrefix {
return name
}
return project + "-" + name
}

// engineForFlag maps a user-facing --engine value to the registry engine key. The
// LocalStack provisioner is keyed by the template name "localstack" (its
// `provides: aws` endpoint), so sqs/sns/aws all resolve there.
func engineForFlag(flag string) string {
switch flag {
case "sqs", "sns", "aws", "localstack":
return "localstack"
case "nats":
return "nats"
case "kafka":
return "kafka"
case "redis":
return "redis"
default:
return flag
}
}

// inferEngine picks the first engine in the preference order that has a live shared
// instance in the workspace (never auto-starting one). Returns "" if none match.
func inferEngine(d orchestrate.UpDeps, order []string) string {
for _, e := range order {
if _, ok := orchestrate.ResolveInstance(d.Model, e); ok {
return e
}
}
return ""
}

// resolveMsgEngine turns the --engine flag (validated against allowed) plus the
// native-default inference order into a concrete registry engine key, or a clear
// error when the flag is unknown or nothing is inferable.
func resolveMsgEngine(d orchestrate.UpDeps, flag string, allowed []string, order []string, domain string) (string, error) {
if flag != "" {
if !contains(allowed, flag) {
return "", fmt.Errorf("invalid --engine %q for %s (want one of %v)", flag, domain, allowed)
}
engine := engineForFlag(flag)
if _, ok := orchestrate.ResolveInstance(d.Model, engine); !ok {
return "", fmt.Errorf("no shared %s engine (for --engine %s) in this workspace; declare one under workspace.shared and run `devstack up` (devstack never auto-starts an engine)", engine, flag)
}
return engine, nil
}
engine := inferEngine(d, order)
if engine == "" {
return "", fmt.Errorf("no active engine to back a %s in this workspace (tried %v); add one under workspace.shared and run `devstack up`, or pass --engine", domain, order)
}
return engine, nil
}

// listMessagingRows returns the project's provisioned rows of the given kind
// (queue|topic|stream) — a lock-free ledger read, prefix-scoped to the project
// unless --all.
func listMessagingRows(cmd *cobra.Command, kind, project string, all bool) ([]state.Provisioned, string, error) {
mgr, closeFn, err := buildManager(cmd)
if err != nil {
return nil, "", err
}
defer closeFn()
proj := project
if proj == "" {
proj = defaultProjectFromModel(mgr.Model)
}
var rows []state.Provisioned
if proj != "" && !all {
rows, err = mgr.DB.ProvisionedFor(proj)
} else {
rows, err = mgr.DB.AllProvisioned()
}
if err != nil {
return nil, proj, err
}
var out []state.Provisioned
for _, r := range rows {
if r.Kind == kind {
out = append(out, r)
}
}
return out, proj, nil
}

// defaultProjectFromModel returns the workspace's single/first project by name.
func defaultProjectFromModel(m *config.Model) string {
names := sortedProjectNames(m)
if len(names) > 0 {
return names[0]
}
return ""
}

func contains(s []string, v string) bool {
for _, x := range s {
if x == v {
return true
}
}
return false
}
164 changes: 164 additions & 0 deletions internal/cli/messaging_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
package cli

import (
"testing"

"github.com/open-source-cloud/devstack/internal/config"
"github.com/open-source-cloud/devstack/internal/orchestrate"
)

func TestMessagingGroupsRegistered(t *testing.T) {
root := NewRootCmd(Options{})
for _, path := range [][]string{
{"queue", "create"}, {"queue", "list"}, {"queue", "rm"},
{"topic", "create"}, {"topic", "list"}, {"topic", "rm"},
{"stream", "create"}, {"stream", "list"}, {"stream", "rm"},
} {
c, _, err := root.Find(path)
if err != nil || c.RunE == nil {
t.Fatalf("%v not registered as a real command: %v", path, err)
}
}
}

func TestMessagingCreateFlags(t *testing.T) {
root := NewRootCmd(Options{})
cases := map[string][]string{
"queue": {"engine", "fifo", "dlq", "max-receive", "no-prefix", "project"},
"topic": {"engine", "subscribe", "no-prefix", "project"},
"stream": {"engine", "partitions", "replicas", "retention", "no-prefix", "project"},
}
for group, flags := range cases {
c, _, err := root.Find([]string{group, "create"})
if err != nil {
t.Fatalf("%s create: %v", group, err)
}
for _, f := range flags {
if c.Flags().Lookup(f) == nil {
t.Errorf("%s create missing --%s", group, f)
}
}
}
}

// loadModelWS builds a config.Model from a temp workspace with the given shared
// block, so engine-inference can be exercised without a live daemon.
func loadModelWS(t *testing.T, sharedYAML string) *config.Model {
t.Helper()
root := writeWS(t,
"apiVersion: devstack/v1\nkind: Workspace\nname: demo\n"+sharedYAML+"projects:\n - { name: web, path: web }\n",
map[string]string{"web": "apiVersion: devstack/v1\nkind: Project\nname: web\nservices:\n app: { template: node.vite }\n"},
)
m, err := config.LoadAt(root)
if err != nil {
t.Fatalf("load: %v", err)
}
return m
}

func TestMessagingEngineInference(t *testing.T) {
tests := []struct {
name string
sharedYAML string
allowed []string
order []string
flag string
wantEngine string
wantErr bool
}{
{
name: "queue infers nats when present (native default)",
sharedYAML: "shared:\n nats: { template: nats }\n redis: { template: redis }\n",
allowed: queueEngines, order: queueOrder, flag: "", wantEngine: "nats",
},
{
name: "queue falls back to redis when no nats",
sharedYAML: "shared:\n redis: { template: redis }\n",
allowed: queueEngines, order: queueOrder, flag: "", wantEngine: "redis",
},
{
name: "explicit --engine sqs maps to localstack instance",
sharedYAML: "shared:\n localstack: { template: localstack }\n",
allowed: queueEngines, order: queueOrder, flag: "sqs", wantEngine: "localstack",
},
{
name: "stream infers nats over kafka",
sharedYAML: "shared:\n kafka: { template: kafka }\n nats: { template: nats }\n",
allowed: streamEngines, order: streamOrder, flag: "", wantEngine: "nats",
},
{
name: "topic infers kafka (native default)",
sharedYAML: "shared:\n kafka: { template: kafka }\n",
allowed: topicEngines, order: topicOrder, flag: "", wantEngine: "kafka",
},
{
name: "unsatisfiable inference errors (no engine, never auto-starts)",
sharedYAML: "shared:\n postgres: { template: postgres }\n",
allowed: queueEngines, order: queueOrder, flag: "", wantErr: true,
},
{
name: "invalid --engine value errors",
sharedYAML: "shared:\n nats: { template: nats }\n",
allowed: queueEngines, order: queueOrder, flag: "pulsar", wantErr: true,
},
{
name: "explicit engine not in workspace errors (never auto-starts)",
sharedYAML: "shared:\n nats: { template: nats }\n",
allowed: queueEngines, order: queueOrder, flag: "sqs", wantErr: true,
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
d := orchestrate.UpDeps{Model: loadModelWS(t, tc.sharedYAML)}
eng, err := resolveMsgEngine(d, tc.flag, tc.allowed, tc.order, "queue")
if tc.wantErr {
if err == nil {
t.Fatalf("want error, got engine %q", eng)
}
return
}
if err != nil {
t.Fatalf("resolveMsgEngine: %v", err)
}
if eng != tc.wantEngine {
t.Errorf("engine = %q, want %q", eng, tc.wantEngine)
}
})
}
}

func TestEngineForFlagAndPrefix(t *testing.T) {
for flag, want := range map[string]string{
"sqs": "localstack", "sns": "localstack", "nats": "nats", "kafka": "kafka", "redis": "redis",
} {
if got := engineForFlag(flag); got != want {
t.Errorf("engineForFlag(%q) = %q, want %q", flag, got, want)
}
}
if got := msgPrefixed("web", "jobs", false); got != "web-jobs" {
t.Errorf("msgPrefixed = %q, want web-jobs", got)
}
if got := msgPrefixed("web", "external", true); got != "external" {
t.Errorf("msgPrefixed --no-prefix = %q, want external", got)
}
}

// TestValidateStreamFlags asserts the parse-time guard: --partitions/--replicas are
// Kafka-only (rejected for NATS) and --retention must parse (spec 29).
func TestValidateStreamFlags(t *testing.T) {
if err := validateStreamFlags("nats", true, false, ""); err == nil {
t.Error("--partitions with --engine nats must error")
}
if err := validateStreamFlags("nats", false, true, ""); err == nil {
t.Error("--replicas with --engine nats must error")
}
if err := validateStreamFlags("kafka", true, true, "168h"); err != nil {
t.Errorf("kafka with partitions/replicas/retention must be valid: %v", err)
}
if err := validateStreamFlags("nats", false, false, "168h"); err != nil {
t.Errorf("nats with only --retention must be valid: %v", err)
}
if err := validateStreamFlags("nats", false, false, "garbage"); err == nil {
t.Error("invalid --retention must error")
}
}
Loading
Loading