Ingress

Status: [Preview] — Grove's inbound event pipeline for pulling events from external systems (SQS today; additional sources over time). The design and core runtime are in place; new source types are added incrementally.

Ingress is how events from outside Grove — SQS queues, webhooks with scheduled polling, third-party APIs — become domain events inside a Grove project. It complements the outbox, which handles events flowing outward.

Prerequisites: Workflows Overview, Outbox, familiarity with the Triggers model. What you'll learn: The ingress architecture (sources, scheduler, ingestor, shedder), how events arrive as Grove domain events, and how the SQS source is configured.


The ingress pipeline

An ingress pipeline has four roles:

External system ──▶ Source ──▶ Ingestor ──▶ IngressWriter ──▶ grove_events
                                  ▲
                                  │
                              Scheduler ───── Shedder

Because events land in grove_events via the normal write path, an ingress pipeline makes external systems indistinguishable from internal modules for the rest of your Grove code: triggers fire, workflows start, outbox rows get created, consumers run.

Sources

A Source is implemented by an SDK crate for each external system. It declares:

Today, the shipped source type is SqsSource (from the grove-ingress-sqs crate). Additional source types are planned; the trait is public so internal teams can add new source types without forking grove-core.

Target keys

Every source has a stable TargetKey composed of:

The scheduler, shedder, and rate limiter all key off this triple. The stability rule (per the ingress architecture doc): the hashed principal element must be stable for at least ten times the longest consuming window — since the shedder's rolling window is 30 seconds, the principal must be stable for at least 5 minutes. In practice this means hashing something like an IAM role ARN or an OAuth client ID, not a short-lived session token.

Typed constructors enforce the stability rule at the API level:

ConstructorStable element
TargetKey::aws_iamIAM role ARN
TargetKey::oauth2client_id:resource_owner_id
TargetKey::long_lived_api_keySHA-256 fingerprint of the key material
TargetKey::hmac_secretHMAC key ID (non-sensitive)
TargetKey::labeled_principalCaller-supplied label (escape hatch for unmoded schemes)

Raw-hash construction is deliberately not part of the public API.

The shedder

The shedder is a per-target circuit breaker that protects Grove and the external system from feedback loops. For each TargetKey it tracks a rolling 30-second window of observations:

A target enters the shedding state when either:

Shedding means the ingestor stops polling that target. The target remains shed until it stays clean for the recovery window (30 seconds by default), at which point polls resume. This is intentionally symmetric: unhealthy targets are isolated quickly, but recovery requires sustained clean signal so flapping doesn't thrash the scheduler.

All shedder thresholds are configurable through ShedderConfig when the ingestor is constructed.

The scheduler

Two scheduler implementations share the same trait:

Swapping schedulers is a construction-time decision; ingestor and source code are unchanged.

The ingestor

The Ingestor<S> runtime drives sources via the scheduler's claim loop. You register each source with a SourceId, and the ingestor:

  1. Asks the scheduler for the next source to poll.
  2. If the shedder reports the target as healthy, invokes the source's poll method.
  3. Records the poll outcome (latency, error) with the shedder.
  4. Returns the claim to the scheduler along with any cursor updates.
  5. Repeats.

A single Ingestor can run any number of heterogeneous sources; placement decisions live on each source.

SQS source

The grove-ingress-sqs crate provides SqsSource, a Source implementation that polls an AWS SQS queue, converts each message to a CloudEvents-shaped payload, and writes it via the ingress writer.

Configuration (SqsSourceConfig):

FieldDefaultMeaning
queue_urlrequiredSQS queue URL to poll
regionfrom envAWS region (e.g. us-east-1)
max_messages10Max messages per ReceiveMessage (SQS caps at 10)
wait_time_seconds20Long-poll wait time (0-20; 20 is SQS's maximum)
project_idrequiredProject ID the events belong to
role_arnrequiredIAM role ARN used as the stable identity for this target key
tenant_idrequiredTenant ID for the target key
schema_version1CloudEvents schema version recorded on each outbox entry

Long-polling with wait_time_seconds = 20 is the recommended setting: SQS holds the connection until a message arrives or the wait time expires, minimizing empty polls and cost. The SQS source deletes each message only after the ingress writer has acknowledged the write — if Grove crashes between pull and write, the message becomes visible again after SQS's visibility timeout and is re-pulled.

Message conversion

Each SQS message is converted to a CloudEvents-shaped payload before it's handed to the writer. The conversion preserves:

Downstream triggers and activities see these as normal Grove events, routed to the configured module by the project's ingress wiring.

Testing

The SQS source is abstracted over a SqsClient trait so tests can inject a fake. The contract test exercises an end-to-end pull → convert → write → event flow without an SQS instance, so project-level integration tests can run fully offline.

Wiring an ingress pipeline

Wiring an ingress pipeline is currently a grove-core-level construction step rather than a project-level configuration surface. Contact your Manzano admin or follow the grove-ingress crate docs if you need to add a new source to a project; a user-facing configuration surface for declaring ingress sources from grove.toml is on the roadmap.

See Also