MCPcopy Create free account
hub / github.com/DoNewsCode/core / TestWriter

Function TestWriter

otkafka/writer_test.go:16–52  ·  view source on GitHub ↗
(t *testing.T)

Source from the content-addressed store, hash-verified

14)
15
16func 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
54func Test_fromWriterConfig(t *testing.T) {
55 writer := fromWriterConfig(WriterConfig{})

Callers

nothing calls this directly

Calls 5

TraceFunction · 0.85
SpanFromMessageFunction · 0.85
WriteMessagesMethod · 0.80
WithLoggerFunction · 0.70
LenMethod · 0.45

Tested by

no test coverage detected