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
2 changes: 1 addition & 1 deletion internal/extensions/extract.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ func ExtractHeaders(ex *openapi3.Example) (map[string]any, bool) {
return ValueHeaders(OpenAPIExampleValue(ex))
}

// EventTrigger is a single x-event-trigger entry (design D8).
// EventTrigger is a single x-event-trigger entry.
type EventTrigger struct {
// Name is the named event fired when the example is selected.
Name string
Expand Down
21 changes: 18 additions & 3 deletions internal/extensions/match.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,11 @@ func getCachedSchema(schema map[string]any) (*gojsonschema.Schema, error) {
}

// EvaluateParamsMatch evaluates whether the given params match the conditions.
func EvaluateParamsMatch(pm ParamsMatch, eval runtime.Evaluator) (bool, error) {
// An evaluation failure (an expression source unavailable in the context, e.g.
// {$event.*} on the reply path) fails closed; when verbose is true the failure
// is logged at warning level (RS.EXT.29), otherwise at debug.
func EvaluateParamsMatch(pm ParamsMatch, eval runtime.Evaluator, verbose ...bool) (bool, error) {
keepVerbose := len(verbose) > 0 && verbose[0]
for expr, condition := range pm {
// Pre-evaluate the condition value when it is itself a runtime
// expression AND the key references the new event/connection contexts
Expand All @@ -58,7 +62,7 @@ func EvaluateParamsMatch(pm ParamsMatch, eval runtime.Evaluator) (bool, error) {
if str, ok := condition.(string); ok && isFullExpression(str) && referencesNewMatchContext(expr) {
resolved, err := eval.Evaluate(str)
if err != nil {
slog.Debug("EvaluateParamsMatch: condition expression evaluation failed", "expr", expr, "condition", str, "err", err)
logMatchEvalFailure(keepVerbose, "EvaluateParamsMatch: condition expression evaluation failed", "expr", expr, "condition", str, "err", err)
return false, nil
}
condition = resolved
Expand All @@ -68,7 +72,7 @@ func EvaluateParamsMatch(pm ParamsMatch, eval runtime.Evaluator) (bool, error) {
value, err := eval.Evaluate(expr)
if err != nil {
// If expression cannot be evaluated, treat as mismatch
slog.Debug("EvaluateParamsMatch: expression evaluation failed", "expr", expr, "err", err)
logMatchEvalFailure(keepVerbose, "EvaluateParamsMatch: expression evaluation failed", "expr", expr, "err", err)
return false, nil
}
slog.Debug("EvaluateParamsMatch: expression evaluated", "expr", expr, "value", value, "condition", condition)
Expand All @@ -94,6 +98,17 @@ func EvaluateParamsMatch(pm ParamsMatch, eval runtime.Evaluator) (bool, error) {
return true, nil
}

// logMatchEvalFailure reports a silent match-evaluation failure. Failures are
// warnings in verbose mode (RS.EXT.29) and debug-level otherwise, so hot paths
// stay quiet by default and operators can see why a condition never matched.
func logMatchEvalFailure(verbose bool, msg string, args ...any) {
if verbose {
slog.Warn(msg, args...)
return
}
slog.Debug(msg, args...)
}

// referencesNewMatchContext reports whether a condition key references the
// event or connection context, the two contexts added by this change. Only
// those conditions pre-resolve full-expression values (design D6); reply-path
Expand Down
2 changes: 1 addition & 1 deletion internal/runtime/expression.go
Original file line number Diff line number Diff line change
Expand Up @@ -233,7 +233,7 @@ type EnvSource struct {
}

// EventSource provides access to the payload of the currently fired event via
// {$event.*} (design D8). Name is the event identity (named-event name or
// {$event.*} (design D3). Name is the event identity (named-event name or
// built-in kind); Data is the event payload. name and data are reserved
// accessor names: {$event.name} returns the identity, {$event.data} the whole
// payload, and any other path resolves within the payload fields.
Expand Down
6 changes: 3 additions & 3 deletions internal/server/add_example_validation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ not reference the event context
When /_mock/examples is invoked
Then the server responds with HTTP 400 and registers nothing

Related spec scenarios: RS.MAPI.34
Related spec scenarios: RS.MAPI.35
*/
func TestAddExampleValidation_NonEventMatchRejected(t *testing.T) {
t.Parallel()
Expand Down Expand Up @@ -186,7 +186,7 @@ Given a POST with an async target and a literal-only match
When /_mock/examples is invoked
Then the server responds with HTTP 400 and registers nothing

Related spec scenarios: RS.MAPI.34
Related spec scenarios: RS.MAPI.35
*/
func TestAddExampleValidation_LiteralOnlyMatchRejected(t *testing.T) {
t.Parallel()
Expand Down Expand Up @@ -306,7 +306,7 @@ Then the server responds with HTTP 400
When the same body is posted with validate: false
Then the server accepts it

Related spec scenarios: RS.MAPI.5
Related spec scenarios: RS.MAPI.5, RS.MAPI.34
*/
func TestAddExampleValidation_ValidateFlag(t *testing.T) {
t.Parallel()
Expand Down
2 changes: 1 addition & 1 deletion internal/server/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,7 @@ func (e *exampleEngine) selectFromBucket(bucket map[string]*openapi3.Example, ke
}
if requireMatch {
pm, _ := extensions.ExtractParamsMatch(ex)
matched, err := extensions.EvaluateParamsMatch(pm, eval)
matched, err := extensions.EvaluateParamsMatch(pm, eval, e.verbose)
if err != nil {
if e.verbose {
slog.Debug("Error evaluating params-match", "example", k, "error", err)
Expand Down
2 changes: 1 addition & 1 deletion internal/server/engine_async.go
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,7 @@ func (e *exampleEngine) SelectAsyncExample(message *loader.MessageSpec, evaluato
continue
}
if match, ok := extensions.ValueMatch(view); ok {
matched, err := extensions.EvaluateParamsMatch(extensions.ParamsMatch(match), evaluator)
matched, err := extensions.EvaluateParamsMatch(extensions.ParamsMatch(match), evaluator, e.verbose)
if err != nil || !matched {
continue
}
Expand Down
2 changes: 1 addition & 1 deletion internal/server/event_broker.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ type delaySchedule struct {
type eventDeliverer func(sub channelSubscription, payload map[string]any)

// eventBroker decouples OpenAPI event triggers from AsyncAPI consumers
// (design D8). Subscriptions are keyed by match identity + schema scope.
// (design D3). Subscriptions are keyed by match identity + schema scope.
type eventBroker struct {
mu sync.RWMutex
byEvent map[string][]channelSubscription // identity -> subscriptions
Expand Down
2 changes: 1 addition & 1 deletion internal/server/event_server.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import (
"github.com/mamonth/oasmock/internal/loader"
)

// eventBus orchestrates the event driver (design D8). It owns the subscription
// eventBus orchestrates the event driver (design D3). It owns the subscription
// broker and the interval scheduler, and delegates message rendering/delivery
// to the messageDelivery engine. It never reaches into Server.
type eventBus struct {
Expand Down
6 changes: 3 additions & 3 deletions internal/server/message_delivery.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ import (

// messageDelivery renders and delivers subscribed AsyncAPI message examples
// through the MessageRenderer and ConsumerBus contracts. It is the cohesive
// delivery engine behind the event driver (design D8): it owns the sources of
// delivery engine behind the event driver (design D6): it owns the sources of
// per-emission rendering (state, env, event, connection), the recipient
// partition, the delayed-emission cancellation and the push/observer side
// effects. It never reaches into Server or the broker/scheduler registries.
Expand Down Expand Up @@ -157,7 +157,7 @@ func (d *messageDelivery) evaluateConnectionBucket(bucket extensions.ParamsMatch
return true, nil
}
eval := d.eventEvaluator(state, env, eventName, payload, connectionSourceFromInfo(candidate))
return extensions.EvaluateParamsMatch(bucket, eval)
return extensions.EvaluateParamsMatch(bucket, eval, d.verbose)
}

// deliverExample runs the shared selection + render + recipient-partition
Expand All @@ -177,7 +177,7 @@ func (d *messageDelivery) deliverExample(sub channelSubscription, examples []*lo
}
evaluator := d.eventEvaluator(state, env, eventName, payload, connSource)
if len(common) > 0 {
ok, cErr := extensions.EvaluateParamsMatch(common, evaluator)
ok, cErr := extensions.EvaluateParamsMatch(common, evaluator, d.verbose)
if cErr != nil || !ok {
continue
}
Expand Down
2 changes: 1 addition & 1 deletion internal/server/registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -146,7 +146,7 @@ func (r *exampleRegistry) exampleEligible(ex dynamicExample, eval runtime.Evalua
if len(ex.conditions) == 0 {
return true
}
matched, err := extensions.EvaluateParamsMatch(extensions.ParamsMatch(ex.conditions), eval)
matched, err := extensions.EvaluateParamsMatch(extensions.ParamsMatch(ex.conditions), eval, r.verbose)
return err == nil && matched
}

Expand Down
2 changes: 1 addition & 1 deletion internal/server/server_http.go
Original file line number Diff line number Diff line change
Expand Up @@ -242,7 +242,7 @@ func (s *Server) selectAndGenerateResponse(r *http.Request, mapping *RouteMappin
}

// fireExampleTriggers dispatches the x-event-trigger events declared on an
// OpenAPI response example against the schema's event broker (design D8).
// OpenAPI response example against the schema's event broker.
func (s *Server) fireExampleTriggers(example *openapi3.Example, prefix string) {
if s.eventBus == nil || example == nil {
return
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ The `POST /_mock/examples` request SHALL reject field combinations that mix or m
- **WHEN** a POST request includes both `interval` and an event-based `match`, or an `interval` that is not a positive integer
- **THEN** the server responds with HTTP 400

#### Scenario RS.MAPI.34: Non-event match on an async target
#### Scenario RS.MAPI.35: Non-event match on an async target
- **WHEN** a POST request includes an AsyncAPI target and a `match` whose conditions reference only `{$connection.*}` (or literal values) with no `{$event.*}` reference
- **THEN** the server responds with HTTP 400 and registers nothing (a runtime example needs a trigger; a connection-only match has none)

Expand Down
46 changes: 33 additions & 13 deletions openspec/specs/asyncapi-management/spec.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
# asyncapi-management Specification

## Purpose
Management API endpoints for driving AsyncAPI mocking: delayed/targeted/broadcast push, consumer discovery, recurring schedules, fire-event, and connection lifecycle control.
Management API endpoints for driving AsyncAPI mocking: delayed/targeted/broadcast push, consumer discovery, management WebSocket event stream, fire-event, and connection lifecycle control.
## Requirements
### Requirement: Delayed example push to consumer
The mock server SHALL extend the management API with an endpoint to push a message example to consumers of an AsyncAPI channel, accepting an optional `delay` (milliseconds) before the push is delivered.
Expand Down Expand Up @@ -38,7 +38,7 @@ The pushed message SHALL be deliverable to a single consumer connection or broad
- **THEN** the server responds with HTTP 404

### Requirement: Connected consumer discovery
The mock server SHALL expose the currently connected consumers per AsyncAPI channel, including open SignalR streams.
The mock server SHALL expose the currently connected consumers per AsyncAPI channel, including open SignalR streams. The `channel` query parameter SHALL be optional: when omitted, the server SHALL return consumers across all channels; when present, it SHALL return only consumers of that channel.

#### Scenario RS.AMG.8: Listing connected consumers
- **WHEN** a management request queries consumers for an AsyncAPI channel with active connections
Expand All @@ -48,6 +48,10 @@ The mock server SHALL expose the currently connected consumers per AsyncAPI chan
- **WHEN** a management request queries consumers for an AsyncAPI channel with no active connections
- **THEN** the server returns an empty list

#### Scenario RS.AMG.22: Listing all consumers without a channel filter
- **WHEN** a management request queries consumers without a `channel` parameter and consumers are connected on multiple channels
- **THEN** the server returns a single flat list of consumers across all channels (raw ws and SignalR), and an empty list when none are connected

### Requirement: Templated push payloads
Pushed message payloads SHALL support runtime expressions ({$state.*}, {$env.*}) evaluated at delivery time, using the schema's state namespace.

Expand All @@ -59,17 +63,6 @@ Pushed message payloads SHALL support runtime expressions ({$state.*}, {$env.*})
- **WHEN** a management request pushes a payload containing an unresolvable or malformed expression
- **THEN** the server rejects the request with HTTP 400

### Requirement: Recurring scheduled push
The mock server SHALL support scheduling repeated pushes of a message example to a channel at a fixed interval.

#### Scenario RS.AMG.12: Scheduling a recurring push
- **WHEN** a management request schedules a push with an `interval` in milliseconds
- **THEN** the message is delivered repeatedly at that interval until stopped (or the server shuts down)

#### Scenario RS.AMG.13: Stopping a recurring push
- **WHEN** a management request stops a previously scheduled recurring push (by its push ID)
- **THEN** no further deliveries occur for that schedule

### Requirement: Connection lifecycle control
The mock server SHALL allow a management request to terminate a connected consumer's connection, with an optional close reason, or to simulate an abrupt client-side drop.

Expand Down Expand Up @@ -108,3 +101,30 @@ The mock server SHALL expose a management endpoint to fire a named event ad-hoc,
- **WHEN** the server shuts down with in-flight delayed pushes/events pending
- **THEN** no delayed delivery occurs after shutdown (pending timers are cancelled)

### Requirement: Management WebSocket event stream
The mock server SHALL expose a general management WebSocket stream at `GET /_mock/stream` (upgrade) through which a client subscribes to runtime notifications. V1 SHALL be notifications-only: the client sets filters at connection time via `events` and `channels` query parameters (comma-separated, `*` wildcard supported), and the server pushes JSON envelopes. A non-upgrade request to `/_mock/stream` SHALL be rejected.

#### Scenario RS.AMG.23: Subscribing with event and channel filters
- **WHEN** a client connects to `/_mock/stream` with `?events=orderCreated&channels=/alerts`
- **THEN** the client receives envelopes only for matching events and channels; an omitted filter matches everything

#### Scenario RS.AMG.24: Receiving an event-fired envelope
- **WHEN** a named event fires (spec-triggered or via the management API) and a subscribed client is connected
- **THEN** the client receives an envelope of type `event` with the event name, payload, schema scope, and global flag

#### Scenario RS.AMG.25: Receiving a push envelope
- **WHEN** a management push delivers a message to a channel
- **THEN** a subscribed client receives an envelope of type `push` with the channel, target connection (when targeted), and payload

#### Scenario RS.AMG.26: Receiving consumer lifecycle envelopes
- **WHEN** a consumer connects to or disconnects from a channel (raw ws or SignalR)
- **THEN** a subscribed client receives an envelope of type `consumer` with a `connected`/`disconnected` action, connection ID, and channel

#### Scenario RS.AMG.27: Receiving schedule start/stop envelopes
- **WHEN** a periodic message example is registered with `interval` via `POST /_mock/examples` (or spec `x-mock-interval`) or removed via `DELETE /_mock/examples/{exampleId}`
- **THEN** a subscribed client receives an envelope of type `schedule` with a `started`/`stopped` action, example ID, channel, and interval

#### Scenario RS.AMG.28: Non-upgrade request to the stream endpoint
- **WHEN** a plain HTTP (non-WebSocket) request is sent to `/_mock/stream`
- **THEN** the server responds with HTTP 405

Loading
Loading