diff --git a/contrib/cloud.google.com/go/pubsub.v1/admin_test.go b/contrib/cloud.google.com/go/pubsub.v1/admin_test.go new file mode 100644 index 00000000000..4e183b178d6 --- /dev/null +++ b/contrib/cloud.google.com/go/pubsub.v1/admin_test.go @@ -0,0 +1,302 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2025 Datadog, Inc. + +package pubsub + +import ( + "context" + "fmt" + "testing" + "time" + + vkit "cloud.google.com/go/pubsub/apiv1" + pubsubpb "cloud.google.com/go/pubsub/apiv1/pubsubpb" + "cloud.google.com/go/pubsub/pstest" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "google.golang.org/api/iterator" + "google.golang.org/api/option" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + + "github.com/DataDog/dd-trace-go/v2/contrib/cloud.google.com/go/pubsubtrace" + "github.com/DataDog/dd-trace-go/v2/ddtrace/ext" + "github.com/DataDog/dd-trace-go/v2/ddtrace/mocktracer" +) + +const adminProjectID = "project" + +func setupAdmin(t *testing.T) (context.Context, mocktracer.Tracer, *vkit.PublisherClient, *vkit.SubscriberClient, *vkit.SchemaClient) { + mt := mocktracer.Start() + t.Cleanup(mt.Stop) + + srv := pstest.NewServer() + t.Cleanup(func() { assert.NoError(t, srv.Close()) }) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + t.Cleanup(cancel) + + // The admin GAPIC clients issue their RPCs over this connection, so + // installing the interceptor here traces their admin operations. + conn, err := grpc.NewClient(srv.Addr, + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithChainUnaryInterceptor(pubsubtrace.UnaryAdminInterceptorV1()), + ) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, conn.Close()) }) + + pc, err := vkit.NewPublisherClient(ctx, option.WithGRPCConn(conn)) + require.NoError(t, err) + sc, err := vkit.NewSubscriberClient(ctx, option.WithGRPCConn(conn)) + require.NoError(t, err) + schema, err := vkit.NewSchemaClient(ctx, option.WithGRPCConn(conn)) + require.NoError(t, err) + + return ctx, mt, pc, sc, schema +} + +func topicName(id string) string { + return fmt.Sprintf("projects/%s/topics/%s", adminProjectID, id) +} + +func subName(id string) string { + return fmt.Sprintf("projects/%s/subscriptions/%s", adminProjectID, id) +} + +func snapshotName(id string) string { + return fmt.Sprintf("projects/%s/snapshots/%s", adminProjectID, id) +} + +func schemaName(id string) string { + return fmt.Sprintf("projects/%s/schemas/%s", adminProjectID, id) +} + +func projectName() string { + return fmt.Sprintf("projects/%s", adminProjectID) +} + +func drain[T any](t *testing.T, next func() (T, error)) { + t.Helper() + for { + if _, err := next(); err == iterator.Done { + return + } else { + require.NoError(t, err) + } + } +} + +func TestTraceAdminTopicOperations(t *testing.T) { + ctx, mt, pc, _, _ := setupAdmin(t) + + _, err := pc.CreateTopic(ctx, &pubsubpb.Topic{Name: topicName("topic")}) + require.NoError(t, err) + + _, err = pc.GetTopic(ctx, &pubsubpb.GetTopicRequest{Topic: topicName("topic")}) + require.NoError(t, err) + + it := pc.ListTopics(ctx, &pubsubpb.ListTopicsRequest{Project: projectName()}) + drain(t, it.Next) + + err = pc.DeleteTopic(ctx, &pubsubpb.DeleteTopicRequest{Topic: topicName("topic")}) + require.NoError(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 4) + + assertAdminSpan(t, spans[0], "CreateTopic", "CreateTopic "+topicName("topic")) + assertAdminSpan(t, spans[1], "GetTopic", "GetTopic "+topicName("topic")) + assertAdminSpan(t, spans[2], "ListTopics", "ListTopics "+projectName()) + assertAdminSpan(t, spans[3], "DeleteTopic", "DeleteTopic "+topicName("topic")) +} + +func TestTraceAdminSubscriptionOperations(t *testing.T) { + ctx, mt, pc, sc, _ := setupAdmin(t) + + _, err := pc.CreateTopic(ctx, &pubsubpb.Topic{Name: topicName("topic")}) + require.NoError(t, err) + + _, err = sc.CreateSubscription(ctx, &pubsubpb.Subscription{ + Name: subName("sub"), + Topic: topicName("topic"), + }) + require.NoError(t, err) + + _, err = sc.GetSubscription(ctx, &pubsubpb.GetSubscriptionRequest{Subscription: subName("sub")}) + require.NoError(t, err) + + it := sc.ListSubscriptions(ctx, &pubsubpb.ListSubscriptionsRequest{Project: projectName()}) + drain(t, it.Next) + + err = sc.DeleteSubscription(ctx, &pubsubpb.DeleteSubscriptionRequest{Subscription: subName("sub")}) + require.NoError(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 5) + + assertAdminSpan(t, spans[0], "CreateTopic", "CreateTopic "+topicName("topic")) + assertAdminSpan(t, spans[1], "CreateSubscription", "CreateSubscription "+subName("sub")) + assertAdminSpan(t, spans[2], "GetSubscription", "GetSubscription "+subName("sub")) + assertAdminSpan(t, spans[3], "ListSubscriptions", "ListSubscriptions "+projectName()) + assertAdminSpan(t, spans[4], "DeleteSubscription", "DeleteSubscription "+subName("sub")) +} + +func TestTraceAdminSnapshotOperations(t *testing.T) { + ctx, mt, pc, sc, _ := setupAdmin(t) + + _, err := pc.CreateTopic(ctx, &pubsubpb.Topic{Name: topicName("topic")}) + require.NoError(t, err) + _, err = sc.CreateSubscription(ctx, &pubsubpb.Subscription{ + Name: subName("sub"), + Topic: topicName("topic"), + }) + require.NoError(t, err) + + // pstest does not implement snapshots, so these RPCs error — but the + // interceptor must still emit spans with the resolved resource path. + _, err = sc.CreateSnapshot(ctx, &pubsubpb.CreateSnapshotRequest{ + Name: snapshotName("snap"), + Subscription: subName("sub"), + }) + require.Error(t, err) + + _, err = sc.GetSnapshot(ctx, &pubsubpb.GetSnapshotRequest{Snapshot: snapshotName("snap")}) + require.Error(t, err) + + it := sc.ListSnapshots(ctx, &pubsubpb.ListSnapshotsRequest{Project: projectName()}) + _, err = it.Next() + require.Error(t, err) + require.NotEqual(t, iterator.Done, err) + + err = sc.DeleteSnapshot(ctx, &pubsubpb.DeleteSnapshotRequest{Snapshot: snapshotName("snap")}) + require.Error(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 6) + + assertAdminSpan(t, spans[0], "CreateTopic", "CreateTopic "+topicName("topic")) + assertAdminSpan(t, spans[1], "CreateSubscription", "CreateSubscription "+subName("sub")) + assertAdminSpan(t, spans[2], "CreateSnapshot", "CreateSnapshot "+snapshotName("snap")) + assert.NotNil(t, spans[2].Tag(ext.ErrorMsg)) + assertAdminSpan(t, spans[3], "GetSnapshot", "GetSnapshot "+snapshotName("snap")) + assert.NotNil(t, spans[3].Tag(ext.ErrorMsg)) + assertAdminSpan(t, spans[4], "ListSnapshots", "ListSnapshots "+projectName()) + assert.NotNil(t, spans[4].Tag(ext.ErrorMsg)) + assertAdminSpan(t, spans[5], "DeleteSnapshot", "DeleteSnapshot "+snapshotName("snap")) + assert.NotNil(t, spans[5].Tag(ext.ErrorMsg)) +} + +func TestTraceAdminSchemaOperations(t *testing.T) { + ctx, mt, _, _, sc := setupAdmin(t) + + const avroDef = `{"type":"record","name":"Test","fields":[{"name":"f","type":"string"}]}` + _, err := sc.CreateSchema(ctx, &pubsubpb.CreateSchemaRequest{ + Parent: projectName(), + Schema: &pubsubpb.Schema{ + Type: pubsubpb.Schema_AVRO, + Definition: avroDef, + }, + SchemaId: "schema", + }) + require.NoError(t, err) + + _, err = sc.GetSchema(ctx, &pubsubpb.GetSchemaRequest{Name: schemaName("schema")}) + require.NoError(t, err) + + it := sc.ListSchemas(ctx, &pubsubpb.ListSchemasRequest{Parent: projectName()}) + drain(t, it.Next) + + err = sc.DeleteSchema(ctx, &pubsubpb.DeleteSchemaRequest{Name: schemaName("schema")}) + require.NoError(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 4) + + assertAdminSpan(t, spans[0], "CreateSchema", "CreateSchema "+projectName()) + assertAdminSpan(t, spans[1], "GetSchema", "GetSchema "+schemaName("schema")) + assertAdminSpan(t, spans[2], "ListSchemas", "ListSchemas "+projectName()) + assertAdminSpan(t, spans[3], "DeleteSchema", "DeleteSchema "+schemaName("schema")) +} + +func TestTraceAdminError(t *testing.T) { + ctx, mt, pc, _, _ := setupAdmin(t) + + _, err := pc.GetTopic(ctx, &pubsubpb.GetTopicRequest{Topic: topicName("missing")}) + require.Error(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 1) + assert.Equal(t, err.Error(), spans[0].Tag(ext.ErrorMsg)) + assertAdminSpan(t, spans[0], "GetTopic", "GetTopic "+topicName("missing")) +} + +func TestTraceAdminMissingResource(t *testing.T) { + ctx, mt, pc, _, _ := setupAdmin(t) + + // Recognized admin RPCs with an empty resource field must still emit a + // span; TraceAdmin falls back to a method-only resource name. + _, createErr := pc.CreateTopic(ctx, &pubsubpb.Topic{}) + _, getErr := pc.GetTopic(ctx, &pubsubpb.GetTopicRequest{}) + + spans := mt.FinishedSpans() + require.Len(t, spans, 2) + + assert.Equal(t, "gcp.pubsub.request", spans[0].OperationName()) + assert.Equal(t, "CreateTopic", spans[0].Tag(ext.ResourceName)) + assert.Equal(t, "CreateTopic", spans[0].Tag("pubsub.method")) + assert.Nil(t, spans[0].Tag(ext.GCPProjectID)) + if createErr != nil { + assert.Equal(t, createErr.Error(), spans[0].Tag(ext.ErrorMsg)) + } + + assert.Equal(t, "gcp.pubsub.request", spans[1].OperationName()) + assert.Equal(t, "GetTopic", spans[1].Tag(ext.ResourceName)) + assert.Equal(t, "GetTopic", spans[1].Tag("pubsub.method")) + assert.Nil(t, spans[1].Tag(ext.GCPProjectID)) + if getErr != nil { + assert.Equal(t, getErr.Error(), spans[1].Tag(ext.ErrorMsg)) + } +} + +func TestTraceAdminWithService(t *testing.T) { + mt := mocktracer.Start() + defer mt.Stop() + + srv := pstest.NewServer() + defer func() { assert.NoError(t, srv.Close()) }() + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + conn, err := grpc.NewClient(srv.Addr, + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithChainUnaryInterceptor(pubsubtrace.UnaryAdminInterceptorV1(WithService("my-admin-service"))), + ) + require.NoError(t, err) + defer func() { assert.NoError(t, conn.Close()) }() + + pc, err := vkit.NewPublisherClient(ctx, option.WithGRPCConn(conn)) + require.NoError(t, err) + + _, err = pc.CreateTopic(ctx, &pubsubpb.Topic{Name: topicName("topic")}) + require.NoError(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 1) + assert.Equal(t, "my-admin-service", spans[0].Tag(ext.ServiceName)) +} + +func assertAdminSpan(t *testing.T, span *mocktracer.Span, method, resource string) { + t.Helper() + assert.Equal(t, "gcp.pubsub.request", span.OperationName()) + assert.Equal(t, resource, span.Tag(ext.ResourceName)) + assert.Equal(t, ext.SpanTypeWorker, span.Tag(ext.SpanType)) + assert.Equal(t, ext.SpanKindClient, span.Tag(ext.SpanKind)) + assert.Equal(t, "cloud.google.com/go/pubsub.v1", span.Tag(ext.Component)) + assert.Equal(t, ext.MessagingSystemGCPPubsub, span.Tag(ext.MessagingSystem)) + assert.Equal(t, method, span.Tag("pubsub.method")) + assert.Equal(t, adminProjectID, span.Tag(ext.GCPProjectID)) + assert.Equal(t, "cloud.google.com/go/pubsub.v1", span.Integration()) +} diff --git a/contrib/cloud.google.com/go/pubsub.v1/orchestrion.yml b/contrib/cloud.google.com/go/pubsub.v1/orchestrion.yml index 55790ac394b..4833a71001a 100644 --- a/contrib/cloud.google.com/go/pubsub.v1/orchestrion.yml +++ b/contrib/cloud.google.com/go/pubsub.v1/orchestrion.yml @@ -23,7 +23,7 @@ aspects: __dd_instr *instrumentation.Instrumentation __dd_pstrace *pubsubtrace.Tracer ) - + func init() { component := instrumentation.PackageGCPPubsub __dd_instr = instrumentation.Load(component) @@ -93,3 +93,30 @@ aspects: defer func() { {{ $publishResult }}.DDCloseSpan = __dd_closeSpan }() + + ## Admin / management operations ## + # In v1, admin RPCs live on the GAPIC PublisherClient (topics), SubscriberClient + # (subscriptions + snapshots) and SchemaClient under cloud.google.com/go/pubsub/apiv1, + # issued over their gRPC connection. Rather than wrapping every method, we append a + # unary client interceptor to each admin client constructor. The interceptor emits a + # gcp.pubsub.request span for admin RPCs and forwards the data-plane RPCs (Publish, + # Pull, Acknowledge, ...) that share the same connection untouched. + # + # Limitation: option.WithGRPCDialOption is ignored when option.WithGRPCConn is set, + # so admin spans are not created for clients that reuse a pre-built grpc.ClientConn. + # WithGRPCConn users should add pubsubtrace.UnaryAdminInterceptorV1 on their dial + - id: Admin client interceptor + join-point: + one-of: + - function-call: cloud.google.com/go/pubsub/apiv1.NewPublisherClient + - function-call: cloud.google.com/go/pubsub/apiv1.NewSubscriberClient + - function-call: cloud.google.com/go/pubsub/apiv1.NewSchemaClient + advice: + - append-args: + type: google.golang.org/api/option.ClientOption + values: + - imports: + option: google.golang.org/api/option + grpc: google.golang.org/grpc + pubsubtrace: github.com/DataDog/dd-trace-go/v2/contrib/cloud.google.com/go/pubsubtrace + template: option.WithGRPCDialOption(grpc.WithChainUnaryInterceptor(pubsubtrace.UnaryAdminInterceptorV1())) diff --git a/contrib/cloud.google.com/go/pubsub.v2/admin_test.go b/contrib/cloud.google.com/go/pubsub.v2/admin_test.go new file mode 100644 index 00000000000..f894d35d950 --- /dev/null +++ b/contrib/cloud.google.com/go/pubsub.v2/admin_test.go @@ -0,0 +1,303 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2025 Datadog, Inc. + +package pubsub + +import ( + "context" + "fmt" + "testing" + "time" + + vkit "cloud.google.com/go/pubsub/v2/apiv1" + "cloud.google.com/go/pubsub/v2/apiv1/pubsubpb" + "cloud.google.com/go/pubsub/v2/pstest" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "google.golang.org/api/iterator" + "google.golang.org/api/option" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + + "github.com/DataDog/dd-trace-go/v2/contrib/cloud.google.com/go/pubsubtrace" + "github.com/DataDog/dd-trace-go/v2/ddtrace/ext" + "github.com/DataDog/dd-trace-go/v2/ddtrace/mocktracer" +) + +const adminProjectID = "project" + +func setupAdmin(t *testing.T) (context.Context, mocktracer.Tracer, *vkit.TopicAdminClient, *vkit.SubscriptionAdminClient, *vkit.SchemaClient) { + mt := mocktracer.Start() + t.Cleanup(mt.Stop) + + srv := pstest.NewServer() + t.Cleanup(func() { assert.NoError(t, srv.Close()) }) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + t.Cleanup(cancel) + + // The admin GAPIC clients issue their RPCs over this connection, so + // installing the interceptor here traces their admin operations. + conn, err := grpc.NewClient(srv.Addr, + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithChainUnaryInterceptor(pubsubtrace.UnaryAdminInterceptorV2()), + ) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, conn.Close()) }) + + tac, err := vkit.NewTopicAdminClient(ctx, option.WithGRPCConn(conn)) + require.NoError(t, err) + sac, err := vkit.NewSubscriptionAdminClient(ctx, option.WithGRPCConn(conn)) + require.NoError(t, err) + schema, err := vkit.NewSchemaClient(ctx, option.WithGRPCConn(conn)) + require.NoError(t, err) + + return ctx, mt, tac, sac, schema +} + +func topicName(id string) string { + return fmt.Sprintf("projects/%s/topics/%s", adminProjectID, id) +} + +func subName(id string) string { + return fmt.Sprintf("projects/%s/subscriptions/%s", adminProjectID, id) +} + +func snapshotName(id string) string { + return fmt.Sprintf("projects/%s/snapshots/%s", adminProjectID, id) +} + +func schemaName(id string) string { + return fmt.Sprintf("projects/%s/schemas/%s", adminProjectID, id) +} + +func projectName() string { + return fmt.Sprintf("projects/%s", adminProjectID) +} + +func drain[T any](t *testing.T, next func() (T, error)) { + t.Helper() + for { + if _, err := next(); err == iterator.Done { + return + } else { + require.NoError(t, err) + } + } +} + +func TestTraceAdminTopicOperations(t *testing.T) { + ctx, mt, tac, _, _ := setupAdmin(t) + + _, err := tac.CreateTopic(ctx, &pubsubpb.Topic{Name: topicName("topic")}) + require.NoError(t, err) + + _, err = tac.GetTopic(ctx, &pubsubpb.GetTopicRequest{Topic: topicName("topic")}) + require.NoError(t, err) + + it := tac.ListTopics(ctx, &pubsubpb.ListTopicsRequest{Project: projectName()}) + drain(t, it.Next) + + err = tac.DeleteTopic(ctx, &pubsubpb.DeleteTopicRequest{Topic: topicName("topic")}) + require.NoError(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 4) + + assertAdminSpan(t, spans[0], "CreateTopic", "CreateTopic "+topicName("topic")) + assertAdminSpan(t, spans[1], "GetTopic", "GetTopic "+topicName("topic")) + assertAdminSpan(t, spans[2], "ListTopics", "ListTopics "+projectName()) + assertAdminSpan(t, spans[3], "DeleteTopic", "DeleteTopic "+topicName("topic")) +} + +func TestTraceAdminSubscriptionOperations(t *testing.T) { + ctx, mt, tac, sac, _ := setupAdmin(t) + + _, err := tac.CreateTopic(ctx, &pubsubpb.Topic{Name: topicName("topic")}) + require.NoError(t, err) + + _, err = sac.CreateSubscription(ctx, &pubsubpb.Subscription{ + Name: subName("sub"), + Topic: topicName("topic"), + }) + require.NoError(t, err) + + _, err = sac.GetSubscription(ctx, &pubsubpb.GetSubscriptionRequest{Subscription: subName("sub")}) + require.NoError(t, err) + + it := sac.ListSubscriptions(ctx, &pubsubpb.ListSubscriptionsRequest{Project: projectName()}) + drain(t, it.Next) + + err = sac.DeleteSubscription(ctx, &pubsubpb.DeleteSubscriptionRequest{Subscription: subName("sub")}) + require.NoError(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 5) + + assertAdminSpan(t, spans[0], "CreateTopic", "CreateTopic "+topicName("topic")) + assertAdminSpan(t, spans[1], "CreateSubscription", "CreateSubscription "+subName("sub")) + assertAdminSpan(t, spans[2], "GetSubscription", "GetSubscription "+subName("sub")) + assertAdminSpan(t, spans[3], "ListSubscriptions", "ListSubscriptions "+projectName()) + assertAdminSpan(t, spans[4], "DeleteSubscription", "DeleteSubscription "+subName("sub")) +} + +func TestTraceAdminSnapshotOperations(t *testing.T) { + ctx, mt, tac, sac, _ := setupAdmin(t) + + _, err := tac.CreateTopic(ctx, &pubsubpb.Topic{Name: topicName("topic")}) + require.NoError(t, err) + _, err = sac.CreateSubscription(ctx, &pubsubpb.Subscription{ + Name: subName("sub"), + Topic: topicName("topic"), + }) + require.NoError(t, err) + + // pstest does not implement snapshots, so these RPCs error — but the + // interceptor must still emit spans with the resolved resource path. + _, err = sac.CreateSnapshot(ctx, &pubsubpb.CreateSnapshotRequest{ + Name: snapshotName("snap"), + Subscription: subName("sub"), + }) + require.Error(t, err) + + _, err = sac.GetSnapshot(ctx, &pubsubpb.GetSnapshotRequest{Snapshot: snapshotName("snap")}) + require.Error(t, err) + + it := sac.ListSnapshots(ctx, &pubsubpb.ListSnapshotsRequest{Project: projectName()}) + _, err = it.Next() + require.Error(t, err) + require.NotEqual(t, iterator.Done, err) + + err = sac.DeleteSnapshot(ctx, &pubsubpb.DeleteSnapshotRequest{Snapshot: snapshotName("snap")}) + require.Error(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 6) + + assertAdminSpan(t, spans[0], "CreateTopic", "CreateTopic "+topicName("topic")) + assertAdminSpan(t, spans[1], "CreateSubscription", "CreateSubscription "+subName("sub")) + assertAdminSpan(t, spans[2], "CreateSnapshot", "CreateSnapshot "+snapshotName("snap")) + assert.NotNil(t, spans[2].Tag(ext.ErrorMsg)) + assertAdminSpan(t, spans[3], "GetSnapshot", "GetSnapshot "+snapshotName("snap")) + assert.NotNil(t, spans[3].Tag(ext.ErrorMsg)) + assertAdminSpan(t, spans[4], "ListSnapshots", "ListSnapshots "+projectName()) + assert.NotNil(t, spans[4].Tag(ext.ErrorMsg)) + assertAdminSpan(t, spans[5], "DeleteSnapshot", "DeleteSnapshot "+snapshotName("snap")) + assert.NotNil(t, spans[5].Tag(ext.ErrorMsg)) +} + +func TestTraceAdminSchemaOperations(t *testing.T) { + ctx, mt, _, _, sc := setupAdmin(t) + + const avroDef = `{"type":"record","name":"Test","fields":[{"name":"f","type":"string"}]}` + _, err := sc.CreateSchema(ctx, &pubsubpb.CreateSchemaRequest{ + Parent: projectName(), + Schema: &pubsubpb.Schema{ + Type: pubsubpb.Schema_AVRO, + Definition: avroDef, + }, + SchemaId: "schema", + }) + require.NoError(t, err) + + _, err = sc.GetSchema(ctx, &pubsubpb.GetSchemaRequest{Name: schemaName("schema")}) + require.NoError(t, err) + + it := sc.ListSchemas(ctx, &pubsubpb.ListSchemasRequest{Parent: projectName()}) + drain(t, it.Next) + + err = sc.DeleteSchema(ctx, &pubsubpb.DeleteSchemaRequest{Name: schemaName("schema")}) + require.NoError(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 4) + + assertAdminSpan(t, spans[0], "CreateSchema", "CreateSchema "+projectName()) + assertAdminSpan(t, spans[1], "GetSchema", "GetSchema "+schemaName("schema")) + assertAdminSpan(t, spans[2], "ListSchemas", "ListSchemas "+projectName()) + assertAdminSpan(t, spans[3], "DeleteSchema", "DeleteSchema "+schemaName("schema")) +} + +func TestTraceAdminError(t *testing.T) { + ctx, mt, tac, _, _ := setupAdmin(t) + + // Getting a topic that does not exist returns an error, which must be recorded on the span. + _, err := tac.GetTopic(ctx, &pubsubpb.GetTopicRequest{Topic: topicName("missing")}) + require.Error(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 1) + assert.Equal(t, err.Error(), spans[0].Tag(ext.ErrorMsg)) + assertAdminSpan(t, spans[0], "GetTopic", "GetTopic "+topicName("missing")) +} + +func TestTraceAdminMissingResource(t *testing.T) { + ctx, mt, tac, _, _ := setupAdmin(t) + + // Recognized admin RPCs with an empty resource field must still emit a + // span; TraceAdmin falls back to a method-only resource name. + _, createErr := tac.CreateTopic(ctx, &pubsubpb.Topic{}) + _, getErr := tac.GetTopic(ctx, &pubsubpb.GetTopicRequest{}) + + spans := mt.FinishedSpans() + require.Len(t, spans, 2) + + assert.Equal(t, "gcp.pubsub.request", spans[0].OperationName()) + assert.Equal(t, "CreateTopic", spans[0].Tag(ext.ResourceName)) + assert.Equal(t, "CreateTopic", spans[0].Tag("pubsub.method")) + assert.Nil(t, spans[0].Tag(ext.GCPProjectID)) + if createErr != nil { + assert.Equal(t, createErr.Error(), spans[0].Tag(ext.ErrorMsg)) + } + + assert.Equal(t, "gcp.pubsub.request", spans[1].OperationName()) + assert.Equal(t, "GetTopic", spans[1].Tag(ext.ResourceName)) + assert.Equal(t, "GetTopic", spans[1].Tag("pubsub.method")) + assert.Nil(t, spans[1].Tag(ext.GCPProjectID)) + if getErr != nil { + assert.Equal(t, getErr.Error(), spans[1].Tag(ext.ErrorMsg)) + } +} + +func TestTraceAdminWithService(t *testing.T) { + mt := mocktracer.Start() + defer mt.Stop() + + srv := pstest.NewServer() + defer func() { assert.NoError(t, srv.Close()) }() + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + conn, err := grpc.NewClient(srv.Addr, + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithChainUnaryInterceptor(pubsubtrace.UnaryAdminInterceptorV2(WithService("my-admin-service"))), + ) + require.NoError(t, err) + defer func() { assert.NoError(t, conn.Close()) }() + + tac, err := vkit.NewTopicAdminClient(ctx, option.WithGRPCConn(conn)) + require.NoError(t, err) + + _, err = tac.CreateTopic(ctx, &pubsubpb.Topic{Name: topicName("topic")}) + require.NoError(t, err) + + spans := mt.FinishedSpans() + require.Len(t, spans, 1) + assert.Equal(t, "my-admin-service", spans[0].Tag(ext.ServiceName)) +} + +func assertAdminSpan(t *testing.T, span *mocktracer.Span, method, resource string) { + t.Helper() + assert.Equal(t, "gcp.pubsub.request", span.OperationName()) + assert.Equal(t, resource, span.Tag(ext.ResourceName)) + assert.Equal(t, ext.SpanTypeWorker, span.Tag(ext.SpanType)) + assert.Equal(t, ext.SpanKindClient, span.Tag(ext.SpanKind)) + assert.Equal(t, "cloud.google.com/go/pubsub.v2", span.Tag(ext.Component)) + assert.Equal(t, ext.MessagingSystemGCPPubsub, span.Tag(ext.MessagingSystem)) + assert.Equal(t, method, span.Tag("pubsub.method")) + assert.Equal(t, adminProjectID, span.Tag(ext.GCPProjectID)) + assert.Equal(t, "cloud.google.com/go/pubsub.v2", span.Integration()) +} diff --git a/contrib/cloud.google.com/go/pubsub.v2/go.mod b/contrib/cloud.google.com/go/pubsub.v2/go.mod index 712bc3ee407..c6556adca56 100644 --- a/contrib/cloud.google.com/go/pubsub.v2/go.mod +++ b/contrib/cloud.google.com/go/pubsub.v2/go.mod @@ -16,6 +16,7 @@ require ( cloud.google.com/go/auth/oauth2adapt v0.2.8 // indirect cloud.google.com/go/compute/metadata v0.9.0 // indirect cloud.google.com/go/iam v1.5.3 // indirect + cloud.google.com/go/pubsub v1.50.1 // indirect github.com/DataDog/datadog-agent/comp/core/tagger/origindetection v0.82.0-rc.2 // indirect github.com/DataDog/datadog-agent/pkg/obfuscate v0.82.0-rc.2 // indirect github.com/DataDog/datadog-agent/pkg/opentelemetry-mapping-go/otlp/attributes v0.82.0-rc.2 // indirect diff --git a/contrib/cloud.google.com/go/pubsub.v2/go.sum b/contrib/cloud.google.com/go/pubsub.v2/go.sum index bcbce923e14..4ec9c7f389f 100644 --- a/contrib/cloud.google.com/go/pubsub.v2/go.sum +++ b/contrib/cloud.google.com/go/pubsub.v2/go.sum @@ -9,6 +9,8 @@ cloud.google.com/go/compute/metadata v0.9.0 h1:pDUj4QMoPejqq20dK0Pg2N4yG9zIkYGdB cloud.google.com/go/compute/metadata v0.9.0/go.mod h1:E0bWwX5wTnLPedCKqk3pJmVgCBSM6qQI1yTBdEb3C10= cloud.google.com/go/iam v1.5.3 h1:+vMINPiDF2ognBJ97ABAYYwRgsaqxPbQDlMnbHMjolc= cloud.google.com/go/iam v1.5.3/go.mod h1:MR3v9oLkZCTlaqljW6Eb2d3HGDGK5/bDv93jhfISFvU= +cloud.google.com/go/pubsub v1.50.1 h1:fzbXpPyJnSGvWXF1jabhQeXyxdbCIkXTpjXHy7xviBM= +cloud.google.com/go/pubsub v1.50.1/go.mod h1:6YVJv3MzWJUVdvQXG081sFvS0dWQOdnV+oTo++q/xFk= cloud.google.com/go/pubsub/v2 v2.0.0 h1:0qS6mRJ41gD1lNmM/vdm6bR7DQu6coQcVwD+VPf0Bz0= cloud.google.com/go/pubsub/v2 v2.0.0/go.mod h1:0aztFxNzVQIRSZ8vUr79uH2bS3jwLebwK6q1sgEub+E= github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU= diff --git a/contrib/cloud.google.com/go/pubsub.v2/orchestrion.yml b/contrib/cloud.google.com/go/pubsub.v2/orchestrion.yml index e87a0f7906f..abe25869e95 100644 --- a/contrib/cloud.google.com/go/pubsub.v2/orchestrion.yml +++ b/contrib/cloud.google.com/go/pubsub.v2/orchestrion.yml @@ -24,7 +24,7 @@ aspects: __dd_instr *instrumentation.Instrumentation __dd_pstrace *pubsubtrace.Tracer ) - + func init() { component := instrumentation.PackageGCPPubsubV2 __dd_instr = instrumentation.Load(component) @@ -94,7 +94,34 @@ aspects: defer func() { {{ $publishResult }}.DDCloseSpan = __dd_closeSpan }() - + + + ## Admin / management operations ## + # Pub/Sub admin (topic, subscription, snapshot and schema management) RPCs are + # issued by the GAPIC clients in cloud.google.com/go/pubsub/v2/apiv1 over their gRPC + # connection. Rather than wrapping every method, we append a unary client interceptor + # to each admin client constructor. The interceptor emits a gcp.pubsub.request span + # for admin RPCs and forwards the data-plane RPCs (Publish, Pull, Acknowledge, ...) + # that share the same connection untouched. + # + # Limitation: option.WithGRPCDialOption is ignored when option.WithGRPCConn is set, + # so admin spans are not created for clients that reuse a pre-built grpc.ClientConn. + # WithGRPCConn users should add pubsubtrace.UnaryAdminInterceptorV2 on their dial + - id: Admin client interceptor + join-point: + one-of: + - function-call: cloud.google.com/go/pubsub/v2/apiv1.NewTopicAdminClient + - function-call: cloud.google.com/go/pubsub/v2/apiv1.NewSubscriptionAdminClient + - function-call: cloud.google.com/go/pubsub/v2/apiv1.NewSchemaClient + advice: + - append-args: + type: google.golang.org/api/option.ClientOption + values: + - imports: + option: google.golang.org/api/option + grpc: google.golang.org/grpc + pubsubtrace: github.com/DataDog/dd-trace-go/v2/contrib/cloud.google.com/go/pubsubtrace + template: option.WithGRPCDialOption(grpc.WithChainUnaryInterceptor(pubsubtrace.UnaryAdminInterceptorV2())) # Note: These aspects are also necessary for v1, they should be present in only one of the orchestrion.yml files, # otherwise they will be applied twice and will cause the build to fail. diff --git a/contrib/cloud.google.com/go/pubsubtrace/admin.go b/contrib/cloud.google.com/go/pubsubtrace/admin.go new file mode 100644 index 00000000000..6d4e0df4fcc --- /dev/null +++ b/contrib/cloud.google.com/go/pubsubtrace/admin.go @@ -0,0 +1,117 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2025 Datadog, Inc. + +// v2 admin request→resource mapping, implemented as a grpc unary client interceptor. +// Uses cloud.google.com/go/pubsub/v2/apiv1/pubsubpb, which is distinct from v1's pubsubpb +// (see admin_v1.go). + +package pubsubtrace + +import ( + "sync" + + "cloud.google.com/go/pubsub/v2/apiv1/pubsubpb" + "google.golang.org/grpc" + + "github.com/DataDog/dd-trace-go/v2/instrumentation" +) + +// UnaryAdminInterceptorV2 returns a grpc.UnaryClientInterceptor that traces +// TopicAdminClient, SubscriptionAdminClient, and SchemaClient admin operations. +// +// When constructing admin clients with option.WithGRPCConn, install this +// interceptor on the dial that creates the connection (WithGRPCDialOption on +// the client constructor is ignored in that case). +func UnaryAdminInterceptorV2(opts ...Option) grpc.UnaryClientInterceptor { + return defaultTracerV2().unaryAdminInterceptor(resolveAdminResourceV2, opts...) +} + +var ( + v2TracerOnce sync.Once + v2Tracer *Tracer +) + +func defaultTracerV2() *Tracer { + v2TracerOnce.Do(func() { + component := instrumentation.PackageGCPPubsubV2 + v2Tracer = NewTracer(instrumentation.Load(component), component) + }) + return v2Tracer +} + +// resolveAdminResourceV2 maps a v2 admin request to its resource path. +// ok is false for non-admin requests (Publish, Pull, Acknowledge, IAM, ...). +func resolveAdminResourceV2(req any) (resourcePath string, ok bool) { + switch r := req.(type) { + // TopicAdminClient + case *pubsubpb.Topic: + return r.GetName(), true + case *pubsubpb.UpdateTopicRequest: + return r.GetTopic().GetName(), true + case *pubsubpb.GetTopicRequest: + return r.GetTopic(), true + case *pubsubpb.ListTopicsRequest: + return r.GetProject(), true + case *pubsubpb.ListTopicSubscriptionsRequest: + return r.GetTopic(), true + case *pubsubpb.ListTopicSnapshotsRequest: + return r.GetTopic(), true + case *pubsubpb.DeleteTopicRequest: + return r.GetTopic(), true + case *pubsubpb.DetachSubscriptionRequest: + return r.GetSubscription(), true + + // SubscriptionAdminClient + case *pubsubpb.Subscription: + return r.GetName(), true + case *pubsubpb.GetSubscriptionRequest: + return r.GetSubscription(), true + case *pubsubpb.UpdateSubscriptionRequest: + return r.GetSubscription().GetName(), true + case *pubsubpb.ListSubscriptionsRequest: + return r.GetProject(), true + case *pubsubpb.DeleteSubscriptionRequest: + return r.GetSubscription(), true + case *pubsubpb.ModifyPushConfigRequest: + return r.GetSubscription(), true + case *pubsubpb.GetSnapshotRequest: + return r.GetSnapshot(), true + case *pubsubpb.ListSnapshotsRequest: + return r.GetProject(), true + case *pubsubpb.CreateSnapshotRequest: + return r.GetName(), true + case *pubsubpb.UpdateSnapshotRequest: + return r.GetSnapshot().GetName(), true + case *pubsubpb.DeleteSnapshotRequest: + return r.GetSnapshot(), true + case *pubsubpb.SeekRequest: + return r.GetSubscription(), true + + // SchemaClient + case *pubsubpb.CreateSchemaRequest: + return r.GetParent(), true + case *pubsubpb.GetSchemaRequest: + return r.GetName(), true + case *pubsubpb.ListSchemasRequest: + return r.GetParent(), true + case *pubsubpb.ListSchemaRevisionsRequest: + return r.GetName(), true + case *pubsubpb.CommitSchemaRequest: + return r.GetName(), true + case *pubsubpb.RollbackSchemaRequest: + return r.GetName(), true + case *pubsubpb.DeleteSchemaRevisionRequest: + return r.GetName(), true + case *pubsubpb.DeleteSchemaRequest: + return r.GetName(), true + case *pubsubpb.ValidateSchemaRequest: + return r.GetParent(), true + case *pubsubpb.ValidateMessageRequest: + return r.GetParent(), true + + default: + return "", false + } +} diff --git a/contrib/cloud.google.com/go/pubsubtrace/admin_interceptor.go b/contrib/cloud.google.com/go/pubsubtrace/admin_interceptor.go new file mode 100644 index 00000000000..9daa945a6f3 --- /dev/null +++ b/contrib/cloud.google.com/go/pubsubtrace/admin_interceptor.go @@ -0,0 +1,40 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2026 Datadog, Inc. + +package pubsubtrace + +import ( + "context" + "strings" + + "google.golang.org/grpc" +) + +// adminResolver returns the resource path and true for recognized admin requests, else false. +type adminResolver func(req any) (resourcePath string, ok bool) + +// unaryAdminInterceptor builds a grpc.UnaryClientInterceptor that emits a +// gcp.pubsub.request span for each admin RPC recognised by resolve. +func (tr *Tracer) unaryAdminInterceptor(resolve adminResolver, opts ...Option) grpc.UnaryClientInterceptor { + return func(ctx context.Context, method string, req, reply any, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, callOpts ...grpc.CallOption) error { + resourcePath, ok := resolve(req) + if !ok { + return invoker(ctx, method, req, reply, cc, callOpts...) + } + ctx, finish := tr.TraceAdmin(ctx, adminMethodName(method), resourcePath, opts...) + err := invoker(ctx, method, req, reply, cc, callOpts...) + finish(err) + return err + } +} + +// adminMethodName returns the RPC method name from a gRPC full-method string, e.g. +// "/google.pubsub.v1.Publisher/CreateTopic" -> "CreateTopic". +func adminMethodName(fullMethod string) string { + if i := strings.LastIndex(fullMethod, "/"); i >= 0 { + return fullMethod[i+1:] + } + return fullMethod +} diff --git a/contrib/cloud.google.com/go/pubsubtrace/admin_v1.go b/contrib/cloud.google.com/go/pubsubtrace/admin_v1.go new file mode 100644 index 00000000000..2f70a8d4fef --- /dev/null +++ b/contrib/cloud.google.com/go/pubsubtrace/admin_v1.go @@ -0,0 +1,117 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2025 Datadog, Inc. + +// v1 admin request→resource mapping. Uses cloud.google.com/go/pubsub/apiv1/pubsubpb, +// which is distinct from v2's pubsubpb (see admin.go). + +package pubsubtrace + +import ( + "sync" + + pubsubpb "cloud.google.com/go/pubsub/apiv1/pubsubpb" + "google.golang.org/grpc" + + "github.com/DataDog/dd-trace-go/v2/instrumentation" +) + +// UnaryAdminInterceptorV1 returns a grpc.UnaryClientInterceptor that traces the admin +// operations of the cloud.google.com/go/pubsub (v1) GAPIC clients (PublisherClient, +// SubscriberClient, SchemaClient). +// +// When constructing admin clients with option.WithGRPCConn, install this +// interceptor on the dial that creates the connection (WithGRPCDialOption on +// the client constructor is ignored in that case). +func UnaryAdminInterceptorV1(opts ...Option) grpc.UnaryClientInterceptor { + return defaultTracerV1().unaryAdminInterceptor(resolveAdminResourceV1, opts...) +} + +var ( + v1TracerOnce sync.Once + v1Tracer *Tracer +) + +func defaultTracerV1() *Tracer { + v1TracerOnce.Do(func() { + component := instrumentation.PackageGCPPubsub + v1Tracer = NewTracer(instrumentation.Load(component), component) + }) + return v1Tracer +} + +// resolveAdminResourceV1 maps a v1 admin request to its resource path. +// ok is false for non-admin requests (Publish, Pull, Acknowledge, ModifyAckDeadline, IAM, ...). +func resolveAdminResourceV1(req any) (resourcePath string, ok bool) { + switch r := req.(type) { + // PublisherClient (topics) + case *pubsubpb.Topic: + return r.GetName(), true + case *pubsubpb.UpdateTopicRequest: + return r.GetTopic().GetName(), true + case *pubsubpb.GetTopicRequest: + return r.GetTopic(), true + case *pubsubpb.ListTopicsRequest: + return r.GetProject(), true + case *pubsubpb.ListTopicSubscriptionsRequest: + return r.GetTopic(), true + case *pubsubpb.ListTopicSnapshotsRequest: + return r.GetTopic(), true + case *pubsubpb.DeleteTopicRequest: + return r.GetTopic(), true + case *pubsubpb.DetachSubscriptionRequest: + return r.GetSubscription(), true + + // SubscriberClient (subscriptions + snapshots) + case *pubsubpb.Subscription: + return r.GetName(), true + case *pubsubpb.GetSubscriptionRequest: + return r.GetSubscription(), true + case *pubsubpb.UpdateSubscriptionRequest: + return r.GetSubscription().GetName(), true + case *pubsubpb.ListSubscriptionsRequest: + return r.GetProject(), true + case *pubsubpb.DeleteSubscriptionRequest: + return r.GetSubscription(), true + case *pubsubpb.ModifyPushConfigRequest: + return r.GetSubscription(), true + case *pubsubpb.GetSnapshotRequest: + return r.GetSnapshot(), true + case *pubsubpb.ListSnapshotsRequest: + return r.GetProject(), true + case *pubsubpb.CreateSnapshotRequest: + return r.GetName(), true + case *pubsubpb.UpdateSnapshotRequest: + return r.GetSnapshot().GetName(), true + case *pubsubpb.DeleteSnapshotRequest: + return r.GetSnapshot(), true + case *pubsubpb.SeekRequest: + return r.GetSubscription(), true + + // SchemaClient + case *pubsubpb.CreateSchemaRequest: + return r.GetParent(), true + case *pubsubpb.GetSchemaRequest: + return r.GetName(), true + case *pubsubpb.ListSchemasRequest: + return r.GetParent(), true + case *pubsubpb.ListSchemaRevisionsRequest: + return r.GetName(), true + case *pubsubpb.CommitSchemaRequest: + return r.GetName(), true + case *pubsubpb.RollbackSchemaRequest: + return r.GetName(), true + case *pubsubpb.DeleteSchemaRevisionRequest: + return r.GetName(), true + case *pubsubpb.DeleteSchemaRequest: + return r.GetName(), true + case *pubsubpb.ValidateSchemaRequest: + return r.GetParent(), true + case *pubsubpb.ValidateMessageRequest: + return r.GetParent(), true + + default: + return "", false + } +} diff --git a/contrib/cloud.google.com/go/pubsubtrace/config.go b/contrib/cloud.google.com/go/pubsubtrace/config.go index be79564b194..ebaa2311b96 100644 --- a/contrib/cloud.google.com/go/pubsubtrace/config.go +++ b/contrib/cloud.google.com/go/pubsubtrace/config.go @@ -19,6 +19,7 @@ type config struct { serviceSource string publishSpanName string receiveSpanName string + requestSpanName string measured bool propagationAsSpanLinks bool } @@ -35,6 +36,7 @@ func (tr *Tracer) defaultConfig() *config { serviceSource: string(tr.component), publishSpanName: tr.instr.OperationName(instrumentation.ComponentProducer, nil), receiveSpanName: tr.instr.OperationName(instrumentation.ComponentConsumer, nil), + requestSpanName: tr.instr.OperationName(instrumentation.ComponentClient, nil), measured: false, propagationAsSpanLinks: propagationAsSpanLinks, } diff --git a/contrib/cloud.google.com/go/pubsubtrace/tracing.go b/contrib/cloud.google.com/go/pubsubtrace/tracing.go index 7aab6c89bad..58405ad4513 100644 --- a/contrib/cloud.google.com/go/pubsubtrace/tracing.go +++ b/contrib/cloud.google.com/go/pubsubtrace/tracing.go @@ -161,9 +161,6 @@ func (tr *Tracer) TraceReceiveFunc(s Subscription, opts ...Option) func(ctx cont if projectID := projectIDFromResourceName(s.String()); projectID != "" { opts = append(opts, tracer.Tag(ext.GCPProjectID, projectID)) } - if projectID := projectIDFromResourceName(s.String()); projectID != "" { - opts = append(opts, tracer.Tag(ext.GCPProjectID, projectID)) - } if cfg.serviceName != "" { opts = append(opts, instrumentation.ServiceNameWithSource(cfg.serviceName, cfg.serviceSource)) } @@ -181,8 +178,47 @@ func (tr *Tracer) TraceReceiveFunc(s Subscription, opts ...Option) func(ctx cont } } -// extracts the GCP project ID from a Pubsub resource name of the form -// "projects/{project}/topics/{topic}" or "projects/{project}/subscriptions/{subscription}" +// TraceAdmin starts a span for a Pub/Sub admin operation (e.g. CreateTopic, ListSubscriptions, DeleteSchema). +// It is driven by the unary client interceptor in admin.go / admin_v1.go, which is the single source of +// truth for the (method, resourcePath) mapping across both the manual and orchestrion instrumentation. +func (tr *Tracer) TraceAdmin(ctx context.Context, method string, resourcePath string, opts ...Option) (context.Context, func(err error)) { + cfg := tr.defaultConfig() + for _, opt := range opts { + opt.apply(cfg) + } + resource := method + if resourcePath != "" { + resource = method + " " + resourcePath + } + spanOpts := []tracer.StartSpanOption{ + tracer.ResourceName(resource), + tracer.SpanType(ext.SpanTypeWorker), + tracer.Tag(ext.Component, tr.component), + tracer.Tag(ext.SpanKind, ext.SpanKindClient), + tracer.Tag(ext.MessagingSystem, ext.MessagingSystemGCPPubsub), + tracer.Tag("pubsub.method", method), + tracer.Measured(), + } + if projectID := projectIDFromResourceName(resourcePath); projectID != "" { + spanOpts = append(spanOpts, tracer.Tag(ext.GCPProjectID, projectID)) + } + if cfg.serviceName != "" { + spanOpts = append(spanOpts, instrumentation.ServiceNameWithSource(cfg.serviceName, cfg.serviceSource)) + } + + span, ctx := tracer.StartSpanFromContext(ctx, cfg.requestSpanName, spanOpts...) + + var once sync.Once + finish := func(err error) { + once.Do(func() { + span.Finish(tracer.WithError(err)) + }) + } + return ctx, finish +} + +// extracts the GCP project ID from a Pubsub resource name starting with +// "projects/{project}. e.g. schemas, snapshots, topics and subscriptions func projectIDFromResourceName(name string) string { const prefix = "projects/" if !strings.HasPrefix(name, prefix) { diff --git a/ddtrace/ext/app_types.go b/ddtrace/ext/app_types.go index d06610f7bac..ec7cd8fcdf1 100644 --- a/ddtrace/ext/app_types.go +++ b/ddtrace/ext/app_types.go @@ -80,6 +80,10 @@ const ( // SpanTypeLLM marks a span as an LLM operation. SpanTypeLLM = "llm" + // SpanTypeWorker marks a span as a worker/management operation, such as a + // Pub/Sub admin (topic, subscription, snapshot, schema management) request. + SpanTypeWorker = "worker" + // SpanTypeAerospike marks a span as an Aerospike operation. SpanTypeAerospike = "aerospike" ) diff --git a/go.mod b/go.mod index 4910b7e4f33..aadb12f22c9 100644 --- a/go.mod +++ b/go.mod @@ -5,6 +5,8 @@ go 1.25.0 godebug x509negativeserial=1 require ( + cloud.google.com/go/pubsub v1.50.1 + cloud.google.com/go/pubsub/v2 v2.0.0 github.com/DataDog/datadog-agent/pkg/obfuscate v0.82.0-rc.2 github.com/DataDog/datadog-agent/pkg/proto v0.82.0-rc.2 github.com/DataDog/datadog-agent/pkg/remoteconfig/state v0.82.0-rc.2 @@ -53,6 +55,7 @@ require ( golang.org/x/time v0.15.0 golang.org/x/tools v0.45.0 golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da + google.golang.org/grpc v1.82.0 google.golang.org/protobuf v1.36.12-0.20260116114154-8c4c4ae446ca ) @@ -103,7 +106,6 @@ require ( golang.org/x/text v0.38.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260618152121-87f3d3e198d3 // indirect - google.golang.org/grpc v1.82.0 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index 2f0d3f716be..964b7b236b1 100644 --- a/go.sum +++ b/go.sum @@ -1,3 +1,7 @@ +cloud.google.com/go/pubsub v1.50.1 h1:fzbXpPyJnSGvWXF1jabhQeXyxdbCIkXTpjXHy7xviBM= +cloud.google.com/go/pubsub v1.50.1/go.mod h1:6YVJv3MzWJUVdvQXG081sFvS0dWQOdnV+oTo++q/xFk= +cloud.google.com/go/pubsub/v2 v2.0.0 h1:0qS6mRJ41gD1lNmM/vdm6bR7DQu6coQcVwD+VPf0Bz0= +cloud.google.com/go/pubsub/v2 v2.0.0/go.mod h1:0aztFxNzVQIRSZ8vUr79uH2bS3jwLebwK6q1sgEub+E= github.com/DataDog/datadog-agent/comp/core/tagger/origindetection v0.82.0-rc.2 h1:D78y0XOvACwT7QsOCx0QDDRe3yMGFdyrIgbZKRtoD6E= github.com/DataDog/datadog-agent/comp/core/tagger/origindetection v0.82.0-rc.2/go.mod h1:6LC1ryDn2VNqF0iNapwcLLdsfoFUMnT4p+JPu6sEkHg= github.com/DataDog/datadog-agent/pkg/obfuscate v0.82.0-rc.2 h1:3uIMYQvZNimspb7EFD66vEz0GteKjeT371QT+gVP/JQ= diff --git a/instrumentation/packages.go b/instrumentation/packages.go index 74eb9067cf1..81c10bcda7f 100644 --- a/instrumentation/packages.go +++ b/instrumentation/packages.go @@ -232,6 +232,12 @@ var packages = map[Package]PackageInfo{ buildOpNameV0: staticName("pubsub.publish"), buildOpNameV1: staticName("gcp.pubsub.send"), }, + ComponentClient: { + useDDServiceV0: false, + buildServiceNameV0: staticName(""), + buildOpNameV0: staticName("gcp.pubsub.request"), + buildOpNameV1: staticName("gcp.pubsub.request"), + }, }, }, PackageGCPPubsubV2: { @@ -250,6 +256,12 @@ var packages = map[Package]PackageInfo{ buildOpNameV0: staticName("pubsub.publish"), buildOpNameV1: staticName("gcp.pubsub.send"), }, + ComponentClient: { + useDDServiceV0: false, + buildServiceNameV0: staticName(""), + buildOpNameV0: staticName("gcp.pubsub.request"), + buildOpNameV1: staticName("gcp.pubsub.request"), + }, }, }, PackageConfluentKafkaGo: { diff --git a/internal/orchestrion/_integration/gcp_pubsub.v1/gcp_pubsub_admin.go b/internal/orchestrion/_integration/gcp_pubsub.v1/gcp_pubsub_admin.go new file mode 100644 index 00000000000..8afe854e3e9 --- /dev/null +++ b/internal/orchestrion/_integration/gcp_pubsub.v1/gcp_pubsub_admin.go @@ -0,0 +1,290 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2023-present Datadog, Inc. + +//go:build linux || !githubci + +package gcppubsub + +import ( + "context" + "fmt" + "testing" + + "cloud.google.com/go/pubsub" + vkit "cloud.google.com/go/pubsub/apiv1" + "cloud.google.com/go/pubsub/apiv1/pubsubpb" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/testcontainers/testcontainers-go" + tclog "github.com/testcontainers/testcontainers-go/log" + "github.com/testcontainers/testcontainers-go/modules/gcloud" + "google.golang.org/api/iterator" + "google.golang.org/api/option" + "google.golang.org/api/option/internaloption" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + + "github.com/DataDog/dd-trace-go/instrumentation/testutils/containers/v2" + + "github.com/DataDog/dd-trace-go/v2/ddtrace/ext" + "github.com/DataDog/dd-trace-go/v2/internal/orchestrion/_integration/internal/trace" +) + +const ( + adminProject = "pstest-orchestrion" + adminTopic = "admin-topic" + adminSubscription = "admin-subscription" + adminSnapshot = "admin-snapshot" + adminSchema = "admin-schema" + avroDefinition = `{"type":"record","name":"Test","fields":[{"name":"f","type":"string"}]}` +) + +type adminBase struct { + container *gcloud.GCloudContainer + uri string + projectID string +} + +// emulatorOptions returns the client options needed to reach the plaintext Pub/Sub +// emulator by dialing internally. Passing these (rather than a pre-built +// option.WithGRPCConn) lets orchestrion append the admin interceptor to the GAPIC +// client constructors it hooks. +func emulatorOptions(uri string) []option.ClientOption { + return []option.ClientOption{ + option.WithEndpoint(uri), + option.WithGRPCDialOption(grpc.WithTransportCredentials(insecure.NewCredentials())), + option.WithoutAuthentication(), + internaloption.SkipDialSettingsValidation(), + } +} + +func (b *adminBase) setup(ctx context.Context, t *testing.T) { + containers.SkipIfProviderIsNotHealthy(t) + + var err error + b.container, err = gcloud.RunPubsub(ctx, + "gcr.io/google.com/cloudsdktool/google-cloud-cli:emulators", + gcloud.WithProjectID(adminProject), + testcontainers.WithLogger(tclog.TestLogger(t)), + containers.WithTestLogConsumer(t), + ) + containers.AssertTestContainersError(t, err) + containers.RegisterContainerCleanup(t, b.container) + + b.projectID = b.container.Settings.ProjectID + b.uri = b.container.URI +} + +func (b *adminBase) projectPath() string { + return "projects/" + b.projectID +} + +func (b *adminBase) topicPath(id string) string { + return fmt.Sprintf("projects/%s/topics/%s", b.projectID, id) +} + +func (b *adminBase) subscriptionPath(id string) string { + return fmt.Sprintf("projects/%s/subscriptions/%s", b.projectID, id) +} + +func (b *adminBase) snapshotPath(id string) string { + return fmt.Sprintf("projects/%s/snapshots/%s", b.projectID, id) +} + +func (b *adminBase) schemaPath(id string) string { + return fmt.Sprintf("projects/%s/schemas/%s", b.projectID, id) +} + +func drain[T any](t *testing.T, next func() (T, error)) { + t.Helper() + for { + if _, err := next(); err == iterator.Done { + return + } else { + require.NoError(t, err) + } + } +} + +func (b *adminBase) adminTrace(method, resourcePath string) *trace.Trace { + return &trace.Trace{ + Tags: map[string]any{ + "name": "gcp.pubsub.request", + "type": "worker", + "resource": method + " " + resourcePath, + "service": "gcp_pubsub.v1.test", + }, + Meta: map[string]string{ + "span.kind": "client", + "component": "cloud.google.com/go/pubsub.v1", + "messaging.system": "googlepubsub", + "pubsub.method": method, + ext.GCPProjectID: adminProject, + }, + } +} + +func (b *adminBase) adminErrorTrace(method, resourcePath, errMsg string) *trace.Trace { + tr := b.adminTrace(method, resourcePath) + tr.Meta["error.message"] = errMsg + return tr +} + +// TestCaseAdminClient exercises admin ops via the high-level pubsub.Client. +// CreateTopic / CreateSubscription are backed by GAPIC clients constructed +// internally by pubsub.NewClient (a variadic-spread call). +type TestCaseAdminClient struct { + adminBase + client *pubsub.Client +} + +func (tc *TestCaseAdminClient) Setup(ctx context.Context, t *testing.T) { + tc.setup(ctx, t) + + var err error + tc.client, err = pubsub.NewClient(ctx, tc.projectID, emulatorOptions(tc.uri)...) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, tc.client.Close()) }) +} + +func (tc *TestCaseAdminClient) Run(ctx context.Context, t *testing.T) { + topic, err := tc.client.CreateTopic(ctx, adminTopic) + require.NoError(t, err) + + sub, err := tc.client.CreateSubscription(ctx, adminSubscription, pubsub.SubscriptionConfig{ + Topic: topic, + }) + require.NoError(t, err) + + require.NoError(t, sub.Delete(ctx)) + require.NoError(t, topic.Delete(ctx)) +} + +func (tc *TestCaseAdminClient) ExpectedTraces() trace.Traces { + return trace.Traces{ + tc.adminTrace("CreateTopic", tc.topicPath(adminTopic)), + tc.adminTrace("CreateSubscription", tc.subscriptionPath(adminSubscription)), + tc.adminTrace("DeleteSubscription", tc.subscriptionPath(adminSubscription)), + tc.adminTrace("DeleteTopic", tc.topicPath(adminTopic)), + } +} + +// TestCaseAdminGAPIC exercises admin ops via directly-constructed GAPIC +// PublisherClient and SubscriberClient. +type TestCaseAdminGAPIC struct { + adminBase + missingErrMsg string +} + +func (tc *TestCaseAdminGAPIC) Setup(ctx context.Context, t *testing.T) { + tc.setup(ctx, t) + + // Fixture resources are created before the tracer starts. + client, err := pubsub.NewClient(ctx, tc.projectID, emulatorOptions(tc.uri)...) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, client.Close()) }) + + topic, err := client.CreateTopic(ctx, adminTopic) + require.NoError(t, err) + + _, err = client.CreateSubscription(ctx, adminSubscription, pubsub.SubscriptionConfig{ + Topic: topic, + }) + require.NoError(t, err) +} + +func (tc *TestCaseAdminGAPIC) Run(ctx context.Context, t *testing.T) { + pc, err := vkit.NewPublisherClient(ctx, emulatorOptions(tc.uri)...) + require.NoError(t, err) + defer func() { assert.NoError(t, pc.Close()) }() + + sc, err := vkit.NewSubscriberClient(ctx, emulatorOptions(tc.uri)...) + require.NoError(t, err) + defer func() { assert.NoError(t, sc.Close()) }() + + drain(t, pc.ListTopics(ctx, &pubsubpb.ListTopicsRequest{Project: tc.projectPath()}).Next) + drain(t, sc.ListSubscriptions(ctx, &pubsubpb.ListSubscriptionsRequest{Project: tc.projectPath()}).Next) + + _, err = sc.CreateSnapshot(ctx, &pubsubpb.CreateSnapshotRequest{ + Name: tc.snapshotPath(adminSnapshot), + Subscription: tc.subscriptionPath(adminSubscription), + }) + require.NoError(t, err) + + drain(t, sc.ListSnapshots(ctx, &pubsubpb.ListSnapshotsRequest{Project: tc.projectPath()}).Next) + + _, err = pc.GetTopic(ctx, &pubsubpb.GetTopicRequest{Topic: tc.topicPath(adminTopic)}) + require.NoError(t, err) + + _, err = pc.GetTopic(ctx, &pubsubpb.GetTopicRequest{Topic: tc.topicPath("missing")}) + require.Error(t, err) + tc.missingErrMsg = err.Error() + + err = sc.DeleteSnapshot(ctx, &pubsubpb.DeleteSnapshotRequest{Snapshot: tc.snapshotPath(adminSnapshot)}) + require.NoError(t, err) + + err = sc.DeleteSubscription(ctx, &pubsubpb.DeleteSubscriptionRequest{Subscription: tc.subscriptionPath(adminSubscription)}) + require.NoError(t, err) + + err = pc.DeleteTopic(ctx, &pubsubpb.DeleteTopicRequest{Topic: tc.topicPath(adminTopic)}) + require.NoError(t, err) +} + +func (tc *TestCaseAdminGAPIC) ExpectedTraces() trace.Traces { + return trace.Traces{ + tc.adminTrace("ListTopics", tc.projectPath()), + tc.adminTrace("ListSubscriptions", tc.projectPath()), + tc.adminTrace("CreateSnapshot", tc.snapshotPath(adminSnapshot)), + tc.adminTrace("ListSnapshots", tc.projectPath()), + tc.adminTrace("GetTopic", tc.topicPath(adminTopic)), + tc.adminErrorTrace("GetTopic", tc.topicPath("missing"), tc.missingErrMsg), + tc.adminTrace("DeleteSnapshot", tc.snapshotPath(adminSnapshot)), + tc.adminTrace("DeleteSubscription", tc.subscriptionPath(adminSubscription)), + tc.adminTrace("DeleteTopic", tc.topicPath(adminTopic)), + } +} + +// TestCaseAdminSchema exercises admin ops via a directly-constructed GAPIC +// SchemaClient (not constructed via pubsub.NewClient). +type TestCaseAdminSchema struct { + adminBase +} + +func (tc *TestCaseAdminSchema) Setup(ctx context.Context, t *testing.T) { + tc.setup(ctx, t) +} + +func (tc *TestCaseAdminSchema) Run(ctx context.Context, t *testing.T) { + schemaClient, err := vkit.NewSchemaClient(ctx, emulatorOptions(tc.uri)...) + require.NoError(t, err) + defer func() { assert.NoError(t, schemaClient.Close()) }() + + _, err = schemaClient.CreateSchema(ctx, &pubsubpb.CreateSchemaRequest{ + Parent: tc.projectPath(), + Schema: &pubsubpb.Schema{ + Type: pubsubpb.Schema_AVRO, + Definition: avroDefinition, + }, + SchemaId: adminSchema, + }) + require.NoError(t, err) + + _, err = schemaClient.GetSchema(ctx, &pubsubpb.GetSchemaRequest{Name: tc.schemaPath(adminSchema)}) + require.NoError(t, err) + + drain(t, schemaClient.ListSchemas(ctx, &pubsubpb.ListSchemasRequest{Parent: tc.projectPath()}).Next) + + err = schemaClient.DeleteSchema(ctx, &pubsubpb.DeleteSchemaRequest{Name: tc.schemaPath(adminSchema)}) + require.NoError(t, err) +} + +func (tc *TestCaseAdminSchema) ExpectedTraces() trace.Traces { + return trace.Traces{ + tc.adminTrace("CreateSchema", tc.projectPath()), + tc.adminTrace("GetSchema", tc.schemaPath(adminSchema)), + tc.adminTrace("ListSchemas", tc.projectPath()), + tc.adminTrace("DeleteSchema", tc.schemaPath(adminSchema)), + } +} diff --git a/internal/orchestrion/_integration/gcp_pubsub.v1/generated_test.go b/internal/orchestrion/_integration/gcp_pubsub.v1/generated_test.go index b50b1e5f8af..dd73eb14ab4 100644 --- a/internal/orchestrion/_integration/gcp_pubsub.v1/generated_test.go +++ b/internal/orchestrion/_integration/gcp_pubsub.v1/generated_test.go @@ -18,3 +18,15 @@ import ( func Test(t *testing.T) { harness.Run(t, new(TestCase)) } + +func TestAdminClient(t *testing.T) { + harness.Run(t, new(TestCaseAdminClient)) +} + +func TestAdminGAPIC(t *testing.T) { + harness.Run(t, new(TestCaseAdminGAPIC)) +} + +func TestAdminSchema(t *testing.T) { + harness.Run(t, new(TestCaseAdminSchema)) +} diff --git a/internal/orchestrion/_integration/gcp_pubsub.v2/gcp_pubsub_admin.go b/internal/orchestrion/_integration/gcp_pubsub.v2/gcp_pubsub_admin.go new file mode 100644 index 00000000000..a939cb264df --- /dev/null +++ b/internal/orchestrion/_integration/gcp_pubsub.v2/gcp_pubsub_admin.go @@ -0,0 +1,292 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2023-present Datadog, Inc. + +//go:build linux || !githubci + +package gcppubsub + +import ( + "context" + "fmt" + "testing" + + "cloud.google.com/go/pubsub/v2" + vkit "cloud.google.com/go/pubsub/v2/apiv1" + "cloud.google.com/go/pubsub/v2/apiv1/pubsubpb" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/testcontainers/testcontainers-go" + tclog "github.com/testcontainers/testcontainers-go/log" + "github.com/testcontainers/testcontainers-go/modules/gcloud" + "google.golang.org/api/iterator" + "google.golang.org/api/option" + "google.golang.org/api/option/internaloption" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + + "github.com/DataDog/dd-trace-go/instrumentation/testutils/containers/v2" + + "github.com/DataDog/dd-trace-go/v2/ddtrace/ext" + "github.com/DataDog/dd-trace-go/v2/internal/orchestrion/_integration/internal/trace" +) + +const ( + adminTopic = "admin-topic" + adminSubscription = "admin-subscription" + adminSnapshot = "admin-snapshot" + adminSchema = "admin-schema" + avroDefinition = `{"type":"record","name":"Test","fields":[{"name":"f","type":"string"}]}` +) + +type adminBase struct { + container *gcloud.GCloudContainer + uri string + projectID string +} + +// emulatorOptions returns the client options needed to reach the plaintext Pub/Sub +// emulator by dialing internally. Passing these (rather than a pre-built +// option.WithGRPCConn) lets orchestrion append the admin interceptor to the GAPIC +// client constructors it hooks. +func emulatorOptions(uri string) []option.ClientOption { + return []option.ClientOption{ + option.WithEndpoint(uri), + option.WithGRPCDialOption(grpc.WithTransportCredentials(insecure.NewCredentials())), + option.WithoutAuthentication(), + internaloption.SkipDialSettingsValidation(), + } +} + +func (b *adminBase) setup(ctx context.Context, t *testing.T) { + containers.SkipIfProviderIsNotHealthy(t) + + var err error + b.container, err = gcloud.RunPubsub(ctx, + "gcr.io/google.com/cloudsdktool/google-cloud-cli:emulators", + gcloud.WithProjectID(testProject), + testcontainers.WithLogger(tclog.TestLogger(t)), + containers.WithTestLogConsumer(t), + ) + containers.AssertTestContainersError(t, err) + containers.RegisterContainerCleanup(t, b.container) + + b.projectID = b.container.Settings.ProjectID + b.uri = b.container.URI +} + +func (b *adminBase) projectPath() string { + return "projects/" + b.projectID +} + +func (b *adminBase) topicPath(id string) string { + return fmt.Sprintf("projects/%s/topics/%s", b.projectID, id) +} + +func (b *adminBase) subscriptionPath(id string) string { + return fmt.Sprintf("projects/%s/subscriptions/%s", b.projectID, id) +} + +func (b *adminBase) snapshotPath(id string) string { + return fmt.Sprintf("projects/%s/snapshots/%s", b.projectID, id) +} + +func (b *adminBase) schemaPath(id string) string { + return fmt.Sprintf("projects/%s/schemas/%s", b.projectID, id) +} + +func drain[T any](t *testing.T, next func() (T, error)) { + t.Helper() + for { + if _, err := next(); err == iterator.Done { + return + } else { + require.NoError(t, err) + } + } +} + +func (b *adminBase) adminTrace(method, resourcePath string) *trace.Trace { + return &trace.Trace{ + Tags: map[string]any{ + "name": "gcp.pubsub.request", + "type": "worker", + "resource": method + " " + resourcePath, + "service": "gcp_pubsub.v2.test", + }, + Meta: map[string]string{ + "span.kind": "client", + "component": "cloud.google.com/go/pubsub.v2", + "messaging.system": "googlepubsub", + "pubsub.method": method, + ext.GCPProjectID: testProject, + }, + } +} + +func (b *adminBase) adminErrorTrace(method, resourcePath, errMsg string) *trace.Trace { + tr := b.adminTrace(method, resourcePath) + tr.Meta["error.message"] = errMsg + return tr +} + +// TestCaseAdminClient exercises admin ops via admin clients exposed by the +// high-level pubsub.Client. These are constructed internally by pubsub.NewClient +// (a variadic-spread call). +type TestCaseAdminClient struct { + adminBase + client *pubsub.Client +} + +func (tc *TestCaseAdminClient) Setup(ctx context.Context, t *testing.T) { + tc.setup(ctx, t) + + var err error + tc.client, err = pubsub.NewClient(ctx, tc.projectID, emulatorOptions(tc.uri)...) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, tc.client.Close()) }) +} + +func (tc *TestCaseAdminClient) Run(ctx context.Context, t *testing.T) { + topic, err := tc.client.TopicAdminClient.CreateTopic(ctx, &pubsubpb.Topic{ + Name: tc.topicPath(adminTopic), + }) + require.NoError(t, err) + + _, err = tc.client.SubscriptionAdminClient.CreateSubscription(ctx, &pubsubpb.Subscription{ + Name: tc.subscriptionPath(adminSubscription), + Topic: topic.Name, + }) + require.NoError(t, err) + + drain(t, tc.client.TopicAdminClient.ListTopics(ctx, &pubsubpb.ListTopicsRequest{ + Project: tc.projectPath(), + }).Next) + drain(t, tc.client.SubscriptionAdminClient.ListSubscriptions(ctx, &pubsubpb.ListSubscriptionsRequest{ + Project: tc.projectPath(), + }).Next) + + _, err = tc.client.SubscriptionAdminClient.CreateSnapshot(ctx, &pubsubpb.CreateSnapshotRequest{ + Name: tc.snapshotPath(adminSnapshot), + Subscription: tc.subscriptionPath(adminSubscription), + }) + require.NoError(t, err) + + drain(t, tc.client.SubscriptionAdminClient.ListSnapshots(ctx, &pubsubpb.ListSnapshotsRequest{ + Project: tc.projectPath(), + }).Next) + + err = tc.client.SubscriptionAdminClient.DeleteSnapshot(ctx, &pubsubpb.DeleteSnapshotRequest{ + Snapshot: tc.snapshotPath(adminSnapshot), + }) + require.NoError(t, err) + + err = tc.client.SubscriptionAdminClient.DeleteSubscription(ctx, &pubsubpb.DeleteSubscriptionRequest{ + Subscription: tc.subscriptionPath(adminSubscription), + }) + require.NoError(t, err) + + err = tc.client.TopicAdminClient.DeleteTopic(ctx, &pubsubpb.DeleteTopicRequest{ + Topic: tc.topicPath(adminTopic), + }) + require.NoError(t, err) +} + +func (tc *TestCaseAdminClient) ExpectedTraces() trace.Traces { + return trace.Traces{ + tc.adminTrace("CreateTopic", tc.topicPath(adminTopic)), + tc.adminTrace("CreateSubscription", tc.subscriptionPath(adminSubscription)), + tc.adminTrace("ListTopics", tc.projectPath()), + tc.adminTrace("ListSubscriptions", tc.projectPath()), + tc.adminTrace("CreateSnapshot", tc.snapshotPath(adminSnapshot)), + tc.adminTrace("ListSnapshots", tc.projectPath()), + tc.adminTrace("DeleteSnapshot", tc.snapshotPath(adminSnapshot)), + tc.adminTrace("DeleteSubscription", tc.subscriptionPath(adminSubscription)), + tc.adminTrace("DeleteTopic", tc.topicPath(adminTopic)), + } +} + +// TestCaseAdminGAPIC exercises admin ops via a directly-constructed GAPIC +// TopicAdminClient. +type TestCaseAdminGAPIC struct { + adminBase + missingErrMsg string +} + +func (tc *TestCaseAdminGAPIC) Setup(ctx context.Context, t *testing.T) { + tc.setup(ctx, t) + + // Fixture topic is created before the tracer starts. + client, err := pubsub.NewClient(ctx, tc.projectID, emulatorOptions(tc.uri)...) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, client.Close()) }) + + _, err = client.TopicAdminClient.CreateTopic(ctx, &pubsubpb.Topic{ + Name: tc.topicPath(adminTopic), + }) + require.NoError(t, err) +} + +func (tc *TestCaseAdminGAPIC) Run(ctx context.Context, t *testing.T) { + tac, err := vkit.NewTopicAdminClient(ctx, emulatorOptions(tc.uri)...) + require.NoError(t, err) + defer func() { assert.NoError(t, tac.Close()) }() + + _, err = tac.GetTopic(ctx, &pubsubpb.GetTopicRequest{Topic: tc.topicPath(adminTopic)}) + require.NoError(t, err) + + _, err = tac.GetTopic(ctx, &pubsubpb.GetTopicRequest{Topic: tc.topicPath("missing")}) + require.Error(t, err) + tc.missingErrMsg = err.Error() +} + +func (tc *TestCaseAdminGAPIC) ExpectedTraces() trace.Traces { + return trace.Traces{ + tc.adminTrace("GetTopic", tc.topicPath(adminTopic)), + tc.adminErrorTrace("GetTopic", tc.topicPath("missing"), tc.missingErrMsg), + } +} + +// TestCaseAdminSchema exercises admin ops via a directly-constructed GAPIC +// SchemaClient (not exposed on pubsub.Client). +type TestCaseAdminSchema struct { + adminBase +} + +func (tc *TestCaseAdminSchema) Setup(ctx context.Context, t *testing.T) { + tc.setup(ctx, t) +} + +func (tc *TestCaseAdminSchema) Run(ctx context.Context, t *testing.T) { + schemaClient, err := vkit.NewSchemaClient(ctx, emulatorOptions(tc.uri)...) + require.NoError(t, err) + defer func() { assert.NoError(t, schemaClient.Close()) }() + + _, err = schemaClient.CreateSchema(ctx, &pubsubpb.CreateSchemaRequest{ + Parent: tc.projectPath(), + Schema: &pubsubpb.Schema{ + Type: pubsubpb.Schema_AVRO, + Definition: avroDefinition, + }, + SchemaId: adminSchema, + }) + require.NoError(t, err) + + _, err = schemaClient.GetSchema(ctx, &pubsubpb.GetSchemaRequest{Name: tc.schemaPath(adminSchema)}) + require.NoError(t, err) + + drain(t, schemaClient.ListSchemas(ctx, &pubsubpb.ListSchemasRequest{Parent: tc.projectPath()}).Next) + + err = schemaClient.DeleteSchema(ctx, &pubsubpb.DeleteSchemaRequest{Name: tc.schemaPath(adminSchema)}) + require.NoError(t, err) +} + +func (tc *TestCaseAdminSchema) ExpectedTraces() trace.Traces { + return trace.Traces{ + tc.adminTrace("CreateSchema", tc.projectPath()), + tc.adminTrace("GetSchema", tc.schemaPath(adminSchema)), + tc.adminTrace("ListSchemas", tc.projectPath()), + tc.adminTrace("DeleteSchema", tc.schemaPath(adminSchema)), + } +} diff --git a/internal/orchestrion/_integration/gcp_pubsub.v2/generated_test.go b/internal/orchestrion/_integration/gcp_pubsub.v2/generated_test.go index b50b1e5f8af..dd73eb14ab4 100644 --- a/internal/orchestrion/_integration/gcp_pubsub.v2/generated_test.go +++ b/internal/orchestrion/_integration/gcp_pubsub.v2/generated_test.go @@ -18,3 +18,15 @@ import ( func Test(t *testing.T) { harness.Run(t, new(TestCase)) } + +func TestAdminClient(t *testing.T) { + harness.Run(t, new(TestCaseAdminClient)) +} + +func TestAdminGAPIC(t *testing.T) { + harness.Run(t, new(TestCaseAdminGAPIC)) +} + +func TestAdminSchema(t *testing.T) { + harness.Run(t, new(TestCaseAdminSchema)) +}