land

package
v0.3.0-20260930175335-... Latest Latest
Warning

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

Go to latest
Published: Sep 30, 2026 License: Apache-2.0 Imports: 16 Imported by: 0

Documentation

Overview

Package land implements the trigger stage for the asynchronous land. It consumes a batch ready to land, builds the full land request from the batch's member requests (one step per request, in Contains order), and publishes it to Runway's merge queue using the batch id as the client-owned correlation id. Runway executes the request as a merge and publishes the result to the merge-signal queue, which the landsignal stage consumes and correlates back to the batch by that id.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Controller

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

Controller handles land queue messages. Implements consumer.Controller.

It loads the batch and its member requests, assembles the full land request (one step per member request, in Contains order, each carrying that request's change and land strategy), and publishes it to Runway's merge queue. Runway executes the request as a merge and returns the result on the merge-signal queue; the landsignal stage consumes it and transitions the batch. This controller therefore performs no state transition itself.

func NewController

func NewController(
	logger *zap.SugaredLogger,
	scope tally.Scope,
	stores orchstorage.Factory,
	registry consumer.TopicRegistry,
	runwayTopicKey consumer.TopicKey,
	topicKey consumer.TopicKey,
	consumerGroup string,
) *Controller

NewController creates a new land controller for the orchestrator. runwayTopicKey is the runway-owned topic this controller publishes land requests to (TopicKeyMerge).

func (*Controller) ConsumerGroup

func (c *Controller) ConsumerGroup() string

ConsumerGroup returns the consumer group for offset tracking.

func (*Controller) Name

func (c *Controller) Name() string

Name returns the controller name for logging and metrics.

func (*Controller) Process

func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) error

Process publishes the full Runway merge request for this land. Returns nil to ack (success), or error to nack/reject.

Error classification: deserialize and storage failures are non-retryable (reject to DLQ). The publish to runway is retryable — it is the hand-off that keeps the land alive, so a transient enqueue blip should replay rather than strand the batch.

func (*Controller) TopicKey

func (c *Controller) TopicKey() consumer.TopicKey

TopicKey returns the topic key this controller subscribes to.

Jump to

Keyboard shortcuts

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