Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 35 additions & 5 deletions pkg/services/resource.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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),
)
Comment on lines +70 to +73

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

hyperfleet.resource_type is set to the Kind (e.g. "Cluster") here and at the other three call sites, but the Tracing Standard and Sentinel's existing spans use the plural form (e.g. "clusters"). This mismatch means a TraceQL query on resource_type won't match across API and Sentinel spans for the same resource. Consider using the descriptor's plural form instead.

resource, err := s.resourceDao.Get(ctx, kind, id)
if err != nil {
return nil, handleGetError(kind, "id", id, err)
Expand All @@ -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
}

Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
}

Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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 {
Expand Down
234 changes: 234 additions & 0 deletions pkg/services/resource_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
}
})
Comment thread
coderabbitai[bot] marked this conversation as resolved.
return tp, exporter
Comment on lines +3150 to +3162

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

setupTestTracer swaps the global otel TracerProvider for the test, then each test starts its span via the global otel.Tracer rather than the tp returned by this helper. Using tp.Tracer(...) directly would avoid touching process-global state, which would remove a potential source of flakiness if t.Parallel() is ever added to these tests.

}

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")
}