cardactiondispatch

package
v1.20.0 Latest Latest
Warning

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

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

Documentation

Index

Constants

View Source
const (
	// DefaultDLQRetention is the fallback dead-letter retention: how long a
	// dead-lettered card-action event stays replayable (via tools/card-action-dlq)
	// before pruning, when the env override is unset or invalid. It matches the
	// value the code shipped with before the retention became configurable, so an
	// upgrade that does not set the override keeps the existing recovery window and
	// never silently prunes older DLQ entries on first deploy. Opt into a shorter
	// window via DLQRetentionEnv.
	DefaultDLQRetention = 30 * 24 * time.Hour
	// DLQRetentionEnv names the retention override, expressed in whole days.
	DLQRetentionEnv = "OCTO_CARD_ACTION_DLQ_RETENTION_DAYS"
)
View Source
const (
	HeaderSignature = octosign.HeaderSignature
	HeaderTimestamp = octosign.HeaderTimestamp
	HeaderEventID   = octosign.HeaderEventID
)

Header names come from pkg/octosign so the wire contract has one definition. Kept exported here because existing callers and tests reference them through this package.

View Source
const MaxDecisionResponseBytes = 64 << 10

Variables

View Source
var ErrServiceAlreadyInstalled = errors.New("cardactiondispatch: service already installed")

Functions

func CanonicalRequest

func CanonicalRequest(method, path, timestamp, eventID string, body []byte) string

func DLQRetentionFromEnv added in v1.12.0

func DLQRetentionFromEnv(getenv func(string) string) time.Duration

DLQRetentionFromEnv resolves the dead-letter retention from OCTO_CARD_ACTION_DLQ_RETENTION_DAYS (whole days). Empty / non-integer / non-positive / over-max values fall back to DefaultDLQRetention so a misconfigured override degrades to a safe window rather than truncating the recovery span (NewRedisQueue rejects a non-positive retention outright). Shared by main.go and tools/card-action-dlq so the two binaries never drift on the CODE value. The server is the pruning authority: it prunes lazily on its own Depths() calls with this resolved window. The CLI only ever applies retention on an explicit `replay`; its read-only `depth` never prunes (see DepthsNoPrune), so merely inspecting the DLQ from a shell that lacks the env var can no longer delete server-retained events.

func Install

func Install(ctx ValueStore, service *Service) error

func MarshalCallbackRequest added in v1.12.0

func MarshalCallbackRequest(event Event, format CallbackFormat) ([]byte, error)

MarshalCallbackRequest encodes either the frozen flat request or the route-opted-in octo-card envelope. response_url is intentionally absent until a signed inbound response contract exists.

func Retryable

func Retryable(err error) bool

func Sign

func Sign(secret, method, path, timestamp, eventID string, body []byte) string

func ValidateAuthoritativeDeciderIDs added in v1.17.0

func ValidateAuthoritativeDeciderIDs(deciderUID, deciderSpaceID string) error

ValidateAuthoritativeDeciderIDs validates the optional identity pair returned by a decision consumer. IDs use byte-length bounds, matching Go's len(string) semantics elsewhere in this wire contract. Empty values are allowed for rolling compatibility, but a Space cannot be attributed without a decider.

func Verify

func Verify(secret, signature, method, path, timestamp, eventID string, body []byte) bool

Types

type CallbackFormat added in v1.12.0

type CallbackFormat string
const (
	CallbackFormatLegacy     CallbackFormat = "legacy"
	CallbackFormatOctoCardV1 CallbackFormat = "octo-card-v1"
)

type CardContext added in v1.12.0

type CardContext struct {
	TemplateID      string `json:"template_id,omitempty"`
	TemplateVersion string `json:"template_version,omitempty"`
	View            string `json:"view,omitempty"`
	PrincipalType   string `json:"principal_type,omitempty"`
	PrincipalID     string `json:"principal_id,omitempty"`
	SpaceID         string `json:"space_id,omitempty"`
}

CardContext is derived from the effective server-authored card frame before enqueue. Zero value means a legacy card without registry metadata.

PrincipalType/PrincipalID/SpaceID are additive PR-C D3 fields: the action ingress copies them from the frame's validated `catalog_provenance` marker after proving consistency with the stored sender and authoritative Space. They are absent on frames sent before provenance existed; consumers must not fall back to guessing a principal from sender/owner when they are set.

type DecisionRequest

type DecisionRequest struct {
	EventID     int64                  `json:"event_id,string"`
	ActionID    string                 `json:"action_id"`
	Decision    string                 `json:"decision"`
	OperatorUID string                 `json:"operator_uid"`
	DocID       string                 `json:"doc_id,omitempty"`
	RequestID   string                 `json:"request_id,omitempty"`
	Inputs      map[string]interface{} `json:"inputs"`
	Data        map[string]interface{} `json:"data,omitempty"`
	MessageID   string                 `json:"message_id"`
	ChannelID   string                 `json:"channel_id"`
	ChannelType uint8                  `json:"channel_type"`
	SpaceID     string                 `json:"space_id,omitempty"`
	// OperatorSpaceID carries the operator's trusted current Space to the
	// decision consumer alongside SpaceID (the card's origin Space), so the
	// consumer can attribute the click without conflating the two.
	OperatorSpaceID string `json:"operator_space_id,omitempty"`
	ActedAt         int64  `json:"acted_at"`
	// contains filtered or unexported fields
}

func DecisionRequestFromEvent

func DecisionRequestFromEvent(event Event) DecisionRequest

type DecisionResult

type DecisionResult struct {
	Disposition  Disposition `json:"disposition"`
	State        State       `json:"state"`
	RequesterUID string      `json:"requester_uid,omitempty"`
	// Authoritative decider identity. Docs is the source of truth for WHO
	// decided; for a duplicate/concurrent click that is the FIRST decider, not
	// the current clicker, so these win over the clicker for display. octo-server
	// resolves the decider's display name and operator Space name from these ids
	// internally (user service + active-membership check) — a caller never sends
	// display copy that overrides them. DecidedAt is the authoritative decision
	// timestamp/label produced by Docs.
	DeciderUID     string            `json:"decider_uid,omitempty"`
	DeciderSpaceID string            `json:"decider_space_id,omitempty"`
	DecidedAt      int64             `json:"decided_at,omitempty"`
	Display        map[string]string `json:"display,omitempty"`
}

func DecodeDecisionResult

func DecodeDecisionResult(reader io.Reader) (DecisionResult, error)

type DeliveryError

type DeliveryError struct {
	Category string
	Status   int
	// contains filtered or unexported fields
}

func (*DeliveryError) Error

func (e *DeliveryError) Error() string

func (*DeliveryError) Unwrap

func (e *DeliveryError) Unwrap() error

type Dispatcher

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

func NewDispatcher

func NewDispatcher(queue dispatchQueue, registry *Registry, deliverer callbackDeliverer, finalizer Finalizer, cfg DispatcherConfig) (*Dispatcher, error)

func (*Dispatcher) ProcessOne

func (d *Dispatcher) ProcessOne(ctx context.Context, now time.Time) (bool, error)

ProcessOne claims and completely handles at most one due event. Callback and finalization failures are converted into a durable retry/DLQ transition; an error is returned only when the queue state itself could not be made safe.

func (*Dispatcher) Start

func (d *Dispatcher) Start(parent context.Context) error

func (*Dispatcher) Stop

func (d *Dispatcher) Stop()

func (*Dispatcher) String

func (d *Dispatcher) String() string

type DispatcherConfig

type DispatcherConfig struct {
	LeaseDuration   time.Duration
	PollInterval    time.Duration
	ReclaimInterval time.Duration
	Metrics         *Metrics
	Logger          interface {
		Warn(string, ...zap.Field)
		Error(string, ...zap.Field)
	}
}

type Disposition

type Disposition string
const (
	DispositionApplied   Disposition = "applied"
	DispositionReplayed  Disposition = "replayed"
	DispositionForbidden Disposition = "forbidden"
	DispositionConflict  Disposition = "conflict"
	DispositionNotFound  Disposition = "not_found"
)

type Event

type Event struct {
	EventID     int64  `json:"event_id"`
	SenderUID   string `json:"sender_uid"`
	Owner       string `json:"owner"`
	ActionType  string `json:"action_type"`
	MessageID   string `json:"message_id"`
	ChannelID   string `json:"channel_id"`
	ChannelType uint8  `json:"channel_type"`
	SpaceID     string `json:"space_id,omitempty"`
	ActionID    string `json:"action_id"`
	OperatorUID string `json:"operator_uid"`
	// OperatorSpaceID is the operator's TRUSTED current Space at click time —
	// the value SpaceMiddleware validated and pinned in the request context. It
	// is captured separately from SpaceID (the card/document authoritative origin
	// Space) and never relabels it: the two answer different questions and can
	// differ when the operator acts on a card that originates in another of their
	// Spaces. Empty when the click carried no server-verified Space.
	OperatorSpaceID string                 `json:"operator_space_id,omitempty"`
	ClientToken     string                 `json:"client_token,omitempty"`
	ActedAt         int64                  `json:"acted_at"`
	Inputs          map[string]interface{} `json:"inputs"`
	Data            map[string]interface{} `json:"data,omitempty"`
	Card            CardContext            `json:"card,omitempty"`
}

type Finalizer

type Finalizer interface {
	Finalize(ctx context.Context, event Event, result DecisionResult) error
}

type FinalizerFunc

type FinalizerFunc func(context.Context, Event, DecisionResult) error

func (FinalizerFunc) Finalize

func (f FinalizerFunc) Finalize(ctx context.Context, event Event, result DecisionResult) error

type FinalizerKey

type FinalizerKey struct {
	Owner      string
	ActionType string
}

FinalizerKey binds a route to an optional specialized terminal renderer. Routes without a binding use the standard approval fallback.

type FinalizerRegistry

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

FinalizerRegistry keeps custom terminal visuals explicit while preserving config-only onboarding for standard approval consumers.

func NewFinalizerRegistry

func NewFinalizerRegistry(fallback Finalizer, bindings map[FinalizerKey]Finalizer) (*FinalizerRegistry, error)

func (*FinalizerRegistry) Finalize

func (r *FinalizerRegistry) Finalize(ctx context.Context, event Event, result DecisionResult) error

type HTTPDeliverer

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

func NewHTTPDeliverer

func NewHTTPDeliverer(transport http.RoundTripper, clock func() time.Time) *HTTPDeliverer

func (*HTTPDeliverer) Deliver

func (d *HTTPDeliverer) Deliver(ctx context.Context, route *Route, request DecisionRequest) (DecisionResult, error)

type Lease

type Lease struct {
	Event   Event
	Token   string
	Attempt int
}

type Metrics

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

func NewMetrics

func NewMetrics(reg prometheus.Registerer) *Metrics

type NackOutcome

type NackOutcome string
const (
	NackRequeued      NackOutcome = "requeued"
	NackDeadLettered  NackOutcome = "dead_lettered"
	NackTokenMismatch NackOutcome = "token_mismatch"
)

type NotifyCapability

type NotifyCapability struct {
	SenderUID string
	Owner     string
}

NotifyCapability is the server-authoritative identity granted to one first-party notification caller. The bearer token never supplies owner or sender metadata; it resolves to this value at the ingress boundary.

type QueueConfig

type QueueConfig struct {
	Prefix       string
	LiveTTL      time.Duration
	DLQRetention time.Duration
}

type QueueDepths

type QueueDepths struct {
	Ready  int64
	Leased int64
	DLQ    int64
}

type RedisQueue

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

func NewRedisQueue

func NewRedisQueue(client *rd.Client, cfg QueueConfig) (*RedisQueue, error)

func (*RedisQueue) Ack

func (q *RedisQueue) Ack(eventID int64, token string) (bool, error)

func (*RedisQueue) Claim

func (q *RedisQueue) Claim(now time.Time, leaseDuration time.Duration) (*Lease, error)

func (*RedisQueue) ClearRouteMissing added in v1.15.0

func (q *RedisQueue) ClearRouteMissing(eventID int64, token string) (bool, error)

ClearRouteMissing removes the durable route-missing first-seen marker for an event whose route has resolved at dispatch. Once the route is present the event is no longer route-missing, so its lease must be treated as a delivery lease by ReclaimExpired: without this, a once-route-missing event that reaches delivery still carries the marker, and a hard crash during delivery (no Ack/Nack runs) would have reclaimScript refund its attempt as though it were still a defer cycle — asymmetrically exempting it from the MaxAttempts bound. Clearing the marker here makes a delivery-phase crash advance the attempt like any delivery lease. Token-protected so only the current lease owner clears it; returns true when this worker owned the lease.

func (*RedisQueue) Defer

func (q *RedisQueue) Defer(eventID int64, token string, due time.Time) (bool, error)

Defer returns a capacity-blocked lease to ready without consuming a delivery attempt. The lease token and leased-set membership are checked atomically, so a stale worker cannot move a lease owned by another replica.

func (*RedisQueue) Depths

func (q *RedisQueue) Depths() (QueueDepths, error)

Depths prunes DLQ entries older than the retention window, then reports queue depths. The running server calls this (via refreshDepthMetrics), so it is the single pruning authority and prunes lazily with its own resolved retention. Read-only inspectors must use DepthsNoPrune instead so observing the queue cannot delete recoverable entries.

func (*RedisQueue) DepthsNoPrune added in v1.12.0

func (q *RedisQueue) DepthsNoPrune() (QueueDepths, error)

DepthsNoPrune reports queue depths WITHOUT pruning the DLQ. Use it for read-only inspection (the card-action-dlq `depth` command) so merely observing the DLQ can never delete recoverable entries — even from a shell whose OCTO_CARD_ACTION_DLQ_RETENTION_DAYS differs from the server's. Pruning stays the running server's job (see Depths). The reported DLQ count therefore includes any not-yet-pruned expired entries, which is the honest current contents for a manual inspection.

func (*RedisQueue) Enqueue

func (q *RedisQueue) Enqueue(event Event, due time.Time) error

func (*RedisQueue) Nack

func (q *RedisQueue) Nack(lease Lease, now time.Time, delay time.Duration, maxAttempts int, reason string) (NackOutcome, error)

func (*RedisQueue) ReclaimExpired

func (q *RedisQueue) ReclaimExpired(now time.Time, limit int) (int, error)

func (*RedisQueue) Renew

func (q *RedisQueue) Renew(eventID int64, token string, now time.Time, leaseDuration time.Duration) (bool, error)

func (*RedisQueue) ReplayDLQ

func (q *RedisQueue) ReplayDLQ(eventID int64, due time.Time) (bool, error)

func (*RedisQueue) RouteMissingSeenAt added in v1.12.0

func (q *RedisQueue) RouteMissingSeenAt(eventID int64, now time.Time) (time.Time, error)

RouteMissingSeenAt records (once) and returns when this event's route was first observed missing at dispatch. The bounded route-missing defer window is measured from this point — NOT from Event.ActedAt — so an event that sat in the durable queue for a long time before its first dispatch attempt (a long restart/outage/backlog window carried by the durable queue) still gets the full self-heal window on its first transient miss, instead of being dead-lettered immediately because its acted-at is already older than the window.

type Registry

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

func NewRegistry

func NewRegistry(specs []RouteSpec, getenv func(string) string) (*Registry, error)

func (*Registry) CanNotify

func (r *Registry) CanNotify(capability NotifyCapability, actionType string) bool

CanNotify keeps a capability scoped to only the action types that explicitly declare its notify_token_env. A token for one owner cannot mint another owner's card or access a callback-only route.

func (*Registry) NotifyProducers

func (r *Registry) NotifyProducers() []NotifyCapability

func (*Registry) Resolve

func (r *Registry) Resolve(senderUID, owner, actionType string) Resolution

func (*Registry) ResolveNotifyToken

func (r *Registry) ResolveNotifyToken(token string) (NotifyCapability, bool)

ResolveNotifyToken performs a constant-time comparison against every configured first-party notification capability. Tokens are unique across capabilities, so at most one result can match.

func (*Registry) Route

func (r *Registry) Route(senderUID, owner, actionType string) (*Route, bool)

func (*Registry) ValidateNotifyTokenExclusions

func (r *Registry) ValidateNotifyTokenExclusions(tokens ...string) error

ValidateNotifyTokenExclusions prevents a route-scoped approval token from accidentally inheriting a broader legacy/docs notify capability.

type Resolution

type Resolution struct {
	Kind  ResolutionKind
	Route *Route
}

type ResolutionKind

type ResolutionKind string
const (
	ResolutionCallback ResolutionKind = "callback"
	ResolutionBotPull  ResolutionKind = "bot_pull"
	ResolutionReject   ResolutionKind = "reject"
)

type Route

type Route struct {
	SenderUID      string
	Owner          string
	ActionType     string
	URL            string
	Timeout        time.Duration
	MaxAttempts    int
	BaseBackoff    time.Duration
	MaxBackoff     time.Duration
	MaxInFlight    int
	CallbackFormat CallbackFormat
	// contains filtered or unexported fields
}

type RouteSpec

type RouteSpec struct {
	SenderUID      string
	Owner          string
	ActionType     string
	URL            string
	SecretEnv      string
	NotifyTokenEnv string
	Timeout        time.Duration
	MaxAttempts    int
	BaseBackoff    time.Duration
	MaxBackoff     time.Duration
	MaxInFlight    int
	CallbackFormat CallbackFormat
}

func LoadRouteSpecs

func LoadRouteSpecs(raw string) ([]RouteSpec, error)

type Service

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

func FromContext

func FromContext(ctx ValueStore) (*Service, bool)

func NewService

func NewService(registry *Registry, queue eventQueue, sequence sequenceGenerator) (*Service, error)

func (*Service) CanNotify

func (s *Service) CanNotify(capability NotifyCapability, actionType string) bool

func (*Service) Enqueue

func (s *Service) Enqueue(event Event) (int64, error)

func (*Service) NotifyProducers

func (s *Service) NotifyProducers() []NotifyCapability

func (*Service) Resolve

func (s *Service) Resolve(senderUID, owner, actionType string) Resolution

func (*Service) ResolveNotifyToken

func (s *Service) ResolveNotifyToken(token string) (NotifyCapability, bool)

type State

type State string
const (
	StatePending   State = "pending"
	StateApproved  State = "approved"
	StateDenied    State = "denied"
	StateCancelled State = "cancelled"
)

type ValueStore

type ValueStore interface {
	SetValue(value interface{}, key string)
	Value(key string) interface{}
}

Jump to

Keyboard shortcuts

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