enrich

package module
v0.0.0-...-7a09e4c Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 7, 2026 License: Apache-2.0 Imports: 25 Imported by: 0

README

Go Reference Code Coverage License

Learn an async API by listening to it.

asyncapi-enrich infers an AsyncAPI specification from observed traffic, the way openapi-enrich does for HTTP.

Status: early, but complete end to end. Recording and enrichment both work and are tested — including on real captures against Finnhub and Yahoo Finance.

The workflow

The same one, with one thing changed. You write down what you want to happen and leave the answer blank; go generate fills the answer in; the specification is enriched from it.

What changes is the unit of work. HTTP gives you a request and its one response, so an interaction is a pair. An async API gives you a connection carrying an ordered sequence of messages in both directions that never ends on its own — so an interaction here is a session: the frames you send, the frames that come back, when each arrived, and a declared condition saying when to stop listening.

You author this — api/sessions.json, named so it does not collide with openapi-enrich's api/interactions.json:

[
  {
    "uri": "wss://ws.finnhub.io?token=$FINNHUB_API_KEY",
    "frames": [
      {"send": {"type": "subscribe", "symbol": "AAPL"}}
    ]
  }
]

and recording fills in the rest:

      {"send": {"type": "subscribe", "symbol": "AAPL"}},
      {"at": "412ms", "receive": {"type": "trade", "data": [{"s": "AAPL", "p": 190.5}]}},
      {"at": "1.08s", "receive": {"type": "ping"}}

Only a receive frame gets an at. A send frame is exactly what you authored — its timing is an artefact of our own scheduling, not something the server did, so there is nothing there worth recording.

A session holds only what is needed to open a connection and what crossed it. Everything else — what the server is called, what the messages mean, why the session exists — belongs in the AsyncAPI document, which is the artefact this file is here to improve. The stop condition is not in the file either: it is what you asked the recorder for, not something the server did, so it lives on the command line.

You can also author an unsubscribe, sent right before the connection closes — the natural end of a session, symmetric to the frames it opened with:

{
  "uri": "wss://ws.finnhub.io?token=$FINNHUB_API_KEY",
  "frames": [{"send": {"type": "subscribe", "symbol": "AAPL"}}],
  "unsubscribe": {"type": "unsubscribe", "symbol": "AAPL"}
}

Recording always ends with the real WebSocket closing handshake (RFC 6455 §7.1.1) rather than just cutting the connection, whether or not a session has one of these.

The URI is what the specification is enriched from, the same way openapi-enrich enriches from a recorded request URL: the scheme gives the server's protocol, the host and port give its host, the path gives the channel address, and a credential in the query gives a httpApiKey security scheme with in: query.

Four things a session has to say that a request does not

  • When to stop. A response ends an HTTP request; nothing ends a feed. So the condition is declared on the command line: a timeout, a number of messages, or — the one worth reaching for — one of each kind. A specification is only complete once every message kind has actually been seen, and a generated reader only has to discriminate when there is more than one kind to tell apart.
  • What did not happen. A timeout with conditions unmet is not a failure. "Sixty seconds, four trades, never a ping" is a fact about the API, and it is reported rather than swallowed.
  • When each frame arrived. A heartbeat interval is only discoverable from timestamps, and a client that does not send one gets dropped.
  • That one recording is not enough. A single frame cannot tell you which of its fields are optional, and a session only sees the kinds that happened to occur. Recordings accumulate and their inferred schemas are merged, which is openapi-merge's job.

Recording

go get -tool github.com/MarkRosemaker/asyncapi-enrich/cmd/asyncapi-record
asyncapi-record -f api/sessions.json -kinds trade=3,ping=1 -timeout 60s

Environment variables in a URI are expanded by the tool rather than the shell, to keep the credential out of shell history — and the URI is written back exactly as authored, so the reference survives and the expansion never reaches disk. A URI that carries a credential outright instead is masked on the way out.

Every session in the file records at once — one session waiting on a quiet feed does not hold up another. A session whose existing frames already satisfy the stop condition is left alone and dials nothing, so rerunning a recording that already succeeded costs nothing:

$ asyncapi-record -f testdata/finnhub/api/sessions.json -kinds trade=2,ping=1 -timeout 60s
sessions[0]: already complete, skipped — 12 received (ping=1, trade=11)

A session that falls short of a stricter condition than last time — -messages 5 where only 3 were captured before, say — is not extended. What it had came from a different connection with its own clock; it is discarded, keeping only what was authored, and recorded again from scratch.

Every frame is written to the file as it arrives, not only once recording finishes, so a crash mid-run loses at most the one frame in flight.

Enriching

go get -tool github.com/MarkRosemaker/asyncapi-enrich/cmd/asyncapi-enrich
asyncapi-enrich -spec api/asyncapi.json -sessions api/sessions.json

This is the half that needs no network, split out from recording for the same reason recording was split out from everything else: the machine that captured the traffic is not necessarily the machine — or the moment — you write the spec on. If api/asyncapi.json does not exist yet, it starts from enrich.NewDocument, a minimal valid AsyncAPI 3.1 document, the same way openapi-enrich does.

For every session, it adds a server for the URI's host, one channel on that server — every WebSocket API recorded so far multiplexes all of its messages over a single connection, so there is no per-topic address to key more than one channel by — a send and a receive operation, and a message with an inferred payload schema for each distinct kind of frame observed. A frame's kind comes from a top-level type field when it has one (Finnhub's {"type": "trade", ...}); failing that, from its one top-level key when it has exactly one (Yahoo's {"subscribe": [...]}, which has no type at all); failing that, every frame going the same direction collapses into one message named after that direction — an honest "we don't know how to tell these apart," not a guess.

A query parameter that looks like a credential — the same field names masking checks against — becomes an httpApiKey security scheme named after it, referenced from the server: the credential the recorder found in Finnhub's URI is what tells enrichment the API needs one, and under what name.

Running it again — after a longer recording, or a second one against a different symbol — extends what is already there instead of duplicating it: a server, channel, message, or schema enrichment already produced is found and merged into, not recreated. That merge is where the tool earns its keep. A schema starts out requiring every field the one sample it came from happened to carry; the moment a second sample is merged in without one of those fields, that field stops being required. No single recording can tell "always present" from "just happened to be there this time" apart — only merging across more than one can, which is the reason two recordings beat one no matter how long either of them ran.

AsyncAPI's schema is JSON Schema, whose type keyword is natively a list, so a field seen as both a real value and null merges into an honest union — ["number", "null"] — rather than needing the workarounds OpenAPI 3.0 schemas required (see merge.go's doc comment for how this compares to openapi-merge, which this package deliberately does not depend on).

Masking

Every frame is masked before it reaches disk — as it is captured, not once at the end, so an incremental save mid-recording never writes a secret either. There is no option to turn masking off, because there is no good reason to want one — a credential committed to a public specification repository is a credential to rotate.

Values are replaced by field name, at any depth, case-insensitively and ignoring separators, so api_key, apiKey and API-KEY are one name. A field whose value is an object or an array is replaced whole.

Masker.URL additionally masks credentials in a URL's query string and user information. That is the case openapi-enrich's masker does not cover and this one must: a feed dialled as wss://host/?token=… puts the secret in the URL itself, where masking headers and bodies never reaches it.

The asyncapi family

Module Purpose
asyncapi Parse, validate, and write AsyncAPI 3.x specifications
asyncapi-enrich (this module) Infer specification content from observed traffic
asyncapi-codegen Generate Go types, clients, and servers from a specification

Additional Information

Contributing

Contributions are welcome — please open an issue or a pull request on GitHub.

License

This project is licensed under the Apache 2.0 License.

Documentation

Overview

Package enrich infers an AsyncAPI specification from observed traffic, the way openapi-enrich does for HTTP.

The workflow is the same one: inside the library you want to generate, you write down what you want to happen and leave the answer blank, run go generate, and the tool fills the answer in and enriches the specification from it.

What differs is the unit of work. HTTP gives you a request and its one response, so an interaction is a pair. An async API gives you a connection that carries an ordered sequence of messages in both directions and never ends on its own, so an interaction here is a Session: the frames you send, the frames that come back, when each of them arrived, and a declared condition — given to Recorder, not stored in the file — that says when to stop listening.

Recording

You author a session's URI, its send frames, and optionally a Session.Unsubscribe frame. Recorder.Record dials every session in a file concurrently, plays each one's send frames in order, and commits every frame that arrives — masked, and saved to disk — until Recorder.Until is met or its timeout expires, then performs the WebSocket closing handshake. A session whose existing frames already satisfy Recorder.Until is left alone and reported as skipped, so a rerun after a successful capture costs nothing.

What a server puts inside a payload is what Masker exists for: every frame is masked before it reaches disk, as it is captured, never after.

Enriching

Enrich turns recorded sessions into an AsyncAPI document: a server per host, one channel per server — every WebSocket API recorded so far multiplexes all of its messages over a single connection, so there is no per-topic address to key more than one channel by — a send and a receive operation, and a payload schema per message kind, inferred from the frames and merged across every one observed.

One recording is not enough to describe an API on its own. A single frame cannot tell you which of its fields are optional, and a session only observes the message kinds that happened to occur while it was listening. Merging is what answers the question the whole tool exists to ask: a field present in every sample merged into a schema is required; a field present in only some of them is not, and no single recording — however carefully read — can tell the two apart alone.

Enrich is safe to call more than once, on more than one recording: a server, channel, message, or schema already in the document is extended rather than duplicated.

Index

Constants

View Source
const Replacement = "***"

Replacement is what a masked value is replaced with.

Variables

View Source
var (
	// ErrNoURI is returned for a session with nothing to dial.
	ErrNoURI = errors.New("uri is required")
	// ErrBothDirections is returned for a frame that is both sent and received.
	ErrBothDirections = errors.New("must not set both send and receive")
	// ErrNoDirection is returned for a frame that is neither sent nor received.
	ErrNoDirection = errors.New("must set either send or receive")
)
View Source
var (
	// ErrNoTimeout is returned for a stop condition with no timeout. Without one
	// a recording of a quiet feed would never return.
	ErrNoTimeout = errors.New("must set a timeout")
	// ErrNoDiscriminator is returned when kinds are counted but no field is named
	// to tell one kind from another.
	ErrNoDiscriminator = errors.New("must set a discriminator when kinds are set")
)

Functions

func Enrich

func Enrich(doc *asyncapi.Document, ss Sessions) error

Enrich updates doc in place from recorded sessions: it adds servers, one channel per server, a send and a receive operation, and payload schemas inferred and merged from every frame observed. It is safe to call more than once, on more than one recording: a server, channel, message, or schema already in doc is extended rather than duplicated.

If a field already carries a ContentSchema declaring a protobuf message — something only a maintainer sets, by hand, from a real .proto; this package only ever infers ContentEncoding on its own (see [detectBinaryEncodings]) — Enrich also fails when the newly recorded examples do not decode as that message: see [validateProtoSchemas]. A maintainer who pastes in the real .proto gets told immediately when a recording stops matching it, rather than finding out from a schema that silently drifted.

func NewDocument

func NewDocument() *asyncapi.Document

NewDocument creates a minimal valid AsyncAPI 3.1.0 document as a starting point.

Types

type Duration

type Duration struct {
	// contains filtered or unexported fields
}

Duration is a time.Duration that reads and writes as a string, e.g. "1.5s", so that a hand-authored interactions file says "60s" rather than 60000000000.

It keeps the text it was read from. A file that says "60s" still says "60s" after a recording rewrites it, rather than drifting to "1m0s" — a recording should show up as a diff of the frames that arrived and nothing else.

func NewDuration

func NewDuration(d time.Duration) Duration

NewDuration returns a Duration of d.

func (Duration) Duration

func (d Duration) Duration() time.Duration

Duration returns the duration.

func (Duration) IsZero

func (d Duration) IsZero() bool

IsZero reports whether the duration is zero, which is what omitzero asks.

func (Duration) MarshalText

func (d Duration) MarshalText() ([]byte, error)

MarshalText implements encoding.TextMarshaler.

func (Duration) String

func (d Duration) String() string

String returns the text the duration was read from, or the form time.Duration.String gives it if it was not read from text.

func (*Duration) UnmarshalText

func (d *Duration) UnmarshalText(b []byte) error

UnmarshalText implements encoding.TextUnmarshaler.

type Frame

type Frame struct {
	// At is how long after the connection opened this frame crossed the wire.
	// It is set by [Recorder.Record], and is what makes a heartbeat interval
	// discoverable.
	At Duration `json:"at,omitzero"`
	// Send is the payload the application writes. You author these.
	Send jsontext.Value `json:"send,omitempty"`
	// Receive is the payload the application reads. Recording appends these.
	//
	// A frame that is not JSON is kept as a JSON string rather than dropped: a
	// feed that answers in base64 or plain text is still a feed worth recording,
	// and what it sent is still the evidence.
	Receive jsontext.Value `json:"receive,omitempty"`
}

Frame is a single message, in one direction.

Exactly one of Send and Receive is set: a frame is either something the application wrote or something it read, never both.

func (*Frame) Validate

func (f *Frame) Validate() error

Validate checks that the frame goes in exactly one direction.

type Masker

type Masker struct {
	// contains filtered or unexported fields
}

Masker replaces the values of fields whose names look like credentials.

It runs before anything is written to disk, never after — including every incremental write while a session is still being recorded, not just the final one. A recording that reaches a file unmasked has already leaked: a specification repository is public, and a credential committed to one is a credential to rotate.

func NewMasker

func NewMasker(extra ...string) *Masker

NewMasker returns a Masker that masks the default field names and any extra names given. Pass nothing for the defaults alone.

func (*Masker) Frame

func (m *Masker) Frame(f *Frame) error

Frame masks a single frame's Send and Receive payloads in place. It is called on every frame as it is captured, not once at the end, so that an incremental save mid-recording never writes a secret to disk even for a moment.

func (*Masker) MaskURI

func (m *Masker) MaskURI(uri string) string

MaskURI parses uri, masks it the way Masker.URL does, and returns the result as a string. A uri that fails to parse is returned unchanged rather than dropped — an unparsable URI is a validation problem, not a secret to hide, and Session.Validate is what catches it.

func (*Masker) URL

func (m *Masker) URL(u *url.URL) *url.URL

URL returns the URL with the value of every query parameter that looks like a credential replaced, and any user information removed.

This is the case openapi-enrich's masker does not cover and this one must: a WebSocket feed authenticated as wss://host/?token=… puts the secret in the URL itself, where masking headers and bodies never reaches it.

func (*Masker) Value

func (m *Masker) Value(v jsontext.Value) (jsontext.Value, error)

Value returns the payload with the value of every field that looks like a credential replaced, at any depth. The order of the remaining fields is kept, so a masked recording still reads like what came off the wire.

type Recorder

type Recorder struct {
	// REQUIRED. Until says when each session stops listening.
	Until *Until
	// Mask is applied to every URI and frame before it is kept. A nil Mask means
	// the default masker, not no masking — there is no way to ask for none,
	// because there is no good reason to want one.
	Mask *Masker
	// Dialer connects to the server. The zero value means
	// [websocket.DefaultDialer].
	Dialer *websocket.Dialer
	// Now returns the current time. The zero value means [time.Now]. Tests set it.
	Now func() time.Time
	// Save, if set, is called after every frame a session records — a send, a
	// receive, or the closing unsubscribe — so that a crash mid-recording loses
	// at most the frame in flight, not the whole session. It is given every
	// session, not just the one that changed, since they share one file.
	//
	// Sessions record concurrently, so Save may be called from several
	// goroutines; Recorder serialises those calls itself, so Save need not be
	// safe for concurrent use on its own.
	Save func(Sessions) error
}

Recorder plays sessions against real servers and fills in what comes back.

func (*Recorder) Record

func (r *Recorder) Record(ctx context.Context, ss Sessions) (*Report, error)

Record plays every session that has not already met the stop condition, and fills in the frames that came back. Sessions record concurrently — recording several costs no more wall-clock time than recording one.

A session whose existing frames already satisfy Recorder.Until is left alone and reported as skipped: rerunning a recording that already succeeded dials nothing and returns at once. A session that falls short — a stricter -messages than a previous run found, say — is recorded from scratch: what it already had came from a different connection with its own clock, so it cannot simply be extended.

Record returns the first error that stopped a session from being recorded. A stop condition that was not met is not one of those: a feed that never sent the message you were waiting for is an answer about the feed, and is reported in the returned Report instead.

func (*Recorder) Session

func (r *Recorder) Session(ctx context.Context, s *Session) (*SessionReport, error)

Session plays one session on its own, without the concurrency or incremental saving Recorder.Record provides for a whole file — see that method for the form actually meant for recording.

type Report

type Report struct {
	// Sessions holds one entry per session, in the order they were recorded.
	Sessions []*SessionReport
}

Report says what a recording actually observed.

It exists because the interesting outcome of a recording is often the thing that did not happen. "Sixty seconds, four trades, never a ping" is what tells you the ping in the documentation is not one this feed sends — and that is a fact about the API, arrived at by observation, which is the whole point.

func (*Report) Complete

func (r *Report) Complete() bool

Complete reports whether every session met every condition it declared.

func (*Report) String

func (r *Report) String() string

String summarises the recording in the form go generate should print.

type Session

type Session struct {
	// REQUIRED. URI is the address to dial, e.g.
	// "wss://ws.finnhub.io?token=$FINNHUB_API_KEY".
	//
	// References to environment variables are expanded when dialling and written
	// back unexpanded, so the file stays runnable without ever holding the
	// credential. A URI that carries a literal credential instead is masked on
	// the way out — see [Masker.URL].
	//
	// The URI is what the specification is enriched from: its scheme gives the
	// server's protocol, its host and port give the host, its path gives the
	// channel address, and a credential in its query gives a security scheme of
	// type httpApiKey with `in: query`.
	URI string `json:"uri"`
	// Frames are the messages of this session, in the order they crossed the wire.
	Frames []*Frame `json:"frames"`
	// Unsubscribe, if set, is sent right before the connection closes — the
	// natural end of a session, symmetric to the subscribe frames it opened
	// with. You author this; it is sent once recording stops, whether the stop
	// condition was met or the timeout ran out, as long as the connection is
	// still open. The closing handshake (RFC 6455 §7.1.1) happens either way,
	// with or without one.
	Unsubscribe jsontext.Value `json:"unsubscribe,omitempty"`
}

Session is one connection, and every frame that crossed it.

You author the URI and the frames you want sent. Recording fills in the frames that came back and the time each of them arrived.

func (*Session) Validate

func (s *Session) Validate() error

Validate checks that the session can be recorded.

type SessionReport

type SessionReport struct {
	// Skipped is true when the session's existing frames already satisfied the
	// stop condition, so it was not dialled at all. A rerun of a recording that
	// already succeeded costs nothing.
	Skipped bool
	// Received is how many frames came back.
	Received int
	// Seen counts the frames of each kind, by the value of the discriminator.
	// It is nil when the session named no discriminator.
	Seen map[string]int
	// Short is how many fewer frames came back than the session asked for.
	Short int
	// Missing maps each kind that came up short to how many more were wanted.
	Missing map[string]int
	// NotJSON counts the received frames that were not JSON and were kept as
	// strings instead. A feed that answers in base64 or plain text shows up here
	// rather than as an error.
	NotJSON int
}

SessionReport says what one session observed.

func (*SessionReport) Complete

func (sr *SessionReport) Complete() bool

Complete reports whether the session met every condition it declared.

func (*SessionReport) String

func (sr *SessionReport) String() string

String summarises what the session observed.

type Sessions

type Sessions []*Session

Sessions is the recorded traffic of an async API, the counterpart of openapi-enrich's interactions file. It lives next to the specification it enriches, e.g. api/sessions.json beside api/asyncapi.json.

It holds only two things: what is needed to open a connection, and what crossed it. Everything else — what the server is called, what the messages mean, why the session exists — belongs in the AsyncAPI document, which is the artefact this one is here to improve.

func LoadFromFile

func LoadFromFile(name string) (Sessions, error)

LoadFromFile reads a sessions file.

func ParseSessions

func ParseSessions(data []byte) (Sessions, error)

ParseSessions parses a sessions file already read into memory — from an embedded filesystem, say, where LoadFromFile cannot reach.

func (Sessions) Validate

func (ss Sessions) Validate() error

Validate checks that every session can be recorded.

func (Sessions) WriteToFile

func (ss Sessions) WriteToFile(name string) error

WriteToFile writes the sessions back, formatted the same way every time.

type Until

type Until struct {
	// REQUIRED. Timeout is how long to listen before giving up.
	Timeout time.Duration
	// Messages is the total number of received frames to wait for.
	Messages int
	// Discriminator is the name of the top-level field that says what kind of
	// message a received frame is, e.g. "type". Required when Kinds is set.
	Discriminator string
	// Kinds maps a value of the discriminator field to the number of frames of
	// that kind to wait for, e.g. {"trade": 3, "ping": 1}.
	//
	// This is the condition worth reaching for: a specification is only complete
	// once every kind of message has actually been seen, and a generated reader
	// only has to discriminate when there is more than one kind to tell apart.
	Kinds map[string]int
}

Until says when a recording stops.

It is not part of a session, because it is not part of the API: it is what you asked the recorder for, not something the server did. An async API does not end a session the way a response ends an HTTP request, so the condition has to be stated somewhere, and it belongs with the recorder rather than in a file that describes a connection.

Recording stops as soon as every condition that was set is satisfied, or when Timeout expires — whichever comes first.

A timeout that expires with conditions unmet is not a failure. It is the useful answer to "does this feed ever send a ping?", and is reported in the SessionReport so that the gap is visible rather than silent.

func (*Until) Validate

func (u *Until) Validate() error

Validate checks that the stop condition can be evaluated.

Directories

Path Synopsis
cmd
asyncapi-enrich command
Command asyncapi-enrich enriches an AsyncAPI specification from a recorded sessions file.
Command asyncapi-enrich enriches an AsyncAPI specification from a recorded sessions file.
asyncapi-record command
Command asyncapi-record plays the sessions of a sessions file against real servers and writes back what came over the wire.
Command asyncapi-record plays the sessions of a sessions file against real servers and writes back what came over the wire.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL