From 34f985d60e1ac100e943e6f4934468a7e768ef35 Mon Sep 17 00:00:00 2001 From: Simon Emms Date: Thu, 24 Sep 2026 17:44:09 +0000 Subject: [PATCH] feat: add emit event definition Signed-off-by: Simon Emms --- definitions.go | 41 +++++ definitions_test.go | 114 ++++++++++++- schema.json | 48 ++++++ schema.yaml | 33 ++++ schema_test.go | 403 ++++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 633 insertions(+), 6 deletions(-) diff --git a/definitions.go b/definitions.go index e4326ad..28b5c2b 100644 --- a/definitions.go +++ b/definitions.go @@ -38,6 +38,7 @@ func buildDefinitions() map[string]*jsonschema.Schema { "eventProperties": eventPropertiesDefinition, propExport: exportDefinition, "externalResource": externalResourceDefinition, + "emitTask": emitTaskDefinition, "flowDirective": flowDirectiveDefinition, "forkTask": forkTaskDefinition, "forTask": forTaskDefinition, @@ -760,6 +761,45 @@ var externalResourceDefinition = &jsonschema.Schema{ }, } +var emitTaskDefinition = &jsonschema.Schema{ + Type: typeObject, + Title: "EmitTask", + Description: "Allows workflows to publish events to event brokers or messaging systems, " + + "facilitating communication and coordination between different components and services.", + Required: []string{"emit"}, + UnevaluatedProperties: falseSchema(), + AllOf: []*jsonschema.Schema{ + {Ref: SchemaRef("taskBase")}, + { + Properties: map[string]*jsonschema.Schema{ + "emit": { + Type: typeObject, + Title: "EmitTaskConfiguration", + Description: "The configuration of an event's emission.", + UnevaluatedProperties: falseSchema(), + Required: []string{"event"}, + Properties: map[string]*jsonschema.Schema{ + "event": { + AdditionalProperties: trueSchema(), + Type: typeObject, + Title: "EmitEventDefinition", + Description: "The definition of the event to emit.", + Properties: map[string]*jsonschema.Schema{ + "with": { + Ref: SchemaRef("eventProperties"), + Title: "EmitEventWith", + Description: "Defines the properties of event to emit.", + Required: []string{"type"}, + }, + }, + }, + }, + }, + }, + }, + }, +} + var flowDirectiveDefinition = &jsonschema.Schema{ Title: "FlowDirective", Description: "Represents different transition options for a workflow.", @@ -1344,6 +1384,7 @@ var taskDefinition = &jsonschema.Schema{ OneOf: []*jsonschema.Schema{ {Ref: SchemaRef("callTask")}, {Ref: SchemaRef("doTask")}, + {Ref: SchemaRef("emitTask")}, {Ref: SchemaRef("forTask")}, {Ref: SchemaRef("forkTask")}, {Ref: SchemaRef("listenTask")}, diff --git a/definitions_test.go b/definitions_test.go index 181d306..2a7f450 100644 --- a/definitions_test.go +++ b/definitions_test.go @@ -17,6 +17,7 @@ package schema import ( + "encoding/json" "slices" "testing" @@ -101,6 +102,7 @@ func TestBuildDefinitionsKeys(t *testing.T) { "doTask", defDocumentMetadata, "duration", + "emitTask", propEndpoint, propError, "eventConsumptionStrategy", @@ -134,19 +136,17 @@ func TestBuildDefinitionsKeys(t *testing.T) { for _, key := range expected { assert.Contains(t, defs, key, "buildDefinitions() should contain %q", key) } - - // emitTask is intentionally unsupported. - assert.NotContains(t, defs, "emitTask", "emitTask must not be present") } -// TestTaskDefinitionOneOf verifies that taskDefinition references exactly the -// supported task types and excludes unsupported ones. +// TestTaskDefinitionOneOf verifies that taskDefinition references each of the +// supported task types. func TestTaskDefinitionOneOf(t *testing.T) { refs := schemaRefs(taskDefinition.OneOf) supported := []string{ SchemaRef("callTask"), SchemaRef("doTask"), + SchemaRef("emitTask"), SchemaRef("forTask"), SchemaRef("forkTask"), SchemaRef("listenTask"), @@ -162,7 +162,109 @@ func TestTaskDefinitionOneOf(t *testing.T) { assert.Contains(t, refs, ref, "taskDefinition.OneOf should reference %q", ref) } - assert.NotContains(t, refs, SchemaRef("emitTask"), "emitTask must not appear in taskDefinition.OneOf") + assert.Len(t, refs, len(supported), "taskDefinition.OneOf must not reference any unexpected task types") +} + +// TestEmitTaskDefinitionShape verifies that emitTask mirrors the upstream Open +// Workflow Specification definition: it inherits taskBase, closes the task and +// emit levels to unknown keys, requires emit, emit.event and (when present) +// emit.event.with.type, and reuses the shared eventProperties definition. +func TestEmitTaskDefinitionShape(t *testing.T) { + e := emitTaskDefinition + + assert.Equal(t, typeObject, e.Type) + assert.Equal(t, []string{"emit"}, e.Required, "emitTask must require 'emit'") + assert.Equal(t, falseSchema(), e.UnevaluatedProperties, "emitTask must reject unevaluated properties") + + require.Len(t, e.AllOf, 2) + assert.Equal(t, SchemaRef("taskBase"), e.AllOf[0].Ref, "emitTask must inherit taskBase") + + emit, ok := e.AllOf[1].Properties["emit"] + require.True(t, ok, "emit property must be present") + assert.Equal(t, typeObject, emit.Type) + assert.Equal(t, []string{testPropEvent}, emit.Required, "emit must require 'event'") + assert.Equal(t, falseSchema(), emit.UnevaluatedProperties, "emit must reject unevaluated properties") + + event, ok := emit.Properties[testPropEvent] + require.True(t, ok, "emit.event property must be present") + assert.Equal(t, typeObject, event.Type) + assert.Empty(t, event.Required, "emit.event must not require any properties (upstream leaves 'with' optional)") + assert.Equal(t, trueSchema(), event.AdditionalProperties, "emit.event must allow additional properties") + + with, ok := event.Properties[propWith] + require.True(t, ok, "emit.event.with property must be present") + assert.Equal(t, SchemaRef("eventProperties"), with.Ref, "emit.event.with must reference eventProperties") + assert.Equal(t, []string{propType}, with.Required, "emit.event.with must require 'type'") +} + +// TestEmitTaskDefinitionJSON verifies that emitTask serialises to the upstream +// JSON Schema keywords within the built schema, in particular that the +// falseSchema/trueSchema helpers marshal to boolean schemas and that emitTask +// is part of the task union. +func TestEmitTaskDefinitionJSON(t *testing.T) { + s, err := BuildSchema("1.0.0", "json") + require.NoError(t, err) + + raw, err := json.Marshal(s) + require.NoError(t, err) + + var doc struct { + Defs map[string]json.RawMessage `json:"$defs"` + } + require.NoError(t, json.Unmarshal(raw, &doc)) + + var task struct { + OneOf []*jsonschema.Schema `json:"oneOf"` + } + require.NoError(t, json.Unmarshal(doc.Defs["task"], &task)) + assert.Contains(t, schemaRefs(task.OneOf), SchemaRef("emitTask"), "$defs.task.oneOf must reference emitTask") + + type withJSON struct { + Ref string `json:"$ref"` + Required []string `json:"required"` + } + type eventJSON struct { + Type string `json:"type"` + Required []string `json:"required"` + AdditionalProperties any `json:"additionalProperties"` + Properties map[string]withJSON `json:"properties"` + } + type emitJSON struct { + Type string `json:"type"` + Required []string `json:"required"` + UnevaluatedProperties any `json:"unevaluatedProperties"` + Properties map[string]eventJSON `json:"properties"` + } + var emitTaskJSON struct { + Type string `json:"type"` + Required []string `json:"required"` + UnevaluatedProperties any `json:"unevaluatedProperties"` + AllOf []struct { + Ref string `json:"$ref"` + Properties map[string]emitJSON `json:"properties"` + } `json:"allOf"` + } + require.NoError(t, json.Unmarshal(doc.Defs["emitTask"], &emitTaskJSON)) + + assert.Equal(t, typeObject, emitTaskJSON.Type) + assert.Equal(t, []string{"emit"}, emitTaskJSON.Required) + assert.Equal(t, false, emitTaskJSON.UnevaluatedProperties, "emitTask.unevaluatedProperties must marshal to false") + require.Len(t, emitTaskJSON.AllOf, 2) + assert.Equal(t, SchemaRef("taskBase"), emitTaskJSON.AllOf[0].Ref) + + emit := emitTaskJSON.AllOf[1].Properties["emit"] + assert.Equal(t, typeObject, emit.Type) + assert.Equal(t, []string{testPropEvent}, emit.Required) + assert.Equal(t, false, emit.UnevaluatedProperties, "emit.unevaluatedProperties must marshal to false") + + event := emit.Properties[testPropEvent] + assert.Equal(t, typeObject, event.Type) + assert.Empty(t, event.Required, "emit.event must not declare required properties") + assert.Equal(t, true, event.AdditionalProperties, "emit.event.additionalProperties must marshal to true") + + with := event.Properties[propWith] + assert.Equal(t, SchemaRef("eventProperties"), with.Ref) + assert.Equal(t, []string{propType}, with.Required) } // TestDurationDefinitionShape verifies that duration supports only the object diff --git a/schema.json b/schema.json index 8bc006c..441a692 100644 --- a/schema.json +++ b/schema.json @@ -563,6 +563,51 @@ } ] }, + "emitTask": { + "type": "object", + "title": "EmitTask", + "description": "Allows workflows to publish events to event brokers or messaging systems, facilitating communication and coordination between different components and services.", + "required": [ + "emit" + ], + "unevaluatedProperties": false, + "allOf": [ + { + "$ref": "#/$defs/taskBase" + }, + { + "properties": { + "emit": { + "type": "object", + "properties": { + "event": { + "type": "object", + "properties": { + "with": { + "$ref": "#/$defs/eventProperties", + "title": "EmitEventWith", + "description": "Defines the properties of event to emit.", + "required": [ + "type" + ] + } + }, + "title": "EmitEventDefinition", + "description": "The definition of the event to emit.", + "additionalProperties": true + } + }, + "title": "EmitTaskConfiguration", + "description": "The configuration of an event's emission.", + "required": [ + "event" + ], + "unevaluatedProperties": false + } + } + } + ] + }, "endpoint": { "title": "Endpoint", "description": "Represents an endpoint.", @@ -1524,6 +1569,9 @@ { "$ref": "#/$defs/doTask" }, + { + "$ref": "#/$defs/emitTask" + }, { "$ref": "#/$defs/forTask" }, diff --git a/schema.yaml b/schema.yaml index 15944dc..90dfbce 100644 --- a/schema.yaml +++ b/schema.yaml @@ -339,6 +339,38 @@ $defs: type: integer type: object unevaluatedProperties: false + emitTask: + allOf: + - $ref: '#/$defs/taskBase' + - properties: + emit: + description: The configuration of an event's emission. + properties: + event: + additionalProperties: true + description: The definition of the event to emit. + properties: + with: + $ref: '#/$defs/eventProperties' + description: Defines the properties of event to emit. + required: + - type + title: EmitEventWith + title: EmitEventDefinition + type: object + required: + - event + title: EmitTaskConfiguration + type: object + unevaluatedProperties: false + description: Allows workflows to publish events to event brokers or messaging + systems, facilitating communication and coordination between different components + and services. + required: + - emit + title: EmitTask + type: object + unevaluatedProperties: false endpoint: description: Represents an endpoint. oneOf: @@ -1045,6 +1077,7 @@ $defs: oneOf: - $ref: '#/$defs/callTask' - $ref: '#/$defs/doTask' + - $ref: '#/$defs/emitTask' - $ref: '#/$defs/forTask' - $ref: '#/$defs/forkTask' - $ref: '#/$defs/listenTask' diff --git a/schema_test.go b/schema_test.go index 6d01839..dfb0b9d 100644 --- a/schema_test.go +++ b/schema_test.go @@ -149,3 +149,406 @@ func TestSchema_RejectsUnknownTopLevelProperties(t *testing.T) { assert.NoError(t, err, "document with only known top-level properties should pass validation") }) } + +const ( + testPropEmit = "emit" + testPropEvent = "event" + testEventType = "com.example.order.placed" + testEventSource = "https://example.com/orders" + testStepName = "step1" + testPropListen = "listen" + testPropIf = "if" + testPropID = "id" + testPropData = "data" + testPropSchema = "dataschema" + testUnknownKey = "unknown" + testUnknownVal = "value" + testOwner = "orders" + testBroker = "kafka" +) + +// emitTaskWorkflow returns a minimal workflow whose only task is the given +// task body. Tests use this to exercise the emit task schema in isolation. +func emitTaskWorkflow(task map[string]any) map[string]any { + doc := minimalWorkflow() + doc[propDo] = []any{ + map[string]any{ + testStepName: task, + }, + } + return doc +} + +// emitTask returns an emit task body whose emit.event.with is the given +// value. +func emitTask(with any) map[string]any { + return map[string]any{ + testPropEmit: map[string]any{ + testPropEvent: map[string]any{ + propWith: with, + }, + }, + } +} + +// emitWorkflow returns a minimal workflow with a single emit task whose +// emit.event.with is the given value. +func emitWorkflow(with any) map[string]any { + return emitTaskWorkflow(emitTask(with)) +} + +// listenOneWorkflow returns a minimal workflow with a single listen task +// whose listen.to.one.with is the given value. This is the existing consumer +// of eventProperties and is used as the baseline for parity checks. +func listenOneWorkflow(with any) map[string]any { + return emitTaskWorkflow(map[string]any{ + testPropListen: map[string]any{ + "to": map[string]any{ + "one": map[string]any{ + propWith: with, + }, + }, + }, + }) +} + +// TestSchema_EmitTask_UpstreamExample verifies that the emit example from the +// Open Workflow Specification (examples/emit.yaml) is accepted. +func TestSchema_EmitTask_UpstreamExample(t *testing.T) { + resolved := resolvedTestSchema(t) + + doc := minimalWorkflow() + doc[propDo] = []any{ + map[string]any{ + "emitEvent": emitTask(map[string]any{ + propSource: "https://petstore.com", + propType: "com.petstore.order.placed.v1", + testPropData: map[string]any{ + "client": map[string]any{ + "firstName": "Cruella", + "lastName": "de Vil", + }, + "items": []any{ + map[string]any{"breed": "dalmatian", "quantity": 101}, + }, + }, + }), + }, + } + + assert.NoError(t, resolved.Validate(doc)) +} + +// TestSchema_EmitTask_Accepted verifies that emit is a valid task type and +// that every property supported by eventProperties is accepted under +// emit.event.with. +func TestSchema_EmitTask_Accepted(t *testing.T) { + resolved := resolvedTestSchema(t) + + tests := []struct { + name string + with map[string]any + }{ + { + name: "type only is accepted", + with: map[string]any{propType: testEventType}, + }, + { + name: "all event properties are accepted", + with: map[string]any{ + testPropID: "evt-1", + propSource: testEventSource, + propType: testEventType, + "time": "2026-01-01T00:00:00Z", + "subject": "order-123", + "datacontenttype": "application/json", + testPropSchema: "https://example.com/schemas/order.json", + testPropData: map[string]any{"orderId": "123"}, + }, + }, + { + name: "expression source is accepted", + with: map[string]any{propType: testEventType, propSource: "${ .source }"}, + }, + { + name: "expression dataschema is accepted", + with: map[string]any{propType: testEventType, testPropSchema: "${ .schema }"}, + }, + { + name: "expression data is accepted", + with: map[string]any{propType: testEventType, testPropData: "${ .payload }"}, + }, + { + name: "non-object data is accepted", + with: map[string]any{propType: testEventType, testPropData: []any{1, "two", true}}, + }, + { + name: "extension attribute is accepted", + with: map[string]any{propType: testEventType, "partitionkey": testOwner}, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + assert.NoError(t, resolved.Validate(emitWorkflow(tc.with))) + }) + } +} + +// TestSchema_EmitTask_RequiredFields verifies the upstream required nesting: +// the task requires emit, emit requires event, and event.with (when present) +// requires type. Upstream does not require event.with itself. +func TestSchema_EmitTask_RequiredFields(t *testing.T) { + resolved := resolvedTestSchema(t) + + tests := []struct { + name string + task map[string]any + expectError bool + }{ + { + // Only taskBase properties, so no task type matches. + name: "task without emit is rejected", + task: map[string]any{testPropIf: "${ true }"}, + expectError: true, + }, + { + name: "null emit is rejected", + task: map[string]any{testPropEmit: nil}, + expectError: true, + }, + { + name: "non-object emit is rejected", + task: map[string]any{testPropEmit: testPropEvent}, + expectError: true, + }, + { + name: "empty emit is rejected", + task: map[string]any{testPropEmit: map[string]any{}}, + expectError: true, + }, + { + name: "non-object emit.event is rejected", + task: map[string]any{testPropEmit: map[string]any{testPropEvent: testPropEvent}}, + expectError: true, + }, + { + name: "empty emit.event is accepted", + task: map[string]any{testPropEmit: map[string]any{testPropEvent: map[string]any{}}}, + }, + { + name: "emit.event with only additional properties is accepted", + task: map[string]any{testPropEmit: map[string]any{testPropEvent: map[string]any{"broker": testBroker}}}, + }, + { + name: "empty emit.event.with is rejected", + task: emitTask(map[string]any{}), + expectError: true, + }, + { + name: "emit.event.with without type is rejected", + task: emitTask(map[string]any{propSource: testEventSource, testPropID: "evt-1"}), + expectError: true, + }, + { + name: "non-object emit.event.with is rejected", + task: emitTask(testEventType), + expectError: true, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + err := resolved.Validate(emitTaskWorkflow(tc.task)) + if tc.expectError { + assert.Error(t, err) + } else { + assert.NoError(t, err) + } + }) + } +} + +// TestSchema_EmitTask_InvalidEventProperties verifies that invalid event +// properties are rejected under emit.event.with, and that the existing listen +// task (which also uses eventProperties) rejects the same values. +func TestSchema_EmitTask_InvalidEventProperties(t *testing.T) { + resolved := resolvedTestSchema(t) + + tests := []struct { + name string + with map[string]any + }{ + { + name: "non-string type is rejected", + with: map[string]any{propType: 123}, + }, + { + name: "non-string id is rejected", + with: map[string]any{propType: testEventType, testPropID: 123}, + }, + { + name: "non-string subject is rejected", + with: map[string]any{propType: testEventType, "subject": true}, + }, + { + name: "non-string datacontenttype is rejected", + with: map[string]any{propType: testEventType, "datacontenttype": 1}, + }, + { + name: "non-URI source is rejected", + with: map[string]any{propType: testEventType, propSource: "not a uri"}, + }, + { + name: "non-string source is rejected", + with: map[string]any{propType: testEventType, propSource: 5}, + }, + { + name: "non-URI dataschema is rejected", + with: map[string]any{propType: testEventType, testPropSchema: "not a uri"}, + }, + { + name: "non-string time is rejected", + with: map[string]any{propType: testEventType, "time": 12345}, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + assert.Error(t, resolved.Validate(emitWorkflow(tc.with)), "emit must reject invalid event properties") + assert.Error(t, resolved.Validate(listenOneWorkflow(tc.with)), "listen must reject the same event properties") + }) + } +} + +// TestSchema_EmitTask_UnevaluatedProperties verifies that unknown keys are +// rejected at the task and emit levels, which set unevaluatedProperties to +// false, while emit.event permits additional properties. +func TestSchema_EmitTask_UnevaluatedProperties(t *testing.T) { + resolved := resolvedTestSchema(t) + + t.Run("unknown key at task level is rejected", func(t *testing.T) { + task := emitTask(map[string]any{propType: testEventType}) + task[testUnknownKey] = testUnknownVal + assert.Error(t, resolved.Validate(emitTaskWorkflow(task))) + }) + + t.Run("unknown key inside emit is rejected", func(t *testing.T) { + task := emitTask(map[string]any{propType: testEventType}) + task[testPropEmit].(map[string]any)[testUnknownKey] = testUnknownVal + assert.Error(t, resolved.Validate(emitTaskWorkflow(task))) + }) + + t.Run("additional key inside emit.event is accepted", func(t *testing.T) { + task := emitTask(map[string]any{propType: testEventType}) + task[testPropEmit].(map[string]any)[testPropEvent].(map[string]any)["broker"] = map[string]any{"name": testBroker} + assert.NoError(t, resolved.Validate(emitTaskWorkflow(task))) + }) + + t.Run("additional key inside emit.event does not bypass with validation", func(t *testing.T) { + task := emitTask(map[string]any{propSource: testEventSource}) + task[testPropEmit].(map[string]any)[testPropEvent].(map[string]any)["broker"] = testBroker + assert.Error(t, resolved.Validate(emitTaskWorkflow(task))) + }) +} + +// TestSchema_EmitTask_TaskBase verifies that the emit task inherits and +// validates the shared taskBase properties. +func TestSchema_EmitTask_TaskBase(t *testing.T) { + resolved := resolvedTestSchema(t) + + tests := []struct { + name string + base map[string]any + expectError bool + }{ + { + name: "all taskBase properties are accepted", + base: map[string]any{ + testPropIf: "${ .enabled }", + propInput: map[string]any{propSchema: map[string]any{"document": map[string]any{propType: "object"}}}, + propOutput: map[string]any{propAs: "${ . }"}, + propExport: map[string]any{propAs: "${ . }"}, + propThen: "end", + propMetadata: map[string]any{"owner": testOwner}, + }, + }, + { + name: "non-string if is rejected", + base: map[string]any{testPropIf: true}, + expectError: true, + }, + { + name: "non-string then is rejected", + base: map[string]any{propThen: 1}, + expectError: true, + }, + { + name: "unknown key inside export is rejected", + base: map[string]any{propExport: map[string]any{testUnknownKey: testUnknownVal}}, + expectError: true, + }, + { + name: "non-object metadata is rejected", + base: map[string]any{propMetadata: testOwner}, + expectError: true, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + task := emitTask(map[string]any{propType: testEventType}) + for k, v := range tc.base { + task[k] = v + } + + err := resolved.Validate(emitTaskWorkflow(task)) + if tc.expectError { + assert.Error(t, err) + } else { + assert.NoError(t, err) + } + }) + } +} + +// TestSchema_EmitTask_TaskOneOf verifies that adding emitTask to the task +// OneOf does not allow emit to be combined with other task types, and that +// existing task types remain valid. +func TestSchema_EmitTask_TaskOneOf(t *testing.T) { + resolved := resolvedTestSchema(t) + + t.Run("existing set task is still accepted", func(t *testing.T) { + assert.NoError(t, resolved.Validate(minimalWorkflow())) + }) + + t.Run("emit alongside other emit tasks is accepted", func(t *testing.T) { + doc := minimalWorkflow() + doc[propDo] = []any{ + map[string]any{testStepName: emitTask(map[string]any{propType: testEventType})}, + map[string]any{"step2": map[string]any{propSet: map[string]any{"x": "y"}}}, + map[string]any{"step3": emitTask(map[string]any{propType: testEventType + ".v2"})}, + } + assert.NoError(t, resolved.Validate(doc)) + }) + + for name, other := range map[string]map[string]any{ + "set": {propSet: map[string]any{"x": "y"}}, + "wait": {propWait: map[string]any{propSeconds: 5}}, + testPropListen: {testPropListen: map[string]any{ + "to": map[string]any{"one": map[string]any{propWith: map[string]any{propType: testEventType}}}, + }}, + } { + t.Run("emit combined with "+name+" is rejected", func(t *testing.T) { + task := emitTask(map[string]any{propType: testEventType}) + for k, v := range other { + task[k] = v + } + assert.Error(t, resolved.Validate(emitTaskWorkflow(task))) + }) + } + + t.Run("task with only an unknown key is rejected", func(t *testing.T) { + assert.Error(t, resolved.Validate(emitTaskWorkflow(map[string]any{testUnknownKey: testUnknownVal}))) + }) +}