Skip to content
Open
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
18 changes: 17 additions & 1 deletion events/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,27 @@ Each "vX" folder contains a "base" folder that contains the base event format, w

For example, an event with "type" == "SAVED_PAYMENT" and "app" == "payments" must have a payload matching schema in the file "payments/SAVED_PAYMENT.yaml".

## Event Envelopes

By default, event schemas are payload schemas composed with `base.yaml`, which
describes the historical Stack envelope (`app`, `version`, `type`, `date`, and
`payload`).

A service version can define a complete, version-specific envelope in
`services/<service>/<version>/base.yaml`. The generator composes each sibling
event schema with that base through JSON Schema `allOf`; `base.yaml` is not
published as an event type.

Ledger v3 uses this mechanism because its NATS and Kafka sinks publish the same
new envelope directly: `type`, `ledger`, `date`, `logSequence`, and the complete
`log`. It does not contain the historical `app`, `version`, or `payload`
properties.

## Payments Versions

We decided to go with stack releases starting at v2.0.x.

Before that, the last payments version was v0.9.7.

This is why we do not have a v1.0.0 directory with events, since payments
was never released with a v1.x version.
was never released with a v1.x version.
95 changes: 80 additions & 15 deletions events/events.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
package events

import (
"embed"
"encoding/json"
"fmt"
"io/fs"
"path/filepath"

"embed"
"sort"

"github.com/pkg/errors"
"github.com/xeipuuv/gojsonschema"
Expand All @@ -19,28 +21,44 @@ var baseEvent string
var services embed.FS

func ComputeSchema(serviceName, eventName string) (*gojsonschema.Schema, error) {
base := map[string]any{}
if err := yaml.Unmarshal([]byte(baseEvent), &base); err != nil {
return nil, err
}

ls, err := services.ReadDir(filepath.Join("services", serviceName))
if err != nil {
return nil, errors.Wrapf(err, "reading events directory for service '%s'", serviceName)
}

var moreRecent string
versions := make([]string, 0, len(ls))
for _, directory := range ls {
if moreRecent == "" || semver.Compare(directory.Name(), moreRecent) > 0 {
moreRecent = directory.Name()
if directory.IsDir() && semver.IsValid(directory.Name()) {
versions = append(versions, directory.Name())
}
}
sort.Slice(versions, func(i, j int) bool {
return semver.Compare(versions[i], versions[j]) > 0
})

if moreRecent == "" {
if len(versions) == 0 {
return nil, fmt.Errorf("error retrieving more recent version directory for service '%s'", serviceName)
}

eventData, err := services.ReadFile(fmt.Sprintf("services/%s/%s/%s.yaml", serviceName, moreRecent, eventName))
for _, version := range versions {
_, err := services.ReadFile(fmt.Sprintf("services/%s/%s/%s.yaml", serviceName, version, eventName))
if err == nil {
return ComputeSchemaForVersion(serviceName, version, eventName)
}
if !errors.Is(err, fs.ErrNotExist) {
return nil, err
}
}

return nil, fmt.Errorf("event schema '%s' not found for service '%s'", eventName, serviceName)
}

// ComputeSchemaForVersion returns the schema for an exact service version.
// Versions with a base.yaml describe a complete event envelope and compose
// the event-specific constraints through allOf. Older versions keep using the
// historical shared envelope with the event schema injected under payload.
func ComputeSchemaForVersion(serviceName, version, eventName string) (*gojsonschema.Schema, error) {
eventData, err := services.ReadFile(fmt.Sprintf("services/%s/%s/%s.yaml", serviceName, version, eventName))
if err != nil {
return nil, err
}
Expand All @@ -50,17 +68,64 @@ func ComputeSchema(serviceName, eventName string) (*gojsonschema.Schema, error)
return nil, err
}

base["properties"].(map[string]any)["payload"] = event
versionBaseData, err := services.ReadFile(fmt.Sprintf("services/%s/%s/base.yaml", serviceName, version))
if err == nil {
base := map[string]any{}
if err := yaml.Unmarshal(versionBaseData, &base); err != nil {
return nil, err
}

allOf, _ := base["allOf"].([]any)
base["allOf"] = append(allOf, event)

return compileSchema(base)
}
if !errors.Is(err, fs.ErrNotExist) {
return nil, err
}

base := map[string]any{}
if err := yaml.Unmarshal([]byte(baseEvent), &base); err != nil {
return nil, err
}

properties, ok := base["properties"].(map[string]any)
if !ok {
return nil, errors.New("base event schema has no properties object")
}
properties["payload"] = event

loader := gojsonschema.NewGoLoader(base)
return gojsonschema.NewSchema(loader)
return compileSchema(base)
}

func compileSchema(schema map[string]any) (*gojsonschema.Schema, error) {
data, err := json.Marshal(schema)
if err != nil {
return nil, errors.Wrap(err, "marshaling schema")
}
return gojsonschema.NewSchema(gojsonschema.NewBytesLoader(data))
}

func Check(data []byte, serviceName, eventName string) error {
schema, err := ComputeSchema(serviceName, eventName)
if err != nil {
return errors.Wrap(err, "computing schema")
}

return validate(data, schema)
}

// CheckForVersion validates an event against an exact service version.
func CheckForVersion(data []byte, serviceName, version, eventName string) error {
schema, err := ComputeSchemaForVersion(serviceName, version, eventName)
if err != nil {
return errors.Wrap(err, "computing schema")
}

return validate(data, schema)
}

func validate(data []byte, schema *gojsonschema.Schema) error {
result, err := schema.Validate(gojsonschema.NewStringLoader(string(data)))
if err != nil {
return errors.Wrap(err, "validating schema")
Expand Down
169 changes: 169 additions & 0 deletions events/events_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,3 +57,172 @@ func TestReconciliationAlertEventSchemas(t *testing.T) {
})
}
}

func TestLedgerV3EventSchemas(t *testing.T) {
t.Parallel()

events := map[string]string{
"CREATED_LEDGER": `{
"type":"CREATED_LEDGER",
"ledger":"orders",
"date":"2026-08-28T22:29:41.133323Z",
"logSequence":1,
"log":{
"sequence":1,
"payload":{"createLedger":{"name":"orders","createdAt":"2026-08-28T22:29:41.133323Z"}},
"responseSignature":{}
}
}`,
"COMMITTED_TRANSACTION": `{
"type":"COMMITTED_TRANSACTION",
"ledger":"orders",
"date":"2026-08-28T22:29:41.133323Z",
"logSequence":2,
"log":{
"sequence":2,
"payload":{"apply":{"ledgerName":"orders","log":{
"type":"NEW_TRANSACTION",
"data":{"createdTransaction":{
"transaction":{
"postings":[{"source":"world","destination":"users:123","amount":1000,"asset":"USD/2","color":""}],
"metadata":{"kind":"sale"},
"timestamp":"2026-08-28T22:29:41.133323Z",
"reference":"order-42",
"id":42,
"reverted":false
},
"accountMetadata":{"users:123":{"tier":"gold"}},
"chapterId":7
}},
"date":"2026-08-28T22:29:41.133323Z",
"id":1
}}},
"responseSignature":{}
}
}`,
"REVERTED_TRANSACTION": `{
"type":"REVERTED_TRANSACTION",
"ledger":"orders",
"date":"2026-08-28T22:29:41.133323Z",
"logSequence":3,
"log":{
"sequence":3,
"payload":{"apply":{"ledgerName":"orders","log":{
"type":"REVERTED_TRANSACTION",
"data":{"revertedTransaction":{
"revertedTransactionId":42,
"revertTransaction":{
"postings":[{"source":"users:123","destination":"world","amount":1000,"asset":"USD/2","color":""}],
"metadata":{},
"timestamp":"2026-08-28T22:29:41.133323Z",
"id":43,
"reverted":false
}
}},
"date":"2026-08-28T22:29:41.133323Z",
"id":2
}}},
"responseSignature":{}
}
}`,
"SAVED_METADATA": `{
"type":"SAVED_METADATA",
"ledger":"orders",
"date":"2026-08-28T22:29:41.133323Z",
"logSequence":4,
"log":{
"sequence":4,
"payload":{"apply":{"ledgerName":"orders","log":{
"type":"SET_METADATA",
"data":{"savedMetadata":{"targetType":"ACCOUNT","accountId":"users:123","metadata":{"tier":"platinum"}}},
"date":"2026-08-28T22:29:41.133323Z",
"id":3
}}},
"responseSignature":{}
}
}`,
"DELETED_METADATA": `{
"type":"DELETED_METADATA",
"ledger":"orders",
"date":"2026-08-28T22:29:41.133323Z",
"logSequence":5,
"log":{
"sequence":5,
"payload":{"apply":{"ledgerName":"orders","log":{
"type":"DELETE_METADATA",
"data":{"deletedMetadata":{"targetType":"TRANSACTION","transactionId":42,"key":"kind"}},
"date":"2026-08-28T22:29:41.133323Z",
"id":4
}}},
"responseSignature":{}
}
}`,
"DELETED_LEDGER": `{
"type":"DELETED_LEDGER",
"ledger":"orders",
"date":"2026-08-28T22:29:41.133323Z",
"logSequence":6,
"log":{
"sequence":6,
"payload":{"deleteLedger":{"name":"orders","deletedAt":"2026-08-28T22:29:41.133323Z"}},
"responseSignature":{}
}
}`,
"SKIPPED_ORDER": `{
"type":"SKIPPED_ORDER",
"ledger":"orders",
"date":"2026-08-28T22:29:41.133323Z",
"logSequence":7,
"log":{
"sequence":7,
"payload":{"apply":{"ledgerName":"orders","log":{
"type":"ORDER_SKIPPED",
"data":{"reason":"TRANSACTION_REFERENCE_CONFLICT","context":{"reference":"order-42","existingTransactionId":"42"}},
"date":"2026-08-28T22:29:41.133323Z",
"id":5
}}},
"responseSignature":{}
}
}`,
}

for eventName, event := range events {
eventName, event := eventName, event
t.Run(eventName, func(t *testing.T) {
t.Parallel()

if err := CheckForVersion([]byte(event), "ledger", "v3.0.0", eventName); err != nil {
t.Fatalf("validate event: %v", err)
}
})
}

if err := CheckForVersion([]byte(events["COMMITTED_TRANSACTION"]), "ledger", "v3.0.0", "SAVED_METADATA"); err == nil {
t.Fatal("expected a committed transaction to be rejected by the saved metadata schema")
}
}

func TestComputeSchemaUsesLatestVersionContainingEvent(t *testing.T) {
t.Parallel()

legacyEvent := []byte(`{
"app":"ledger",
"version":"v2",
"date":"2026-08-28T22:29:41Z",
"type":"COMMITTED_TRANSACTIONS",
"payload":{
"ledger":"orders",
"transactions":[{
"postings":[],
"metadata":{},
"id":42,
"timestamp":"2026-08-28T22:29:41Z",
"reverted":false
}]
}
}`)

if err := Check(legacyEvent, "ledger", "COMMITTED_TRANSACTIONS"); err != nil {
t.Fatalf("validate legacy event from newest version containing it: %v", err)
}
}
Loading
Loading