Documentation
¶
Index ¶
- Constants
- Variables
- func CanonicalRequest(method, path, timestamp, eventID string, body []byte) string
- func DLQRetentionFromEnv(getenv func(string) string) time.Duration
- func Install(ctx ValueStore, service *Service) error
- func MarshalCallbackRequest(event Event, format CallbackFormat) ([]byte, error)
- func Retryable(err error) bool
- func Sign(secret, method, path, timestamp, eventID string, body []byte) string
- func ValidateAuthoritativeDeciderIDs(deciderUID, deciderSpaceID string) error
- func Verify(secret, signature, method, path, timestamp, eventID string, body []byte) bool
- type CallbackFormat
- type CardContext
- type DecisionRequest
- type DecisionResult
- type DeliveryError
- type Dispatcher
- type DispatcherConfig
- type Disposition
- type Event
- type Finalizer
- type FinalizerFunc
- type FinalizerKey
- type FinalizerRegistry
- type HTTPDeliverer
- type Lease
- type Metrics
- type NackOutcome
- type NotifyCapability
- type QueueConfig
- type QueueDepths
- type RedisQueue
- func (q *RedisQueue) Ack(eventID int64, token string) (bool, error)
- func (q *RedisQueue) Claim(now time.Time, leaseDuration time.Duration) (*Lease, error)
- func (q *RedisQueue) ClearRouteMissing(eventID int64, token string) (bool, error)
- func (q *RedisQueue) Defer(eventID int64, token string, due time.Time) (bool, error)
- func (q *RedisQueue) Depths() (QueueDepths, error)
- func (q *RedisQueue) DepthsNoPrune() (QueueDepths, error)
- func (q *RedisQueue) Enqueue(event Event, due time.Time) error
- func (q *RedisQueue) Nack(lease Lease, now time.Time, delay time.Duration, maxAttempts int, ...) (NackOutcome, error)
- func (q *RedisQueue) ReclaimExpired(now time.Time, limit int) (int, error)
- func (q *RedisQueue) Renew(eventID int64, token string, now time.Time, leaseDuration time.Duration) (bool, error)
- func (q *RedisQueue) ReplayDLQ(eventID int64, due time.Time) (bool, error)
- func (q *RedisQueue) RouteMissingSeenAt(eventID int64, now time.Time) (time.Time, error)
- type Registry
- func (r *Registry) CanNotify(capability NotifyCapability, actionType string) bool
- func (r *Registry) NotifyProducers() []NotifyCapability
- func (r *Registry) Resolve(senderUID, owner, actionType string) Resolution
- func (r *Registry) ResolveNotifyToken(token string) (NotifyCapability, bool)
- func (r *Registry) Route(senderUID, owner, actionType string) (*Route, bool)
- func (r *Registry) ValidateNotifyTokenExclusions(tokens ...string) error
- type Resolution
- type ResolutionKind
- type Route
- type RouteSpec
- type Service
- func (s *Service) CanNotify(capability NotifyCapability, actionType string) bool
- func (s *Service) Enqueue(event Event) (int64, error)
- func (s *Service) NotifyProducers() []NotifyCapability
- func (s *Service) Resolve(senderUID, owner, actionType string) Resolution
- func (s *Service) ResolveNotifyToken(token string) (NotifyCapability, bool)
- type State
- type ValueStore
Constants ¶
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" )
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.
const MaxDecisionResponseBytes = 64 << 10
Variables ¶
var ErrServiceAlreadyInstalled = errors.New("cardactiondispatch: service already installed")
Functions ¶
func CanonicalRequest ¶
func DLQRetentionFromEnv ¶ added in v1.12.0
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 ValidateAuthoritativeDeciderIDs ¶ added in v1.17.0
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.
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 ¶
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 ¶
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) Stop ¶
func (d *Dispatcher) Stop()
func (*Dispatcher) String ¶
func (d *Dispatcher) String() string
type DispatcherConfig ¶
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 ¶
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 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 ¶
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 QueueDepths ¶
type RedisQueue ¶
type RedisQueue struct {
// contains filtered or unexported fields
}
func NewRedisQueue ¶
func NewRedisQueue(client *rd.Client, cfg QueueConfig) (*RedisQueue, 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 ¶
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) Nack ¶
func (q *RedisQueue) Nack(lease Lease, now time.Time, delay time.Duration, maxAttempts int, reason string) (NackOutcome, error)
func (*RedisQueue) ReclaimExpired ¶
func (*RedisQueue) RouteMissingSeenAt ¶ added in v1.12.0
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 (*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) ValidateNotifyTokenExclusions ¶
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 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 ¶
type Service ¶
type Service struct {
// contains filtered or unexported fields
}
func FromContext ¶
func FromContext(ctx ValueStore) (*Service, bool)
func NewService ¶
func (*Service) CanNotify ¶
func (s *Service) CanNotify(capability NotifyCapability, actionType string) bool
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)