From 412472943d82ba0fc81e31aae830a0e1c845b6aa Mon Sep 17 00:00:00 2001 From: Sammy Kerata Oina <44265300+SammyOina@users.noreply.github.com> Date: Tue, 14 Jul 2026 13:51:01 +0300 Subject: [PATCH] feat: add CVM ID to gRPC metadata and implement log-forwarder stream receiving (#611) Signed-off-by: Sammy Oina --- cmd/agent/main.go | 10 +++++++++- cmd/log-forwarder/main.go | 18 +++++++++++++++++- 2 files changed, 26 insertions(+), 2 deletions(-) diff --git a/cmd/agent/main.go b/cmd/agent/main.go index faf20c7d..7ff48504 100644 --- a/cmd/agent/main.go +++ b/cmd/agent/main.go @@ -39,6 +39,7 @@ import ( runnerclient "github.com/ultravioletrs/cocos/pkg/clients/grpc/runner" "github.com/ultravioletrs/cocos/pkg/ingress" "golang.org/x/sync/errgroup" + "google.golang.org/grpc/metadata" ) const ( @@ -170,6 +171,9 @@ func main() { } // Don't defer close here as we want to keep the connection open + if cfg.CVMId != "" { + ctx = metadata.AppendToOutgoingContext(ctx, "job-id", cfg.CVMId) + } pc, err := newClient.Process(ctx) if err != nil { grpcClient.Close() @@ -236,7 +240,11 @@ func main() { } ingressProxy := ingress.NewProxyServer(logger, backendURL, certProvider) - pc, err := cvmsClient.Process(ctx) + agentCtx := ctx + if cfg.CVMId != "" { + agentCtx = metadata.AppendToOutgoingContext(ctx, "job-id", cfg.CVMId) + } + pc, err := cvmsClient.Process(agentCtx) if err != nil { logger.Error(fmt.Sprintf("failed to connect to cvm server: %s", err)) exitCode = 1 diff --git a/cmd/log-forwarder/main.go b/cmd/log-forwarder/main.go index bfc2b19a..dbeb627b 100644 --- a/cmd/log-forwarder/main.go +++ b/cmd/log-forwarder/main.go @@ -20,6 +20,7 @@ import ( cvmsgrpc "github.com/ultravioletrs/cocos/pkg/clients/grpc/cvm" "golang.org/x/sync/errgroup" "google.golang.org/grpc" + "google.golang.org/grpc/metadata" ) const ( @@ -30,6 +31,7 @@ const ( type config struct { LogLevel string `env:"LOG_FORWARDER_LOG_LEVEL" envAlternate:"AGENT_LOG_LEVEL" envDefault:"debug"` + CVMId string `env:"AGENT_CVM_ID" envDefault:""` } func main() { @@ -100,6 +102,9 @@ func main() { defer cvmClient.Close() // Create stream to Manager + if cfg.CVMId != "" { + ctx = metadata.AppendToOutgoingContext(ctx, "job-id", cfg.CVMId, "connection-type", "log-forwarder") + } stream, err := cvmsClient.Process(ctx) if err != nil { logger.Error(fmt.Sprintf("failed to create stream to manager: %s", err)) @@ -122,12 +127,23 @@ func main() { case msg := <-logQueue: if err := stream.Send(msg); err != nil { logger.Error(fmt.Sprintf("failed to send log to manager: %s", err)) - // Reconnect logic would go here + return err } } } }) + // Stream Receiver Goroutine + g.Go(func() error { + for { + _, err := stream.Recv() + if err != nil { + logger.Error(fmt.Sprintf("stream connection lost: %s", err)) + return err + } + } + }) + g.Go(func() error { ch := make(chan os.Signal, 1) signal.Notify(ch, syscall.SIGINT, syscall.SIGTERM)