intake

package
v0.2.3 Latest Latest
Warning

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

Go to latest
Published: Oct 3, 2026 License: Apache-2.0 Imports: 23 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

Why a poller

Webhooks need a public ingress route to the receiving process; that is a trust-model change this service does not make. The intake 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 'owner/name#12=reviewer,owner/name=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 memory are 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 wakes the agent that owns a pull request or issue when something actionable happens to it: a comment, a review, an inline review comment, or a failed CI run.

It is a tailnet-side poller, not a webhook receiver. Accepting GitHub's webhooks would need public ingress, which is a trust-model change this service deliberately does not make; instead 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.

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 ipchttp.IHTTPClient the binary grants — normally one built by ipchttp.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 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 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 ValidateRepository

func ValidateRepository(repository string) error

ValidateRepository reports whether repository is an owner/name pair.

Types

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 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 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 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 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 ipchttp.IHTTPClient) Option

WithHTTPClient grants the HTTP capability the GitHub API is reached through — normally ipchttp.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 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 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 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.

Jump to

Keyboard shortcuts

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