diff --git a/internal/extensions/extract.go b/internal/extensions/extract.go index 3f172b5..2b37dc4 100644 --- a/internal/extensions/extract.go +++ b/internal/extensions/extract.go @@ -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 diff --git a/internal/extensions/match.go b/internal/extensions/match.go index 28dbc23..070524b 100644 --- a/internal/extensions/match.go +++ b/internal/extensions/match.go @@ -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 @@ -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 @@ -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) @@ -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 diff --git a/internal/runtime/expression.go b/internal/runtime/expression.go index f78e5a2..84f486b 100644 --- a/internal/runtime/expression.go +++ b/internal/runtime/expression.go @@ -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. diff --git a/internal/server/add_example_validation_test.go b/internal/server/add_example_validation_test.go index 53b236d..d2beaf0 100644 --- a/internal/server/add_example_validation_test.go +++ b/internal/server/add_example_validation_test.go @@ -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() @@ -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() @@ -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() diff --git a/internal/server/engine.go b/internal/server/engine.go index b4b38ef..356a630 100644 --- a/internal/server/engine.go +++ b/internal/server/engine.go @@ -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) diff --git a/internal/server/engine_async.go b/internal/server/engine_async.go index 6c975e9..128b7cf 100644 --- a/internal/server/engine_async.go +++ b/internal/server/engine_async.go @@ -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 } diff --git a/internal/server/event_broker.go b/internal/server/event_broker.go index a1b8ad7..38e29b0 100644 --- a/internal/server/event_broker.go +++ b/internal/server/event_broker.go @@ -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 diff --git a/internal/server/event_server.go b/internal/server/event_server.go index b7f3faf..649c163 100644 --- a/internal/server/event_server.go +++ b/internal/server/event_server.go @@ -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 { diff --git a/internal/server/message_delivery.go b/internal/server/message_delivery.go index 0e204ed..f5457dd 100644 --- a/internal/server/message_delivery.go +++ b/internal/server/message_delivery.go @@ -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. @@ -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 @@ -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 } diff --git a/internal/server/registry.go b/internal/server/registry.go index 51c41ce..72a4be7 100644 --- a/internal/server/registry.go +++ b/internal/server/registry.go @@ -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 } diff --git a/internal/server/server_http.go b/internal/server/server_http.go index de91e0d..3883ead 100644 --- a/internal/server/server_http.go +++ b/internal/server/server_http.go @@ -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 diff --git a/openspec/changes/async-management-api-extensions/.openspec.yaml b/openspec/changes/archive/2026-09-06-async-management-api-extensions/.openspec.yaml similarity index 100% rename from openspec/changes/async-management-api-extensions/.openspec.yaml rename to openspec/changes/archive/2026-09-06-async-management-api-extensions/.openspec.yaml diff --git a/openspec/changes/async-management-api-extensions/design.md b/openspec/changes/archive/2026-09-06-async-management-api-extensions/design.md similarity index 100% rename from openspec/changes/async-management-api-extensions/design.md rename to openspec/changes/archive/2026-09-06-async-management-api-extensions/design.md diff --git a/openspec/changes/async-management-api-extensions/proposal.md b/openspec/changes/archive/2026-09-06-async-management-api-extensions/proposal.md similarity index 100% rename from openspec/changes/async-management-api-extensions/proposal.md rename to openspec/changes/archive/2026-09-06-async-management-api-extensions/proposal.md diff --git a/openspec/changes/async-management-api-extensions/specs/asyncapi-management/spec.md b/openspec/changes/archive/2026-09-06-async-management-api-extensions/specs/asyncapi-management/spec.md similarity index 100% rename from openspec/changes/async-management-api-extensions/specs/asyncapi-management/spec.md rename to openspec/changes/archive/2026-09-06-async-management-api-extensions/specs/asyncapi-management/spec.md diff --git a/openspec/changes/async-management-api-extensions/specs/event-driver/spec.md b/openspec/changes/archive/2026-09-06-async-management-api-extensions/specs/event-driver/spec.md similarity index 100% rename from openspec/changes/async-management-api-extensions/specs/event-driver/spec.md rename to openspec/changes/archive/2026-09-06-async-management-api-extensions/specs/event-driver/spec.md diff --git a/openspec/changes/async-management-api-extensions/specs/extensions/spec.md b/openspec/changes/archive/2026-09-06-async-management-api-extensions/specs/extensions/spec.md similarity index 100% rename from openspec/changes/async-management-api-extensions/specs/extensions/spec.md rename to openspec/changes/archive/2026-09-06-async-management-api-extensions/specs/extensions/spec.md diff --git a/openspec/changes/async-management-api-extensions/specs/management-api/spec.md b/openspec/changes/archive/2026-09-06-async-management-api-extensions/specs/management-api/spec.md similarity index 98% rename from openspec/changes/async-management-api-extensions/specs/management-api/spec.md rename to openspec/changes/archive/2026-09-06-async-management-api-extensions/specs/management-api/spec.md index d4699fe..7e55e6e 100644 --- a/openspec/changes/async-management-api-extensions/specs/management-api/spec.md +++ b/openspec/changes/archive/2026-09-06-async-management-api-extensions/specs/management-api/spec.md @@ -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) diff --git a/openspec/changes/async-management-api-extensions/tasks.md b/openspec/changes/archive/2026-09-06-async-management-api-extensions/tasks.md similarity index 100% rename from openspec/changes/async-management-api-extensions/tasks.md rename to openspec/changes/archive/2026-09-06-async-management-api-extensions/tasks.md diff --git a/openspec/specs/asyncapi-management/spec.md b/openspec/specs/asyncapi-management/spec.md index 4cd367f..81a8347 100644 --- a/openspec/specs/asyncapi-management/spec.md +++ b/openspec/specs/asyncapi-management/spec.md @@ -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. @@ -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 @@ -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. @@ -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. @@ -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 + diff --git a/openspec/specs/event-driver/spec.md b/openspec/specs/event-driver/spec.md index db4a3b6..89566ab 100644 --- a/openspec/specs/event-driver/spec.md +++ b/openspec/specs/event-driver/spec.md @@ -1,7 +1,7 @@ # event-driver Specification ## Purpose -Named event bus decoupling OpenAPI example triggers (x-event-trigger) from AsyncAPI message subscriptions (x-send-events), with {$event.*} payload templating and schema-local/global scoping. +Named event bus decoupling OpenAPI example triggers (x-event-trigger) from AsyncAPI message subscriptions (matched via `x-mock-match` against an event context), with {$event.*} payload templating and schema-local/global scoping. ## Requirements ### Requirement: Event trigger on an OpenAPI example The mock server SHALL support firing named events from an OpenAPI response example via the `x-event-trigger` extension, triggered whenever that example is selected for a response. Multiple triggers per example SHALL be supported (list form). @@ -33,48 +33,64 @@ Event names SHALL be schema-local by default and server-wide only when declared - **WHEN** an `x-event-trigger` entry sets `global: true` - **THEN** the event is broadcast over all loaded schemas and any matching subscription anywhere receives it -### Requirement: Event subscription on an AsyncAPI message example -The mock server SHALL support subscribing an AsyncAPI message example to events via the `x-send-events` extension. Each entry references a named event or a built-in trigger (`receive`, `connect`, `cron`). - -#### Scenario RS.EVT.7: Subscribing to a named event -- **WHEN** a message example has `x-send-events` containing `{on: }` -- **THEN** the message is emitted to the channel's consumers whenever that event fires - -#### Scenario RS.EVT.8: Event payload in consumer template -- **WHEN** an event fires with a payload and the subscribed message example references `{$event.}` -- **THEN** the expression resolves to the event payload value at emission time - -#### Scenario RS.EVT.9: Built-in connect trigger -- **WHEN** a message example's `x-send-events` contains `{on: connect}` -- **THEN** the message is emitted to a consumer when it connects (with an optional `wait` delay) - -#### Scenario RS.EVT.10: Built-in cron trigger -- **WHEN** a message example's `x-send-events` contains `{on: cron, wait: }` -- **THEN** the message is emitted repeatedly to the channel's consumers at the given interval - -#### Scenario RS.EVT.11: Built-in receive trigger -- **WHEN** a message example's `x-send-events` contains a flat `receive` entry -- **THEN** the message is emitted when the channel receives a matching client message - ### Requirement: Broadcast delivery with client-side filtering -Event-driven messages SHALL be broadcast to the consuming channel's connected consumers; the mock does NOT route by session or account — consumers filter by payload. +Event-driven messages SHALL be broadcast to the consuming channel's connected consumers by default; consumers may filter by payload. When the example's `x-mock-match` additionally references `{$connection.*}`, delivery SHALL be narrowed to the consumers satisfying those per-connection conditions (two-phase recipient partition); the mock does not otherwise route by session or account. #### Scenario RS.EVT.12: Broadcasting an event-driven message -- **WHEN** an event fires and a message example subscribed to it exists on a channel with active consumers +- **WHEN** an event fires and a message example matching it exists on a channel with active consumers and no `{$connection.*}` conditions - **THEN** the templated message is emitted to all consumers connected to that channel #### Scenario RS.EVT.13: Emitting into an open SignalR stream -- **WHEN** an event fires and the subscribed channel is a SignalR hub stream with open invocation handles +- **WHEN** an event fires and the matching channel is a SignalR hub stream with open invocation handles - **THEN** the templated message is pushed as a `StreamItem` into the channel's open streams (per `signalr-hub-runtime`) #### Scenario RS.EVT.14: Event with no subscribers -- **WHEN** a named event fires but no message example subscribes to it +- **WHEN** a named event fires but no message example matches it - **THEN** the event is accepted with no delivery (no error) #### Scenario RS.EVT.15: Event with no consumers -- **WHEN** an event fires, a subscription exists, but the channel has no connected consumers +- **WHEN** an event fires, a matching example exists, but the channel has no connected consumers - **THEN** the event is accepted without error and no message is delivered +### Requirement: Event-driven emission on an AsyncAPI message example +The mock server SHALL emit an AsyncAPI message example in response to events when its `x-mock-match` references the event context. The event identity is matched via `{$event.name}` (named-event name or built-in kind `connect`/`receive`); payload conditions use `{$event.}` and the whole payload via `{$event.data}`. The example's payload SHALL be templated at emission time against `{$event.*}`, `{$state.*}`, and `{$env.*}`. Periodic emission SHALL be declared with `x-mock-interval` rather than an event condition. + +#### Scenario RS.EVT.22: Emitting on a named event +- **WHEN** a message example has `x-mock-match: {'{$event.name}': }` +- **THEN** the message is emitted to the channel's consumers whenever that event fires + +#### Scenario RS.EVT.23: Event payload in consumer template +- **WHEN** an event fires with a payload and the message example references `{$event.}` +- **THEN** the expression resolves to the event payload value at emission time + +#### Scenario RS.EVT.24: Built-in connect trigger +- **WHEN** a message example has `x-mock-match: {'{$event.name}': connect}` +- **THEN** the message is emitted to the (just-connected) consumer when it connects, subject to an optional `x-mock-delay` + +#### Scenario RS.EVT.25: Built-in receive trigger +- **WHEN** a message example has `x-mock-match: {'{$event.name}': receive}` and the channel receives a client message +- **THEN** the message is emitted with the inbound client message exposed in the event context + +#### Scenario RS.EVT.26: Periodic emission via x-mock-interval +- **WHEN** a message example declares `x-mock-interval: ` instead of any event condition or match +- **THEN** the message is emitted repeatedly to the channel's consumers at the given interval until removed or the server shuts down +- **AND** the example SHALL NOT carry an `x-mock-match` (a periodically driven example has exactly one trigger); a spec declaring both is rejected at load + +### Requirement: Per-connection event delivery +The mock server SHALL narrow event-driven delivery to consumers whose connection context satisfies the example's `{$connection.*}` conditions (e.g., `'{$connection.id}': '{$event.connectionId}'`), evaluating non-connection conditions once per emission and connection conditions per candidate. + +#### Scenario RS.EVT.19: Targeted event delivery by connection id +- **WHEN** an example has `x-mock-match: {'{$event.name}': orderCreated, '{$connection.id}': '{$event.connectionId}'}` and the event fires with a `connectionId` payload +- **THEN** only the consumer whose connection id equals the payload value receives the message + +### Requirement: x-send-events deprecation mapping +The mock server SHALL accept legacy `x-send-events` entries by mapping each `{on, wait}` to the unified form during loading, writing a deprecation note in verbose mode: `on` → `x-mock-match: {'{$event.name}': on}` for named/`connect`/`receive`, and `{on: cron, wait: N}` → `x-mock-interval: N`. + +#### Scenario RS.EVT.18: Mapping legacy x-send-events to match +- **WHEN** a spec still uses `x-send-events: [{on: orderCreated}]` or `[{on: cron, wait: 1000}]` +- **THEN** the server behaves as if the example declared `x-mock-match: {'{$event.name}': orderCreated}` (respectively `x-mock-interval: 1000`) and logs a deprecation note in verbose mode +- **AND** a `{on: cron}` entry without a positive `wait` SHALL be rejected at load with an error naming the missing interval (an interval of 0 would otherwise silently register a dead reply example) + ### Requirement: Management fire-event endpoint The mock server SHALL expose a management API endpoint to fire a named event ad-hoc, reusing the event broker and its delay semantics. diff --git a/openspec/specs/extensions/spec.md b/openspec/specs/extensions/spec.md index 3bba526..7d2d405 100644 --- a/openspec/specs/extensions/spec.md +++ b/openspec/specs/extensions/spec.md @@ -1,6 +1,6 @@ ## Purpose -OpenAPI extensions that enable conditional example selection, state management, dynamic headers, and runtime expression evaluation for the OASMock server. +OpenAPI extensions that enable conditional example selection, state management, dynamic headers, and runtime expression evaluation for the OASMock server. `x-mock-match` is the single selection extension across sync and async examples: it selects against an HTTP/request context for OpenAPI and against an event context (`{$event.*}`) for async-driven examples, filters recipients per connection (`{$connection.*}`), and is complemented by timing-only sibling extensions (`x-mock-interval`, `x-mock-delay`). ## Requirements @@ -103,3 +103,91 @@ The mock server SHALL support value modifiers after a `|` sign in runtime expres #### Scenario RS.EXT.17: Escaping dots in property names - **WHEN** expression `{$request.cookie.dot\.dot}` appears - **THEN** the dot is treated as part of the property name, not a path separator + +### Requirement: Event-context example matching +The mock server SHALL evaluate `x-mock-match` against an event context for async-driven examples. The event context SHALL expose the event identity as `{$event.name}` (the named-event name or the built-in kind), the whole event payload as `{$event.data}`, and each payload field as `{$event.}`. `{$event.name}`/`{$event.data}` are reserved metadata names; payload fields with those names are shadowed. + +#### Scenario RS.EXT.18: Matching the event identity +- **WHEN** an async-driven example has `x-mock-match: {'{$event.name}': orderCreated}` and the `orderCreated` event fires +- **THEN** the server emits the example to the channel's consumers + +#### Scenario RS.EXT.19: Matching the event payload +- **WHEN** an async-driven example has `x-mock-match: {'{$event.accountId}': 'acc-1'}` (or a JSON-schema condition on `{$event.data}`) +- **AND** the fired event's payload satisfies the condition +- **THEN** the server emits the example; otherwise it does not + +#### Scenario RS.EXT.21: Matching built-in events +- **WHEN** an async-driven example has `x-mock-match: {'{$event.name}': connect}` or `{'{$event.name}': receive}` +- **THEN** the example is emitted when a consumer connects to the channel, or when the channel receives a client message, respectively + +### Requirement: Event-driven example classification +The mock server SHALL classify a spec example as event-driven when its `x-mock-match` references `{$event.*}`, as periodically driven when it declares `x-mock-interval`, and otherwise as a sync/async reply. An example SHALL be rejected at load when it mixes `{$event.*}` match conditions with `{$request.*}`/`{$message.*}`/`{$channel.*}` conditions, or when it declares both `x-mock-interval` and any `x-mock-match` conditions (a periodic emission is single-trigger and has no match context to honor). An event-driven example whose `{$event.name}` condition value is itself a runtime expression SHALL be rejected at load: the identity must be a literal string so the subscription key can always match a fired identity. + +#### Scenario RS.EXT.20: Rejecting mixed match contexts +- **WHEN** a spec example's `x-mock-match` contains both `{$event.name}` and `{$message.payload.kind}` conditions in the same map +- **THEN** the server rejects the spec at load with a clear error + +#### Scenario RS.EXT.28: Rejecting dual triggers +- **WHEN** a spec example declares both `x-mock-interval` and an `{$event.name}` match condition +- **THEN** the server rejects the spec at load with a clear error (any `x-mock-match` alongside `x-mock-interval`, event-driven or not, is a load error) + +#### Scenario RS.EXT.33: Rejecting a non-literal event identity +- **WHEN** a spec example has `x-mock-match: {'{$event.name}': '{$state.envName}'}` +- **THEN** the server rejects the spec at load with a clear error instead of registering a subscription keyed by the literal expression string + +#### Scenario RS.EXT.34: Wildcard identity without an {$event.name} pin +- **WHEN** an event-driven match references the event context only through condition values (no `{$event.name}` condition), e.g. `'{$connection.id}': '{$event.connectionId}'` +- **THEN** the example is registered as event-driven with a wildcard identity that evaluates against every fired event + +### Requirement: Per-connection recipient matching +The mock server SHALL partition `x-mock-match` at delivery time: conditions referencing `{$connection.*}` on either side form the recipient filter, evaluated per candidate connection; all other conditions are evaluated once per emission. Only candidates that satisfy both phases receive the message. When no condition references `{$connection.*}`, the server SHALL broadcast to all consumers of the channel as today. + +#### Scenario RS.EXT.24: Two-phase recipient partition +- **WHEN** an async-driven example has `x-mock-match` with a common condition and `'{$connection.id}': '{$event.connectionId}'`, and the event fires +- **THEN** the server evaluates the common condition once, then evaluates only the `{$connection.*}` condition against each candidate connection and delivers the payload to the connections whose id matches the event's `connectionId` + +#### Scenario RS.EXT.25: Broadcast fast path +- **WHEN** an async-driven example has no `{$connection.*}` conditions and no `{$connection.*}` references +- **THEN** the server broadcasts the emitted message to all consumers of the channel (unchanged behavior, no per-connection evaluation) + +#### Scenario RS.EXT.26: Single-recipient connect built-in +- **WHEN** a `connect` event fires and the connected consumer satisfies the example's `{$connection.*}` conditions +- **THEN** the server delivers the message to that single consumer only + +#### Scenario RS.EXT.27: Connection context exposure +- **WHEN** a condition references `{$connection.id}`, `{$connection.channel}`, `{$connection.query.}`, or `{$connection.header.}` +- **THEN** the values resolve from the consumer's connection id, channel address, and metadata captured at upgrade + +### Requirement: Timing extensions x-mock-interval and x-mock-delay +The mock server SHALL support `x-mock-interval` (positive integer milliseconds) on an async example to emit it repeatedly at that cadence, and `x-mock-delay` (integer milliseconds, default 0) to delay emission after an event fire. Neither is an event identity; `x-mock-interval` marks a periodically driven example. Timing values SHALL be integral milliseconds: a fractional value SHALL be rejected at load rather than silently truncated, and a periodically driven example SHALL honor `x-mock-skip` like every other example. + +#### Scenario RS.EXT.22: Interval-driven periodic emission +- **WHEN** an async example declares `x-mock-interval: 1000` +- **THEN** the server emits the message to the channel's consumers roughly every 1000 ms until the example is removed or the server shuts down + +#### Scenario RS.EXT.35: Rejecting a fractional x-mock-interval +- **WHEN** an async example declares `x-mock-interval: 2.5` (a fractional millisecond value) +- **THEN** the server rejects the spec at load with a clear error instead of truncating to 2 ms + +#### Scenario RS.EXT.36: Rejecting a fractional x-mock-delay +- **WHEN** an async example declares `x-mock-delay: 2.5` +- **THEN** the server rejects the spec at load with a clear error instead of silently dropping the delay + +#### Scenario RS.EXT.37: Skipping a periodically driven example +- **WHEN** a periodically driven example declares `x-mock-skip` and a consumer is connected to its channel +- **THEN** the server never emits the example's message while the skip flag is set + +#### Scenario RS.EXT.23: Delayed event emission +- **WHEN** an async-driven example declares `x-mock-delay: 150` and its event fires +- **THEN** the server emits the message 150 ms after the fire + +#### Scenario RS.EXT.30: Reply-path condition values stay literal +- **WHEN** an `x-mock-match` condition key references the reply context (`{$request.*}`/`{$message.*}`/`{$channel.*}`) and its value is a full runtime-expression string +- **THEN** the value is compared as a literal string, never pre-resolved (only conditions whose key references `{$event.*}` or `{$connection.*}` pre-resolve full-expression values) + +### Requirement: Fail-closed match evaluation +An `x-mock-match` condition that references an expression source unavailable in the evaluation context SHALL fail closed (the example does not match) and, in verbose mode, SHALL log a warning. + +#### Scenario RS.EXT.29: Event context unavailable in reply path +- **WHEN** a sync example references `{$event.*}` or a reply-path async example references `{$connection.*}` +- **THEN** the condition never matches and the server logs a verbose-mode warning rather than erroring diff --git a/openspec/specs/management-api/spec.md b/openspec/specs/management-api/spec.md index 190a63b..a1909ea 100644 --- a/openspec/specs/management-api/spec.md +++ b/openspec/specs/management-api/spec.md @@ -1,6 +1,6 @@ ## Purpose -HTTP management API for runtime control of the OASMock server, allowing dynamic addition of mock examples and retrieval of request history. +HTTP management API for runtime control of the OASMock server, allowing dynamic addition of mock examples (OpenAPI and AsyncAPI targets with match/interval/delay), retrieval of request history, and a type-discriminated event resource to fire events. ## Requirements ### Requirement: Management API availability The mock server SHALL provide an HTTP API for runtime management under the `/_mock` path prefix according to [openapi](../../../../openapi.yaml) spec @@ -113,13 +113,66 @@ The mock server SHALL accept a route identifier that resolves to an AsyncAPI cha - **THEN** the server responds with HTTP 400 (no matching route) ### Requirement: Fire an event on the event bus -The management API SHALL expose an endpoint to fire a named event ad-hoc, reusing the event broker and its delay semantics (per `event-driver`). +The management API SHALL expose `POST /_mock/events` to fire a named event ad-hoc. The request SHALL carry a required `type` discriminator (`"fire"` for V1, extensible), along with `event`, `payload`, `delay`, and `global` fields, reusing the event broker and its delay semantics (per `event-driver`). #### Scenario RS.MAPI.22: Firing an event via management API -- **WHEN** a management request fires a named event with a payload and optional delay -- **THEN** the server delivers it like a spec-triggered event (immediately or after the delay) to matching `x-send-events` consumers +- **WHEN** a `POST /_mock/events` request fires a named event with `type: fire`, a payload, and an optional delay +- **THEN** the server delivers it like a spec-triggered event (immediately or after the delay) to matching event-driven message examples #### Scenario RS.MAPI.23: Fire-event payload templating - **WHEN** an ad-hoc fired event payload contains `{$state.*}` or `{$env.*}` expressions - **THEN** they are evaluated against the schema's isolated state namespace and environment before delivery +#### Scenario RS.MAPI.32: Invalid event type +- **WHEN** a `POST /_mock/events` request omits `type` or uses a type other than `fire` +- **THEN** the server responds with HTTP 400 + +### Requirement: Adding a runtime async-driven example +The `POST /_mock/examples` request SHALL accept, for AsyncAPI targets, an optional `match` object (mirroring `x-mock-match` against the event and connection contexts), an optional `interval` (positive integer ms for periodic emission), and an optional `delay` (integer ms). The mock server SHALL register the added message example (payload = `response.body`, headers = `response.headers`) as a live async-driven subscription delivered to the channel's consumers according to its match/interval, templating the payload at emission time against `{$event.*}`, `{$connection.*}`, `{$state.*}`, and `{$env.*}`. + +#### Scenario RS.MAPI.24: Registering a named-event runtime example +- **WHEN** a POST request is sent to `/_mock/examples` with an AsyncAPI target, `response.body`, and `match: {'{$event.name}': orderCreated}` +- **THEN** the server registers the message as a live subscription, responds with success and an example ID, and delivers the message when the `orderCreated` event fires + +#### Scenario RS.MAPI.25: Scheduling repeated delivery via interval +- **WHEN** a POST request includes `interval: 1000` for an AsyncAPI target +- **THEN** the message is delivered repeatedly at the 1000 ms interval until removed (or the server shuts down) + +#### Scenario RS.MAPI.26: Subscribing to the connect and receive built-ins +- **WHEN** a POST request includes `match: {'{$event.name}': connect}` or `{'{$event.name}': receive}` +- **THEN** the message is delivered to a consumer when it connects to the channel, or when the channel receives a client message (with the inbound message payload available to templates), respectively + +#### Scenario RS.MAPI.33: Targeting delivery by connection +- **WHEN** a POST request includes `match` with a `{$connection.*}` condition alongside an event condition +- **THEN** the registered message is delivered only to the channel's consumers satisfying that connection condition when the event fires + +### Requirement: Strict example target validation +The `POST /_mock/examples` request SHALL reject field combinations that mix or misplace sync and async targeting with HTTP 400. An OpenAPI target requires `path` (and uses `response`); an AsyncAPI target requires `channel` (optionally `protocol`); `match`/`interval`/`delay` are only valid on AsyncAPI targets; a runtime example SHALL have exactly one trigger — `interval` OR an `{$event.*}`-based `match`, never both — and `interval` SHALL be a positive integer. + +#### Scenario RS.MAPI.27: Mixing sync and async targeting +- **WHEN** a POST request includes both `path` and `channel` +- **THEN** the server responds with HTTP 400 + +#### Scenario RS.MAPI.28: match or interval on an OpenAPI target +- **WHEN** a POST request includes `path` with `match` (or `interval`) but no AsyncAPI target +- **THEN** the server responds with HTTP 400 + +#### Scenario RS.MAPI.29: Dual or invalid triggers +- **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.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) + +### Requirement: Removing a dynamic example +The mock server SHALL provide `DELETE /_mock/examples/{exampleId}` to remove a dynamically added example and cancel any recurring delivery registered under that example ID. + +#### Scenario RS.MAPI.30: Removing a dynamic example +- **WHEN** a DELETE request is sent to `/_mock/examples/{exampleId}` for an existing example +- **THEN** the server removes the example, stops any recurring delivery, and responds with success + +#### Scenario RS.MAPI.31: Removing an unknown example +- **WHEN** a DELETE request is sent to `/_mock/examples/{unknownId}` that does not exist +- **THEN** the server responds with HTTP 404 + diff --git a/test/asyncapi/management-api/management_api_test.go b/test/asyncapi/management-api/management_api_test.go index 243f336..efcd3e3 100644 --- a/test/asyncapi/management-api/management_api_test.go +++ b/test/asyncapi/management-api/management_api_test.go @@ -5,6 +5,7 @@ import ( "fmt" "net/http" "strings" + "sync" "testing" "time" @@ -36,6 +37,104 @@ func wsConnect(t *testing.T, port int, path string) *websocket.Conn { return conn } +// streamCollector collects envelopes read from a _mock/stream connection so +// tests can assert on them asynchronously as the server pushes them. +type streamCollector struct { + mu sync.Mutex + envs []map[string]any +} + +// collectStream starts a deadline-free reader collecting envelopes from a +// _mock/stream connection. The reader exits when the connection is closed +// (any read error ends the loop); the test closes the connection via defer. +// No read deadlines are used because a gorilla ReadMessage error is permanent +// on a connection. +func collectStream(conn *websocket.Conn) *streamCollector { + c := &streamCollector{} + go func() { + for { + _, raw, err := conn.ReadMessage() + if err != nil { + return + } + var env map[string]any + if json.Unmarshal(raw, &env) == nil { + c.mu.Lock() + c.envs = append(c.envs, env) + c.mu.Unlock() + } + } + }() + return c +} + +// find returns the first envelope matching the predicate. +func (c *streamCollector) find(match func(map[string]any) bool) (map[string]any, bool) { + c.mu.Lock() + defer c.mu.Unlock() + for _, env := range c.envs { + if match(env) { + return env, true + } + } + return nil, false +} + +// waitForEnv polls until an envelope matching the predicate arrives or the +// test times out. +func waitForEnv(t *testing.T, c *streamCollector, what string, match func(map[string]any) bool) map[string]any { + t.Helper() + var got map[string]any + require.Eventuallyf(t, func() bool { + env, ok := c.find(match) + got = env + return ok + }, 3*time.Second, 50*time.Millisecond, "timed out waiting for envelope: %s", what) + return got +} + +// waitQuietEnv asserts no matching envelope arrives within a window. +func waitQuietEnv(t *testing.T, c *streamCollector, what string, dur time.Duration, match func(map[string]any) bool) { + t.Helper() + end := time.Now().Add(dur) + for time.Now().Before(end) { + time.Sleep(30 * time.Millisecond) + if _, ok := c.find(match); ok { + t.Fatalf("expected no envelope during quiet window: %s", what) + } + } +} + +// envOfType returns the per-type payload of an envelope of the given type, or +// nil when the envelope is of a different type. +func envOfType(env map[string]any, typ string) map[string]any { + if env["type"] != typ { + return nil + } + sub, _ := env[typ].(map[string]any) + return sub +} + +// readUntilPayload reads frames from a consumer connection until a frame +// containing the given substring is seen, or the deadline passes. +func readUntilPayload(t *testing.T, conn *websocket.Conn, want string) string { + t.Helper() + _ = conn.SetReadDeadline(time.Now().Add(3 * time.Second)) + deadline := time.Now().Add(3 * time.Second) + var got string + for time.Now().Before(deadline) { + _, raw, rerr := conn.ReadMessage() + if rerr != nil { + break + } + got = string(raw) + if strings.Contains(got, want) { + break + } + } + return got +} + /* Scenario: Runtime match on {$event.name} fired via POST /_mock/events Given a running server with a match on {$event.name} and a connected consumer @@ -280,3 +379,194 @@ func TestIntegration_ScheduleGone410(t *testing.T) { _ = resp.Body.Close() assert.Equal(t, 410, resp.StatusCode) } + +/* +Scenario: A _mock/stream subscriber receives filtered event/push/consumer/schedule envelopes +Given a stream subscriber filtering events=levelUp and channels=/alerts, and an unfiltered subscriber +When an interval example starts and stops, a consumer connects, and the levelUp event fires +Then the filtered subscriber receives schedule/consumer/event/push envelopes but the connect event stays filtered out, while the unfiltered subscriber sees everything + +Related spec scenarios: RS.AMG.23, RS.AMG.24, RS.AMG.25, RS.AMG.26, RS.AMG.27 +*/ +func TestIntegration_ManageStream_Envelopes(t *testing.T) { + t.Parallel() + port, stop := startManagementServer(t) + defer stop() + + filtered := wsConnect(t, port, "/_mock/stream?events=levelUp&channels=/alerts") + defer filtered.Close() //nolint:errcheck + fCol := collectStream(filtered) + + // An omitted filter matches everything (RS.AMG.23). + all := wsConnect(t, port, "/_mock/stream") + defer all.Close() //nolint:errcheck + allCol := collectStream(all) + + // Register an interval example -> schedule started (RS.AMG.27). + addBody := `{"channel":"/alerts","interval":300,"response":{"code":200,"body":{"mytick":true}}}` + addResp, err := http.Post(fmt.Sprintf("http://localhost:%d/_mock/examples", port), "application/json", strings.NewReader(addBody)) + require.NoError(t, err) + var add map[string]any + require.NoError(t, json.NewDecoder(addResp.Body).Decode(&add)) + addResp.Body.Close() //nolint:errcheck + require.Equal(t, 200, addResp.StatusCode) + exampleID, _ := add["id"].(string) + require.NotEmpty(t, exampleID) + + started := waitForEnv(t, fCol, "schedule started", func(env map[string]any) bool { + sub := envOfType(env, "schedule") + return sub != nil && sub["action"] == "started" && sub["channel"] == "/alerts" && sub["exampleId"] == exampleID + }) + assert.EqualValues(t, 300, started["schedule"].(map[string]any)["interval"]) + + // Connect a consumer -> consumer lifecycle envelope (RS.AMG.26); the + // connect built-in event (name "connect") is filtered out on the events + // filter but reaches the unfiltered subscriber (RS.AMG.23). + conn := wsConnect(t, port, "/alerts") + defer conn.Close() //nolint:errcheck + _ = conn.SetReadDeadline(time.Now().Add(2 * time.Second)) + _, _, _ = conn.ReadMessage() // connect built-in welcome + + waitForEnv(t, fCol, "consumer connected", func(env map[string]any) bool { + sub := envOfType(env, "consumer") + return sub != nil && sub["action"] == "connected" && sub["channel"] == "/alerts" + }) + // The connect event violates the events=levelUp filter, so it can never + // appear on the filtered subscriber even though the brand-new consumer + // fired it (the broadcast is synchronous with the connect). + waitQuietEnv(t, fCol, "filtered connect event", 400*time.Millisecond, func(env map[string]any) bool { + sub := envOfType(env, "event") + return sub != nil && sub["name"] == "connect" + }) + // The unfiltered subscriber sees it. + waitForEnv(t, allCol, "unfiltered connect event", func(env map[string]any) bool { + sub := envOfType(env, "event") + return sub != nil && sub["name"] == "connect" + }) + + // Fire levelUp -> event and push envelopes (RS.AMG.24, RS.AMG.25). + _, err = http.Post(fmt.Sprintf("http://localhost:%d/_mock/events", port), "application/json", + strings.NewReader(`{"type":"fire","event":"levelUp","payload":{"level":"warn","msg":"boom"}}`)) + require.NoError(t, err) + + waitForEnv(t, fCol, "levelUp event", func(env map[string]any) bool { + sub := envOfType(env, "event") + return sub != nil && sub["name"] == "levelUp" + }) + waitForEnv(t, fCol, "levelUp push", func(env map[string]any) bool { + sub := envOfType(env, "push") + if sub == nil || sub["channel"] != "/alerts" { + return false + } + payload, _ := sub["payload"].(map[string]any) + return payload != nil && payload["msg"] == "boom" + }) + + // DELETE the interval example -> schedule stopped (RS.AMG.27). + delReq, err := http.NewRequest(http.MethodDelete, fmt.Sprintf("http://localhost:%d/_mock/examples/%s", port, exampleID), nil) + require.NoError(t, err) + delResp, err := http.DefaultClient.Do(delReq) + require.NoError(t, err) + _ = delResp.Body.Close() + assert.Equal(t, 200, delResp.StatusCode) + + waitForEnv(t, fCol, "schedule stopped", func(env map[string]any) bool { + sub := envOfType(env, "schedule") + return sub != nil && sub["action"] == "stopped" && sub["exampleId"] == exampleID + }) +} + +/* +Scenario: The schema's connect built-in fires on consumer connection +Given a channel whose spec declares an x-mock-match on {$event.name}: connect +When a ws consumer connects +Then the templated welcome example is delivered + +Related spec scenarios: RS.EVT.24, RS.EXT.21, RS.EXT.26 +*/ +func TestIntegration_ConnectBuiltIn_Spec(t *testing.T) { + t.Parallel() + port, stop := startManagementServer(t) + defer stop() + + conn := wsConnect(t, port, "/alerts") + defer conn.Close() //nolint:errcheck + + got := readUntilPayload(t, conn, `"msg":"welcome"`) + assert.Contains(t, got, `"msg":"welcome"`) +} + +/* +Scenario: A runtime receive built-in example echoes inbound traffic +Given a POST /_mock/examples example matching {$event.name}: receive +When a consumer sends a client message on the channel +Then the templated message is emitted with the inbound payload in the context + +Related spec scenarios: RS.EVT.25, RS.EVT.23, RS.MAPI.26, RS.EXT.21 +*/ +func TestIntegration_ReceiveBuiltIn_Runtime(t *testing.T) { + t.Parallel() + port, stop := startManagementServer(t) + defer stop() + + addBody := `{"channel":"/alerts","match":{"{$event.name}":"receive"},"response":{"code":200,"body":{"echoed":"{$event.text}"}}}` + resp, err := http.Post(fmt.Sprintf("http://localhost:%d/_mock/examples", port), "application/json", strings.NewReader(addBody)) + require.NoError(t, err) + _ = resp.Body.Close() + assert.Equal(t, 200, resp.StatusCode) + + conn := wsConnect(t, port, "/alerts") + defer conn.Close() //nolint:errcheck + _ = conn.SetReadDeadline(time.Now().Add(2 * time.Second)) + _, _, _ = conn.ReadMessage() // connect built-in welcome + + require.NoError(t, conn.WriteMessage(websocket.TextMessage, []byte(`{"text":"hi"}`))) + + got := readUntilPayload(t, conn, `"echoed":"hi"`) + assert.Contains(t, got, `"echoed":"hi"`) +} + +/* +Scenario: The legacy x-send-events shim still emits on a deprecated subscription +Given a spec example declared with x-send-events: [{on: legacyAlert}] +When the legacyAlert event fires via POST /_mock/events +Then the shim-mapped example is emitted with the event payload + +Related spec scenarios: RS.EVT.18 +*/ +func TestIntegration_XSendEventsShim_LegacyEmission(t *testing.T) { + t.Parallel() + port, stop := startManagementServer(t) + defer stop() + + conn := wsConnect(t, port, "/alerts") + defer conn.Close() //nolint:errcheck + _ = conn.SetReadDeadline(time.Now().Add(2 * time.Second)) + _, _, _ = conn.ReadMessage() // connect built-in welcome + + _, err := http.Post(fmt.Sprintf("http://localhost:%d/_mock/events", port), "application/json", + strings.NewReader(`{"type":"fire","event":"legacyAlert","payload":{"level":"warn"}}`)) + require.NoError(t, err) + + got := readUntilPayload(t, conn, `"legacy":"warn"`) + assert.Contains(t, got, `"legacy":"warn"`) +} + +/* +Scenario: A non-upgrade request to the management stream is rejected +Given a plain HTTP GET to /_mock/stream +When the request is made +Then the server responds with HTTP 405 + +Related spec scenarios: RS.AMG.28 +*/ +func TestIntegration_ManageStream_NonUpgrade405(t *testing.T) { + t.Parallel() + port, stop := startManagementServer(t) + defer stop() + + resp, err := http.Get(fmt.Sprintf("http://localhost:%d/_mock/stream", port)) + require.NoError(t, err) + _ = resp.Body.Close() + assert.Equal(t, 405, resp.StatusCode) +}