How the External Connector Runtime Works for Event Integration in LoopX
The LoopX external connector runtime provides a provider-neutral, file-based system that safely stores, orders, and acknowledges external events for Agents, enforcing durability and privacy guarantees through cursor-based state management and exclusive file locking.
LoopX separates event integration into two distinct layers: providers that generate external events and a runtime that handles durable processing. This architecture lives primarily in loopx/extensions/external_connector_runtime.py and enables secure, ordered event ingestion without exposing sensitive payloads to the provider layer. Understanding this runtime is essential for building reliable Agent integrations with external services.
Core Concepts of the Connector Runtime
The runtime implements six fundamental primitives that work together to guarantee safe event processing.
Connector Binding
A connector binding is a concise, owner-local description that links an Agent to a single external source. It captures the goal reference, provider kind, source type, capture policies, and storage locations.
The build_external_connector_binding() function validates all required fields, including token formats and enumerated policy values. Located at line 205 of the runtime file, this factory ensures every binding conforms to the expected schema before any events flow through the system.
Inbox and Cursor State
Each binding receives dedicated storage paths constructed by _binding_paths() at line 558. This includes:
- An inbox directory for pending event JSON files
- A cursor state file tracking the last committed cursor and pending acknowledgements
These paths remain owner-local, ensuring no other Agent or process can access the binding's private data.
Event Capture
The capture_external_connector_events() function (line 591) implements the main ingestion logic. Providers deliver events as pages containing up to 500 events plus cursor metadata. The runtime:
- Validates the page against the expected cursor
- Filters duplicate events
- Applies capture policies (e.g.,
addressed_onlyfor group messages) - Writes accepted events to
inbox/events/<sha256(event_id)>.json - Updates cursor state only after verification—without advancing until events are acknowledged
Inbox Drain
Agents retrieve pending work through drain_external_connector_inbox() at line 735. This returns up to 20 events (configurable) in FIFO order based on internal sequence numbers. Critically, the drain operation returns only owner-private event data—no provider payload details leak through this interface.
Inbox Inspection
For monitoring without exposure, inspect_external_connector_inbox() (line 765) provides a content-free status projection. This reveals counts, freshness metrics, and failure statistics without disclosing private identifiers or event contents.
Acknowledgement Decision
Before an event can be settled, decide_external_event_ack() at line 331 verifies that:
- A durable effect receipt exists (e.g., a committed turn or TODO update)
- If the connector's response policy requires it, a verified response receipt is present
This prevents premature acknowledgements that could lose work or violate delivery guarantees.
Event Settlement
Finally, settle_external_connector_event() at line 864 performs the atomic acknowledgement. It removes the event from inbox storage, records the acknowledgement in cursor state, and optionally advances to page_next_cursor if provided by the provider.
All state-mutating operations use exclusive_file_lock from loopx/file_lock.py to guarantee atomicity across concurrent processes.
Complete Event Integration Workflow
The following sequence illustrates a typical integration from binding creation through settlement.
Step 1: Create the Binding
from loopx.extensions.external_connector_runtime import build_external_connector_binding
binding = build_external_connector_binding(
goal_ref="customer_support",
agent_ref="support_agent_001",
provider_kind="slack",
source_kind="group_message",
source_ref="#triage-channel",
capture_policy="addressed_only",
ingress_policy="async_inbox",
response_policy="topic_reply",
cursor_ref="slack_cursor.json",
lifecycle="connected",
capabilities=["realtime_receive"],
inbox_ref="slack_inbox",
)
This configures the runtime to capture only messages addressed to the Agent, store them in a local inbox, and expect threaded replies as responses.
Step 2: Provider Constructs Event Page
On the provider side, raw events transform into runtime-compatible pages:
from loopx.extensions.external_connector_provider import build_external_connector_provider_page
provider_page = build_external_connector_provider_page(
events=[
{
"schema_version": "agent_external_connector_event_v0",
"event_id": "evt-20240904-001",
"content": "Help needed with invoice #4421",
"occurred_at": "2024-09-04T14:30:00Z",
"addressed": True,
"event_cursor": "cursor_abc",
}
],
expected_cursor="cursor_prev",
next_cursor="cursor_def",
has_more=False,
)
Step 3: Runtime Captures the Page
from loopx.extensions.external_connector_runtime import capture_external_connector_events
result = capture_external_connector_events(
project="/var/loopx/projects/acme_corp",
binding=binding,
events=provider_page["events"],
expected_cursor=provider_page["expected_cursor"],
next_cursor=provider_page["next_cursor"],
execute=True,
)
print(f"Captured {result['accepted_count']} events, {result['duplicate_count']} duplicates filtered")
The runtime writes accepted events to the inbox directory and updates checkpoint state atomically under file lock.
Step 4: Agent Drains Pending Events
from loopx.extensions.external_connector_runtime import drain_external_connector_inbox
pending = drain_external_connector_inbox(
project="/var/loopx/projects/acme_corp",
binding=binding,
limit=10,
)
for event in pending["items"]:
print(f"Processing {event['event_id']}: {event['content'][:50]}...")
Step 5: Generate Durable Effect and Settle
from loopx.extensions.external_connector_runtime import (
decide_external_event_ack,
settle_external_connector_event
)
effect_receipt = {
"schema_version": "agent_external_event_effect_receipt_v0",
"event_id": event["event_id"],
"effect_id": "turn-20240904-883",
"status": "committed",
"effect_kind": "working_session_turn",
}
# Verify acknowledgement is permitted
ack_decision = decide_external_event_ack(
project="/var/loopx/projects/acme_corp",
binding=binding,
event_id=event["event_id"],
effect_receipt=effect_receipt,
)
if ack_decision["allowed"]:
settlement = settle_external_connector_event(
project="/var/loopx/projects/acme_corp",
binding=binding,
event_id=event["event_id"],
effect_receipt=effect_receipt,
response_receipt=None, # Required if response_policy demands it
execute=True,
)
print(f"Settled: cursor advanced to {settlement['new_cursor']}")
Policy Configuration Reference
Capture and response policies control runtime behavior at the binding level.
| Capture Policy | Behavior |
|---|---|
addressed_only |
Accept only events explicitly addressed to the Agent |
configured_source_all |
Accept all events from the configured source |
incremental |
Accept events incrementally, tracking per-event cursors |
| Response Policy | Required for Ack |
|---|---|
no_response |
Effect receipt only |
topic_reply |
Effect receipt + verified response receipt |
| Ingress Policy | Storage Target |
|---|---|
async_inbox |
Runtime-managed inbox (file-based) |
sync_direct |
Direct synchronous delivery (bypasses runtime) |
Concurrency and Durability Guarantees
The LoopX external connector runtime provides specific guarantees through its implementation:
- At-least-once delivery: Events persist in inbox until explicitly acknowledged
- Exactly-once processing: Duplicate detection via
event_idprevents reprocessing - Ordered consumption: FIFO drain based on internal sequence numbers
- Crash recovery: Cursor state and inbox files survive process restarts
- Cross-process safety: Exclusive file locks prevent corruption during concurrent capture/drain operations
These properties make the runtime suitable for production Agent deployments requiring reliable external integrations.
Key Source Files
| Path | Purpose |
|---|---|
loopx/extensions/external_connector_runtime.py |
Core runtime implementation: binding, capture, drain, inspect, ack decision, settlement |
loopx/extensions/external_connector_provider.py |
Provider-side helpers for permission management and page construction |
loopx/file_lock.py |
Cross-process exclusive locking primitive |
tests/test_external_connector_runtime.py |
Comprehensive test coverage demonstrating all API behaviors |
Summary
- The external connector runtime in LoopX decouples event providers from Agent processing through a file-based, cursor-governed inbox system
build_external_connector_binding()creates validated, owner-local connector configurations with configurable policiescapture_external_connector_events()ingests provider pages with duplicate filtering, policy enforcement, and atomic cursor updatesdrain_external_connector_inbox()returns pending events in FIFO order while preserving privacy of provider payloadsdecide_external_event_ack()andsettle_external_connector_event()implement conditional, atomic acknowledgements that advance cursor state only after durable effects commit- All state mutations are protected by
exclusive_file_lockfor cross-process safety
Frequently Asked Questions
What is the maximum number of events per capture page?
The runtime accepts pages containing up to 500 events. Providers should batch appropriately and use the has_more flag to indicate additional pages available for subsequent capture calls.
How does the runtime prevent duplicate event processing?
Each event carries a unique event_id which the runtime hashes to generate the inbox filename. The capture_external_connector_events() function filters events whose IDs already exist in the pending set, ensuring idempotent capture even if providers retry delivery.
Can multiple Agents drain the same inbox concurrently?
No—bindings are owner-local and tied to a specific Agent reference. The runtime path construction includes the Agent identifier, creating namespace isolation. Within a single Agent, concurrent drain operations are serialized through the exclusive file lock mechanism.
What happens if an Agent crashes during event processing?
Unacknowledged events remain in the inbox across restarts. The cursor state preserves the last committed position with the provider. On restart, the Agent drains pending events from the inbox and reprocesses them—the runtime provides at-least-once delivery guarantees, so Agents must implement idempotent effect generation for true exactly-once semantics.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →