diff --git a/go.mod b/go.mod index 136887cd55..ceea56a1a3 100644 --- a/go.mod +++ b/go.mod @@ -244,6 +244,7 @@ require ( github.com/opencontainers/selinux v1.13.1 // indirect github.com/pelletier/go-toml/v2 v2.2.4 // indirect github.com/peterbourgon/diskv v2.0.1+incompatible // indirect + github.com/pierrec/lz4/v4 v4.1.22 // indirect github.com/pjbgf/sha1cd v0.5.0 // indirect github.com/planetscale/vtprotobuf v0.6.1-0.20240319094008-0393e58bdf10 // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect @@ -261,6 +262,7 @@ require ( github.com/sagikazarmark/locafero v0.12.0 // indirect github.com/segmentio/asm v1.1.3 // indirect github.com/segmentio/encoding v0.5.3 // indirect + github.com/segmentio/kafka-go v0.4.47 // indirect github.com/sergi/go-diff v1.4.0 // indirect github.com/sirupsen/logrus v1.9.4 // indirect github.com/skeema/knownhosts v1.3.2 // indirect diff --git a/go.sum b/go.sum index 69240a0923..197dd28e01 100644 --- a/go.sum +++ b/go.sum @@ -699,6 +699,7 @@ github.com/kisielk/errcheck v1.1.0/go.mod h1:EZBBE59ingxPouuu3KfxchcWSUPOHkagtvW github.com/kisielk/errcheck v1.2.0/go.mod h1:/BMXB+zMLi60iA8Vv6Ksmxu/1UDYcXs4uQLJ+jE2L00= github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= +github.com/klauspost/compress v1.15.9/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU= github.com/klauspost/compress v1.18.5 h1:/h1gH5Ce+VWNLSWqPzOVn6XBO+vJbCNGvjoaGBFW2IE= github.com/klauspost/compress v1.18.5/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y= @@ -895,6 +896,10 @@ github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0 github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= github.com/peterbourgon/diskv v2.0.1+incompatible h1:UBdAOUP5p4RWqPBg048CAvpKN+vxiaj6gdUUzhl4XmI= github.com/peterbourgon/diskv v2.0.1+incompatible/go.mod h1:uqqh8zWWbv1HBMNONnaR/tNboyR3/BZd58JJSHlUSCU= +github.com/pierrec/lz4 v2.6.1+incompatible h1:9UY3+iC23yxF0UfGaYrGplQ+79Rg+h/q9FV9ix19jjM= +github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= +github.com/pierrec/lz4/v4 v4.1.22 h1:cKFw6uJDK+/gfw5BcDL0JL5aBsAFdsIT18eRtLj7VIU= +github.com/pierrec/lz4/v4 v4.1.22/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= github.com/pjbgf/sha1cd v0.5.0 h1:a+UkboSi1znleCDUNT3M5YxjOnN1fz2FhN48FlwCxs0= github.com/pjbgf/sha1cd v0.5.0/go.mod h1:lhpGlyHLpQZoxMv8HcgXvZEhcGs0PG/vsZnEJ7H0iCM= github.com/pkg/diff v0.0.0-20210226163009-20ebb0f2a09e/go.mod h1:pJLUxLENpZxwdsKMEsNbx1VGcRFpLqf3715MtcvvzbA= @@ -996,6 +1001,8 @@ github.com/segmentio/asm v1.1.3 h1:WM03sfUOENvvKexOLp+pCqgb/WDjsi7EK8gIsICtzhc= github.com/segmentio/asm v1.1.3/go.mod h1:Ld3L4ZXGNcSLRg4JBsZ3//1+f/TjYl0Mzen/DQy1EJg= github.com/segmentio/encoding v0.5.3 h1:OjMgICtcSFuNvQCdwqMCv9Tg7lEOXGwm1J5RPQccx6w= github.com/segmentio/encoding v0.5.3/go.mod h1:HS1ZKa3kSN32ZHVZ7ZLPLXWvOVIiZtyJnO1gPH1sKt0= +github.com/segmentio/kafka-go v0.4.47 h1:IqziR4pA3vrZq7YdRxaT3w1/5fvIH5qpCwstUanQQB0= +github.com/segmentio/kafka-go v0.4.47/go.mod h1:HjF6XbOKh0Pjlkr5GVZxt6CsjjwnmhVOfURM5KMd8qg= github.com/sergi/go-diff v1.0.0/go.mod h1:0CfEIISq7TuYL3j771MWULgwwjU+GofnZX9QAmXWZgo= github.com/sergi/go-diff v1.4.0 h1:n/SP9D5ad1fORl+llWyN+D6qoUETXNZARKjyY2/KVCw= github.com/sergi/go-diff v1.4.0/go.mod h1:A0bzQcvG0E7Rwjx0REVgAGH58e96+X0MeOfepqsbeW4= @@ -1111,6 +1118,9 @@ github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/xanzy/ssh-agent v0.3.3 h1:+/15pJfg/RsTxqYcX6fHqOXZwwMP+2VyYWJeWM2qQFM= github.com/xanzy/ssh-agent v0.3.3/go.mod h1:6dzNDKs0J9rVPHPhaGCukekBHKqfl+L3KghI1Bc68Uw= +github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= +github.com/xdg-go/scram v1.1.2/go.mod h1:RT/sEzTbU5y00aCK8UOx6R7YryM0iF1N2MOmC3kKLN4= +github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM= github.com/xeipuuv/gojsonpointer v0.0.0-20180127040702-4e3ac2762d5f/go.mod h1:N2zxlSyiKSe5eX1tZViRH5QA0qijqEDrYZiPEAiq3wU= github.com/xeipuuv/gojsonpointer v0.0.0-20190905194746-02993c407bfb h1:zGWFAtiMcyryUHoUjUJX0/lt1H2+i2Ka2n+D3DImSNo= github.com/xeipuuv/gojsonpointer v0.0.0-20190905194746-02993c407bfb/go.mod h1:N2zxlSyiKSe5eX1tZViRH5QA0qijqEDrYZiPEAiq3wU= @@ -1531,6 +1541,7 @@ golang.org/x/text v0.3.4/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.5/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= +golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ= golang.org/x/text v0.4.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= golang.org/x/text v0.5.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= golang.org/x/text v0.6.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= diff --git a/pkg/functions/function.go b/pkg/functions/function.go index 86c96321ba..d22d79c9bf 100644 --- a/pkg/functions/function.go +++ b/pkg/functions/function.go @@ -107,7 +107,7 @@ type Function struct { // Invoke defines hints for use when invoking this function. // See Client.Invoke for usage. - Invoke string `yaml:"invoke,omitempty" jsonschema:"enum=http,enum=cloudevent"` + Invoke string `yaml:"invoke,omitempty" jsonschema:"enum=http,enum=cloudevent,enum=kafka"` // Build defines the build properties for a function Build BuildSpec `yaml:"build,omitempty"` diff --git a/pkg/functions/invoke.go b/pkg/functions/invoke.go index 45d08e16ba..deb9005131 100644 --- a/pkg/functions/invoke.go +++ b/pkg/functions/invoke.go @@ -8,12 +8,14 @@ import ( "io" "net/http" "net/url" + "os" + "strings" "time" cloudevents "github.com/cloudevents/sdk-go/v2" cehttp "github.com/cloudevents/sdk-go/v2/protocol/http" - "github.com/google/uuid" + "github.com/segmentio/kafka-go" ) const ( @@ -55,29 +57,13 @@ func NewInvokeMessage() InvokeMessage { // invocation message. Returned is metadata (such as HTTP headers or // CloudEvent fields) and a stringified version of the payload. func invoke(ctx context.Context, c *Client, f Function, target string, m InvokeMessage, verbose bool) (metadata map[string][]string, body string, err error) { - // Get the first available route from 'local', 'remote', a named environment - // or treat target - route, err := invocationRoute(ctx, c, f, target) // choose instance to invoke - if err != nil { - return - } - - // Format" either 'http' or 'cloudevent' + // Format" either 'http', 'cloudevent' or 'kafka' // TODO: discuss if providing a Format on Message should a) update the // function to use the new format if none is defined already (backwards // compatibility fix) or b) always update the function, even if it was already // set. Once decided, codify in a test. format := DefaultInvokeFormat - // RequestType is expected GET or POST - if m.RequestType == "" { - m.RequestType = DefaultInvokeRequestType - } - - if verbose { - fmt.Printf("Invoking '%v' function at %v\n", f.Invoke, route) - } - if f.Invoke != "" { // Prefer the format set during function creation if defined. format = f.Invoke @@ -85,12 +71,35 @@ func invoke(ctx context.Context, c *Client, f Function, target string, m InvokeM if m.Format != "" { // Use the override specified on the message if provided format = m.Format - if verbose { + } + + var route string + if format != "kafka" { + // Get the first available route from 'local', 'remote', a named environment + // or treat target + route, err = invocationRoute(ctx, c, f, target) // choose instance to invoke + if err != nil { + return + } + } + + // RequestType is expected GET or POST + if m.RequestType == "" { + m.RequestType = DefaultInvokeRequestType + } + + if verbose { + if format == "kafka" { + fmt.Printf("Invoking '%v' function\n", f.Invoke) + } else { + fmt.Printf("Invoking '%v' function at %v\n", f.Invoke, route) + } + if m.Format != "" { fmt.Printf("Invoking '%v' function using '%v' format\n", f.Invoke, m.Format) } } - if m.RequestType != "POST" && m.RequestType != "GET" { + if format != "kafka" && m.RequestType != "POST" && m.RequestType != "GET" { err = fmt.Errorf("http request type '%v' not supported, expected GET or POST", m.RequestType) return } @@ -107,6 +116,29 @@ func invoke(ctx context.Context, c *Client, f Function, target string, m InvokeM // This will be used most likely only for very special cases body, err = sendGetEvent(ctx, route, m, c.transport, verbose) } + case "kafka": + // Get brokers and topics from Function configuration + funcEnvs, err := Interpolate(f.Run.Envs) + if err != nil { + return nil, "", err + } + brokers := funcEnvs["KAFKA_BROKERS"] + if brokers == "" { + brokers = os.Getenv("KAFKA_BROKERS") + } + topics := funcEnvs["KAFKA_TOPICS"] + if topics == "" { + topics = os.Getenv("KAFKA_TOPICS") + } + + if brokers == "" { + return nil, "", errors.New("KAFKA_BROKERS not found in func.yaml (run.envs) or environment") + } + if topics == "" { + return nil, "", errors.New("KAFKA_TOPICS not found in func.yaml (run.envs) or environment") + } + + return sendKafka(ctx, brokers, topics, m, verbose) default: err = fmt.Errorf("format '%v' not supported", format) } @@ -290,3 +322,52 @@ func sendHttp(ctx context.Context, route string, m InvokeMessage, t http.RoundTr b, err := io.ReadAll(resp.Body) return resp.Header, string(b), err } + +type kafkaWriter interface { + WriteMessages(ctx context.Context, msgs ...kafka.Message) error + Close() error +} + +var newKafkaWriter = func(brokers []string, topic string) kafkaWriter { + return &kafka.Writer{ + Addr: kafka.TCP(brokers...), + Topic: topic, + } +} + +// sendKafka produces a message to the specified Kafka topics using the brokers +func sendKafka(ctx context.Context, brokers string, topics string, m InvokeMessage, verbose bool) (map[string][]string, string, error) { + brokerList := strings.Split(brokers, ",") + for i, b := range brokerList { + brokerList[i] = strings.TrimSpace(b) + } + + topicList := strings.Split(topics, ",") + for i, t := range topicList { + topicList[i] = strings.TrimSpace(t) + } + + if verbose { + fmt.Printf("Connecting to Kafka brokers: %v\n", brokerList) + } + + for _, topic := range topicList { + if topic == "" { + continue + } + if verbose { + fmt.Printf("Producing test message to topic: %q\n", topic) + } + writer := newKafkaWriter(brokerList, topic) + defer writer.Close() + + err := writer.WriteMessages(ctx, kafka.Message{ + Value: m.Data, + }) + if err != nil { + return nil, "", fmt.Errorf("failed to write message to Kafka topic %q: %w", topic, err) + } + } + + return nil, fmt.Sprintf("Message sent to Kafka topic(s): %s\n", topics), nil +} diff --git a/pkg/functions/invoke_test.go b/pkg/functions/invoke_test.go new file mode 100644 index 0000000000..cc42e8b489 --- /dev/null +++ b/pkg/functions/invoke_test.go @@ -0,0 +1,144 @@ +package functions + +import ( + "context" + "os" + "testing" + + "github.com/segmentio/kafka-go" +) + +type mockKafkaWriter struct { + writeMessagesFunc func(ctx context.Context, msgs ...kafka.Message) error + closeFunc func() error +} + +func (m *mockKafkaWriter) WriteMessages(ctx context.Context, msgs ...kafka.Message) error { + if m.writeMessagesFunc != nil { + return m.writeMessagesFunc(ctx, msgs...) + } + return nil +} + +func (m *mockKafkaWriter) Close() error { + if m.closeFunc != nil { + return m.closeFunc() + } + return nil +} + +func TestInvokeKafka(t *testing.T) { + // Setup test directories/files for a mock function + tempDir, err := os.MkdirTemp("", "func-test") + if err != nil { + t.Fatal(err) + } + defer os.RemoveAll(tempDir) + + client := New() + + // Initialize function + f, err := client.Init(Function{ + Root: tempDir, + Runtime: "go", + Name: "my-kafka-func", + }) + if err != nil { + t.Fatal(err) + } + + // Add KAFKA_BROKERS and KAFKA_TOPICS to function environment + f.Run.Envs.Add("KAFKA_BROKERS", "localhost:9092,localhost:9093") + f.Run.Envs.Add("KAFKA_TOPICS", "test-topic-1,test-topic-2") + f.Invoke = "kafka" + if err := f.Write(); err != nil { + t.Fatal(err) + } + + // Mock newKafkaWriter + var writtenMessages []kafka.Message + var closedCount int + origNewKafkaWriter := newKafkaWriter + t.Cleanup(func() { + newKafkaWriter = origNewKafkaWriter + }) + + newKafkaWriter = func(brokers []string, topic string) kafkaWriter { + return &mockKafkaWriter{ + writeMessagesFunc: func(ctx context.Context, msgs ...kafka.Message) error { + writtenMessages = append(writtenMessages, msgs...) + return nil + }, + closeFunc: func() error { + closedCount++ + return nil + }, + } + } + + m := InvokeMessage{ + Data: []byte("test kafka message"), + } + + metadata, body, err := client.Invoke(context.Background(), tempDir, "", m) + if err != nil { + t.Fatalf("expected nil error, got %v", err) + } + + if metadata != nil { + t.Errorf("expected nil metadata, got %v", metadata) + } + + expectedBody := "Message sent to Kafka topic(s): test-topic-1,test-topic-2\n" + if body != expectedBody { + t.Errorf("expected body %q, got %q", expectedBody, body) + } + + if len(writtenMessages) != 2 { + t.Errorf("expected 2 messages to be written (one for each topic), got %d", len(writtenMessages)) + } + + for _, msg := range writtenMessages { + if string(msg.Value) != "test kafka message" { + t.Errorf("expected message value 'test kafka message', got %q", string(msg.Value)) + } + } + + if closedCount != 2 { + t.Errorf("expected close to be called 2 times, got %d", closedCount) + } +} + +func TestInvokeKafkaMissingConfig(t *testing.T) { + tempDir, err := os.MkdirTemp("", "func-test") + if err != nil { + t.Fatal(err) + } + defer os.RemoveAll(tempDir) + + client := New() + + // Initialize function + f, err := client.Init(Function{ + Root: tempDir, + Runtime: "go", + Name: "my-kafka-func-missing", + }) + if err != nil { + t.Fatal(err) + } + + f.Invoke = "kafka" + if err := f.Write(); err != nil { + t.Fatal(err) + } + + m := InvokeMessage{ + Data: []byte("test kafka message"), + } + + _, _, err = client.Invoke(context.Background(), tempDir, "", m) + if err == nil { + t.Fatal("expected error due to missing Kafka configuration, got nil") + } +} diff --git a/schema/func_yaml-schema.json b/schema/func_yaml-schema.json index 47e2e1b2b2..4783f6a58b 100644 --- a/schema/func_yaml-schema.json +++ b/schema/func_yaml-schema.json @@ -210,7 +210,8 @@ "invoke": { "enum": [ "http", - "cloudevent" + "cloudevent", + "kafka" ], "type": "string", "description": "Invoke defines hints for use when invoking this function.\nSee Client.Invoke for usage."