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 ¶
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.