(t *testing.T)
| 14 | ) |
| 15 | |
| 16 | func TestWriter(t *testing.T) { |
| 17 | if os.Getenv("KAFKA_ADDR") == "" { |
| 18 | t.Skip("set KAFKA_ADDR to run TestModule_ProvideRunGroup") |
| 19 | return |
| 20 | } |
| 21 | addrs := strings.Split(os.Getenv("KAFKA_ADDR"), ",") |
| 22 | |
| 23 | { |
| 24 | ctx := context.Background() |
| 25 | kw := kafka.Writer{ |
| 26 | Addr: kafka.TCP(addrs...), |
| 27 | Topic: "trace", |
| 28 | } |
| 29 | tracer := mocktracer.New() |
| 30 | w := Trace(&kw, tracer, WithLogger(log.NewNopLogger())) |
| 31 | span, ctx := opentracing.StartSpanFromContextWithTracer(ctx, tracer, "test") |
| 32 | span.SetBaggageItem("foo", "bar") |
| 33 | err := w.WriteMessages(ctx, kafka.Message{Value: []byte(`hello`)}) |
| 34 | assert.NoError(t, err) |
| 35 | assert.Len(t, tracer.FinishedSpans(), 1) |
| 36 | span.Finish() |
| 37 | } |
| 38 | |
| 39 | { |
| 40 | ctx := context.Background() |
| 41 | kr := kafka.NewReader(kafka.ReaderConfig{Brokers: addrs, Topic: "trace", GroupID: "test", MinBytes: 1, MaxBytes: 1}) |
| 42 | tracer := mocktracer.New() |
| 43 | msg, err := kr.ReadMessage(ctx) |
| 44 | assert.NoError(t, err) |
| 45 | assert.Equal(t, "hello", string(msg.Value)) |
| 46 | span, _, err := SpanFromMessage(ctx, tracer, &msg) |
| 47 | assert.NoError(t, err) |
| 48 | foo := span.BaggageItem("foo") |
| 49 | assert.Equal(t, "bar", foo) |
| 50 | span.Finish() |
| 51 | } |
| 52 | } |
| 53 | |
| 54 | func Test_fromWriterConfig(t *testing.T) { |
| 55 | writer := fromWriterConfig(WriterConfig{}) |
nothing calls this directly
no test coverage detected