Member Junction
    Preparing search index...

    Base for integration connectors that ingest from a shared External Data Source (EDS) — the "dbconnection" heart of the ingestion-connector hierarchy.

    It does NOT open its own connection, re-implement per-engine introspection, or write any dialect SQL. Instead it resolves the SAME first-class MJ: External Data Sources row that EDS live-read/materialize use (connection config + credential via CredentialEngine + the engine driver), through ExternalDataSourceRouter, and on top of that shared connection provides the whole integration contract: TestConnection, schema introspection (mapping EDS's ExternalSchemaDescriptor → the framework's SourceSchemaInfo), object/field discovery, AND incremental delta ingestion (FetchChanges).

    FetchChanges is generic across every EDS driver: it passes a structured incrementalSince watermark bound plus raw ordering columns to driver.RunView, and the EDS driver renders the dialect predicate, identifier quoting, and literal formatting itself. So this connector carries NO dialect knowledge — the only per-family variance is whether discovery is authoritative (SQL introspection enumerates the full column set; document introspection samples — see the family subclasses). A future family whose driver can't take these RunView params simply overrides FetchChanges.

    This supersedes the SQL-Server-hardcoded, inline-mssql RelationalDBConnector: every engine's connect/introspect/read is single-sourced in the EDS drivers, so this class is engine-agnostic.

    The MJ: Company Integrations connection names its shared EDS row via Configuration.externalDataSourceID (JSON). Credentials live on the EDS row (CredentialID → CredentialEngine), NEVER in the integration's Configuration.

    Hierarchy (View Summary)

    Index

    Constructors

    Properties

    ExternalDataSourceIDKey: "externalDataSourceID" = 'externalDataSourceID'

    Configuration key on CompanyIntegration.Configuration naming the shared EDS data-source row.

    Accessors

    • get DiscoveryIsAuthoritative(): boolean

      §7 — does this connector's discovery (DiscoverObjects/DiscoverFields → IntrospectSchema) return the AUTHORITATIVE, COMPLETE gamut of objects/fields the credentials expose? Default false (safe): a connector must explicitly affirm this. Override to true ONLY when DiscoverObjects hits a real list/describe endpoint that returns EVERYTHING accessible — so an object/field absent from a refresh genuinely means the source dropped it (and may be deactivated). Leave false for stubbed discovery (returns nothing → static metadata is all we have), a cache-driven IntrospectSchema, or any partial/scoped enumeration — there, absence proves nothing and MUST NOT deactivate.

      Returns boolean

    • get IntegrationName(): string

      The canonical integration name (e.g., "HubSpot", "Rasa.io"). Used by GetActionGeneratorConfig() and IntegrationActionExecutor to match connectors to action Config.IntegrationName.

      Override in subclasses. Defaults to the class name.

      Returns string

    • get MaxConcurrencyHint(): number

      Highest SAFE per-layer concurrency the source tolerates (plan.md §7 peak parallelization) — the ceiling the engine's adaptive controller ramps toward. null → use configured syncConcurrency.

      Returns number

    • get MonotonicWatermark(): boolean

      We fetch strictly in watermark order, so the last batch's max watermark IS the true high-water mark and an updated row always re-surfaces at a new, higher watermark — the engine can safely narrow the next incremental to it (§ MonotonicWatermark).

      Returns boolean

    • get RateLimitPolicy(): RateLimitPolicy

      Token-bucket rate-limit policy for this connector's source API (plan.md §7 peak-aware rate limiting). null → the engine derives a conservative rate from Integration.BatchRequestWaitTime. Override to push to the source's real limits.

      Returns RateLimitPolicy

    • get SupportsBatchWrite(): boolean

      Whether this connector supports batched target writes (plan.md §7 aggressive batching).

      Returns boolean

    • get SupportsCreate(): boolean

      Whether this connector supports creating new records in the external system.

      Returns boolean

    • get SupportsDelete(): boolean

      Whether this connector supports deleting records from the external system.

      Returns boolean

    • get SupportsGet(): boolean

      Whether this connector supports reading/fetching records. Always true.

      Returns boolean

    • get SupportsListing(): boolean

      Whether this connector supports paginated listing of records.

      Returns boolean

    • get SupportsSearch(): boolean

      Whether this connector supports searching/querying records with filters.

      Returns boolean

    • get SupportsUpdate(): boolean

      Whether this connector supports updating existing records in the external system.

      Returns boolean

    • get SupportsUpsert(): boolean

      Whether this connector supports idempotent upserts (create-or-update keyed by a unique business property). Connectors override this AND Upsert to enable it.

      Returns boolean

    Methods

    • Builds a CRUDResult for a record CREATE, failing LOUDLY when the external system returned no usable record ID. A 2xx response with an empty/undefined ID means the create did not durably produce a record we can track — returning Success:true there silently loses the record and causes duplicate creates on the next sync (the HubSpot-association class of bug, fixed in next commit 9f718a7e). This makes that failure explicit at the connector boundary.

      Parameters

      • externalID: string
      • statusCode: number
      • objectName: string

      Returns CRUDResult

    • Coerce a raw watermark value into a Date for ExternalRecord.ModifiedAt (undefined when unparseable). Delegates full ISO-8601 date-times to the shared parseIso8601AsUtc helper — the same one the EDS drivers use for incremental-watermark literal formatting and Mongo date coercion — so a ZONELESS ISO string is interpreted as UTC consistently everywhere, not via a locally-reimplemented regex. Falls back to Date's own parsing for non-ISO shapes (e.g. a date-only string) that helper intentionally rejects. (In practice the driver returns Dates from RunView; this string path is the fallback.)

      Parameters

      • value: unknown

      Returns Date

    • Deterministic content key for a genuinely PK-less row (rare); lets such tables still dedupe. Excludes watermarkField when given — the watermark column changes on every update by definition, so including it would mint a new ExternalID (and therefore a new inserted record) on every update to the same PK-less row instead of matching its prior identity.

      Parameters

      • row: Record<string, unknown>
      • OptionalwatermarkField: string

      Returns string

    • Best-effort default watermark column, matched by NAME only (engine-neutral — never inferred from a native type, since e.g. Postgres timestamp is a plain datetime, not a rowversion). Only sets the DEFAULT IncrementalWatermarkField; an operator can override it on the persisted IntegrationObject. Returns undefined when no conventional column is present (that object then syncs full, not incremental).

      Parameters

      Returns string

    • Discovery via the connector's READ PATH (FetchChanges), TIME-BOUNDED — the way to gather a statistically-significant sample when a single DiscoverFields sample is too small to PROVE a key. "Discovery is the sync read path with the save removed."

      Loops FetchChanges as a read-only FULL fetch (WatermarkValue=null, nothing persisted), threading pagination/keyset cursors across batches, and streams every record through the data-informed field + provable-PK inference. It stops at the discovery TIME BUDGET (default 5 min), or a record cap, or source exhaustion — whichever comes first — so the provable-PK decision is made on as much real data as the budget allows. It NEVER fabricates a key: if even this larger sample yields no provable single/composite PK, the field set comes back PK-less and the object is honestly not added.

      Falls back to the single-sample DiscoverFields if the read path can't run for this object (e.g. a connector whose FetchChanges needs an already-persisted IO row that doesn't exist yet).

      Parameters

      Returns Promise<ExternalFieldSchema[]>

    • Stage-2 field discovery for sources WITHOUT a describe/introspection endpoint (file feeds, undocumented JSON list endpoints): stream the source's actual records — READ-ONLY, no save, no ack — and derive the full field set + data-informed PK/uniqueness/nullability from the gathered statistics. The connector supplies whatever read-only fetch yields the records; this helper turns that stream into ExternalFieldSchema[].

      Why data-informed: streaming the real values lets pickPrimaryKeyFromStats pick the PK from evidence (uniqueness/non-null statistics) COMBINED with the naming convention, rather than a name guess alone. The PK is a SOFT key, so the pick is best-available, not strict-significance: a confident unique+non-null column wins outright; otherwise a near-unique / convention-named column is taken as a soft key (a PK-less object would stall CodeGen). The scan is time-bounded — it stops on exhaustion OR opts.Discovery.TimeBudgetMs; more rows simply mean stronger claims.

      Provable-only encoding into the standard flags:

      • IsPrimaryKey — set ONLY on the single statistics-first pick. Multiple equally-ranked unique columns leave PK unset here (ambiguous → the pipeline's SoftPKClassifier LLM tiebreaker decides, fed these same stats). Zero unique columns → no PK is fabricated.
      • IsUniqueKey — set when the column was all-distinct over the scan AND uniqueness was provable (the distinct-cap wasn't hit).
      • AllowsNull — asserted true ONLY when a null/absent value was actually observed; otherwise left undefined (permissive default). Never fabricates NOT NULL — critical under a time-capped partial scan where unseen rows could still be null.

      The PK emitted here is SOFT (it rides additionalSchemaInfo via the persist + DDL path; it is NEVER a hard DB key), so a wrong inference can never reject a valid row — the engine dedupes via the record-map. IsReadOnly defaults to true (stream discovery targets read feeds); a writable source overrides via opts.ReadOnly.

      Parameters

      • records:
            | AsyncIterable<Record<string, unknown>, any, any>
            | Iterable<Record<string, unknown>, any, any>

        A read-only sync/async iterable of source records (the caller's fetch yields them).

      • Optionalopts: { Discovery?: StreamDiscoveryOptions; Pk?: PkPickOptions; ReadOnly?: boolean }

      Returns Promise<ExternalFieldSchema[]>

    • The record source DiscoverFieldsViaFetch streams for field/PK inference. Default: loop FetchChanges (full fetch), yielding each record's fields until maxRecords. A protocol subclass (e.g. REST) overrides this to sample a template-var CHILD with the correct record-constrained, recursive stream. Yields plain field maps; the caller stops it at maxRecords.

      Parameters

      Returns AsyncGenerator<Record<string, unknown>>

    • Parse a Retry-After / rate-limit signal out of a failed response or thrown error into milliseconds so the engine can back off precisely. Return undefined when the error is not a throttle (or carries no hint).

      Parameters

      • _error: unknown

      Returns number

    • Returns the ActionGeneratorConfig for this connector, combining the integration name, category, icon, and objects into a ready-to-use configuration for ActionMetadataGenerator.Generate().

      Override in subclasses to customize the config (e.g., icon, category). Returns null by default if GetIntegrationObjects() returns empty.

      Returns ActionGeneratorConfig

    • Returns suggested default field mappings for an external object to MJ entity. Override in subclasses to provide intelligent defaults.

      Parameters

      • _objectName: string

        Name of the external object

      • _entityName: string

        Name of the target MJ entity

      Returns DefaultFieldMapping[]

      Array of default field mappings (empty by default)

    • Returns the integration objects and their fields that this connector supports, for use by the ActionMetadataGenerator. This is static metadata that does NOT require a live connection — it describes the connector's known object model.

      Override in subclasses to provide connector-specific objects/fields. Returns an empty array by default (no action generation available).

      Returns IntegrationObjectInfo[]

    • The batch's max watermark = the last row's watermark value (rows are watermark-ordered ascending).

      Parameters

      • rows: Record<string, unknown>[]
      • watermarkField: string

      Returns string

    • Type-driven post-processing hook (plan.md §10): a connector may normalize/enforce a record's values to the resolved column formats AFTER transform/normalize and BEFORE write. Default returns the record unchanged. (Named for this system — NOT MCP, not take.) The engine ALSO applies target-type constraint enforcement; this is the connector-side complement.

      Parameters

      Returns ExternalRecord

    • Name of a stable, monotonic ordering key (PK/identity) usable for KEYSET/seek resume on watermark-less objects (plan.md §7 — resume from last-seen key, robust to mid-stream insert/delete). null → keyset resume unavailable for this object.

      Parameters

      • _objectName: string

      Returns string

    • Upserts a record — a single idempotent create-or-update keyed by a unique business property (e.g. email), eliminating the search-then-create race window. Override in subclasses whose external system exposes a keyed upsert primitive. Check SupportsUpsert before calling.

      Parameters

      Returns Promise<CRUDResult>