Skip to content
Draft
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
36 changes: 35 additions & 1 deletion .github/workflows/checks.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -532,7 +532,41 @@ jobs:

- name: Run BDD Tests
run: |
CUKES_LOG_HANDLER=console go test ./tests-bdd -v --tags=cukes --godog.random --godog.format="cucumber:$(pwd)/cukes_platform_report.json,pretty:$(pwd)/cukes_platform_report.log,pretty" ./features
CUKES_LOG_HANDLER=console go test ./tests-bdd -v --tags=cukes --godog.random --godog.format="cucumber:$(pwd)/cukes_platform_report.json,pretty:$(pwd)/cukes_platform_report.log,pretty" ./features 2>&1 | tee cukes_test_output.log

- name: Summarize authorization performance
if: ${{ !cancelled() }}
run: |
PERFORMANCE_ROWS=$(awk '
/authorization v2 concurrent multi-resource performance/ {
delete values
for (i = 1; i <= NF; i++) {
split($i, pair, "=")
values[pair[1]] = pair[2]
}
printf "%d\t%s\t%s\t%s\t%s\t%d\n",
values["concurrency"],
values["wall_duration"],
values["median_request_duration"],
values["p95_request_duration"],
values["maximum_request_duration"],
values["failed_request_count"]
}
' cukes_test_output.log | sort -n | awk -F '\t' '{
printf "| %s | `%s` | `%s` | `%s` | `%s` | %s |\n", $1, $2, $3, $4, $5, $6
}')

{
echo "### Authorization v2 concurrency performance"
echo
if [[ -z "$PERFORMANCE_ROWS" ]]; then
echo "No authorization performance scenarios completed."
else
echo "| Concurrency | Wall time | Median request | p95 request | Maximum request | Failures |"
echo "| ---: | ---: | ---: | ---: | ---: | ---: |"
echo "$PERFORMANCE_ROWS"
fi
} >> "$GITHUB_STEP_SUMMARY"

- name: Check for undefined steps
run: |
Expand Down
110 changes: 100 additions & 10 deletions tests-bdd/cukes/steps_attributes.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,10 @@ import (
"context"
"errors"
"fmt"
"log/slog"
"strings"
"sync"
"time"

"github.com/cucumber/godog"
"github.com/opentdf/platform/protocol/go/policy"
Expand All @@ -19,6 +22,19 @@ type AttributesStepDefinitions struct {
PlatformCukesContext *PlatformTestSuiteContext
}

func parseAttributeRule(rule string) (policy.AttributeRuleTypeEnum, error) {
switch strings.TrimSpace(rule) {
case "anyOf":
return policy.AttributeRuleTypeEnum_ATTRIBUTE_RULE_TYPE_ENUM_ANY_OF, nil
case "allOf":
return policy.AttributeRuleTypeEnum_ATTRIBUTE_RULE_TYPE_ENUM_ALL_OF, nil
case "hierarchy":
return policy.AttributeRuleTypeEnum_ATTRIBUTE_RULE_TYPE_ENUM_HIERARCHY, nil
default:
return policy.AttributeRuleTypeEnum_ATTRIBUTE_RULE_TYPE_ENUM_UNSPECIFIED, fmt.Errorf("unknown attribute rule type %s", rule)
}
}

func (s *AttributesStepDefinitions) aAttributeDef(ctx context.Context, _ string, _ string) (context.Context, error) {
return ctx, nil
}
Expand Down Expand Up @@ -98,16 +114,9 @@ func (s *AttributesStepDefinitions) iSendARequestToCreateAnAttributeWithGenerate
return ctx, fmt.Errorf("unable to get namespace id for %s", namespaceRef)
}

var ruleType policy.AttributeRuleTypeEnum
switch strings.TrimSpace(rule) {
case "anyOf":
ruleType = policy.AttributeRuleTypeEnum_ATTRIBUTE_RULE_TYPE_ENUM_ANY_OF
case "allOf":
ruleType = policy.AttributeRuleTypeEnum_ATTRIBUTE_RULE_TYPE_ENUM_ALL_OF
case "hierarchy":
ruleType = policy.AttributeRuleTypeEnum_ATTRIBUTE_RULE_TYPE_ENUM_HIERARCHY
default:
return ctx, fmt.Errorf("unknown attribute rule type %s", rule)
ruleType, err := parseAttributeRule(rule)
if err != nil {
return ctx, err
}

values := make([]string, 0, valueCount)
Expand All @@ -128,11 +137,92 @@ func (s *AttributesStepDefinitions) iSendARequestToCreateAnAttributeWithGenerate
return ctx, nil
}

func (s *AttributesStepDefinitions) iSendARequestToCreateAnAttributeWithBatchedGeneratedValues(ctx context.Context, referenceID, namespaceRef, name, rule string, valueCount, batchSize int) (context.Context, error) {
scenarioContext := GetPlatformScenarioContext(ctx)
scenarioContext.ClearError()
if valueCount < 1 {
return ctx, errors.New("generated value count must be positive")
}
if batchSize < 1 {
return ctx, errors.New("generated value batch size must be positive")
}

namespaceID, ok := scenarioContext.GetObject(strings.TrimSpace(namespaceRef)).(string)
if !ok {
return ctx, fmt.Errorf("unable to get namespace id for %s", namespaceRef)
}
ruleType, err := parseAttributeRule(rule)
if err != nil {
return ctx, err
}

created, err := scenarioContext.SDK.Attributes.CreateAttribute(ctx, &attributes.CreateAttributeRequest{
NamespaceId: namespaceID,
Name: strings.TrimSpace(name),
Rule: ruleType,
})
if err != nil {
scenarioContext.SetError(err)
return ctx, nil
}
if created.GetAttribute() == nil {
return ctx, errors.New("create attribute returned no attribute")
}

started := time.Now()
values := make([]*policy.Value, valueCount)
for batchStart := 0; batchStart < valueCount; batchStart += batchSize {
batchEnd := min(batchStart+batchSize, valueCount)
batchCtx, cancel := context.WithCancel(ctx)
errCh := make(chan error, batchEnd-batchStart)
var wg sync.WaitGroup
for i := batchStart; i < batchEnd; i++ {
wg.Add(1)
go func(index int) {
defer wg.Done()
resp, createErr := scenarioContext.SDK.Attributes.CreateAttributeValue(batchCtx, &attributes.CreateAttributeValueRequest{
AttributeId: created.GetAttribute().GetId(),
Value: fmt.Sprintf("v%04d", index),
})
if createErr != nil {
errCh <- fmt.Errorf("create generated attribute value v%04d: %w", index, createErr)
cancel()
return
}
if resp.GetValue() == nil {
errCh <- fmt.Errorf("create generated attribute value v%04d returned no value", index)
cancel()
return
}
values[index] = resp.GetValue()
}(i)
}
wg.Wait()
cancel()
close(errCh)
if batchErr, hasBatchErr := <-errCh; hasBatchErr {
scenarioContext.SetError(batchErr)
return ctx, nil
}
}

created.GetAttribute().Values = values
scenarioContext.RecordObject(strings.TrimSpace(referenceID), created.GetAttribute())
scenarioContext.TestSuiteContext.Logger.Info(
"created generated attribute values in batches",
slog.Int("value_count", valueCount),
slog.Int("batch_size", batchSize),
slog.Duration("duration", time.Since(started)),
)
return ctx, nil
}

func RegisterAttributeStepDefinitions(ctx *godog.ScenarioContext, x *PlatformTestSuiteContext) {
stepDefinitions := AttributesStepDefinitions{
PlatformCukesContext: x,
}
ctx.Step(`^a (anyOf|allOf|hierarchy) attribute definition with values: "([^"]*)"$`, stepDefinitions.aAttributeDef)
ctx.Step(`^I send a request to create an attribute with:$`, stepDefinitions.iSendARequestToCreateAnAttributeWith)
ctx.Step(`^I send a request to create an attribute referenced as "([^"]*)" in namespace "([^"]*)" named "([^"]*)" with rule "([^"]*)" and (\d+) generated values$`, stepDefinitions.iSendARequestToCreateAnAttributeWithGeneratedValues)
ctx.Step(`^I send a request to create an attribute referenced as "([^"]*)" in namespace "([^"]*)" named "([^"]*)" with rule "([^"]*)" and (\d+) generated values in batches of (\d+)$`, stepDefinitions.iSendARequestToCreateAnAttributeWithBatchedGeneratedValues)
}
161 changes: 161 additions & 0 deletions tests-bdd/cukes/steps_authorization.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,20 @@ import (
"context"
"errors"
"fmt"
"log/slog"
"net"
"sort"
"strconv"
"strings"
"sync"
"time"

"github.com/cucumber/godog"
authzV2 "github.com/opentdf/platform/protocol/go/authorization/v2"
"github.com/opentdf/platform/protocol/go/entity"
"github.com/opentdf/platform/protocol/go/policy"
"google.golang.org/protobuf/encoding/protojson"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/anypb"
"google.golang.org/protobuf/types/known/structpb"
)
Expand Down Expand Up @@ -214,6 +219,160 @@ func (s *AuthorizationServiceStepDefinitions) iSendAMultiResourceDecisionRequest
return ctx, nil
}

func (s *AuthorizationServiceStepDefinitions) iSendAMultiResourceDecisionRequestWithin(ctx context.Context, entityChainID, action, maximumDuration string, tbl *godog.Table) (context.Context, error) {
limit, err := time.ParseDuration(maximumDuration)
if err != nil {
return ctx, fmt.Errorf("parse maximum decision duration %q: %w", maximumDuration, err)
}
if limit <= 0 {
return ctx, fmt.Errorf("maximum decision duration must be positive, got %s", limit)
}

requestCtx, cancel := context.WithTimeout(ctx, limit)
defer cancel()
started := time.Now()
_, err = s.iSendAMultiResourceDecisionRequestForEntityChainForActionOnResources(requestCtx, entityChainID, action, tbl)
elapsed := time.Since(started)

scenarioContext := GetPlatformScenarioContext(ctx)
scenarioContext.TestSuiteContext.Logger.Info(
"authorization v2 multi-resource performance",
slog.Duration("duration", elapsed),
slog.Duration("maximum_duration", limit),
slog.Int("resource_count", len(tbl.Rows)-1),
)
if err != nil {
return ctx, err
}
if requestErr := scenarioContext.GetError(); requestErr != nil {
return ctx, fmt.Errorf("multi-resource decision failed after %s: %w", elapsed, requestErr)
}
if elapsed > limit {
return ctx, fmt.Errorf("multi-resource decision took %s, exceeding the %s limit", elapsed, limit)
}

return ctx, nil
}

func (s *AuthorizationServiceStepDefinitions) iSendConcurrentMultiResourceDecisionRequestsWithin(ctx context.Context, concurrency int, entityChainID, action, maximumDuration string, tbl *godog.Table) (context.Context, error) {
if concurrency <= 0 {
return ctx, fmt.Errorf("concurrency must be positive, got %d", concurrency)
}

limit, err := time.ParseDuration(maximumDuration)
if err != nil {
return ctx, fmt.Errorf("parse maximum decision duration %q: %w", maximumDuration, err)
}
if limit <= 0 {
return ctx, fmt.Errorf("maximum decision duration must be positive, got %s", limit)
}

scenarioContext := GetPlatformScenarioContext(ctx)
scenarioContext.ClearError()
entityChain, err := buildEntityChainFromIDs(scenarioContext, entityChainID)
if err != nil {
return ctx, err
}
resources, resourceFQNMap, err := buildResourcesFromTable(tbl)
if err != nil {
return ctx, err
}

req := &authzV2.GetDecisionMultiResourceRequest{
EntityIdentifier: &authzV2.EntityIdentifier{
Identifier: &authzV2.EntityIdentifier_EntityChain{EntityChain: entityChain},
},
Action: &policy.Action{Name: strings.ToLower(action)},
Resources: resources,
FulfillableObligationFqns: getAllObligationsFromScenario(scenarioContext),
}

// Keep the measurement gate separate from the timeout so slow requests still
// produce useful latency results before the scenario fails.
const minimumRequestTimeout = 5 * time.Second
requestCtx, cancel := context.WithTimeout(ctx, max(limit, minimumRequestTimeout))
defer cancel()
start := make(chan struct{})
responses := make([]*authzV2.GetDecisionMultiResourceResponse, concurrency)
requestErrors := make([]error, concurrency)
durations := make([]time.Duration, concurrency)

var wg sync.WaitGroup
wg.Add(concurrency)
for i := range concurrency {
go func() {
defer wg.Done()
<-start
started := time.Now()
clonedReq, ok := proto.Clone(req).(*authzV2.GetDecisionMultiResourceRequest)
if !ok {
requestErrors[i] = errors.New("clone multi-resource decision request")
return
}
responses[i], requestErrors[i] = scenarioContext.SDK.AuthorizationV2.GetDecisionMultiResource(requestCtx, clonedReq)
durations[i] = time.Since(started)
}()
}

wallStarted := time.Now()
close(start)
wg.Wait()
wallElapsed := time.Since(wallStarted)
sortedDurations := append([]time.Duration(nil), durations...)
sort.Slice(sortedDurations, func(i, j int) bool { return sortedDurations[i] < sortedDurations[j] })
maximumRequestDuration := sortedDurations[len(sortedDurations)-1]
medianRequestDuration := sortedDurations[(len(sortedDurations)-1)/2]
p95RequestDuration := sortedDurations[(95*len(sortedDurations)-1)/100]
failedRequestCount := 0
var firstRequestError error
for i, requestErr := range requestErrors {
if requestErr != nil {
failedRequestCount++
if firstRequestError == nil {
firstRequestError = fmt.Errorf("concurrent multi-resource decision request %d failed after %s: %w", i+1, durations[i], requestErr)
}
continue
}
if responses[i] == nil {
failedRequestCount++
if firstRequestError == nil {
firstRequestError = fmt.Errorf("concurrent multi-resource decision request %d returned no response", i+1)
}
continue
}
if i > 0 && !proto.Equal(responses[0], responses[i]) {
failedRequestCount++
if firstRequestError == nil {
firstRequestError = fmt.Errorf("concurrent multi-resource decision request %d returned a different response", i+1)
}
}
}

scenarioContext.TestSuiteContext.Logger.Info(
"authorization v2 concurrent multi-resource performance",
slog.Int("concurrency", concurrency),
slog.Duration("wall_duration", wallElapsed),
slog.Duration("median_request_duration", medianRequestDuration),
slog.Duration("p95_request_duration", p95RequestDuration),
slog.Duration("maximum_request_duration", maximumRequestDuration),
slog.Duration("maximum_duration", limit),
slog.Int("resource_count", len(resources)),
slog.Int("failed_request_count", failedRequestCount),
)

if firstRequestError != nil {
return ctx, firstRequestError
}
if maximumRequestDuration > limit {
return ctx, fmt.Errorf("slowest concurrent multi-resource decision took %s, exceeding the %s limit", maximumRequestDuration, limit)
}

scenarioContext.RecordObject(multiDecisionResponseKey, responses[0])
scenarioContext.RecordObject(decisionResponse, responses[0])
scenarioContext.RecordObject("resourceFQNMap", resourceFQNMap)
return ctx, nil
}

func (s *AuthorizationServiceStepDefinitions) iSendAMultiResourceDecisionRequestForEntityChainForActionOnResourcesWithNoFulfillableObligations(ctx context.Context, entityChainID string, action string, tbl *godog.Table) (context.Context, error) {
return s.iSendAMultiResourceDecisionRequestForEntityChainForActionOnResourcesWithFulfillableObligations(ctx, entityChainID, action, "[]", tbl)
}
Expand Down Expand Up @@ -556,6 +715,8 @@ func RegisterAuthorizationStepDefinitions(ctx *godog.ScenarioContext) {
ctx.Step(`^I send a decision request for entity chain "([^"]*)" for "([^"]*)" action on resource "([^"]*)" with fulfillable obligations "([^"]*)"$`, stepDefinitions.iSendADecisionRequestForEntityChainForActionOnResourceWithFulfillableObligations)
ctx.Step(`^I send a decision request for entity chain "([^"]*)" for "([^"]*)" action on resource "([^"]*)" with no fulfillable obligations$`, stepDefinitions.iSendADecisionRequestForEntityChainForActionOnResourceWithNoFulfillableObligations)
ctx.Step(`^I send a multi-resource decision request for entity chain "([^"]*)" for "([^"]*)" action on resources:$`, stepDefinitions.iSendAMultiResourceDecisionRequestForEntityChainForActionOnResources)
ctx.Step(`^I send a multi-resource decision request for entity chain "([^"]*)" for "([^"]*)" action on resources within "([^"]*)":$`, stepDefinitions.iSendAMultiResourceDecisionRequestWithin)
ctx.Step(`^I send (\d+) concurrent multi-resource decision requests for entity chain "([^"]*)" for "([^"]*)" action on resources each within "([^"]*)":$`, stepDefinitions.iSendConcurrentMultiResourceDecisionRequestsWithin)
ctx.Step(`^I send a multi-resource decision request for entity chain "([^"]*)" for "([^"]*)" action on resources with no fulfillable obligations:$`, stepDefinitions.iSendAMultiResourceDecisionRequestForEntityChainForActionOnResourcesWithNoFulfillableObligations)
ctx.Step(`^I send a multi-resource decision request for entity chain "([^"]*)" for "([^"]*)" action on resources with fulfillable obligations "([^"]*)":$`, stepDefinitions.iSendAMultiResourceDecisionRequestForEntityChainForActionOnResourcesWithFulfillableObligations)
ctx.Step(`^I should get a "([^"]*)" decision response$`, stepDefinitions.iShouldGetADecisionResponse)
Expand Down
Loading
Loading