diff --git a/pkg/services/resource.go b/pkg/services/resource.go index 33783391..1c5d724d 100644 --- a/pkg/services/resource.go +++ b/pkg/services/resource.go @@ -7,6 +7,9 @@ import ( "strings" "time" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" + "github.com/openshift-hyperfleet/hyperfleet-api/pkg/api" "github.com/openshift-hyperfleet/hyperfleet-api/pkg/dao" "github.com/openshift-hyperfleet/hyperfleet-api/pkg/db" @@ -64,6 +67,10 @@ func (s *sqlResourceService) Get(ctx context.Context, kind, id string) (*api.Res if svcErr := validateKind(kind); svcErr != nil { return nil, svcErr } + trace.SpanFromContext(ctx).SetAttributes( + attribute.String("hyperfleet.resource_id", id), + attribute.String("hyperfleet.resource_type", registry.MustGet(kind).Plural), + ) resource, err := s.resourceDao.Get(ctx, kind, id) if err != nil { return nil, handleGetError(kind, "id", id, err) @@ -83,7 +90,14 @@ func (s *sqlResourceService) Create( ) (*api.Resource, *errors.ServiceError) { resource.Kind = kind - if svcErr := validateResourceName(kind, resource.Name); svcErr != nil { + if svcErr := validateKind(kind); svcErr != nil { + return nil, svcErr + } + trace.SpanFromContext(ctx).SetAttributes( + attribute.String("hyperfleet.resource_type", registry.MustGet(kind).Plural), + ) + + if svcErr := validateName(kind, resource.Name); svcErr != nil { return nil, svcErr } @@ -115,6 +129,7 @@ func (s *sqlResourceService) Create( if err != nil { return nil, handleCreateError(kind, err) } + trace.SpanFromContext(ctx).SetAttributes(attribute.String("hyperfleet.resource_id", resource.ID)) if len(resource.Labels) > 0 { if labelErr := s.resourceLabelDao.ReplaceLabels(ctx, resource.ID, resource.Labels); labelErr != nil { @@ -151,6 +166,10 @@ func (s *sqlResourceService) Patch( if svcErr := validateKind(kind); svcErr != nil { return nil, svcErr } + trace.SpanFromContext(ctx).SetAttributes( + attribute.String("hyperfleet.resource_id", id), + attribute.String("hyperfleet.resource_type", registry.MustGet(kind).Plural), + ) resource, err := s.resourceDao.GetForUpdate(ctx, kind, id) if err != nil { return nil, handleGetError(kind, "id", id, err) @@ -226,6 +245,10 @@ func (s *sqlResourceService) Delete(ctx context.Context, kind, id string) (*api. if svcErr := validateKind(kind); svcErr != nil { return nil, svcErr } + trace.SpanFromContext(ctx).SetAttributes( + attribute.String("hyperfleet.resource_id", id), + attribute.String("hyperfleet.resource_type", registry.MustGet(kind).Plural), + ) resource, err := s.resourceDao.GetForUpdate(ctx, kind, id) if err != nil { return nil, handleSoftDeleteError(kind, err) @@ -386,9 +409,11 @@ func (s *sqlResourceService) checkCanDelete( func (s *sqlResourceService) GetByOwner( ctx context.Context, kind, id, ownerID string, ) (*api.Resource, *errors.ServiceError) { + trace.SpanFromContext(ctx).SetAttributes(attribute.String("hyperfleet.resource_id", id)) if svcErr := validateKind(kind); svcErr != nil { return nil, svcErr } + trace.SpanFromContext(ctx).SetAttributes(attribute.String("hyperfleet.resource_type", registry.MustGet(kind).Plural)) resource, err := s.resourceDao.GetByOwner(ctx, kind, id, ownerID) if err != nil { return nil, handleGetError(kind, "id", id, err) @@ -464,10 +489,14 @@ func (s *sqlResourceService) ListByOwner( } func (s *sqlResourceService) GetByID(ctx context.Context, id string) (*api.Resource, *errors.ServiceError) { + trace.SpanFromContext(ctx).SetAttributes(attribute.String("hyperfleet.resource_id", id)) resource, err := s.resourceDao.GetByID(ctx, id) if err != nil { return nil, handleGetError("Resource", "id", id, err) } + trace.SpanFromContext(ctx).SetAttributes( + attribute.String("hyperfleet.resource_type", registry.MustGet(resource.Kind).Plural), + ) return resource, nil } @@ -501,9 +530,11 @@ func (s *sqlResourceService) ListAll( func (s *sqlResourceService) ProcessAdapterStatus( ctx context.Context, kind, resourceID string, adapterStatus *api.AdapterStatus, ) (*api.AdapterStatus, *errors.ServiceError) { + trace.SpanFromContext(ctx).SetAttributes(attribute.String("hyperfleet.resource_id", resourceID)) if svcErr := validateKind(kind); svcErr != nil { return nil, svcErr } + trace.SpanFromContext(ctx).SetAttributes(attribute.String("hyperfleet.resource_type", registry.MustGet(kind).Plural)) // Step 1: Acquire a row-level exclusive lock on the resource. Concurrent // adapter status updates for the same resource are serialized here. @@ -797,10 +828,7 @@ func validateKind(kind string) *errors.ServiceError { } // Name format/length validation is handled by OpenAPI spec validation middleware. -func validateResourceName(kind, name string) *errors.ServiceError { - if svcErr := validateKind(kind); svcErr != nil { - return svcErr - } +func validateName(kind, name string) *errors.ServiceError { if name == "" { return errors.Validation("%s name cannot be empty", kind) } @@ -861,9 +889,11 @@ func applyResourcePatch(resource *api.Resource, patch *api.ResourcePatch) error } func (s *sqlResourceService) ForceDelete(ctx context.Context, kind, id, reason string) *errors.ServiceError { + trace.SpanFromContext(ctx).SetAttributes(attribute.String("hyperfleet.resource_id", id)) if svcErr := validateKind(kind); svcErr != nil { return svcErr } + trace.SpanFromContext(ctx).SetAttributes(attribute.String("hyperfleet.resource_type", registry.MustGet(kind).Plural)) resource, err := s.resourceDao.GetForUpdate(ctx, kind, id) if err != nil { diff --git a/pkg/services/resource_test.go b/pkg/services/resource_test.go index 13a7bb55..81cf2662 100644 --- a/pkg/services/resource_test.go +++ b/pkg/services/resource_test.go @@ -9,6 +9,8 @@ import ( "time" . "github.com/onsi/gomega" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" "gorm.io/gorm" "github.com/openshift-hyperfleet/hyperfleet-api/pkg/api" @@ -3142,3 +3144,235 @@ func TestProcessAdapterStatus_FinalizedTrue_RecomputesConditions_WhenHardDeleteB Expect(*recon.Reason).To(ContainSubstring("Children"), "Reason should indicate waiting for child resources") } + +// --- Span attribute tests --- + +func setupTestTracer(t *testing.T) (*sdktrace.TracerProvider, *tracetest.InMemoryExporter) { + t.Helper() + exporter := tracetest.NewInMemoryExporter() + tp := sdktrace.NewTracerProvider( + sdktrace.WithSampler(sdktrace.AlwaysSample()), + sdktrace.WithSyncer(exporter), + ) + t.Cleanup(func() { + if err := tp.Shutdown(context.Background()); err != nil { + t.Errorf("failed to shutdown tracer: %v", err) + } + }) + return tp, exporter +} + +func findSpanAttribute(spans tracetest.SpanStubs, attrKey string) (string, bool) { + for _, span := range spans { + for _, attr := range span.Attributes { + if string(attr.Key) == attrKey { + return attr.Value.AsString(), true + } + } + } + return "", false +} + +func TestResourceService_SetsSpanAttributes(t *testing.T) { + setupTestDescriptors() + + tests := []struct { + invoke func(ctx context.Context, svc ResourceService) error + name string + seedID string + expectedResourceID string + expectedType string + expectError bool + }{ + { + name: "Get", + seedID: "ch-1", + expectedResourceID: "ch-1", + expectedType: "channels", + invoke: func(ctx context.Context, svc ResourceService) error { + _, err := svc.Get(ctx, "Channel", "ch-1") + return svcErrOrNil(err) + }, + }, + { + name: "Get not found still tags span", + expectedResourceID: "nonexistent", + expectedType: "channels", + expectError: true, + invoke: func(ctx context.Context, svc ResourceService) error { + _, err := svc.Get(ctx, "Channel", "nonexistent") + return svcErrOrNil(err) + }, + }, + { + name: "Create", + expectedResourceID: "ch-new", + expectedType: "channels", + invoke: func(ctx context.Context, svc ResourceService) error { + _, err := svc.Create(ctx, "Channel", testResource("Channel", "ch-new", "beta"), nil) + return svcErrOrNil(err) + }, + }, + { + name: "Create failure still tags resource_type", + expectedType: "channels", + expectError: true, + invoke: func(ctx context.Context, svc ResourceService) error { + r := testResource("Channel", "", "") + r.Name = "" // triggers name validation error + _, err := svc.Create(ctx, "Channel", r, nil) + return svcErrOrNil(err) + }, + }, + { + name: "Patch", + seedID: "ch-1", + expectedResourceID: "ch-1", + expectedType: "channels", + invoke: func(ctx context.Context, svc ResourceService) error { + _, err := svc.Patch(ctx, "Channel", "ch-1", &api.ResourcePatch{ + Spec: map[string]interface{}{"key": "updated"}, + }) + return svcErrOrNil(err) + }, + }, + { + name: "Delete", + seedID: "ch-1", + expectedResourceID: "ch-1", + expectedType: "channels", + invoke: func(ctx context.Context, svc ResourceService) error { + _, err := svc.Delete(ctx, "Channel", "ch-1") + return svcErrOrNil(err) + }, + }, + { + name: "GetByID", + seedID: "ch-1", + expectedResourceID: "ch-1", + expectedType: "channels", + invoke: func(ctx context.Context, svc ResourceService) error { + _, err := svc.GetByID(ctx, "ch-1") + return svcErrOrNil(err) + }, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + RegisterTestingT(t) + tp, exporter := setupTestTracer(t) + + mockDao := newMockResourceDao() + svc, _, _ := newTestResourceService(mockDao) + if tc.seedID != "" { + mockDao.addResource(testResource("Channel", tc.seedID, "stable")) + } + + ctx, span := tp.Tracer("test").Start(context.Background(), "test") + err := tc.invoke(ctx, svc) + span.End() + if tc.expectError { + Expect(err).ToNot(BeNil()) + } else { + Expect(err).To(BeNil()) + } + + if flushErr := tp.ForceFlush(context.Background()); flushErr != nil { + t.Fatalf("failed to flush spans: %v", flushErr) + } + spans := exporter.GetSpans() + + if tc.expectedResourceID != "" { + resourceID, found := findSpanAttribute(spans, "hyperfleet.resource_id") + Expect(found).To(BeTrue(), "hyperfleet.resource_id attribute not found") + Expect(resourceID).To(Equal(tc.expectedResourceID)) + } + + resourceType, found := findSpanAttribute(spans, "hyperfleet.resource_type") + Expect(found).To(BeTrue(), "hyperfleet.resource_type attribute not found") + Expect(resourceType).To(Equal(tc.expectedType)) + }) + } +} + +func svcErrOrNil(err *errors.ServiceError) error { + if err != nil { + return err + } + return nil +} + +func assertSpanAttributes( + t *testing.T, tp *sdktrace.TracerProvider, exporter *tracetest.InMemoryExporter, + expectedID, expectedType string, +) { + t.Helper() + if err := tp.ForceFlush(context.Background()); err != nil { + t.Fatalf("failed to flush spans: %v", err) + } + spans := exporter.GetSpans() + + resourceID, found := findSpanAttribute(spans, "hyperfleet.resource_id") + Expect(found).To(BeTrue(), "hyperfleet.resource_id attribute not found") + Expect(resourceID).To(Equal(expectedID)) + + resourceType, found := findSpanAttribute(spans, "hyperfleet.resource_type") + Expect(found).To(BeTrue(), "hyperfleet.resource_type attribute not found") + Expect(resourceType).To(Equal(expectedType)) +} + +func TestResourceService_GetByOwner_SetsSpanAttributes(t *testing.T) { + RegisterTestingT(t) + setupTestDescriptors() + tp, exporter := setupTestTracer(t) + + mockDao := newMockResourceDao() + svc, _, _ := newTestResourceService(mockDao) + r := testResource("Version", "v-1", "1.0") + r.OwnerID = strPtr("ch-1") + mockDao.addResource(r) + + ctx, span := tp.Tracer("test").Start(context.Background(), "test") + _, svcErr := svc.GetByOwner(ctx, "Version", "v-1", "ch-1") + span.End() + Expect(svcErr).To(BeNil()) + assertSpanAttributes(t, tp, exporter, "v-1", "versions") +} + +func TestResourceService_ForceDelete_SetsSpanAttributes(t *testing.T) { + RegisterTestingT(t) + setupTestDescriptors() + tp, exporter := setupTestTracer(t) + + mockDao := newMockResourceDao() + svc, _, _ := newTestResourceService(mockDao) + r := testResource("Channel", "ch-1", "stable") + now := time.Now().UTC() + r.DeletedTime = &now + mockDao.addResource(r) + + ctx, span := tp.Tracer("test").Start(context.Background(), "test") + svcErr := svc.ForceDelete(ctx, "Channel", "ch-1", "test cleanup") + span.End() + Expect(svcErr).To(BeNil()) + assertSpanAttributes(t, tp, exporter, "ch-1", "channels") +} + +func TestResourceService_ProcessAdapterStatus_SetsSpanAttributes(t *testing.T) { + RegisterTestingT(t) + setupAdapterStatusDescriptors() + tp, exporter := setupTestTracer(t) + + mockDao := newMockResourceDao() + svc, _, _, _ := newTestResourceServiceWithAdapterStatus(mockDao) + r := testResource("TestResource", "r-1", "test") + r.Generation = 1 + mockDao.addResource(r) + + ctx, span := tp.Tracer("test").Start(context.Background(), "test") + _, svcErr := svc.ProcessAdapterStatus(ctx, "TestResource", "r-1", testAdapterStatusRequest(1)) + span.End() + Expect(svcErr).To(BeNil()) + assertSpanAttributes(t, tp, exporter, "r-1", "testresources") +}