Skip to content
Merged
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
25 changes: 14 additions & 11 deletions dispatch.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,11 +43,17 @@ func newDispatcher(
}

// dispatch looks up an enabled ID, re-binds, authorizes, invokes, and converts typed output.
//
// remainingIntermediateBytes is the unused native-result value-body budget.
// ConvertOutput uses min(maxValueBytes, remainingIntermediateBytes) as its
// node and materialization limit. ErrValueLimit maps to resource failure;
// other conversion failures map to capability failure.
func (dispatch *dispatcher) dispatch(
ctx context.Context,
subject authz.Subject,
id string,
args map[string]any,
remainingIntermediateBytes int,
) (any, error) {
if dispatch == nil || dispatch.catalog == nil || dispatch.authorizer == nil {
return nil, execution.ErrInternal
Expand Down Expand Up @@ -98,19 +104,16 @@ func (dispatch *dispatcher) dispatch(
if outcome.err != nil {
return nil, outcome.err
}
converted, conversionErr := entry.Plan.ConvertOutput(outcome.output)
if conversionErr != nil {
return nil, fmt.Errorf("%w: %w", execution.ErrCapabilityFailure, conversionErr)
}
if validationErr := binding.ValidateValue(
converted,
converted, conversionErr := entry.Plan.ConvertOutput(
outcome.output,
dispatch.maxValueDepth,
dispatch.maxValueBytes,
); validationErr != nil {
if errors.Is(validationErr, binding.ErrValueLimit) {
return nil, fmt.Errorf("%w: %w", execution.ErrResourceLimit, validationErr)
min(dispatch.maxValueBytes, remainingIntermediateBytes),
)
if conversionErr != nil {
if errors.Is(conversionErr, binding.ErrValueLimit) {
return nil, fmt.Errorf("%w: %w", execution.ErrResourceLimit, conversionErr)
}
return nil, fmt.Errorf("%w: %w", execution.ErrCapabilityFailure, validationErr)
return nil, fmt.Errorf("%w: %w", execution.ErrCapabilityFailure, conversionErr)
}
return converted, nil
}
Expand Down
141 changes: 125 additions & 16 deletions dispatch_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -130,8 +130,7 @@ func TestDispatchBindsAuthorizesThenInvokes(t *testing.T) {
"limit": int64(25),
"enabled": true,
"weight": 2.5,
},
)
}, 64*1024)

require.NoError(t, err)
assert.Equal(t, []string{"authorize", "handler"}, events)
Expand Down Expand Up @@ -197,8 +196,7 @@ func TestDispatchTranslatesEveryBindValueFailureInternally(t *testing.T) {
t.Context(),
authz.Subject{ID: "subject-1"},
"cap.lookup",
tt.arguments,
)
tt.arguments, 64*1024)

require.ErrorIs(t, err, worker.ErrProtocol)
require.NotErrorIs(t, err, execution.ErrInvalidArguments)
Expand All @@ -220,8 +218,7 @@ func TestDispatchRejectsUnknownIDsBeforeAuthorization(t *testing.T) {
t.Context(),
authz.Subject{ID: "subject-1"},
"cap.missing",
map[string]any{"value": "alpha"},
)
map[string]any{"value": "alpha"}, 64*1024)

require.ErrorIs(t, err, worker.ErrProtocol)
assert.Zero(t, handlerCalls.Load())
Expand Down Expand Up @@ -336,8 +333,7 @@ func TestDispatchClassifiesPolicyAndHandlerFailures(t *testing.T) {
t.Context(),
authz.Subject{ID: "subject-1"},
"cap.lookup",
map[string]any{"value": "alpha"},
)
map[string]any{"value": "alpha"}, 64*1024)

require.ErrorIs(t, err, tt.target)
})
Expand Down Expand Up @@ -380,8 +376,7 @@ func TestDispatchReturnsFreshCanonicalMaps(t *testing.T) {
t.Context(),
authz.Subject{ID: "subject-1"},
"cap.lookup",
decoded,
)
decoded, 64*1024)

require.NoError(t, err)
require.NotNil(t, authorized)
Expand Down Expand Up @@ -451,8 +446,7 @@ func TestDispatchCancellationAfterAllowPreventsInvoke(t *testing.T) {
ctx,
authz.Subject{ID: "subject-1"},
"cap.lookup",
map[string]any{"value": "alpha"},
)
map[string]any{"value": "alpha"}, 64*1024)
result <- dispatchOutcome{value: value, err: err}
}()

Expand Down Expand Up @@ -510,8 +504,7 @@ func TestDispatchClassifiesParentOutputLimits(t *testing.T) {
t.Context(),
authz.Subject{ID: "subject-1"},
"cap.lookup",
map[string]any{"value": "alpha"},
)
map[string]any{"value": "alpha"}, 64*1024)

require.ErrorIs(t, err, execution.ErrResourceLimit)
})
Expand Down Expand Up @@ -540,8 +533,7 @@ func TestDispatchCancellationDuringHandlerReturnsPromptly(t *testing.T) {
ctx,
authz.Subject{ID: "subject-1"},
"cap.lookup",
map[string]any{"value": "alpha"},
)
map[string]any{"value": "alpha"}, 64*1024)
result <- dispatchOutcome{value: value, err: err}
}()

Expand Down Expand Up @@ -606,3 +598,120 @@ func newWidenedDispatchSubject(t *testing.T, authorizer authz.Authorizer, invoke
dispatch: newDispatcher(capabilityCatalog, authorizer, 16, 64*1024),
}
}

// overflowDispatchOutput covers unsigned values above MaxInt64.
type overflowDispatchOutput struct {
// Count is an unsigned 64-bit integer.
Count uint64 `json:"count"`
}

// nanDispatchOutput covers non-finite floating-point results.
type nanDispatchOutput struct {
// Score is a finite floating-point field.
Score float64 `json:"score"`
}

// TestDispatchClassifiesInvalidCompositeOutputs proves invalid runtime values stay capability failures.
func TestDispatchClassifiesInvalidCompositeOutputs(t *testing.T) {
tests := []struct {
// name identifies the invalid handler output.
name string

// plan compiles the handler contract.
plan func(*testing.T) *binding.Plan

// invoke returns an invalid registered output.
invoke catalog.Invoker
}{
{
name: "unsigned overflow",
plan: func(t *testing.T) *binding.Plan {
t.Helper()
plan, err := binding.CompileFor[dispatchInput, overflowDispatchOutput]()
require.NoError(t, err)
return plan
},
invoke: func(context.Context, authz.Subject, any) (any, error) {
return overflowDispatchOutput{Count: uint64(1) << 63}, nil
},
},
{
name: "NaN",
plan: func(t *testing.T) *binding.Plan {
t.Helper()
plan, err := binding.CompileFor[dispatchInput, nanDispatchOutput]()
require.NoError(t, err)
return plan
},
invoke: func(context.Context, authz.Subject, any) (any, error) {
return nanDispatchOutput{Score: math.NaN()}, nil
},
},
{
name: "infinity",
plan: func(t *testing.T) *binding.Plan {
t.Helper()
plan, err := binding.CompileFor[dispatchInput, nanDispatchOutput]()
require.NoError(t, err)
return plan
},
invoke: func(context.Context, authz.Subject, any) (any, error) {
return nanDispatchOutput{Score: math.Inf(1)}, nil
},
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
authorizer := authzmocks.NewMockAuthorizer(t)
authorizer.EXPECT().Authorize(mock.Anything, mock.Anything).Return(nil).Once()
capabilityCatalog, err := catalog.Build([]catalog.Registration{{
ID: "cap.lookup",
Name: "records.lookup",
Summary: "Return one record.",
Description: "Returns the supplied record value.",
Plan: tt.plan(t),
Invoke: tt.invoke,
}}, catalog.Options{
MaxSearchQueryBytes: 256,
MaxSearchResults: 20,
})
require.NoError(t, err)
subject := &dispatchSubject{dispatch: newDispatcher(capabilityCatalog, authorizer, 16, 64*1024)}

_, err = subject.dispatch.dispatch(
t.Context(),
authz.Subject{ID: "subject-1"},
"cap.lookup",
map[string]any{"value": "alpha"},
64*1024,
)

require.ErrorIs(t, err, execution.ErrCapabilityFailure)
require.NotErrorIs(t, err, execution.ErrResourceLimit)
})
}
}

// TestDispatchMapsValueLimitToResourceFailure proves remaining-byte conversion limits are resource failures.
func TestDispatchMapsValueLimitToResourceFailure(t *testing.T) {
authorizer := authzmocks.NewMockAuthorizer(t)
authorizer.EXPECT().Authorize(mock.Anything, mock.Anything).Return(nil).Once()
subject := newDispatchSubject(
t,
authorizer,
func(context.Context, authz.Subject, any) (any, error) {
return dispatchOutput{Value: "alpha"}, nil
},
)

_, err := subject.dispatch.dispatch(
t.Context(),
authz.Subject{ID: "subject-1"},
"cap.lookup",
map[string]any{"value": "alpha"},
1,
)

require.ErrorIs(t, err, execution.ErrResourceLimit)
}
3 changes: 2 additions & 1 deletion internal/binding/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,5 +6,6 @@
// ValidateValue, FromStarlark, and ToStarlark enforce that matrix plus positive
// depth and materialization limits. [json.Number] and other numeric types are
// rejected. Plan.InputShape remains the only descriptor source for compiled
// input fields.
// input fields. Plan.OutputShape remains a flat FieldShape slice whose Type
// strings carry nested list, dict, struct, and nullable notation.
package binding
Loading