messagequeue

package
v0.3.0-20261005223632-... Latest Latest
Warning

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

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

README

SubmitQueue internal message-queue contract

Wire payloads for the queues internal to the SubmitQueue pipeline (gateway and orchestrator). It is internal — used only within the SubmitQueue domain — so it lives under submitqueue/core rather than api/ (Bazel visibility keeps it domain-scoped).

Payloads are defined in proto3 (proto/, generated into protopb/) and serialized as protobuf JSON (protojson), so the MySQL-backed queue keeps storing self-describing JSON. The contract package adds protojson glue (Marshal/Unmarshal), TopicKeys, the pipeline TopicKey constants, and helpers that map generated payloads to submitqueue/entity types. MarshalID/UnmarshalID take the topic key so the bound message is used; unmarshalling through a different type would drop fields added later (DiscardUnknown). Each payload declares the topic key that carries it via the topic_keys proto option (defined in api/base/messagequeue); a contract test round-trips every payload and asserts each topic key is bound to exactly one message.

Shared field types Change and Strategy come from api/base/change and api/base/mergestrategy.

Stages

Each topic key has its own message, even when the first version is only an id and a queue, so a stage can grow fields without touching others.

  • start (TopicKeyStart, Start) — gateway publishes the minted request id and land inputs; start persists a Request. Full payload: this seam crosses services.
  • cancel (TopicKeyCancel, Cancel) — gateway publishes the request id to cancel; cancel reloads the Request. Full payload across the gateway/orchestrator seam.
  • validate (TopicKeyValidate, Validate) — start publishes the request id; validate reloads the Request.
  • batch (TopicKeyBatch, Batch) — landconflictsignal publishes the request id; batch reloads the Request.
  • dependency-analysis (TopicKeyDependencyAnalysis, DependencyAnalysis) — batch publishes the batch id; partitioned by queue.
  • speculate (TopicKeySpeculate, Speculate) — dependency-analysis (and later stages) publish the batch id.
  • build (TopicKeyBuild, Build in submitqueuebuild.proto) — speculate publishes a batch id; build reloads the Batch. The proto filename is not build.proto so it does not collide with Stovepipe's stovepipe/core/messagequeue/proto/build.proto in the protobuf filename registry.
  • buildsignal (TopicKeyBuildSignal, BuildSignal in submitqueuebuildsignal.proto) — build publishes a build id; buildsignal polls and may hold the delivery. Same filename-registry reason as build.
  • submitqueue-land (TopicKeyLand, Merge in submitqueuemerge.proto) — speculate publishes a batch id; land reloads the Batch before handing work to Runway. The proto filename is not merge.proto/land.proto so it does not collide with Runway's api/runway/messagequeue/proto/merge.proto in the protobuf filename registry.
  • conclude (TopicKeyConclude, Conclude) — speculate/landsignal publish a batch id. A failed batch's reason travels in message metadata (MetadataKeyFailureReason), not the payload.
  • log (TopicKeyLog, Log) — orchestrator publishes a full request-log entry; the gateway materializes it. type, status, and event are open strings matching the domain vocabularies.

In-boundary stages (validate through conclude, except start/cancel/log) put only an id on the queue because producer and consumer share storage.

Documentation

Overview

Package messagequeue holds SubmitQueue's internal message-queue contract: the wire payloads for the pipeline queues SubmitQueue owns, defined by the proto files in proto/ and generated into protopb/. The proto is the language-neutral authority; the generated Go types in protopb are the binding for Go callers.

It is internal — used only within the SubmitQueue domain — so it lives under submitqueue/core rather than api/. The message types are generated into protopb; this package adds generic protojson glue (Marshal/Unmarshal), the topic-key reflection lookup (TopicKeys), the pipeline TopicKey constants, and helpers that map between generated payloads and submitqueue/entity types at the controller edge. Payloads are serialized as protobuf JSON, not binary, so the MySQL-backed queue keeps storing self-describing JSON. The topic key that carries each payload is declared on the message itself via the topic_keys proto option (see api/base/messagequeue). Proto filenames that would collide with another domain's contract in the protobuf filename registry are prefixed (submitqueuemerge.proto, submitqueuebuild.proto, submitqueuebuildsignal.proto).

Index

Constants

View Source
const MetadataKeyFailureReason = "failure_reason"

MetadataKeyFailureReason is the conclude message's metadata attribute carrying a failed batch's human-readable reason. Set by the failure sites (land and speculate) on the conclude publish and read by conclude to stamp the request's terminal log; absent on the landed and cancelled paths.

Variables

This section is empty.

Functions

func CancelToEntity

func CancelToEntity(m *Cancel) entity.CancelRequest

CancelToEntity copies a cancel payload onto the domain cancellation.

func LandRequestFromStart

func LandRequestFromStart(m *Start) entity.LandRequest

LandRequestFromStart copies a start payload onto the gateway-owned land request.

func LogToEntity

func LogToEntity(m *Log) entity.RequestLog

LogToEntity copies a log payload onto a request-log entry. An empty type is treated as a status entry, matching entries written before the type field existed. A nil metadata map becomes an empty map.

func Marshal

func Marshal(m proto.Message) ([]byte, error)

Marshal serializes any contract message to protojson bytes for the queue payload, keeping the proto field names (snake_case) on the wire.

func MarshalID

func MarshalID(key TopicKey, id, queue string) ([]byte, error)

MarshalID serializes the id-only payload bound to key. Start, cancel, and log are not id-only; callers marshal those messages directly.

func TopicKeys

func TopicKeys(m proto.Message) []string

TopicKeys returns the stable logical topic keys bound to a message via the topic_keys proto option — not concrete wire names; a caller maps each key to its backend's topic name. Returns nil for a message that declares no keys.

func Unmarshal

func Unmarshal[T proto.Message](b []byte, m T) error

Unmarshal deserializes protojson bytes into the contract message m, tolerating unknown fields so an additive contract change is ignored rather than rejected.

func UnmarshalBatchID

func UnmarshalBatchID(key TopicKey, b []byte) (entity.BatchID, error)

UnmarshalBatchID reads a batch-scoped id-only payload.

func UnmarshalBuildID

func UnmarshalBuildID(key TopicKey, b []byte) (entity.BuildID, error)

UnmarshalBuildID reads a buildsignal payload.

func UnmarshalCancelRequest

func UnmarshalCancelRequest(b []byte) (entity.CancelRequest, error)

UnmarshalCancelRequest reads a cancel payload into the domain cancellation.

func UnmarshalID

func UnmarshalID(key TopicKey, b []byte) (id, queue string, err error)

UnmarshalID reads id and queue from the id-only payload bound to key. Consumers pass the topic they subscribe to so a field added to that message is decoded rather than discarded as unknown on a different type.

func UnmarshalLandRequest

func UnmarshalLandRequest(b []byte) (entity.LandRequest, error)

UnmarshalLandRequest reads a start payload into the gateway-owned land request.

func UnmarshalRequestID

func UnmarshalRequestID(key TopicKey, b []byte) (entity.RequestID, error)

UnmarshalRequestID reads a request-scoped id-only payload (validate or batch).

func UnmarshalRequestLog

func UnmarshalRequestLog(b []byte) (entity.RequestLog, error)

UnmarshalRequestLog reads a log payload into a request-log entry.

Types

type Batch

type Batch = protopb.Batch

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

type Build

type Build = protopb.Build

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

type BuildSignal

type BuildSignal = protopb.BuildSignal

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

type Cancel

type Cancel = protopb.Cancel

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

func CancelFromEntity

func CancelFromEntity(r entity.CancelRequest) *Cancel

CancelFromEntity copies a cancellation onto the cancel payload.

type Conclude

type Conclude = protopb.Conclude

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

type DependencyAnalysis

type DependencyAnalysis = protopb.DependencyAnalysis

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

type Log

type Log = protopb.Log

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

func LogFromEntity

func LogFromEntity(r entity.RequestLog) *Log

LogFromEntity copies a request-log entry onto the log payload.

type Merge

type Merge = protopb.Merge

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

type Speculate

type Speculate = protopb.Speculate

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

type Start

type Start = protopb.Start

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

func StartFromLandRequest

func StartFromLandRequest(r entity.LandRequest) *Start

StartFromLandRequest copies a gateway-owned land request onto the start payload.

type TopicKey

type TopicKey = consumer.TopicKey

TopicKey is the typed identifier used to look up a queue backend, topic name, and subscription config in a consumer.TopicRegistry. The constants below are the logical topic keys for SubmitQueue's internal pipeline stages; they are the same strings each message lists in its topic_keys option.

const (
	// TopicKeyStart carries new land requests from the gateway to start.
	TopicKeyStart TopicKey = "start"
	// TopicKeyCancel carries cancellation requests from the gateway to cancel.
	TopicKeyCancel TopicKey = "cancel"
	// TopicKeyValidate carries request ids from start to validate.
	TopicKeyValidate TopicKey = "validate"
	// TopicKeyBatch carries request ids from landconflictsignal to batch.
	TopicKeyBatch TopicKey = "batch"
	// TopicKeyDependencyAnalysis carries newly created batch ids for conflict
	// analysis. Messages must be partitioned by queue: analysis reads the
	// queue's dependency-eligible batches, so two batches of one queue analyzed
	// concurrently would each miss the other.
	TopicKeyDependencyAnalysis TopicKey = "dependency-analysis"
	// TopicKeySpeculate carries batch ids for speculation.
	TopicKeySpeculate TopicKey = "speculate"
	// TopicKeyBuild carries batch ids whose speculated heads should be built.
	TopicKeyBuild TopicKey = "build"
	// TopicKeyBuildSignal carries build ids to poll. The consumer calls
	// BuildRunner.Status, persists the latest status, publishes the batch id to
	// TopicKeySpeculate so the state machine re-evaluates, and holds the
	// delivery for the next poll when the build has not yet reached a terminal
	// state.
	TopicKeyBuildSignal TopicKey = "buildsignal"
	// TopicKeyLand carries batch ids to the internal land stage before Runway.
	TopicKeyLand TopicKey = "submitqueue-land"
	// TopicKeyConclude carries batch ids for terminal request reconciliation.
	TopicKeyConclude TopicKey = "conclude"
	// TopicKeyLog carries per-request log entries from the orchestrator to the gateway.
	TopicKeyLog TopicKey = "log"
)

type Validate

type Validate = protopb.Validate

Wire payload types. These alias the generated protobuf bindings so callers reference the contract through this curated package rather than protopb.

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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