Skip to content

Go/sarama: consumer group ID is recorded as a Kafka topic (false ASYNC_CALLS), and real topics from ProducerMessage{Topic} / ConsumerGroup.Consume are not extracted #2539

Description

@akynte

Version

Version: 0.11.0 (release). Also reproduced on main @ 96c3f41 built from source.

Platform

Linux (x64)

Install channel

GitHub release archive / install.sh / install.ps1

Binary variant

standard

What happened, and what did you expect?

For an idiomatic sarama producer and consumer of the same topic, in two repositories:

  1. False edge. In the consumer, sarama.NewConsumerGroup(brokers, "ledger-group", cfg)
    produces ASYNC_CALLS main -> ledger-group with
    {"callee":"sarama.NewConsumerGroup","url_path":"ledger-group","broker":"kafka"}.
    "ledger-group" is the consumer group ID, not a topic.
  2. Missed consumer topic. g.Consume(ctx, []string{"payments.charged"}, h) produces no
    ASYNC_CALLS edge for payments.charged.
  3. Missed producer topic. p.SendMessage(&sarama.ProducerMessage{Topic: "payments.charged", ...})
    produces no ASYNC_CALLS edge.

Result: cross_async_calls: 0 between producer and consumer. Before investigating, we
had read edge 1 as "consumer detected", which is how the false edge can mislead users.

Expected: producer ASYNC_CALLS to topic payments.charged, consumer ASYNC_CALLS to topic
payments.charged, no edge for the group ID, and one CROSS_ASYNC_CALLS between the
repos.

Reproduction

  1. Code:
// producer/main.go   (go.mod requires github.com/IBM/sarama)
package main

import "github.com/IBM/sarama"

func publishCharged(p sarama.SyncProducer) error {
	_, _, err := p.SendMessage(&sarama.ProducerMessage{Topic: "payments.charged", Value: sarama.StringEncoder("{}")})
	return err
}

func main() {
	p, _ := sarama.NewSyncProducer([]string{"kafka:9092"}, nil)
	_ = publishCharged(p)
}
// consumer/main.go
package main

import (
	"context"

	"github.com/IBM/sarama"
)

type handler struct{}

func (handler) Setup(sarama.ConsumerGroupSession) error   { return nil }
func (handler) Cleanup(sarama.ConsumerGroupSession) error { return nil }
func (handler) ConsumeClaim(s sarama.ConsumerGroupSession, c sarama.ConsumerGroupClaim) error {
	for range c.Messages() {
	}
	return nil
}

func consumeCharged(g sarama.ConsumerGroup) error {
	return g.Consume(context.Background(), []string{"payments.charged"}, handler{})
}

func main() {
	g, _ := sarama.NewConsumerGroup([]string{"kafka:9092"}, "ledger-group", nil)
	_ = consumeCharged(g)
}
  1. Commands:
codebase-memory-mcp cli query_graph '{"project":"<consumer>","query":"MATCH (a)-[r:ASYNC_CALLS]->(b) RETURN a.name, b.name, r"}'
codebase-memory-mcp cli get_graph_schema '{"project":"<producer>"}'
codebase-memory-mcp cli index_repository '{"repo_path":"/path/to/producer","name":"<producer>","mode":"cross-repo-intelligence","target_projects":["<consumer>"]}'
  1. Output: the consumer has one row, main ledger-group {"callee":"sarama.NewConsumerGroup","url_path":"ledger-group",...};
    the producer has no ASYNC_CALLS; and cross_async_calls: 0.

Attached repro.sh, section "Issue 3".

What we observed while debugging (instrumented build of 96c3f41)

  • The extracted first_string_arg values were:
    • sarama.NewConsumerGroup -> "ledger-group", the first string literal among the
      arguments, which is the group ID;
    • g.Consume -> none (the topic is inside a []string{...} literal);
    • p.SendMessage -> none (the topic is a struct field).
  • extract_composite_queue_field() already recovers topics from composite literals, but
    is_queue_topic_field() (internal/cbm/extract_calls.c:2244) only accepts
    QueueUrl/QueueURL/TopicArn/TopicARN/QueueName/TopicName/QueueArn/QueueARN/Destination.
    It does not accept Topic, which is the field used by sarama ProducerMessage, kafka-go
    Message/Writer/ReaderConfig, franz-go kgo.Record and confluent TopicPartition.
    Adding Topic alone made first_string_arg become payments.charged. The edge was
    still not emitted, because p.SendMessage (a method on a sarama-typed variable)
    resolves to nothing and its callee text does not contain a service-pattern keyword.
    This looks like the same external-package classification limit as the separate
    net/http report.
  • sarama.NewConsumerGroup / NewConsumerGroupFromClient take a group ID, not a topic.
    Excluding them, or reading topics from Consume's []string argument instead, would
    remove the false edge.

Full self-contained script (creates throwaway repos in a temp dir, isolated cache): see attached repro.sh.txt, section "Issue N". Run with bash repro.sh.txt

repro.sh.txt

Logs

"== Issue 3: Go sarama topics not extracted; consumer group ID recorded as a topic"
repo i3-producer go.mod <<<'module example.com/producer

go 1.22

require github.com/IBM/sarama v1.45.0'
repo i3-producer main.go <<'GO'
package main

import "github.com/IBM/sarama"

func publishCharged(p sarama.SyncProducer) error {
	_, _, err := p.SendMessage(&sarama.ProducerMessage{Topic: "payments.charged", Value: sarama.StringEncoder("{}")})
	return err
}

func main() {
	p, _ := sarama.NewSyncProducer([]string{"kafka:9092"}, nil)
	_ = publishCharged(p)
}
GO
repo i3-consumer go.mod <<<'module example.com/consumer

go 1.22

require github.com/IBM/sarama v1.45.0'
repo i3-consumer main.go <<'GO'
package main

import (
	"context"

	"github.com/IBM/sarama"
)

type handler struct{}

func (handler) Setup(sarama.ConsumerGroupSession) error   { return nil }
func (handler) Cleanup(sarama.ConsumerGroupSession) error { return nil }
func (handler) ConsumeClaim(s sarama.ConsumerGroupSession, c sarama.ConsumerGroupClaim) error {
	for range c.Messages() {
	}
	return nil
}

func consumeCharged(g sarama.ConsumerGroup) error {
	return g.Consume(context.Background(), []string{"payments.charged"}, handler{})
}

func main() {
	g, _ := sarama.NewConsumerGroup([]string{"kafka:9092"}, "ledger-group", nil)
	_ = consumeCharged(g)
}
GO
commit i3-producer; commit i3-consumer; index i3-producer; index i3-consumer
check "ASYNC_CALLS in producer (topic payments.charged)" "$(count i3-producer ASYNC_CALLS)" 1
echo "  consumer ASYNC_CALLS edges (expected: one to topic payments.charged, none to the group ID):"
rows i3-consumer 'MATCH (a)-[r:ASYNC_CALLS]->(b) RETURN a.name, b.name, r' | sed 's/^/    /'
check "CROSS_ASYNC_CALLS producer -> consumer" "$(cross i3-producer i3-consumer cross_async_calls)" 1

Diagnostics trajectory (memory / performance / leak issues)


Project scale (if relevant)

No response

Confirmations

  • I searched existing issues and this is not a duplicate.
  • My reproduction uses shareable code (a dummy snippet or a public OSS repository), not proprietary code.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    bugSomething isn't workingparsing/qualityGraph extraction bugs, false positives, missing edgeswindowsWindows-specific issues

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions