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
6 changes: 3 additions & 3 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -81,13 +81,13 @@ jobs:
cache-dependency-path: go.sum

- name: Test Runtime B
run: go test -count=1 ./runtime-b/...
run: go test -count=1 ./runtime-b/... ./internal/nacosregistration/...

- name: Race test Runtime B
run: go test -race ./runtime-b/...
run: go test -race ./runtime-b/... ./internal/nacosregistration/...

- name: Vet Runtime B
run: go vet ./runtime-b/...
run: go vet ./runtime-b/... ./internal/nacosregistration/...

images:
runs-on: ubuntu-latest
Expand Down
3 changes: 2 additions & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ module github.com/NeKiro-project/NeKiro-Samples
go 1.26.0

require (
github.com/NeKiro-project/NeKiro v0.0.0-20260804142931-aad73c450435
github.com/NeKiro-project/NeKiro v0.0.0-20260810043416-3e815b89cb37
github.com/NeKiro-project/nekiro-sdk-go v0.0.0-20260804145402-39c37c8929b5
github.com/a2aproject/a2a-go v0.3.15
github.com/golang-jwt/jwt/v5 v5.3.1
Expand All @@ -17,6 +17,7 @@ require (
github.com/go-logr/stdr v1.2.2 // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.22.0 // indirect
github.com/nacos-group/nacos-sdk-go/v2 v2.3.5 // indirect
github.com/santhosh-tekuri/jsonschema/v6 v6.0.2 // indirect
go.opentelemetry.io/auto/sdk v1.1.0 // indirect
go.opentelemetry.io/otel v1.35.0 // indirect
Expand Down
6 changes: 4 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
github.com/Masterminds/semver/v3 v3.5.0 h1:kQceYJfbupGfZOKZQg0kou0DgAKhzDg2NZPAwZ/2OOE=
github.com/Masterminds/semver/v3 v3.5.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM=
github.com/NeKiro-project/NeKiro v0.0.0-20260804142931-aad73c450435 h1:Qwd8L31Yto8HZgDvekdMi53jEMNAt1uwqiXPh8IWkVA=
github.com/NeKiro-project/NeKiro v0.0.0-20260804142931-aad73c450435/go.mod h1:ceTCf86U3pNqSjqJXWgkOhfCv+BbEUicGD1IfY+R8Js=
github.com/NeKiro-project/NeKiro v0.0.0-20260810043416-3e815b89cb37 h1:ai0eN+G2k6rwtF9h1dDhDKzloueH6AyuJp2UOIoyQec=
github.com/NeKiro-project/NeKiro v0.0.0-20260810043416-3e815b89cb37/go.mod h1:JCIEeiLu52WC/Q5QlcAKmWKRtW7CNLkZ3lV3BAn92Oo=
github.com/NeKiro-project/nekiro-sdk-go v0.0.0-20260804145402-39c37c8929b5 h1:wx8nuBluyNb7MhMOURtX/9/Q+FBNcniJqoCMpxBnrwE=
github.com/NeKiro-project/nekiro-sdk-go v0.0.0-20260804145402-39c37c8929b5/go.mod h1:lxAQLsSVXE3lmWoz+SMDKsDOc1b0Vj7okAGNkE7KFRI=
github.com/a2aproject/a2a-go v0.3.15 h1:h5YpCiPq3jxQ5rIns7oDjPag3ivP8u817AzdA4F+NiI=
Expand Down Expand Up @@ -46,6 +46,8 @@ github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
github.com/mattn/go-sqlite3 v1.14.32 h1:JD12Ag3oLy1zQA+BNn74xRgaBbdhbNIDYvQUEuuErjs=
github.com/mattn/go-sqlite3 v1.14.32/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y=
github.com/nacos-group/nacos-sdk-go/v2 v2.3.5 h1:Hux7C4N4rWhwBF5Zm4yyYskrs9VTgrRTA8DZjoEhQTs=
github.com/nacos-group/nacos-sdk-go/v2 v2.3.5/go.mod h1:ygUBdt7eGeYBt6Lz2HO3wx7crKXk25Mp80568emGMWU=
github.com/oasdiff/yaml v0.1.1 h1:6nHx+pn9gBRM6YpBlFZFQGCCd1nuvqOBtTD3KKTgGxY=
github.com/oasdiff/yaml v0.1.1/go.mod h1:EYJNoyktvWMJ0Hmhx+6qTaqMOsalUaRGT8Sj1hNcegU=
github.com/oasdiff/yaml3 v0.0.14 h1:aLJee3hxBK2H5wdXd9iPcIXb93Nty1Ge0pT171eHtkw=
Expand Down
76 changes: 62 additions & 14 deletions internal/nacosregistration/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ import (
"strconv"
"strings"
"time"

"github.com/NeKiro-project/NeKiro/registry"
)

const (
Expand All @@ -23,31 +25,46 @@ var identifierPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$`

type Config struct {
Mode string
AgentID string
InstanceID string
AgentCardVersion string
ReleaseID string
CardDigest string
CanonicalEndpoint string
Audience string
APIOrigin string
NamespaceID string
GroupName string
ServiceName string
ClusterName string
PortName string
AdvertisedIP string
AdvertisedPort int
Weight float64
HeartbeatInterval time.Duration
HeartbeatTimeout time.Duration
IPDeleteTimeout time.Duration
RequestTimeout time.Duration
AuthMode string
AccessToken string
InstanceID string
}

func Load(lookup func(string) (string, bool), prefix, instanceID string) (Config, error) {
if lookup == nil || !identifierPattern.MatchString(instanceID) || prefix != "RUNTIME_A" && prefix != "RUNTIME_B" {
func Load(lookup func(string) (string, bool), prefix, agentID, instanceID string) (Config, error) {
if lookup == nil || !identifierPattern.MatchString(agentID) || !identifierPattern.MatchString(instanceID) || prefix != "RUNTIME_A" && prefix != "RUNTIME_B" {
return Config{}, errorsFor(prefix, "registration dependencies are invalid")
}
name := func(suffix string) string { return prefix + "_" + suffix }
mode, err := required(lookup, name("REGISTRATION_MODE"))
if err != nil {
return Config{}, err
}
config := Config{Mode: mode, InstanceID: instanceID}
nacosSuffixes := []string{"NACOS_API_ORIGIN", "NACOS_NAMESPACE_ID", "NACOS_GROUP_NAME", "NACOS_SERVICE_NAME", "NACOS_CLUSTER_NAME", "NACOS_ADVERTISED_IP", "NACOS_ADVERTISED_PORT", "NACOS_HEARTBEAT_INTERVAL_MS", "NACOS_REQUEST_TIMEOUT_MS", "NACOS_AUTH_MODE", "NACOS_ACCESS_TOKEN"}
config := Config{Mode: mode, AgentID: agentID, InstanceID: instanceID}
nacosSuffixes := []string{
"AGENT_CARD_VERSION", "RELEASE_ID", "CARD_DIGEST", "CANONICAL_ENDPOINT", "AUDIENCE",
"NACOS_API_ORIGIN", "NACOS_NAMESPACE_ID", "NACOS_GROUP_NAME", "NACOS_SERVICE_NAME", "NACOS_CLUSTER_NAME", "NACOS_PORT_NAME",
"NACOS_ADVERTISED_IP", "NACOS_ADVERTISED_PORT", "NACOS_WEIGHT", "NACOS_HEARTBEAT_INTERVAL_MS", "NACOS_HEARTBEAT_TIMEOUT_MS",
"NACOS_IP_DELETE_TIMEOUT_MS", "NACOS_REQUEST_TIMEOUT_MS", "NACOS_AUTH_MODE", "NACOS_ACCESS_TOKEN",
}
if mode == ModeDisabled {
for _, suffix := range nacosSuffixes {
if _, exists := lookup(name(suffix)); exists {
Expand All @@ -59,6 +76,18 @@ func Load(lookup func(string) (string, bool), prefix, instanceID string) (Config
if mode != ModeNacos {
return Config{}, fmt.Errorf("%s is unsupported", name("REGISTRATION_MODE"))
}
for environment, destination := range map[string]*string{
name("AGENT_CARD_VERSION"): &config.AgentCardVersion,
name("RELEASE_ID"): &config.ReleaseID,
name("CARD_DIGEST"): &config.CardDigest,
name("CANONICAL_ENDPOINT"): &config.CanonicalEndpoint,
name("AUDIENCE"): &config.Audience,
} {
*destination, err = required(lookup, environment)
if err != nil {
return Config{}, err
}
}
if config.APIOrigin, err = required(lookup, name("NACOS_API_ORIGIN")); err != nil {
return Config{}, err
}
Expand All @@ -70,6 +99,7 @@ func Load(lookup func(string) (string, bool), prefix, instanceID string) (Config
name("NACOS_GROUP_NAME"): &config.GroupName,
name("NACOS_SERVICE_NAME"): &config.ServiceName,
name("NACOS_CLUSTER_NAME"): &config.ClusterName,
name("NACOS_PORT_NAME"): &config.PortName,
} {
*destination, err = requiredIdentifier(lookup, environment)
if err != nil {
Expand All @@ -86,16 +116,23 @@ func Load(lookup func(string) (string, bool), prefix, instanceID string) (Config
return Config{}, err
}
config.AdvertisedPort = int(port)
heartbeat, err := requiredUnsigned(lookup, name("NACOS_HEARTBEAT_INTERVAL_MS"), minimumMillis, maximumMillis)
weight, err := requiredUnsigned(lookup, name("NACOS_WEIGHT"), 1, 10000)
if err != nil {
return Config{}, err
}
config.HeartbeatInterval = time.Duration(heartbeat) * time.Millisecond
timeout, err := requiredUnsigned(lookup, name("NACOS_REQUEST_TIMEOUT_MS"), minimumMillis, maximumMillis)
if err != nil {
config.Weight = float64(weight)
if config.HeartbeatInterval, err = requiredDuration(lookup, name("NACOS_HEARTBEAT_INTERVAL_MS"), 1000, maximumMillis); err != nil {
return Config{}, err
}
if config.HeartbeatTimeout, err = requiredDuration(lookup, name("NACOS_HEARTBEAT_TIMEOUT_MS"), 1001, 300000); err != nil {
return Config{}, err
}
if config.IPDeleteTimeout, err = requiredDuration(lookup, name("NACOS_IP_DELETE_TIMEOUT_MS"), 1002, 600000); err != nil {
return Config{}, err
}
if config.RequestTimeout, err = requiredDuration(lookup, name("NACOS_REQUEST_TIMEOUT_MS"), minimumMillis, maximumMillis); err != nil {
return Config{}, err
}
config.RequestTimeout = time.Duration(timeout) * time.Millisecond
config.AuthMode, err = required(lookup, name("NACOS_AUTH_MODE"))
if err != nil {
return Config{}, err
Expand All @@ -117,17 +154,23 @@ func Load(lookup func(string) (string, bool), prefix, instanceID string) (Config
}

func (config Config) Validate() error {
if !identifierPattern.MatchString(config.InstanceID) {
return errorsFor("runtime", "instance ID is invalid")
if !identifierPattern.MatchString(config.AgentID) || !identifierPattern.MatchString(config.InstanceID) {
return errorsFor("runtime", "registration identity is invalid")
}
if config.Mode == ModeDisabled {
return nil
}
if config.Mode != ModeNacos || validateOrigin(config.APIOrigin, "Nacos API origin") != nil || !identifierPattern.MatchString(config.NamespaceID) || !identifierPattern.MatchString(config.GroupName) || !identifierPattern.MatchString(config.ServiceName) || !identifierPattern.MatchString(config.ClusterName) {
if config.Mode != ModeNacos || validateOrigin(config.APIOrigin, "Nacos API origin") != nil || !identifierPattern.MatchString(config.NamespaceID) || !identifierPattern.MatchString(config.GroupName) || !identifierPattern.MatchString(config.ServiceName) || !identifierPattern.MatchString(config.ClusterName) || !identifierPattern.MatchString(config.PortName) {
return errorsFor("runtime", "Nacos registration tuple is invalid")
}
if _, err := registry.NewReleaseTarget(registry.ReleaseTargetInput{
AgentID: config.AgentID, AgentCardVersion: config.AgentCardVersion, ReleaseID: config.ReleaseID,
CardDigest: config.CardDigest, CanonicalEndpoint: config.CanonicalEndpoint, Audience: config.Audience,
}); err != nil {
return errorsFor("runtime", "exact Release target is invalid")
}
parsedIP := net.ParseIP(config.AdvertisedIP)
if parsedIP == nil || parsedIP.String() != config.AdvertisedIP || config.AdvertisedPort < 1 || config.AdvertisedPort > 65535 || config.HeartbeatInterval < minimumMillis*time.Millisecond || config.HeartbeatInterval > maximumMillis*time.Millisecond || config.RequestTimeout < minimumMillis*time.Millisecond || config.RequestTimeout > maximumMillis*time.Millisecond {
if parsedIP == nil || parsedIP.String() != config.AdvertisedIP || config.AdvertisedPort < 1 || config.AdvertisedPort > 65535 || config.Weight < 1 || config.Weight > 10000 || config.Weight != float64(int(config.Weight)) || config.HeartbeatInterval < time.Second || config.HeartbeatInterval > time.Minute || config.HeartbeatTimeout <= config.HeartbeatInterval || config.HeartbeatTimeout > 5*time.Minute || config.IPDeleteTimeout <= config.HeartbeatTimeout || config.IPDeleteTimeout > 10*time.Minute || config.RequestTimeout < minimumMillis*time.Millisecond || config.RequestTimeout > maximumMillis*time.Millisecond {
return errorsFor("runtime", "Nacos registration endpoint or timing is invalid")
}
if config.AuthMode != AuthNone && config.AuthMode != AuthAccessToken || config.AuthMode == AuthNone && config.AccessToken != "" || config.AuthMode == AuthAccessToken && strings.TrimSpace(config.AccessToken) == "" {
Expand Down Expand Up @@ -169,6 +212,11 @@ func requiredUnsigned(lookup func(string) (string, bool), name string, minimum,
return parsed, nil
}

func requiredDuration(lookup func(string) (string, bool), name string, minimum, maximum int64) (time.Duration, error) {
value, err := requiredUnsigned(lookup, name, minimum, maximum)
return time.Duration(value) * time.Millisecond, err
}

func validateOrigin(value, name string) error {
parsed, err := url.Parse(value)
if err != nil || parsed.Scheme != "http" && parsed.Scheme != "https" || parsed.Host == "" || parsed.User != nil || parsed.Path != "/nacos" || parsed.RawPath != "" || parsed.RawQuery != "" || parsed.ForceQuery || parsed.Fragment != "" || parsed.RawFragment != "" {
Expand Down
60 changes: 60 additions & 0 deletions internal/nacosregistration/config_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
package nacosregistration

import "testing"

func TestLoadRequiresExactReleaseAndExplicitFreshness(t *testing.T) {
values := validEnvironment()
config, err := Load(mapLookup(values), "RUNTIME_B", "runtime-b", "runtime-b-primary")
if err != nil {
t.Fatal(err)
}
if config.ReleaseID != "rel_runtime_b_1" || config.PortName != "a2a" || config.HeartbeatTimeout.Milliseconds() != 5000 || config.IPDeleteTimeout.Milliseconds() != 10000 {
t.Fatalf("config=%#v", config)
}
for _, name := range []string{
"RUNTIME_B_RELEASE_ID", "RUNTIME_B_CARD_DIGEST", "RUNTIME_B_CANONICAL_ENDPOINT", "RUNTIME_B_AUDIENCE",
"RUNTIME_B_NACOS_PORT_NAME", "RUNTIME_B_NACOS_WEIGHT", "RUNTIME_B_NACOS_HEARTBEAT_TIMEOUT_MS", "RUNTIME_B_NACOS_IP_DELETE_TIMEOUT_MS",
} {
invalid := validEnvironment()
delete(invalid, name)
if _, err := Load(mapLookup(invalid), "RUNTIME_B", "runtime-b", "runtime-b-primary"); err == nil {
t.Errorf("missing %s was accepted", name)
}
}
}

func TestLoadRejectsMismatchedTargetAndFreshnessOrder(t *testing.T) {
for name, mutate := range map[string]func(map[string]string){
"audience": func(values map[string]string) { values["RUNTIME_B_AUDIENCE"] = "http://runtime-a:8091" },
"heartbeat timeout": func(values map[string]string) { values["RUNTIME_B_NACOS_HEARTBEAT_TIMEOUT_MS"] = "1000" },
"delete timeout": func(values map[string]string) { values["RUNTIME_B_NACOS_IP_DELETE_TIMEOUT_MS"] = "5000" },
} {
t.Run(name, func(t *testing.T) {
values := validEnvironment()
mutate(values)
if _, err := Load(mapLookup(values), "RUNTIME_B", "runtime-b", "runtime-b-primary"); err == nil {
t.Fatal("invalid registration config was accepted")
}
})
}
}

func validEnvironment() map[string]string {
return map[string]string{
"RUNTIME_B_REGISTRATION_MODE": "nacos", "RUNTIME_B_AGENT_CARD_VERSION": "1.0.0", "RUNTIME_B_RELEASE_ID": "rel_runtime_b_1",
"RUNTIME_B_CARD_DIGEST": "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef",
"RUNTIME_B_CANONICAL_ENDPOINT": "http://runtime-b:8092/", "RUNTIME_B_AUDIENCE": "http://runtime-b:8092",
"RUNTIME_B_NACOS_API_ORIGIN": "http://nacos:8848/nacos", "RUNTIME_B_NACOS_NAMESPACE_ID": "public",
"RUNTIME_B_NACOS_GROUP_NAME": "NEKIRO", "RUNTIME_B_NACOS_SERVICE_NAME": "runtime-b", "RUNTIME_B_NACOS_CLUSTER_NAME": "DEFAULT",
"RUNTIME_B_NACOS_PORT_NAME": "a2a", "RUNTIME_B_NACOS_ADVERTISED_IP": "127.0.0.1", "RUNTIME_B_NACOS_ADVERTISED_PORT": "8092",
"RUNTIME_B_NACOS_WEIGHT": "1", "RUNTIME_B_NACOS_HEARTBEAT_INTERVAL_MS": "1000", "RUNTIME_B_NACOS_HEARTBEAT_TIMEOUT_MS": "5000",
"RUNTIME_B_NACOS_IP_DELETE_TIMEOUT_MS": "10000", "RUNTIME_B_NACOS_REQUEST_TIMEOUT_MS": "1000", "RUNTIME_B_NACOS_AUTH_MODE": "none",
}
}

func mapLookup(values map[string]string) func(string) (string, bool) {
return func(name string) (string, bool) {
value, ok := values[name]
return value, ok
}
}
Loading
Loading