intake

package
v0.3.0 Latest Latest
Warning

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

Go to latest
Published: Oct 7, 2026 License: Apache-2.0 Imports: 38 Imported by: 0

README

intake — GitHub events wake the owning agent

services/intake turns actionable GitHub activity on a pull request or issue into a typed candace.intake.v1.Event and delivers it through the relay to the agent that owns that pull request or issue. An idle in-process agent blocked in Messenger.Receive wakes when the event arrives; a host or network agent finds it at its next inbox fetch.

Kind GitHub source
EVENT_KIND_COMMENT IssueCommentEvent (created) on a pull request or an issue
EVENT_KIND_REVIEW PullRequestReviewEvent (submitted review, any state)
EVENT_KIND_REVIEW_COMMENT PullRequestReviewCommentEvent (inline diff comment)
EVENT_KIND_CHECK_FAILURE a workflow run whose conclusion is failure, one event per attempt

The webhook receiver

csf serve mounts WebhookReceiver at POST /api/intake/github. GitHub delivers there as things happen; nothing polls.

  1. Verify: X-Hub-Signature-256 must be the HMAC-SHA256 of the body under github.webhook_secret in <state>/providers.json (owner-only), compared in constant time.
  2. Keep once: the raw body is kept as deliveries/<X-GitHub-Delivery>.json; a redelivery of a kept GUID is a duplicate.
  3. Record: each delivery and each typed event is a record of the GitHub event stream.
  4. Route: SessionRoute and MergeRoute take each event to whoever acts on it.
  5. Recover: at start, every delivery since the last kept one that never arrived is redelivered.

A missing secret refuses every delivery with no_webhook_secret; a missing, malformed or wrong signature is bad_signature.

The GitHub event stream is the run directory <state>/00000000-0000-0000-0000-000000000401: an events.jsonl and the kept deliveries, and no run.json, so it is never resumed or shown as a session. csf events -assignment 00000000-0000-0000-0000-000000000401 prints it, and the mining loop mines it with every other run's log.

Kind Delivery Action Routed to
EVENT_KIND_PULL_REQUEST_OPENED pull_request opened
EVENT_KIND_PULL_REQUEST_MERGED pull_request closed, merged dispatcher
EVENT_KIND_PULL_REQUEST_CLOSED pull_request closed
EVENT_KIND_COMMENT issue_comment created owning session
EVENT_KIND_REVIEW pull_request_review submitted owning session
EVENT_KIND_REVIEW_COMMENT pull_request_review_comment created owning session
EVENT_KIND_CHECK_RUN_COMPLETED check_run completed owning session
EVENT_KIND_CHECK_SUITE_COMPLETED check_suite completed owning session
EVENT_KIND_PUSH push
EVENT_KIND_ISSUE_OPENED issues opened
EVENT_KIND_ISSUE_CLOSED issues closed
EVENT_KIND_RELEASE_PUBLISHED release published Workbench

The dispatcher takes a merge only when its base is main, and the owning session takes a comment only on a pull request and a check only when it failed. The owning session is the live session whose branch is the event's head branch, or whose recorded pull request is the event's. A csf release published after the host started shows on the Workbench as the notice that csf upgrade installs it. /metrics exports csf_github_deliveries_total, csf_github_events_total and csf_github_event_latency_seconds, each with a panel in the CSF dashboard's GitHub row.

Why a poller as well

Webhooks need a public ingress route to the receiving process; that is a trust-model change the operator makes, not this service. On a host without it, the poller dials out to the GitHub REST API instead and needs no listener at all:

  • each feed (/repos/{owner}/{name}/events, /repos/{owner}/{name}/actions/runs?status=failure) is requested with If-None-Match, so an unchanged feed answers 304 and does not count against the rate limit;
  • a round never runs sooner than GitHub's X-Poll-Interval;
  • a spent primary limit (X-RateLimit-Remaining: 0) backs off until X-RateLimit-Reset; a secondary limit backs off for Retry-After; no wait exceeds WithMaxBackoff (default one hour).

Mounting it (bare binary first)

client, err := ipchttp.NewHTTPClient(ipcnet.NewHostNetwork()) // the only way out
core, err := relay.NewRelay[*intakev1.Event]()
messenger, err := relay.NewMessenger[*intakev1.Event](core)
routes, err := intake.NewStaticRoutes(intake.Route{Repository: "owner/name", Number: 12, Agent: "reviewer"})
service, err := intake.NewEventIntake(
    intake.WithHTTPClient(client),
    intake.WithRelay(core, messenger),
    intake.WithRoutes(routes),
    intake.WithRepositories(routes.Repositories()...),
    intake.WithToken(token), // read by the binary, never by the service
)
host.Mount("relay", core)     // before the intake: it registers on Start
host.Mount("intake", service) // starts the poller and the dispatcher

app/intake is the runnable composition:

CSF_INTAKE_GITHUB_TOKEN=… go run ./app/intake/cmd --routes 'candacelabs/example#12=reviewer,candacelabs/example=triage'

With no agent attached it registers a journal in each routed agent's place and logs one JSON line per delivered event. It is not deployed anywhere.

Contracts

  • IEventQueue — deduplication and hand-off between poller and dispatcher. MemoryEventQueue is the implementation; a durable one in csfpg must keep the semantics written on the interface.
  • IRoutes — which agent owns a subject. StaticRoutes is the configured table until agents' ownership claims are stored.
  • IAgentDirectory — the relay's Register/Resolve; *relay.Relay[*intakev1.Event] satisfies it.

Limits

  • The queue and deduplication are kept in memory. Events older than the service's start are ignored by default (WithBacklogSince), so a restart does not replay the feed — and events that arrived while it was down are not delivered either.
  • Delivery is at most once: an event whose agent is not registered when it is dispatched is counted in Status().Undeliverable and dropped.
  • Only the first page of each feed is read per round; a repository busier than 100 events per poll interval loses the overflow.
  • The counters are exposed through Status() only; there are no Prometheus series or dashboard yet.

Documentation

Overview

Package intake turns GitHub activity into the typed candace.intake.v1.Event and takes each to whoever acts on it. It has two sources.

The WebhookReceiver is the one csf serve mounts: GitHub delivers to WebhookPath as things happen. Each delivery's X-Hub-Signature-256 is checked against the webhook secret in constant time; a verified delivery is kept once by its X-GitHub-Delivery GUID (its raw body is the corpus copy), recorded with its typed events on the GitHub event stream (Stream), and every event is offered to the [EventRoute]s: SessionRoute queues pull request feedback in the session that owns the branch, MergeRoute marks a merge into main on the slice dispatcher. At Start, recovery asks GitHub to redeliver every delivery since the last one kept that never arrived. The receiver needs public ingress to its one route, which is the operator's decision; until then it is proven with signed replays.

The EventIntake is a tailnet-side poller for a host without ingress: it polls the GitHub REST API with conditional requests (an unchanged feed answers 304 and costs no rate limit), honours the X-Poll-Interval GitHub asks for, and backs off until the rate-limit window resets when it is exhausted. The rest of this comment describes the poller.

Flow

poll → normalize → deduplicate → queue → route → relay.

Each feed (CS-6: a registry in data, see feeds.go) turns one GitHub response into typed candace.intake.v1 Events. The IEventQueue drops an event whose identifier it has already accepted and holds the rest; the IRoutes table names the agent registered for the event's pull request or issue; the relay delivers the event into that agent's inbox. An in-process agent is woken by its blocked Receive returning; a host or network agent finds it at its next inbox fetch.

Capabilities and goroutines

The GitHub API is reached only through the iohttp.IHTTPClient the binary grants — normally one built by iohttp.NewHTTPClient from an ipc/net dialer. The service reads no environment and opens no socket of its own. It is a runtime service: EventIntake.Start starts exactly two goroutines on the scope it is given, the poller and the dispatcher, and both return when that scope is canceled.

Limits

The queue and the deduplication memory are in memory (a durable queue in csfdb is a follow-up), so a restart forgets what was delivered; the default backlog cut-off — events older than the service's start are ignored — is what keeps a restart from re-waking every agent. Routing is a static table until agents' ownership claims are stored. Delivery is at most once: an event whose agent is not registered when it is dispatched is counted and dropped.

Index

Constants

View Source
const (
	// DefaultAPIBaseURL is GitHub's REST API.
	DefaultAPIBaseURL = "https://api.github.com"
	// DefaultPollInterval is the wait between polling rounds when GitHub asks
	// for no longer one.
	DefaultPollInterval = time.Minute
	// DefaultMaxBackoff caps any single wait, however far away GitHub says the
	// rate-limit reset is, so a skewed clock cannot park the poller for hours.
	DefaultMaxBackoff = time.Hour
	// DefaultSourceAgent is the relay agent identifier the intake registers
	// itself under and sends from.
	DefaultSourceAgent relay.AgentID = "intake"
)
View Source
const (
	// DefaultQueueCapacity is how many accepted events wait for dispatch
	// before Offer blocks the poller.
	DefaultQueueCapacity = 256
	// DefaultDeduplicationMemory is how many recent event identifiers the
	// memory queue remembers. A GitHub feed page holds at most 100 events,
	// so this spans many pages of every feed.
	DefaultDeduplicationMemory = 4096
)
View Source
const (
	// StreamAssignment names the stream's run directory.
	StreamAssignment = "00000000-0000-0000-0000-000000000401"
	// StreamSession is the session identifier its records carry.
	StreamSession = "github"
	// DeliveriesDirectory, inside the stream's directory, keeps each
	// accepted delivery's raw body as <delivery>.json: the corpus copy, and
	// the record that the delivery was received.
	DeliveriesDirectory = "deliveries"
)

The GitHub event stream is a run directory of its own under the state directory, <state>/<StreamAssignment>, holding an events.jsonl and no run record. So the readers of every run's event log read it unchanged: csf events -assignment <StreamAssignment> prints it, and the mining loop, which mines <corpus>/*/events.jsonl, mines it. Having no run record, it is never resumed as a session and never shown as a session's card.

View Source
const (
	// EventTypeDelivery records one delivery and its outcome.
	EventTypeDelivery = "github_delivery"
	// EventTypeEvent records one typed event a delivery carried.
	EventTypeEvent = "github_event"
	// EventTypeRouted records what a route did with one event.
	EventTypeRouted = "github_event_routed"
)

The stream's record types.

View Source
const (
	KeyDelivery   = "delivery"
	KeyEvent      = "github_event"
	KeyOutcome    = "outcome"
	KeyRefusal    = "refusal"
	KeyKind       = "kind"
	KeyEventID    = "event_id"
	KeyRepository = "repository"
	KeySubject    = "subject"
	KeyNumber     = "number"
	KeyBranch     = "branch"
	KeyBase       = "base"
	KeyConclusion = "conclusion"
	KeyActor      = "actor"
	KeyURL        = "url"
	KeySummary    = "summary"
	KeyOccurredAt = "occurred_at"
	KeyLatency    = "latency_seconds"
	KeyPayload    = "payload"
	KeyTarget     = "target"
	KeyDetail     = "detail"
)

The stream records' attribute keys.

View Source
const (
	// WebhookPath is the one route GitHub's webhooks are delivered to.
	WebhookPath = "/api/intake/github"
	// HeaderSignature carries the HMAC-SHA256 of the body under the
	// webhook's secret, as sha256=<hex>.
	HeaderSignature = "X-Hub-Signature-256"
	// HeaderEvent names the delivery's event: pull_request, push, ...
	HeaderEvent = "X-GitHub-Event"
	// HeaderDelivery is the delivery's GUID, the same on every redelivery.
	HeaderDelivery = "X-GitHub-Delivery"

	// MaxDeliveryBytes is GitHub's own cap on a webhook payload.
	MaxDeliveryBytes = 25 << 20
)

The webhook route and the headers GitHub sends on every delivery.

View Source
const (
	WebhookPullRequest              = "pull_request"
	WebhookPullRequestReview        = "pull_request_review"
	WebhookPullRequestReviewComment = "pull_request_review_comment"
	WebhookIssueComment             = "issue_comment"
	WebhookIssues                   = "issues"
	WebhookCheckRun                 = "check_run"
	WebhookCheckSuite               = "check_suite"
	WebhookPush                     = "push"
	WebhookRelease                  = "release"
	// WebhookPing is the delivery GitHub sends when a hook is created.
	WebhookPing = "ping"
)

The X-GitHub-Event names of the webhook deliveries the receiver turns into events. A delivery of any other name is acknowledged and recorded as ignored.

View Source
const (
	// RecoveryRetention is how far back GitHub keeps a hook's deliveries:
	// three days. A host that never kept a delivery looks back this far.
	RecoveryRetention = 72 * time.Hour
	// RecoverySlack is how far before the newest kept delivery recovery
	// still looks. Deliveries are answered concurrently, so one GitHub sent
	// earlier can be kept after a later one; the slack covers that overlap
	// many times over, since a delivery is answered within GitHub's own
	// ten-second timeout.
	RecoverySlack = 10 * time.Minute

	// EventTypeRedelivery records one redelivery recovery asked GitHub for.
	EventTypeRedelivery = "github_redelivery"
	// KeyHook is a hook's identifier on a redelivery record.
	KeyHook = "hook"
)
View Source
const (
	TargetSession    = "session"
	TargetDispatcher = "dispatcher"
)

The targets the routes record.

View Source
const MainBranch = "main"

MainBranch is the branch whose merges move the dispatcher's frontier.

View Source
const ProviderName = "intake"

ProviderName is the Provider of every SourceAddress.

Variables

View Source
var (
	// ErrMissingOption reports a required grant that was not supplied.
	ErrMissingOption = errors.New("intake: a required option is missing")
	// ErrAlreadyStarted reports a second Start.
	ErrAlreadyStarted = errors.New("intake: already started")
)
View Source
var (
	// ErrUnrouted reports an event whose subject no route names.
	ErrUnrouted = errors.New("intake: no agent owns this subject")
	// ErrInvalidRoute reports a route without a repository or a valid agent.
	ErrInvalidRoute = errors.New("intake: invalid route")
)
View Source
var (
	// ErrNoSecret is what a secret source returns when no webhook secret is
	// configured.
	ErrNoSecret = errors.New("intake: no GitHub webhook secret is configured")
	// ErrNoStateDirectory reports a receiver built without its state
	// directory.
	ErrNoStateDirectory = errors.New("intake: WithStateDirectory is required")
	// ErrNoSecretSource reports a receiver built without a secret source.
	ErrNoSecretSource = errors.New("intake: WithWebhookSecret is required")
)
View Source
var ErrInvalidEvent = errors.New("intake: invalid event")

ErrInvalidEvent reports an event that fails its contract's validation.

View Source
var ErrNoSource = errors.New("intake: a source address needs a source name")

ErrNoSource reports a source address built without a source name.

Functions

func EventLine added in v0.3.0

func EventLine(event *intakev1.Event) string

EventLine is one line naming an event: "pull_request_merged candacelabs/csf#12 by octocat: title".

func Latency added in v0.3.0

func Latency(event *intakev1.Event, now time.Time) time.Duration

Latency is how long after the event happened it was recorded at now; zero for an event that names no time.

func NormalizeDelivery added in v0.3.0

func NormalizeDelivery(name string, delivery string, body []byte) (events []*intakev1.Event, handled bool, err error)

NormalizeDelivery turns one webhook delivery into its typed events, each stamped with the delivery's identifier. handled is false for an event name the receiver does not act on.

func ReadWebhookSecret added in v0.3.0

func ReadWebhookSecret(path string) ([]byte, error)

ReadWebhookSecret reads github.webhook_secret from the providers.json at path, which must be readable by its owner only. No file, or no key, is ErrNoSecret.

func SessionMessage added in v0.3.0

func SessionMessage(event *intakev1.Event) string

SessionMessage is the message a session is sent for an event.

func Sign added in v0.3.0

func Sign(body []byte, secret []byte) string

Sign is the X-Hub-Signature-256 value GitHub sends for body under secret: what a replay of a recorded delivery signs it with.

func ValidateRepository

func ValidateRepository(repository string) error

ValidateRepository reports whether repository is an owner/name pair.

Types

type Answer added in v0.3.0

type Answer struct {
	Delivery string   `json:"delivery,omitempty"`
	Event    string   `json:"event,omitempty"`
	Outcome  Outcome  `json:"outcome"`
	Refusal  Refusal  `json:"refusal,omitempty"`
	Events   []string `json:"events,omitempty"`
}

Answer is the receiver's JSON answer to one delivery.

type Clock

type Clock struct {
	Now   func() time.Time
	After func(duration time.Duration) <-chan time.Time
}

Clock is the time source the poller waits on. It is data, not behaviour: SystemClock fills it from package time, and a spec fills it with a fixed time and an After that records what it was asked to wait.

func SystemClock

func SystemClock() Clock

SystemClock is the wall clock.

type Counts added in v0.3.0

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

Counts is the receiver's counters. The key sets are fixed when it is built and only the values change, so the maps are read without a lock.

func NewDeliveryCounts added in v0.3.0

func NewDeliveryCounts() *Counts

NewDeliveryCounts builds the counters at zero: what a receiver starts from, and what the catalog spec describes.

func (*Counts) Collect added in v0.3.0

func (counts *Counts) Collect(metrics chan<- prometheus.Metric)

Collect sends every series; a kind never seen reports no latency.

func (*Counts) Deliveries added in v0.3.0

func (counts *Counts) Deliveries(outcome Outcome, refusal Refusal) int64

Deliveries is how many deliveries ended with outcome and refusal.

func (*Counts) Describe added in v0.3.0

func (counts *Counts) Describe(descriptions chan<- *prometheus.Desc)

Describe sends the receiver's descriptors.

func (*Counts) Events added in v0.3.0

func (counts *Counts) Events(kind intakev1.EventKind) int64

Events is how many events of kind were recorded.

type EventIntake

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

EventIntake polls GitHub and delivers actionable events to their owning agents through the relay. It is a runtime service: mount it after the relay.

func NewEventIntake

func NewEventIntake(options ...Option) (*EventIntake, error)

NewEventIntake validates its options and builds an intake that has not started. WithHTTPClient, WithRelay, WithRoutes and WithRepositories are required.

func (*EventIntake) Source

func (intake *EventIntake) Source() SourceAddress

Source is the address the intake registers and sends from.

func (*EventIntake) Start

func (intake *EventIntake) Start(scope *runtime.Scope) error

Start registers the intake with the relay — which must already be started — and starts the poller and the dispatcher on scope. Both return when scope is canceled; neither outlives it.

func (*EventIntake) Status

func (intake *EventIntake) Status() Status

Status returns a snapshot of the counters.

type EventRoute added in v0.3.0

type EventRoute func(ctx context.Context, event *intakev1.Event) (Routed, error)

EventRoute takes one accepted event to whoever acts on it. Routes are registered in data (WithEventRoutes); the receiver runs every route on every event and records what each did.

func MergeRoute added in v0.3.0

func MergeRoute(merged func(ctx context.Context, pullRequestURL string) ([]string, error)) EventRoute

MergeRoute records a pull request merged into main on the dispatcher's frontier: merged marks merged every slice that recorded the pull request and returns their identifiers.

func SessionRoute added in v0.3.0

func SessionRoute(sessions ISessionMessenger) EventRoute

SessionRoute queues a comment, a review, an inline review comment or a failed check on a pull request as a message in the virtual session that owns it: the live session whose branch is the event's head branch or whose recorded pull request is the event's subject.

type GitHubHookDeliveries added in v0.3.0

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

GitHubHookDeliveries is recovery's IHookDeliveries over CSF's GitHub protocol client.

func NewGitHubHookDeliveries added in v0.3.0

func NewGitHubHookDeliveries(client *github.GitHubClient) (*GitHubHookDeliveries, error)

NewGitHubHookDeliveries builds recovery's view of client.

func (*GitHubHookDeliveries) ListHookDeliveries added in v0.3.0

func (hooks *GitHubHookDeliveries) ListHookDeliveries(ctx context.Context, repository string, hook int64, cursor string) ([]HookDelivery, string, error)

ListHookDeliveries lists one page of a hook's deliveries, newest first.

func (*GitHubHookDeliveries) ListHooks added in v0.3.0

func (hooks *GitHubHookDeliveries) ListHooks(ctx context.Context, repository string) ([]Hook, error)

ListHooks lists the repository's webhooks, up to one page of 100.

func (*GitHubHookDeliveries) RedeliverHookDelivery added in v0.3.0

func (hooks *GitHubHookDeliveries) RedeliverHookDelivery(ctx context.Context, repository string, hook int64, delivery int64) error

RedeliverHookDelivery asks GitHub to deliver one attempt again.

type Hook added in v0.3.0

type Hook struct {
	ID  int64
	URL string
}

Hook is one repository webhook as GitHub lists it.

type HookDelivery added in v0.3.0

type HookDelivery struct {
	ID          int64
	GUID        string
	Event       string
	DeliveredAt time.Time
	StatusCode  int
	Redelivery  bool
}

HookDelivery is one attempt GitHub made to deliver to a hook; a redelivery is an attempt of its own with the same GUID.

type IAgentDirectory

type IAgentDirectory interface {
	Register(ctx context.Context, registration relay.Registration) error
	Resolve(ctx context.Context, agent relay.AgentID) (relay.Registration, error)
}

IAgentDirectory is the part of the relay the intake registers itself with and resolves owning agents through. *relay.Relay[*intakev1.Event] satisfies it.

type IEventQueue

type IEventQueue interface {
	Offer(ctx context.Context, event *intakev1.Event) (bool, error)
	Take(ctx context.Context) (*intakev1.Event, error)
}

IEventQueue holds accepted events between the poller and the dispatcher, and is where deduplication happens. A durable implementation (in csfdb) must keep these semantics:

  • Offer accepts an event at most once per identifier: offering an identifier it has accepted before reports false and queues nothing.
  • Offer reports ErrInvalidEvent for an event without an identifier.
  • An Offer that fails — its context ends before the event is queued — leaves no trace, so offering the same event again can succeed.
  • Take returns accepted events oldest first and blocks until one is available or its context ends.

type IHookDeliveries added in v0.3.0

type IHookDeliveries interface {
	// ListHooks lists the repository's webhooks.
	ListHooks(ctx context.Context, repository string) ([]Hook, error)
	// ListHookDeliveries lists one page of the hook's deliveries, newest
	// first, from cursor ("" for the newest); next is "" after the last page.
	ListHookDeliveries(ctx context.Context, repository string, hook int64, cursor string) (deliveries []HookDelivery, next string, err error)
	// RedeliverHookDelivery asks GitHub to deliver one attempt again.
	RedeliverHookDelivery(ctx context.Context, repository string, hook int64, delivery int64) error
}

IHookDeliveries is the part of the GitHub capability recovery reads a repository's hooks and their deliveries through, and asks redeliveries of.

type IRoutes

type IRoutes interface {
	// Route returns the owning agent, or an error wrapping [ErrUnrouted].
	Route(ctx context.Context, subject *intakev1.Subject) (relay.AgentID, error)
}

IRoutes names the agent that owns an event's subject. The static table is the implementation until agents' ownership claims are stored; a claims implementation answers the same question from them.

type ISessionMessenger added in v0.3.0

ISessionMessenger is the part of the harness that lists the virtual sessions and queues a message into one; the agent session service satisfies it.

type MemoryEventQueue

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

MemoryEventQueue is the in-process IEventQueue: a bounded channel of pending events and a bounded memory of recent identifiers.

func NewMemoryEventQueue

func NewMemoryEventQueue(options ...QueueOption) (*MemoryEventQueue, error)

NewMemoryEventQueue validates its options and returns an empty queue.

func (*MemoryEventQueue) Offer

func (queue *MemoryEventQueue) Offer(ctx context.Context, event *intakev1.Event) (bool, error)

Offer queues event unless its identifier was accepted before; see IEventQueue.

func (*MemoryEventQueue) Take

func (queue *MemoryEventQueue) Take(ctx context.Context) (*intakev1.Event, error)

Take returns the oldest queued event; see IEventQueue.

type Option

type Option func(configuration *configuration) error

Option configures an EventIntake before NewEventIntake builds anything.

func WithAPIBaseURL

func WithAPIBaseURL(base string) Option

WithAPIBaseURL points the poller at a GitHub Enterprise server or a double.

func WithBacklogSince

func WithBacklogSince(since time.Time) Option

WithBacklogSince accepts events that happened at or after since. The default is the moment the intake starts, so a restart does not re-wake every agent with the feed's history.

func WithClock

func WithClock(clock Clock) Option

WithClock replaces the wall clock.

func WithHTTPClient

func WithHTTPClient(client iohttp.IHTTPClient) Option

WithHTTPClient grants the HTTP capability the GitHub API is reached through — normally iohttp.NewHTTPClient over an ipc/net dialer. Required.

func WithIgnoredActors

func WithIgnoredActors(logins ...string) Option

WithIgnoredActors drops events caused by these logins — typically the agents' own bot account, so an agent is not woken by its own comment.

func WithLogger

func WithLogger(logger *slog.Logger) Option

WithLogger receives one record per failed poll, rate-limit backoff and undeliverable event.

func WithMaxBackoff

func WithMaxBackoff(backoff time.Duration) Option

WithMaxBackoff caps any single wait.

func WithPollInterval

func WithPollInterval(interval time.Duration) Option

WithPollInterval sets the wait between rounds. GitHub's X-Poll-Interval lengthens it, never shortens it.

func WithQueue

func WithQueue(queue IEventQueue) Option

WithQueue replaces the in-memory queue. Optional.

func WithRelay

func WithRelay(directory IAgentDirectory, messenger relay.IMessenger[*intakev1.Event]) Option

WithRelay supplies the relay events are delivered through: the directory the intake registers with and resolves owners in, and the messenger it sends on. Required.

func WithRepositories

func WithRepositories(repositories ...string) Option

WithRepositories names the owner/name repositories to poll. Required.

func WithRoutes

func WithRoutes(routes IRoutes) Option

WithRoutes supplies the table naming each subject's owning agent. Required.

func WithSourceAgent

func WithSourceAgent(agent relay.AgentID) Option

WithSourceAgent sets the relay agent identifier the intake sends from.

func WithToken

func WithToken(token string) Option

WithToken authenticates every request. Without one GitHub allows 60 requests an hour, which is too few for more than one repository.

type Outcome added in v0.3.0

type Outcome string

Outcome is how the receiver ended one delivery.

const (
	// OutcomeAccepted: verified, recorded, and its events routed.
	OutcomeAccepted Outcome = "accepted"
	// OutcomeDuplicate: a delivery with this GUID was accepted before.
	OutcomeDuplicate Outcome = "duplicate"
	// OutcomeIgnored: verified and recorded, but of an event or action CSF
	// does not act on.
	OutcomeIgnored Outcome = "ignored"
	// OutcomeRefused: not verified; see the Refusal.
	OutcomeRefused Outcome = "refused"
	// OutcomeFailed: verified, but it could not be recorded; GitHub sees a
	// failure and the delivery is recovered by redelivery.
	OutcomeFailed Outcome = "failed"
)

type QueueOption

type QueueOption func(configuration *queueConfiguration)

QueueOption configures a MemoryEventQueue.

func WithDeduplicationMemory

func WithDeduplicationMemory(memory int) QueueOption

WithDeduplicationMemory sets how many recent identifiers are remembered. An identifier older than that is forgotten and would be accepted again.

func WithQueueCapacity

func WithQueueCapacity(capacity int) QueueOption

WithQueueCapacity sets how many events wait for dispatch before Offer blocks.

type RecoveryReport added in v0.3.0

type RecoveryReport struct {
	Redelivered, Failed int
}

RecoveryReport is what one recovery did.

type Refusal added in v0.3.0

type Refusal string

Refusal is why the receiver refused a delivery. Each is a stable value the answer, the stream and the refused-deliveries series carry.

const (
	// RefusedNoSecret: providers.json holds no github.webhook_secret, so no
	// delivery can be verified and every one is refused.
	RefusedNoSecret Refusal = "no_webhook_secret"
	// RefusedSignature: X-Hub-Signature-256 is missing, malformed, or not
	// the body's HMAC under the secret.
	RefusedSignature Refusal = "bad_signature"
	// RefusedDelivery: X-GitHub-Delivery or X-GitHub-Event is missing or
	// malformed.
	RefusedDelivery Refusal = "bad_delivery"
	// RefusedTooLarge: the body exceeds MaxDeliveryBytes.
	RefusedTooLarge Refusal = "too_large"
)

type Route

type Route struct {
	Repository string
	// Number is a pull request or issue number; zero routes every subject in
	// the repository that has no route of its own.
	Number uint64
	Agent  relay.AgentID
}

Route assigns the events of one repository, or of one pull request or issue in it, to an agent.

func ParseRoutes

func ParseRoutes(specification string) ([]Route, error)

ParseRoutes reads a comma-separated route list in which each entry is owner/name=agent (every subject in the repository) or owner/name#N=agent (pull request or issue N). It validates the syntax only; NewStaticRoutes validates the table.

type Routed added in v0.3.0

type Routed struct {
	// Target names what received the event: a session, the dispatcher, ...
	Target string
	// Detail is the target's own identifier: an assignment, a slice.
	Detail string
}

Routed is what one route did with one event; the zero value means the route does not apply to it.

type SourceAddress

type SourceAddress struct {
	model.InProcessTier
	Source string
}

SourceAddress is where an intake's envelopes come from. The relay only carries envelopes between registered addresses, so the intake registers itself at this address before it sends; nothing is ever delivered to it. It is in-process: the intake runs in the runtime that holds the relay.

func NewSourceAddress

func NewSourceAddress(source string) (SourceAddress, error)

NewSourceAddress addresses the intake named source.

func (SourceAddress) Key

func (address SourceAddress) Key() string

Key is intake/in_process/<source>.

func (SourceAddress) Provider

func (address SourceAddress) Provider() string

Provider is ProviderName.

type StaticRoutes

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

StaticRoutes is a fixed routing table from configuration.

func NewStaticRoutes

func NewStaticRoutes(routes ...Route) (*StaticRoutes, error)

NewStaticRoutes validates routes and builds the table. Two routes for the same subject are rejected rather than resolved by order.

func (*StaticRoutes) Repositories

func (routes *StaticRoutes) Repositories() []string

Repositories lists every repository the table routes, sorted and each once: the repositories worth polling.

func (*StaticRoutes) Route

func (routes *StaticRoutes) Route(ctx context.Context, subject *intakev1.Subject) (relay.AgentID, error)

Route prefers the subject's own route, then its repository's.

type Status

type Status struct {
	// Polls counts requests sent, NotModified the 304s among them,
	// RateLimited the responses that started a backoff, and PollFailures the
	// requests that failed or answered something unusable.
	Polls, NotModified, RateLimited, PollFailures int64
	// Accepted counts events queued, Duplicates events the queue had already
	// accepted, and Skipped events dropped as backlog, ignored actors or
	// invalid.
	Accepted, Duplicates, Skipped int64
	// Delivered counts events handed to their agent's inbox, Unrouted events
	// no route owns, and Undeliverable events whose agent could not be
	// resolved or reached.
	Delivered, Unrouted, Undeliverable int64
	// BackoffUntil is when the latest rate-limit backoff ends; zero if none.
	BackoffUntil time.Time
}

Status is a snapshot of what the intake has done since it started. Every counter only grows.

type Stream added in v0.3.0

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

Stream appends the GitHub event stream and keeps the deliveries' raw bodies. Every record is one whole line appended to the log, so concurrent deliveries interleave whole records.

func NewStream added in v0.3.0

func NewStream(stateDirectory string, source clock.IClock) (*Stream, error)

NewStream builds the stream of the state directory, creating its directory.

func (*Stream) Directory added in v0.3.0

func (stream *Stream) Directory() string

Directory is the stream's run directory.

func (*Stream) Keep added in v0.3.0

func (stream *Stream) Keep(delivery string, body []byte) (kept bool, err error)

Keep stores a delivery's raw body unless one with its GUID is kept already; kept is false for a duplicate.

func (*Stream) Kept added in v0.3.0

func (stream *Stream) Kept(delivery string) bool

Kept reports whether a delivery with this GUID was received.

func (*Stream) LastKept added in v0.3.0

func (stream *Stream) LastKept() (time.Time, error)

LastKept is when the newest kept delivery was received; zero when none is.

func (*Stream) RecordDelivery added in v0.3.0

func (stream *Stream) RecordDelivery(ctx context.Context, answer Answer, cause error)

RecordDelivery appends one delivery's outcome.

func (*Stream) RecordEvent added in v0.3.0

func (stream *Stream) RecordEvent(ctx context.Context, event *intakev1.Event)

RecordEvent appends one typed event, with how long after it happened it was recorded.

func (*Stream) RecordRouted added in v0.3.0

func (stream *Stream) RecordRouted(ctx context.Context, event *intakev1.Event, routed Routed, cause error)

RecordRouted appends what one route did with one event.

type WebhookOption added in v0.3.0

type WebhookOption func(receiver *WebhookReceiver) error

WebhookOption configures a WebhookReceiver.

func WithEventRoutes added in v0.3.0

func WithEventRoutes(routes ...EventRoute) WebhookOption

WithEventRoutes adds routes every accepted event is offered to.

func WithRecovery added in v0.3.0

func WithRecovery(deliveries IHookDeliveries, repositories ...string) WebhookOption

WithRecovery grants the GitHub capability recovery works through and names the owner/name repositories whose hooks deliver here. At Start the receiver asks GitHub to redeliver every delivery since the newest one it kept that it never received: what arrived while this host was down or unreachable.

func WithStateDirectory added in v0.3.0

func WithStateDirectory(directory string) WebhookOption

WithStateDirectory names the harness state directory the GitHub event stream lives in. Required.

func WithWebhookClock added in v0.3.0

func WithWebhookClock(source clock.IClock) WebhookOption

WithWebhookClock replaces the system clock records and latencies are measured on.

func WithWebhookLogger added in v0.3.0

func WithWebhookLogger(logger *slog.Logger) WebhookOption

WithWebhookLogger receives one record per refused or failed delivery.

func WithWebhookSecret added in v0.3.0

func WithWebhookSecret(source func() ([]byte, error)) WebhookOption

WithWebhookSecret grants the source of the webhook secret, read at every delivery so a secret added later takes effect without a restart. A source returning ErrNoSecret or an empty secret makes the receiver refuse every delivery. Required.

type WebhookReceiver added in v0.3.0

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

WebhookReceiver verifies GitHub's webhook deliveries, records each on the GitHub event stream, and routes the events they carry.

func NewWebhookReceiver added in v0.3.0

func NewWebhookReceiver(options ...WebhookOption) (*WebhookReceiver, error)

NewWebhookReceiver validates the whole option set before building the receiver.

func (*WebhookReceiver) Counts added in v0.3.0

func (receiver *WebhookReceiver) Counts() *Counts

Counts is what the receiver has done since it was built.

func (*WebhookReceiver) Receive added in v0.3.0

func (receiver *WebhookReceiver) Receive(ctx context.Context, header http.Header, body io.Reader) (int, Answer)

Receive handles one delivery: its headers and body as GitHub sent them. It returns the HTTP status and the answer GitHub is sent.

func (*WebhookReceiver) Recover added in v0.3.0

func (receiver *WebhookReceiver) Recover(ctx context.Context) (RecoveryReport, error)

Recover runs recovery once, now; Start runs it after a restart.

func (*WebhookReceiver) Register added in v0.3.0

func (receiver *WebhookReceiver) Register(router gin.IRouter)

Register mounts the webhook route on the caller's router.

func (*WebhookReceiver) Start added in v0.3.0

func (receiver *WebhookReceiver) Start(scope *runtime.Scope) error

Start recovers the deliveries missed while this host was down, when recovery was granted, in one goroutine on scope.

Jump to

Keyboard shortcuts

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