Documentation
¶
Overview ¶
Package signal carries external signals to waiting executions. Signals are published to a durable JetStream stream so they survive restarts and are redelivered until acknowledged; the engine applies them idempotently using each message's stream sequence number (spec §7).
Index ¶
- type Delivery
- type Signals
- func (s *Signals) Consume(ctx context.Context, durable string, ...) (jetstream.ConsumeContext, error)
- func (s *Signals) EnsureStream(ctx context.Context) error
- func (s *Signals) Publish(ctx context.Context, execID, name string, payload []byte) error
- func (s *Signals) Subject(execID, name string) string
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Signals ¶
type Signals struct {
// contains filtered or unexported fields
}
Signals publishes and consumes external signals within one namespace.
func (*Signals) Consume ¶
func (s *Signals) Consume( ctx context.Context, durable string, handler func(context.Context, Delivery) error, ) (jetstream.ConsumeContext, error)
Consume sets up a durable consumer and invokes handler for every signal. The handler must persist state before returning nil; only then is the message acked (CAS-before-ack). A returned error triggers redelivery. The returned ConsumeContext must be stopped by the caller.
func (*Signals) EnsureStream ¶
EnsureStream creates the signals stream if it does not exist.