Documentation
¶
Overview ¶
Package notify delivers operational events (a tripped watchdog, an unhealthy device, a lost worker) to an external sink such as a webhook.
It is deliberately dependency-free of the rest of this module: nothing here imports another internal package, so its delivery semantics — retry, backoff, drop-on-full-queue, graceful close — can be tested in isolation from the controller that will eventually call it.
The property every caller depends on is that Notify never blocks and never fails its caller. Notify is meant to be called from the scheduler's hot path, from the reaper's sweep, and from HTTP handlers; a wedged or unreachable webhook must never stall a job, a device release, or a sweep. When the delivery queue is full, the event is dropped and the drop is counted, not blocked on.
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Event ¶
type Event struct {
Kind Kind `json:"event"`
Device string `json:"device"`
Job string `json:"job"`
Reason string `json:"reason"`
At time.Time `json:"at"`
}
Event is one notification. It is the payload a Sink delivers, and its JSON tags are a public wire format: a later task documents this shape in the README, so field names and tags must not change casually.
type Kind ¶
type Kind string
Kind names the class of event being reported. Values are part of the public wire format documented for consumers of the webhook sink.
const ( KindWatchdogTrip Kind = "watchdog_trip" KindDeviceUnhealthy Kind = "device_unhealthy" KindWorkerLost Kind = "worker_lost" KindJobLost Kind = "job_lost" KindVerifyFailed Kind = "verify_failed" KindLeaseExpired Kind = "lease_expired" // KindDeviceRecovered is the counterpart to device_unhealthy: a device // that left the pool because a process might still have been holding it // has come back, on proof from its worker rather than an operator's // clear. It is the one event here that reports something going right, // and it exists because an unexplained recovery is as alarming as an // unexplained quarantine — Reason names both what the device was out for // and what proof brought it back. KindDeviceRecovered Kind = "device_recovered" )
type Notifier ¶
type Notifier struct {
// contains filtered or unexported fields
}
Notifier delivers Events to a Sink through a bounded, buffered queue served by a single background goroutine. The zero value is not usable; construct one with New. A nil *Notifier is explicitly supported by every method here — it is the "no webhook configured" deployment, which is the default, so call sites never need a nil check before calling Notify or Close.
func New ¶
New creates a Notifier delivering events to sink, and starts the single background goroutine that serves its queue. Callers must eventually call Close to release that goroutine.
If sink is nil, New returns nil instead of a Notifier that would crash its background goroutine on the first delivery attempt. This matters because a later task constructs the notifier conditionally on whether a webhook is configured — exactly the shape that would otherwise hand New a nil Sink. A nil *Notifier is already the documented "no webhook configured" state that every method here handles safely, so the caller gets correct behaviour for free instead of a crash.
func (*Notifier) Close ¶
Close stops accepting new events and drains and attempts to deliver whatever is already queued. It is safe to call on a nil receiver and safe to call more than once.
When ctx is done (including already-cancelled), Close cancels the shared delivery context, which unsticks any pending backoff sleep and any in-flight or subsequent Sink.Deliver call that itself honours ctx cancellation as required by the Sink contract. Close then unconditionally waits for the background goroutine to exit before returning, so the goroutine is never leaked — but this means ctx only bounds how long Close waits when the Sink cooperates. A Sink that ignores its context can still make Close block past ctx's deadline, because Close will not return while the goroutine is still running a Deliver call. "Bounded by ctx" is therefore a best-effort promise contingent on the Sink, not a hard guarantee independent of it.
func (*Notifier) Dropped ¶
Dropped reports how many events have been dropped so far because the queue was full or because Notify was called after Close. Safe on a nil receiver, returning 0.
Only the full-queue case is logged (in Notify, at the point of the drop); a post-Close drop is counted here but not logged, since by then the Notifier is shutting down and there is no delivery goroutine left to usefully report through. Also see Notify's doc comment for the one drop path this counter cannot see at all: an event that loses its race against a concurrent Close.
func (*Notifier) Notify ¶
Notify enqueues e for delivery and returns immediately. If the queue is full, or the Notifier is nil, closed, or otherwise unable to accept e, the event is dropped and the drop is counted — Notify must never block its caller and must never panic, since it is called from the scheduler's hot path, the reaper's sweep, and HTTP handlers.
If e.At is the zero time, Notify stamps it with time.Now().UTC() before enqueueing. At is part of the documented wire format an operator's consumer may sort or dedupe on, so it must never ship as the Go zero value just because a caller forgot to set it.
Note on shutdown: the check for "has Close been called" and the send to the queue are two separate steps, not one atomic operation. An event whose Notify call is in flight at the exact moment Close runs can lose the race after passing the closed check but before its send is received by drainRemaining — that event is silently lost and NOT counted in Dropped, unlike every other drop path. This is a deliberate accepted gap, not a bug: it is bounded by QueueSize, only possible in the brief shutdown window, and closing it would mean either blocking Notify (which this package exists to avoid) or closing a channel with concurrent senders (a panic risk we specifically designed around, see the queue field comment above).
type Options ¶
type Options struct {
// QueueSize bounds how many events may be buffered awaiting delivery.
// Once full, Notify drops new events rather than blocking its caller.
QueueSize int
// Attempts is the number of delivery attempts made per event before
// giving up on it. Must be at least 1.
Attempts int
// Backoff is the delay before the second attempt; it doubles between
// each subsequent attempt.
Backoff time.Duration
}
Options configures a Notifier.
type Sink ¶
Sink delivers a single Event, returning an error if delivery failed.
Deliver MUST respect ctx: it must stop trying and return once ctx is done. This is not a nicety — Close cancels the shared delivery context and then unconditionally waits for the background goroutine to exit, so a Sink that ignores ctx turns a bounded Close into one that hangs until the Sink itself returns, regardless of the caller's own deadline. The webhook Sink in this package honours ctx via http.NewRequestWithContext; any other Sink implementation must do the same.
func NewWebhook ¶
NewWebhook returns a Sink that POSTs each Event as JSON to url. If client is nil, a default client with a bounded timeout is used. If client is non-nil but has no timeout configured, each request is still individually bounded by defaultWebhookTimeout so an unreachable or hanging endpoint cannot stall a delivery attempt forever.