Skip to content
Draft
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
185 changes: 143 additions & 42 deletions agent.proto
Original file line number Diff line number Diff line change
@@ -1,57 +1,110 @@
syntax = "proto3";

import "google/protobuf/struct.proto";
import "google/protobuf/timestamp.proto";

option go_package = "github.com/formancehq/membership/internal/grpc/generated";

package server;

// Legacy bidirectional streaming agent service.
// Kept alongside AgentService for ~3 months to support old agent binaries that
// haven't yet migrated to the pull model. Sunset target: 2026-09-09. New
// agents must use AgentService. Membership-api keeps Server registered behind
// a feature flag.
service Server {
rpc Join(stream Message) returns (stream Order) {}
}

message ConnectRequest {
string id = 1;
map<string, string> tags = 2;
string baseUrl = 3;
bool production = 4;
// Pull-based agent service. The agent polls for stacks and reports status via unary RPCs.
service AgentService {
// Agent pulls stacks that need syncing, paginated by (updated_at, id) cursor
rpc ListStacks(ListStacksRequest) returns (ListStacksResponse);

// Agent reports observed state back to membership
rpc ReportStackStatus(ReportStackStatusRequest) returns (ReportStackStatusResponse);
rpc ReportStackDeleted(ReportStackDeletedRequest) returns (ReportStackDeletedResponse);
rpc ReportModuleStatus(ReportModuleStatusRequest) returns (ReportModuleStatusResponse);
rpc ReportModuleDeleted(ReportModuleDeletedRequest) returns (ReportModuleDeletedResponse);

// Version management
rpc UpsertVersion(UpsertVersionRequest) returns (UpsertVersionResponse);
rpc DeleteVersion(DeleteVersionRequest) returns (DeleteVersionResponse);

// Agent heartbeat (replaces ping/pong)
rpc Heartbeat(HeartbeatRequest) returns (HeartbeatResponse);

// Graceful disconnect — marks the region as inactive
rpc Disconnect(DisconnectRequest) returns (DisconnectResponse);
}

message Order {
reserved 5;
oneof message {
Connected connected = 1;
Stack existingStack = 2;
DeletedStack deletedStack = 3;
Ping ping = 4;
DisabledStack disabledStack = 6;
EnabledStack enabledStack = 7;
}
map<string, string> metadata = 8;
// ─── ListStacks ───

message ListStacksRequest {
string region_id = 1;
int32 page_size = 2;
string cursor = 3; // opaque, base64-encoded (updated_at, id). Empty = start from beginning.
}

message Message {
oneof message {
StatusChanged statusChanged = 1;
Pong pong = 2;

AddedVersion addedVersion = 3;
DeletedVersion deletedVersion = 4;
UpdatedVersion updatedVersion = 5;

ModuleStatusChanged moduleStatusChanged = 6;
ModuleDeleted moduleDeleted = 7;
message ListStacksResponse {
repeated Stack stacks = 1;
string next_cursor = 2; // empty if last page
bool has_more = 3;
}

DeletedStack stackDeleted = 8;
}
map<string, string> metadata = 9;
// ─── Report messages ───

message ReportStackStatusRequest {
StatusChanged status_changed = 1;
}
message ReportStackStatusResponse {}

message Connected {}
message ReportStackDeletedRequest {
DeletedStack stack_deleted = 1;
}
message ReportStackDeletedResponse {}

message Ping {}
message ReportModuleStatusRequest {
ModuleStatusChanged module_status_changed = 1;
}
message ReportModuleStatusResponse {}

message Pong {}
message ReportModuleDeletedRequest {
ModuleDeleted module_deleted = 1;
}
message ReportModuleDeletedResponse {}

// ─── Version management ───

message UpsertVersionRequest {
string name = 1;
map<string, string> versions = 2;
bool deprecated = 3;
}
message UpsertVersionResponse {}

message DeleteVersionRequest {
string name = 1;
}
message DeleteVersionResponse {}

// ─── Heartbeat ───

message HeartbeatRequest {
string region_id = 1;
string base_url = 2;
repeated string additional_base_urls = 3;
string version = 4;
bool production = 5;
repeated string capabilities = 6;
repeated string modules = 7;
}
message HeartbeatResponse {}

message DisconnectRequest {}
message DisconnectResponse {}

// ─── Shared data types ───

message Stack {
string clusterName = 1;
Expand All @@ -66,6 +119,8 @@ message Stack {
map<string, string> additionalLabels = 10;
map<string, string> additionalAnnotations = 11;
repeated Module modules = 12;
string expectedStatus = 13;
google.protobuf.Timestamp updated_at = 14;
}

message Module {
Expand Down Expand Up @@ -112,14 +167,6 @@ message DeletedStack {
string clusterName = 1;
}

message DisabledStack {
string clusterName = 1;
}

message EnabledStack {
string clusterName = 1;
}

message AuthConfig {
string clientId = 1;
string clientSecret = 2;
Expand All @@ -145,4 +192,58 @@ message UpdatedVersion {

message DeletedVersion {
string name = 1;
}
}

// ─── Legacy streaming protocol messages (Server.Join) ───
// Kept for backwards compatibility with old agents. Remove on 2026-09-09 sunset.

message ConnectRequest {
string id = 1;
map<string, string> tags = 2;
string baseUrl = 3;
bool production = 4;
}

message Order {
reserved 5;
oneof message {
Connected connected = 1;
Stack existingStack = 2;
DeletedStack deletedStack = 3;
Ping ping = 4;
DisabledStack disabledStack = 6;
EnabledStack enabledStack = 7;
}
map<string, string> metadata = 8;
}

message Message {
oneof message {
StatusChanged statusChanged = 1;
Pong pong = 2;

AddedVersion addedVersion = 3;
DeletedVersion deletedVersion = 4;
UpdatedVersion updatedVersion = 5;

ModuleStatusChanged moduleStatusChanged = 6;
ModuleDeleted moduleDeleted = 7;

DeletedStack stackDeleted = 8;
}
map<string, string> metadata = 9;
}

message Connected {}

message Ping {}

message Pong {}

message DisabledStack {
string clusterName = 1;
}

message EnabledStack {
string clusterName = 1;
}
9 changes: 6 additions & 3 deletions cmd/root.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ const (
productionFlag = "production"
outdatedFlag = "outdated"
resyncPeriodFlag = "resync-period"
pollIntervalFlag = "poll-interval"
)

var rootCmd = &cobra.Command{
Expand Down Expand Up @@ -89,6 +90,7 @@ func init() {
rootCmd.Flags().Bool(productionFlag, false, "Is a production agent")
rootCmd.Flags().Bool(outdatedFlag, false, "Set the region as outdated when connecting")
rootCmd.Flags().Duration(resyncPeriodFlag, 5*time.Minute, "Resync period of K8S resources")
rootCmd.Flags().Duration(pollIntervalFlag, 10*time.Second, "Poll interval for membership API")
rootCmd.Flags().BoolP("toggle", "t", false, "Help message for toggle")
}

Expand All @@ -107,13 +109,13 @@ func runAgent(cmd *cobra.Command, _ []string) error {
return errors.New("missing id")
}

credentials, err := createGRPCTransportCredentials(cmd)
creds, err := createGRPCTransportCredentials(cmd)
if err != nil {
return err
}

dialOptions := make([]grpc.DialOption, 0)
dialOptions = append(dialOptions, grpc.WithTransportCredentials(credentials))
dialOptions = append(dialOptions, grpc.WithTransportCredentials(creds))

baseUrlString, _ := cmd.Flags().GetString(baseUrlFlag)
if baseUrlString == "" {
Expand Down Expand Up @@ -152,6 +154,7 @@ func runAgent(cmd *cobra.Command, _ []string) error {
resyncPeriod, _ := cmd.Flags().GetDuration(resyncPeriodFlag)
outdated, _ := cmd.Flags().GetBool(outdatedFlag)
additionalBaseUrls, _ := cmd.Flags().GetStringSlice(additionalBaseUrlsFlag)
pollInterval, _ := cmd.Flags().GetDuration(pollIntervalFlag)

options := []fx.Option{
fx.Supply(restConfig),
Expand All @@ -160,7 +163,6 @@ func runAgent(cmd *cobra.Command, _ []string) error {
return logging.ContextWithLogger(cmd.Context(), l)
}),
internal.NewModule(
service.IsDebug(cmd),
serverAddress,
authenticator,
internal.ClientInfo{
Expand All @@ -171,6 +173,7 @@ func runAgent(cmd *cobra.Command, _ []string) error {
Outdated: outdated,
Version: Version,
}, resyncPeriod,
pollInterval,
dialOptions...,
),
otlp.FXModuleFromFlags(cmd, otlp.WithServiceVersion(Version)),
Expand Down
8 changes: 3 additions & 5 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@ module github.com/formancehq/stack/components/agent
go 1.25.0

require (
github.com/alitto/pond v1.9.2
github.com/formancehq/go-libs/v2 v2.2.4
github.com/formancehq/operator/v3 v3.10.0
github.com/google/uuid v1.6.0
Expand All @@ -14,11 +13,9 @@ require (
github.com/spf13/cobra v1.10.2
github.com/stretchr/testify v1.11.1
github.com/zitadel/oidc/v3 v3.45.5
go.opentelemetry.io/otel v1.43.0
go.opentelemetry.io/otel/trace v1.43.0
go.uber.org/fx v1.24.0
go.uber.org/mock v0.6.0
golang.org/x/oauth2 v0.36.0
golang.org/x/sync v0.20.0
google.golang.org/grpc v1.80.0
google.golang.org/protobuf v1.36.11
k8s.io/apiextensions-apiserver v0.35.3
Expand Down Expand Up @@ -100,13 +97,15 @@ require (
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.67.0 // indirect
go.opentelemetry.io/contrib/propagators/b3 v1.42.0 // indirect
go.opentelemetry.io/otel v1.43.0 // indirect
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.43.0 // indirect
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.42.0 // indirect
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.43.0 // indirect
go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.42.0 // indirect
go.opentelemetry.io/otel/log v0.18.0 // indirect
go.opentelemetry.io/otel/metric v1.43.0 // indirect
go.opentelemetry.io/otel/sdk v1.43.0 // indirect
go.opentelemetry.io/otel/trace v1.43.0 // indirect
go.opentelemetry.io/proto/otlp v1.10.0 // indirect
go.uber.org/dig v1.19.0 // indirect
go.uber.org/multierr v1.11.0 // indirect
Expand All @@ -116,7 +115,6 @@ require (
golang.org/x/crypto v0.52.0 // indirect
golang.org/x/mod v0.35.0 // indirect
golang.org/x/net v0.55.0 // indirect
golang.org/x/sync v0.20.0 // indirect
golang.org/x/sys v0.45.0 // indirect
golang.org/x/term v0.43.0 // indirect
golang.org/x/text v0.37.0 // indirect
Expand Down
4 changes: 0 additions & 4 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,6 @@ github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1
github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM=
github.com/ThreeDotsLabs/watermill v1.5.1 h1:t5xMivyf9tpmU3iozPqyrCZXHvoV1XQDfihas4sV0fY=
github.com/ThreeDotsLabs/watermill v1.5.1/go.mod h1:Uop10dA3VeJWsSvis9qO3vbVY892LARrKAdki6WtXS4=
github.com/alitto/pond v1.9.2 h1:9Qb75z/scEZVCoSU+osVmQ0I0JOeLfdTDafrbcJ8CLs=
github.com/alitto/pond v1.9.2/go.mod h1:xQn3P/sHTYcU/1BR3i86IGIrilcrGC2LiS+E2+CJWsI=
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
github.com/bmatcuk/doublestar/v4 v4.10.0 h1:zU9WiOla1YA122oLM6i4EXvGW62DvKZVxIe6TYWexEs=
Expand Down Expand Up @@ -266,8 +264,6 @@ go.uber.org/fx v1.24.0 h1:wE8mruvpg2kiiL1Vqd0CC+tr0/24XIB10Iwp2lLWzkg=
go.uber.org/fx v1.24.0/go.mod h1:AmDeGyS+ZARGKM4tlH4FY2Jr63VjbEDJHtqXTGP5hbo=
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
go.uber.org/mock v0.6.0 h1:hyF9dfmbgIX5EfOdasqLsWD6xqpNZlXblLB/Dbnwv3Y=
go.uber.org/mock v0.6.0/go.mod h1:KiVJ4BqZJaMj4svdfmHM0AUx4NJYO8ZNpPnZn1Z+BBU=
go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0=
go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y=
go.uber.org/zap v1.27.1 h1:08RqriUEv8+ArZRYSTXy1LeBScaMpVSTBhCeaZYfMYc=
Expand Down
Loading