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
- Source — knows how to pull events from one external system (one SQS queue, one webhook endpoint). Emits events into an
IngressWriter. - Scheduler — decides which source to poll next and hands out work to ingestor workers.
- Ingestor — the runtime that drives sources through the scheduler's claim loop.
- Shedder — per-target circuit breaker that stops hammering unhealthy targets.
- IngressWriter — writes incoming events to
grove_events(and therefore the outbox) so downstream triggers and activities see them like any domain event.
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:
- Scheme — a short identifier like
"sqs"or"webhook". - Placement — where it runs: on every replica (
PerReplica), on a single cluster-elected replica (ClusterSingleton), or on a scheduled cadence (Scheduled). - Target key — the stable identity used for scheduling, shedding, and rate-limit accounting (see below).
- Poll — the actual work: pull events, emit them via the writer, return a
PollOutcome(Emitted(n),Empty, orThrottled(duration)).
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:
- Domain — the external service (e.g.
"sqs.us-east-1.amazonaws.com"). - Principal hash — a SHA-256 hash of a stable identity element.
- Tenant id — the tenant the source belongs to.
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:
| Constructor | Stable element |
|---|---|
TargetKey::aws_iam | IAM role ARN |
TargetKey::oauth2 | client_id:resource_owner_id |
TargetKey::long_lived_api_key | SHA-256 fingerprint of the key material |
TargetKey::hmac_secret | HMAC key ID (non-sensitive) |
TargetKey::labeled_principal | Caller-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:
- Success(latency) — the poll succeeded with the measured latency.
- Error — the poll failed.
A target enters the shedding state when either:
- P50 emit latency exceeds 5 seconds (default), or
- Error rate exceeds 25% (default)
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:
NaiveScheduler— local, no coordination. Suitable for development, tests, and single-replica deployments.- Tiller scheduler — database-backed, claim-lease coordination across replicas for production. Uses the same claim-lease pattern as the outbox to ensure each source is polled by one worker at a time.
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:
- Asks the scheduler for the next source to poll.
- If the shedder reports the target as healthy, invokes the source's
pollmethod. - Records the poll outcome (latency, error) with the shedder.
- Returns the claim to the scheduler along with any cursor updates.
- 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):
| Field | Default | Meaning |
|---|---|---|
queue_url | required | SQS queue URL to poll |
region | from env | AWS region (e.g. us-east-1) |
max_messages | 10 | Max messages per ReceiveMessage (SQS caps at 10) |
wait_time_seconds | 20 | Long-poll wait time (0-20; 20 is SQS's maximum) |
project_id | required | Project ID the events belong to |
role_arn | required | IAM role ARN used as the stable identity for this target key |
tenant_id | required | Tenant ID for the target key |
schema_version | 1 | CloudEvents 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:
- SQS message ID, receipt handle, and message attributes.
- The body (JSON by default; raw bytes otherwise).
- Timestamps (enqueue time, first-receive time).
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
- Outbox — the outbound companion: how Grove events leave the project
- Triggers — what fires when an ingress-delivered event lands in
grove_events - Activities — consuming ingress-sourced events in a workflow
- Workflows Overview — the bigger picture: ingress feeds workflows