feat(connectors): add MongoDB connector (#305) (#310)

* refactor(connectors): split KtxDialect into core and KtxSqlDialect

Separate the dialect contract into a driver-agnostic core (display/ref
formatting and type mapping) and a SQL-only extension (query generators).
The catalog and entity-details paths resolve the core dialect for any
snapshot driver, so it must stay free of SQL generation; this is the
prerequisite refactor for adding non-SQL primary sources.

- KtxDialect keeps type, formatDisplayRef, parseDisplayRef,
  columnDisplayTablePartCount, mapDataType, mapToDimensionType
- KtxSqlDialect extends it with quoteIdentifier, formatTableName, and the
  query/sample/statistics generators; the 7 SQL dialects implement it
- add getSqlDialectForDriver for SQL drivers; the 7 connectors and the
  relationship-benchmark harness consume it
- thread the relationship pipeline (profiling/validation/composite/
  discovery) as KtxSqlDialect | null so a non-SQL source skips coverage SQL
  and its candidates stay in review; local-enrichment builds the SQL
  dialect only when the connector advertises readOnlySql

Pure extraction: no behavior change for the existing 7 drivers.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* feat(connectors): add MongoDB connector for issue #305

Add a read-only MongoDB connector that treats a database as a primary
context source: collections map to tables and inferred top-level fields to
columns. MongoDB is the first non-SQL source (readOnlySql: false), so
ktx sql and metric compilation do not apply, but its collections flow
through ingest, descriptions, and relationship discovery.

- schema-inference: infer a flat column schema from the most recent
  sample_size documents (by _id desc, or order_by for non-ObjectId keys).
  Union BSON types per field, mark multi-type fields mixed (string), keep
  sub-documents/arrays as a single opaque json column, derive nullability
  from presence, treat _id as the primary key
- connector: KtxMongoDbScanConnector behind an injectable client seam;
  strictly read-only (find/listCollections/estimatedDocumentCount only),
  no executeReadOnly; resolves env:/file: via resolveKtxConfigReference
- core-only KtxMongoDbDialect and a live-database introspection adapter
- wire the mongodb driver: driver union, dialect registry, driver
  registration (scopeConfigKey databases), mongodbConnectionSchema,
  connection-drivers, normalizeDriver, the live-database route, and the
  ktx setup picker. ktx sql is refused by the read-only SQL capability gate
- tests: schema inference, connector snapshot via a fake client, dialect,
  driver-schema parsing, and the ktx sql rejection

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* docs(integrations): document the MongoDB primary source

Add a MongoDB section to the primary-sources reference: connection config
(url, databases, enabled_tables, sample_size, order_by), mongodb+srv/TLS/
Atlas notes, the schema-inference explainer, a features matrix, and the
non-SQL caveat. Update the frontmatter and connection field reference.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix(connectors): address review blockers on the MongoDB connector

- introspect: skip estimatedDocumentCount for views. The count command is
  rejected on a MongoDB view (CommandNotSupportedOnView), so counting a view
  aborted introspect for the whole connection; compute estimatedRows only for
  real collections, as ClickHouse does.
- sl: refuse a semantic-layer query against a non-SQL connection instead of
  defaulting it to the Postgres dialect. compileLocalSlQuery (the shared CLI +
  MCP path) now rejects a driver with no SQL dialect via the new
  isSqlQueryableDriver authority, keeping MongoDB context-only per issue #305.
- tests: cover input.tableScope and the empty-scope skip for the Mongo
  connector (the scan layer does not post-filter), the view no-count path, and
  the ktx sl query refusal for a mongodb connection.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* polish(mongodb): compute sampled nullCount and document sampling caveats

Address the non-blocking review notes:

- sampleColumn now counts null/absent values over the sampled window instead of
  returning nullCount: null, since the documents are already in hand
- warn that a custom order_by must be indexed (an unindexed sort hits MongoDB's
  in-memory sort limit on large collections) in the connection schema and docs
- note that sampled values for nested fields are stringified, not faithfully
  serialized, so the json opacity is deliberate

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* docs(examples): add a MongoDB connector example

A manual, container-backed example mirroring examples/postgres-historic:

- docker-compose.yml + init/seed.js seed a representative dataset (nested
  documents, arrays, a Decimal128, a mixed-type field, a nullable field, an
  ObjectId reference, and a view) on first container start
- scripts/smoke.sh + introspect-smoke.mjs assert the connector's inferred
  schema with no LLM credentials — the same introspection entry point ktx
  ingest's database-schema stage uses, including the view-no-count path
- README.md documents the smoke and a full keyless ktx ingest run
  (claude-code LLM + managed sentence-transformers embeddings)

Works with Docker Compose or podman compose. Verified end to end.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* chore: ignore examples/** in knip to fix dead-code false positives

The MongoDB connector example files (examples/mongodb/init/seed.js and
examples/mongodb/scripts/introspect-smoke.mjs) are used at runtime but were
flagged as unused by knip. Add examples/** to the ignore array, matching the
existing .context/** entry.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_0114qQV8fJ5a5ME3XbMVRzbL

* fix(mongodb): refuse non-SQL connections before SQL analysis

`ktx sql` and the MCP sql_execution tool resolved a SQL-analysis dialect
(falling back to Postgres for a non-SQL driver) and ran read-only
validation before the connector capability gate refused the connection.
For a MongoDB connection that spun up the parser/daemon and produced
Postgres parser diagnostics instead of a clean non-SQL refusal.

Route both entry points through a shared assertSqlQueryableConnection
guard before dialect selection, mirroring compileLocalSlQuery. The
federated duckdb path has no driver and is exempted at each call site.
Add CLI and MCP regression tests asserting validation/connector work
never starts for a MongoDB connection.

* fix(mongodb): pass CI gates (dialect boundary, secrets, setup test)

Three latent failures in the connector surfaced once CI ran on the branch:

- connector.ts imported the concrete KtxMongoDbDialect, which the connector
  dialect-import boundary forbids. Route it through getDialectForDriver('mongodb')
  and widen inferKtxMongoCollectionColumns to the base KtxDialect (it only uses
  mapDataType/mapToDimensionType).
- detect-secrets flagged a test ObjectId hex and the mongodb+srv example URL;
  annotate both with allowlist pragmas.
- the "shows every supported database" setup test omitted the new MongoDB option.

---------

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Co-authored-by: Luca Martial <48870843+luca-martial@users.noreply.github.com>
Co-authored-by: Luca Martial <lucamrtl@gmail.com>
Co-authored-by: Andrey Avtomonov <andreybavt@gmail.com>
This commit is contained in:
Pintouch 2026-06-29 15:17:56 +02:00 committed by GitHub
parent 4f084186f1
commit 2afab61417
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
59 changed files with 1971 additions and 129 deletions

View file

@ -1,6 +1,6 @@
import pLimit from 'p-limit';
import type { KtxLlmRuntimePort } from '../../context/llm/runtime-port.js';
import { getDialectForDriver } from '../connections/dialects.js';
import { getSqlDialectForDriver } from '../connections/dialects.js';
import { buildDefaultKtxProjectConfig, type KtxScanRelationshipConfig } from '../project/config.js';
import { KtxDescriptionGenerator } from './description-generation.js';
import { buildKtxColumnEmbeddingText } from './embedding-text.js';
@ -486,7 +486,9 @@ export async function runLocalScanEnrichment(
snapshot,
connectionId: input.connectionId,
});
const dialect = getDialectForDriver(snapshot.driver);
const dialect = input.connector.capabilities.readOnlySql
? getSqlDialectForDriver(snapshot.driver)
: null;
const now = input.now ?? (() => new Date());
const state = completedKtxScanEnrichmentStateSummary();
const syncId = input.syncId ?? input.context.runId;

View file

@ -131,12 +131,13 @@ function normalizeDriver(driver: string | undefined): KtxConnectionDriver {
normalized === 'clickhouse' ||
normalized === 'sqlserver' ||
normalized === 'bigquery' ||
normalized === 'snowflake'
normalized === 'snowflake' ||
normalized === 'mongodb'
) {
return normalized;
}
throw new Error(
`Standalone ktx scan supports postgres/sqlite/mysql/clickhouse/sqlserver/bigquery/snowflake in this phase, received "${driver ?? 'unknown'}"`,
`Standalone ktx scan supports postgres/sqlite/mysql/clickhouse/sqlserver/bigquery/snowflake/mongodb in this phase, received "${driver ?? 'unknown'}"`,
);
}

View file

@ -6,7 +6,7 @@ import { gunzipSync } from 'node:zlib';
import Database from 'better-sqlite3';
import YAML from 'yaml';
import { z } from 'zod';
import { getDialectForDriver } from '../connections/dialects.js';
import { getSqlDialectForDriver } from '../connections/dialects.js';
import type { KtxLlmRuntimePort } from '../llm/runtime-port.js';
import type { KtxEnrichedRelationship, KtxEnrichedSchema, KtxRelationshipType } from './enrichment-types.js';
import { snapshotToKtxEnrichedSchema } from './local-enrichment.js';
@ -537,7 +537,7 @@ export function ktxRelationshipBenchmarkDetectorWithLlm(
const formalLinks = formalMetadata.accepted.map((relationship) => relationshipToBenchmarkLink(relationship));
const acceptedKeys = new Set(formalLinks.map(fkKey));
const sqliteDataAvailable = Boolean(input.dataPath && input.snapshot.driver === 'sqlite');
const dialect = getDialectForDriver(input.snapshot.driver);
const dialect = getSqlDialectForDriver(input.snapshot.driver);
const profilingExecutor =
sqliteDataAvailable && input.mode !== 'profiling_disabled'
? new KtxRelationshipBenchmarkSqliteExecutor(input.dataPath as string)
@ -552,6 +552,7 @@ export function ktxRelationshipBenchmarkDetectorWithLlm(
})
: await profileKtxRelationshipSchema({
connectionId: input.snapshot.connectionId,
driver: input.snapshot.driver,
dialect,
schema: input.schema,
executor: profilingExecutor,
@ -673,7 +674,7 @@ export function currentKtxRelationshipBenchmarkDetector(): KtxRelationshipBenchm
const formalLinks = formalMetadata.accepted.map((relationship) => relationshipToBenchmarkLink(relationship));
const acceptedKeys = new Set(formalLinks.map(fkKey));
const sqliteDataAvailable = Boolean(input.dataPath && input.snapshot.driver === 'sqlite');
const dialect = getDialectForDriver(input.snapshot.driver);
const dialect = getSqlDialectForDriver(input.snapshot.driver);
const profilingExecutor =
sqliteDataAvailable && input.mode !== 'profiling_disabled'
? new KtxRelationshipBenchmarkSqliteExecutor(input.dataPath as string)
@ -688,6 +689,7 @@ export function currentKtxRelationshipBenchmarkDetector(): KtxRelationshipBenchm
})
: await profileKtxRelationshipSchema({
connectionId: input.snapshot.connectionId,
driver: input.snapshot.driver,
dialect,
schema: input.schema,
executor: profilingExecutor,

View file

@ -1,4 +1,4 @@
import type { KtxDialect } from '../connections/dialects.js';
import type { KtxSqlDialect } from '../connections/dialects.js';
import type { KtxEnrichedColumn, KtxEnrichedSchema, KtxEnrichedTable, KtxRelationshipType } from './enrichment-types.js';
import {
type KtxRelationshipProfileArtifact,
@ -56,7 +56,7 @@ export interface KtxCompositeRelationshipCandidate {
export interface DiscoverKtxCompositeRelationshipsInput {
connectionId: string;
dialect: KtxDialect;
dialect: KtxSqlDialect;
schema: KtxEnrichedSchema;
profiles: KtxRelationshipProfileArtifact;
executor: KtxRelationshipReadOnlyExecutor | null;
@ -227,11 +227,11 @@ function sqlSuffix(fragment: string): string {
return fragment ? ` ${fragment}` : '';
}
function aliasedTupleSelect(dialect: KtxDialect, columns: readonly string[]): string {
function aliasedTupleSelect(dialect: KtxSqlDialect, columns: readonly string[]): string {
return columns.map((column, index) => `${dialect.quoteIdentifier(column)} AS c${index}`).join(', ');
}
function nonNullPredicate(dialect: KtxDialect, columns: readonly string[]): string {
function nonNullPredicate(dialect: KtxSqlDialect, columns: readonly string[]): string {
return columns.map((column) => `${dialect.quoteIdentifier(column)} IS NOT NULL`).join(' AND ');
}
@ -242,7 +242,7 @@ function tupleEquality(columns: number): string {
}
function buildTupleDistinctSql(input: {
dialect: KtxDialect;
dialect: KtxSqlDialect;
table: KtxTableRef;
columns: readonly string[];
}): string {
@ -257,7 +257,7 @@ function buildTupleDistinctSql(input: {
}
function buildCompositeCoverageSql(input: {
dialect: KtxDialect;
dialect: KtxSqlDialect;
childTable: KtxTableRef;
childColumns: readonly string[];
parentTable: KtxTableRef;
@ -322,7 +322,7 @@ function hasAcceptedSubset(
async function detectCompositePrimaryKeys(input: {
connectionId: string;
dialect: KtxDialect;
dialect: KtxSqlDialect;
table: KtxEnrichedTable;
profiles: KtxRelationshipProfileArtifact;
executor: KtxRelationshipReadOnlyExecutor;
@ -426,7 +426,7 @@ function compatibleTuple(sourceColumns: readonly KtxEnrichedColumn[], targetColu
async function validateCompositeRelationship(input: {
connectionId: string;
dialect: KtxDialect;
dialect: KtxSqlDialect;
sourceTable: KtxEnrichedTable;
sourceColumns: readonly KtxEnrichedColumn[];
targetKey: KtxCompositePrimaryKeyCandidate;

View file

@ -1,5 +1,5 @@
import type { KtxLlmRuntimePort } from '../../context/llm/runtime-port.js';
import type { KtxDialect } from '../connections/dialects.js';
import type { KtxSqlDialect } from '../connections/dialects.js';
import type { KtxScanRelationshipConfig } from '../project/config.js';
import type { KtxEnrichedRelationship, KtxEnrichedSchema, KtxRelationshipUpdate } from './enrichment-types.js';
import {
@ -34,7 +34,7 @@ import type {
export interface DiscoverKtxRelationshipsInput {
connectionId: string;
dialect: KtxDialect;
dialect: KtxSqlDialect | null;
connector: KtxScanConnector;
schema: KtxEnrichedSchema;
context: KtxScanContext;
@ -122,20 +122,21 @@ function compositeSummary(relationships: readonly KtxCompositeRelationshipCandid
async function detectCompositeRelationships(input: {
connectionId: string;
dialect: KtxDialect;
dialect: KtxSqlDialect | null;
schema: KtxEnrichedSchema;
profile: KtxRelationshipProfileArtifact;
executor: KtxRelationshipReadOnlyExecutor | null;
context: DiscoverKtxRelationshipsInput['context'];
warnings: KtxScanWarning[];
}): Promise<KtxCompositeRelationshipCandidate[]> {
if (!input.executor || !input.profile.sqlAvailable) {
if (!input.executor || !input.profile.sqlAvailable || !input.dialect) {
return [];
}
const dialect = input.dialect;
try {
const compositeDetection = await discoverKtxCompositeRelationships({
connectionId: input.connectionId,
dialect: input.dialect,
dialect,
schema: input.schema,
profiles: input.profile,
executor: input.executor,
@ -223,6 +224,7 @@ export async function discoverKtxRelationships(
const profileCache = createKtxRelationshipProfileCache();
const profile = await profileKtxRelationshipSchema({
connectionId: input.connectionId,
driver: input.connector.driver,
dialect: input.dialect,
schema: input.schema,
executor,

View file

@ -1,4 +1,4 @@
import type { KtxDialect } from '../connections/dialects.js';
import type { KtxSqlDialect } from '../connections/dialects.js';
import type { KtxEnrichedColumn, KtxEnrichedSchema, KtxEnrichedTable } from './enrichment-types.js';
import { mapWithConcurrency } from './relationship-validation.js';
import type {
@ -56,7 +56,8 @@ export interface KtxRelationshipProfileCache {
export interface ProfileKtxRelationshipSchemaInput {
connectionId: string;
dialect: KtxDialect;
driver: KtxConnectionDriver;
dialect: KtxSqlDialect | null;
schema: KtxEnrichedSchema;
executor: KtxRelationshipReadOnlyExecutor | null;
ctx: KtxScanContext;
@ -123,7 +124,7 @@ function columnKey(table: KtxEnrichedTable, column: KtxEnrichedColumn): string {
function tableProfileCacheKey(input: {
connectionId: string;
dialect: KtxDialect;
dialect: KtxSqlDialect;
ctx: KtxScanContext;
table: KtxTableRef;
sampleValuesPerColumn: number;
@ -149,7 +150,7 @@ function sqlSuffix(fragment: string): string {
return fragment ? ` ${fragment}` : '';
}
function sampledTableSql(dialect: KtxDialect, tableSql: string, limit: number): string {
function sampledTableSql(dialect: KtxSqlDialect, tableSql: string, limit: number): string {
const top = dialect.getTopClause(limit);
if (top) {
return `(SELECT ${top} * FROM ${tableSql}) AS relationship_profile_sample`;
@ -158,7 +159,7 @@ function sampledTableSql(dialect: KtxDialect, tableSql: string, limit: number):
}
function sampleValuesSql(input: {
dialect: KtxDialect;
dialect: KtxSqlDialect;
tableSql: string;
columnSql: string;
limit: number;
@ -175,7 +176,7 @@ function sampleValuesSql(input: {
}
function columnProfileSelectSql(input: {
dialect: KtxDialect;
dialect: KtxSqlDialect;
tableSql: string;
profileTableSql: string;
column: KtxEnrichedColumn;
@ -218,7 +219,7 @@ function splitSampleValues(value: unknown): string[] {
async function queryCount(input: {
connectionId: string;
dialect: KtxDialect;
dialect: KtxSqlDialect;
table: KtxTableRef;
executor: KtxRelationshipReadOnlyExecutor;
ctx: KtxScanContext;
@ -233,7 +234,7 @@ async function queryCount(input: {
async function queryTableProfile(input: {
connectionId: string;
dialect: KtxDialect;
dialect: KtxSqlDialect;
table: KtxEnrichedTable;
executor: KtxRelationshipReadOnlyExecutor;
ctx: KtxScanContext;
@ -320,10 +321,10 @@ type TableProfileResult =
export async function profileKtxRelationshipSchema(
input: ProfileKtxRelationshipSchemaInput,
): Promise<KtxRelationshipProfileArtifact> {
if (!input.executor) {
if (!input.executor || !input.dialect) {
return {
connectionId: input.connectionId,
driver: input.dialect.type,
driver: input.driver,
sqlAvailable: false,
queryCount: 0,
tables: [],
@ -337,6 +338,7 @@ export async function profileKtxRelationshipSchema(
const columns: Record<string, KtxRelationshipColumnProfile> = {};
const warnings: string[] = [];
const executor = input.executor;
const dialect = input.dialect;
const enabledTables = input.schema.tables.filter((candidate) => candidate.enabled);
const tableResults = await mapWithConcurrency<KtxEnrichedTable, TableProfileResult>(
@ -347,7 +349,7 @@ export async function profileKtxRelationshipSchema(
const profileSampleRows = input.profileSampleRows ?? 10000;
const cacheKey = tableProfileCacheKey({
connectionId: input.connectionId,
dialect: input.dialect,
dialect,
ctx: input.ctx,
table: table.ref,
sampleValuesPerColumn,
@ -361,7 +363,7 @@ export async function profileKtxRelationshipSchema(
try {
const tableProfile = await queryTableProfile({
connectionId: input.connectionId,
dialect: input.dialect,
dialect,
table,
executor,
ctx: input.ctx,
@ -403,7 +405,7 @@ export async function profileKtxRelationshipSchema(
return {
connectionId: input.connectionId,
driver: input.dialect.type,
driver: input.driver,
sqlAvailable: true,
queryCount: queryTotal,
tables,

View file

@ -1,4 +1,4 @@
import type { KtxDialect } from '../connections/dialects.js';
import type { KtxSqlDialect } from '../connections/dialects.js';
import type { KtxRelationshipEndpoint } from './enrichment-types.js';
import { applyKtxRelationshipValidationBudget, type KtxRelationshipValidationBudget } from './relationship-budget.js';
import type { KtxRelationshipDiscoveryCandidate } from './relationship-candidates.js';
@ -44,7 +44,7 @@ export interface KtxValidatedRelationshipDiscoveryCandidate
export interface ValidateKtxRelationshipDiscoveryCandidatesInput {
connectionId: string;
dialect: KtxDialect;
dialect: KtxSqlDialect | null;
candidates: readonly KtxRelationshipDiscoveryCandidate[];
profiles: KtxRelationshipProfileArtifact;
executor: KtxRelationshipReadOnlyExecutor | null;
@ -108,7 +108,7 @@ function sqlSuffix(fragment: string): string {
}
function buildCoverageSql(input: {
dialect: KtxDialect;
dialect: KtxSqlDialect;
childTable: KtxTableRef;
childColumn: string;
parentTable: KtxTableRef;
@ -237,13 +237,14 @@ export async function validateKtxRelationshipDiscoveryCandidates(
input: ValidateKtxRelationshipDiscoveryCandidatesInput,
): Promise<KtxValidatedRelationshipDiscoveryCandidate[]> {
const settings = mergeSettings(input.settings);
if (!input.executor || !input.profiles.sqlAvailable) {
if (!input.executor || !input.profiles.sqlAvailable || !input.dialect) {
return input.candidates.map((candidate) =>
reviewWithoutValidation(candidate, input.profiles, 'validation_unavailable'),
);
}
const executor = input.executor;
const dialect = input.dialect;
async function validateCandidate(
candidate: KtxRelationshipDiscoveryCandidate,
@ -260,7 +261,7 @@ export async function validateKtxRelationshipDiscoveryCandidates(
{
connectionId: input.connectionId,
sql: buildCoverageSql({
dialect: input.dialect,
dialect,
childTable: candidate.from.table,
childColumn: sourceColumn,
parentTable: candidate.to.table,

View file

@ -7,7 +7,8 @@ export type KtxConnectionDriver =
| 'bigquery'
| 'snowflake'
| 'mysql'
| 'clickhouse';
| 'clickhouse'
| 'mongodb';
export type KtxScanMode = 'structural' | 'relationships' | 'enriched';