From 890f8e7cfd6911fcbd8c7afffd26d02cd9fdaddf Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 23 Aug 2026 14:04:34 +0000 Subject: [PATCH] fix(service-datasource): one live introspection per datasource per validation sweep MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit validateAll/validateDatasource used to call config.introspect(datasource) once per federated OBJECT, all concurrently — M objects on one datasource dialled the same remote M times per sweep. The sweep now threads a per-call memo (sweepScopedIntrospect) through the validation body: the in-flight introspection promise is memoised by datasource for the lifetime of ONE validateEach call and discarded when the call returns, so a long-lived service never serves a stale schema to a later sweep. Direct validateObject calls keep reading live on every call. Pins: counting fakes assert one read per datasource per sweep (both sweep spellings), a second sweep reads live again and sees a remote change (per-call, not per-instance), direct calls stay live, and one unreachable remote costs one connection attempt while every object on it still gets its failure row. The IExternalDatasourceService.validateAll docstring stops promising "parallelised per datasource" for an implementation that parallelised per object. Fixes #10962 Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01APWX2AwT3a4xDcjPCe8bk4 --- .../validate-sweep-introspection-memo.md | 6 + .../external-datasource-service.test.ts | 148 ++++++++++++++++++ .../src/external-datasource-service.ts | 55 ++++++- .../contracts/external-datasource-service.ts | 2 +- 4 files changed, 208 insertions(+), 3 deletions(-) create mode 100644 .changeset/validate-sweep-introspection-memo.md diff --git a/.changeset/validate-sweep-introspection-memo.md b/.changeset/validate-sweep-introspection-memo.md new file mode 100644 index 0000000000..17c99bb613 --- /dev/null +++ b/.changeset/validate-sweep-introspection-memo.md @@ -0,0 +1,6 @@ +--- +'@objectstack/service-datasource': patch +'@objectstack/spec': patch +--- + +`validateAll`/`validateDatasource` now read each datasource's live schema once per sweep instead of once per federated object: the sweep threads a per-call introspection memo through the validation body, so M objects on one datasource cost one remote introspection round-trip (a rejected read is shared the same way — one connection attempt, M failure rows). The memo lives and dies inside a single call, so a long-lived service never serves a stale schema to a later sweep, and direct `validateObject` calls still read live every time. The `IExternalDatasourceService.validateAll` docstring, which promised "parallelised per datasource" while the implementation parallelised per object, now states the actual behaviour. diff --git a/packages/services/service-datasource/src/__tests__/external-datasource-service.test.ts b/packages/services/service-datasource/src/__tests__/external-datasource-service.test.ts index 3a17d834b2..c587934e3a 100644 --- a/packages/services/service-datasource/src/__tests__/external-datasource-service.test.ts +++ b/packages/services/service-datasource/src/__tests__/external-datasource-service.test.ts @@ -598,6 +598,154 @@ describe('validateDatasource', () => { }); }); +/** + * [#10962] Per-sweep introspection memo — the CALL COUNT is the deliverable. + * + * `validateObject` reads the live remote schema on every call, so a sweep over + * M federated objects on one datasource used to perform M concurrent + * `introspect(datasource)` round-trips against the same remote. The fix + * threads a per-call memo through `validateAll`/`validateDatasource`: one live + * read per datasource per sweep. + * + * The memo's LIFETIME is the other half of the contract, and the harder one: + * per call, never per instance. A `Map` cached on the service would pass every + * counting assertion a per-call memo passes — only the second-sweep cases + * below (a fresh read per sweep, and a remote change visible to the next + * sweep) tell them apart. Do not weaken those to "at least once". + */ +describe('per-sweep introspection memo [#10962]', () => { + const M_DATASOURCE = 'wh'; + const SIDE_DATASOURCE = 'wh_b'; + + /** M > 1 objects on one datasource — the population the memo collapses. */ + const OBJECTS: ObjectLike[] = [ + ...['wh_orders_a', 'wh_orders_b', 'wh_orders_c'].map((name) => ({ + name, + datasource: M_DATASOURCE, + external: { remoteName: 'orders' }, + fields: { order_id: { type: 'text' } }, + })), + { + name: 'side_orders', + datasource: SIDE_DATASOURCE, + external: { remoteName: 'orders' }, + fields: { order_id: { type: 'text' } }, + }, + ]; + + function makeCounting(opts: { unreachable?: readonly string[] } = {}) { + const introspected: string[] = []; + const unreachable = new Set(opts.unreachable ?? []); + let tables: IntrospectedSchema['tables'] = { + orders: { + name: 'orders', + indexes: [], + columns: [{ name: 'order_id', type: 'text', nullable: false, primaryKey: true }], + }, + }; + const svc = new ExternalDatasourceService({ + introspect: async (datasource: string) => { + introspected.push(datasource); + if (unreachable.has(datasource)) throw new Error(`connect ECONNREFUSED (${datasource})`); + return { dialect: 'postgres', introspectedAt: '2026-08-23T00:00:00.000Z', tables }; + }, + getDatasource: async (name: string) => + [M_DATASOURCE, SIDE_DATASOURCE].includes(name) ? { name, schemaMode: 'external' } : undefined, + getObject: async (name: string) => OBJECTS.find((o) => o.name === name), + listObjects: async () => OBJECTS, + logger: { warn: () => {} }, + }); + return { + svc, + introspected, + /** Simulate the remote dropping its tables between sweeps. */ + dropRemoteTables: () => { + tables = {}; + }, + }; + } + + it('validateAll reads each datasource once per sweep — M objects, one live read', async () => { + // The claim is only non-vacuous when M really exceeds 1. + expect(OBJECTS.filter((o) => o.datasource === M_DATASOURCE).length).toBeGreaterThan(1); + const { svc, introspected } = makeCounting(); + + const report = await svc.validateAll(); + + expect([...introspected].sort()).toEqual([M_DATASOURCE, SIDE_DATASOURCE]); + expect(report.ok).toBe(true); + expect(report.results).toHaveLength(OBJECTS.length); + }); + + it('validateDatasource reads its datasource once for M objects', async () => { + const { svc, introspected } = makeCounting(); + + const report = await svc.validateDatasource(M_DATASOURCE); + + expect(introspected).toEqual([M_DATASOURCE]); + expect(report.ok).toBe(true); + expect(report.results).toHaveLength(3); + }); + + it('a second sweep reads live again — the memo is per call, not a service-instance cache', async () => { + const { svc, introspected, dropRemoteTables } = makeCounting(); + + const first = await svc.validateAll(); + expect(first.ok).toBe(true); + expect(introspected.filter((d) => d === M_DATASOURCE)).toHaveLength(1); + + // The remote changes between sweeps. A per-instance cache would keep the + // counting pin above green while answering this sweep from last sweep's + // schema — stale `ok: true` — which is exactly what this case refuses. + dropRemoteTables(); + const second = await svc.validateAll(); + + expect(introspected.filter((d) => d === M_DATASOURCE)).toHaveLength(2); + expect(second.ok).toBe(false); + for (const r of second.results) { + expect(r.ok).toBe(false); + expect(r.diffs[0]).toMatchObject({ kind: 'missing_table', severity: 'error' }); + } + }); + + it('direct validateObject stays live: two calls are two reads, and a remote change is seen', async () => { + const { svc, introspected, dropRemoteTables } = makeCounting(); + + const before = await svc.validateObject('wh_orders_a'); + expect(before.ok).toBe(true); + + dropRemoteTables(); + const after = await svc.validateObject('wh_orders_a'); + + expect(introspected).toEqual([M_DATASOURCE, M_DATASOURCE]); + expect(after.ok).toBe(false); + expect(after.diffs[0]).toMatchObject({ kind: 'missing_table', severity: 'error' }); + }); + + it('one unreachable remote costs ONE connection attempt and still yields a failure row per object', async () => { + const { svc, introspected } = makeCounting({ unreachable: [M_DATASOURCE] }); + + const report = await svc.validateDatasource(M_DATASOURCE); + + expect(introspected).toEqual([M_DATASOURCE]); + expect(report.ok).toBe(false); + expect(report.results).toHaveLength(3); + for (const r of report.results) { + expect(r).toMatchObject({ + ok: false, + datasource: M_DATASOURCE, + diffs: [ + expect.objectContaining({ + kind: 'missing_table', + severity: 'error', + actual: `connect ECONNREFUSED (${M_DATASOURCE})`, + }), + ], + }); + } + }); +}); + describe('refreshCatalog', () => { it('produces a snapshot with suggested field types', async () => { const svc = makeService(); diff --git a/packages/services/service-datasource/src/external-datasource-service.ts b/packages/services/service-datasource/src/external-datasource-service.ts index a88f4e6402..a62ba9d855 100644 --- a/packages/services/service-datasource/src/external-datasource-service.ts +++ b/packages/services/service-datasource/src/external-datasource-service.ts @@ -593,6 +593,25 @@ export class ExternalDatasourceService implements IExternalDatasourceService { } async validateObject(objectName: string): Promise { + // A direct call performs its own live read — no memo. Read reuse is the + // sweep's per-call concern ({@link validateEach}), never this method's: a + // long-lived service must answer every direct call from the remote's + // schema as it is NOW (pinned: two direct calls are two live reads). + return this.validateObjectUsing(objectName, (ds) => this.config.introspect(ds)); + } + + /** + * [#10962] The body of {@link validateObject}, with the live-schema read + * abstracted behind `readSchema` so one sweep can share a single read per + * datasource across all of its objects. `readSchema` is either + * `config.introspect` itself (the public single-object path above) or the + * per-sweep memoised reader from {@link sweepScopedIntrospect} — never a + * cache that outlives one call. + */ + private async validateObjectUsing( + objectName: string, + readSchema: (datasource: string) => Promise, + ): Promise { const obj = await this.config.getObject(objectName); if (!obj) { throw new Error(`Object '${objectName}' not found.`); @@ -605,7 +624,7 @@ export class ExternalDatasourceService implements IExternalDatasourceService { return { ok: true, datasource, object: objectName, diffs: [] }; } - const schema = await this.config.introspect(datasource); + const schema = await readSchema(datasource); const dialect = schema.dialect as SqlDialect | undefined; const remoteName = obj.external?.remoteName ?? obj.name; const table = this.findTable(schema, remoteName); @@ -684,6 +703,34 @@ export class ExternalDatasourceService implements IExternalDatasourceService { return o.external !== undefined || Boolean(o.datasource && o.datasource !== 'default'); } + /** + * [#10962] One live schema read per datasource per SWEEP. + * + * Returns a reader that memoises `config.introspect` by datasource name for + * the lifetime of ONE {@link validateEach} call. The memo is a local of that + * call — deliberately NOT an instance field — so a long-lived service can + * never serve a stale schema to a later sweep: the next `validateAll()` / + * `validateDatasource()` builds a fresh memo and reads live again (both + * directions pinned in `__tests__/external-datasource-service.test.ts`). + * + * The PROMISE is memoised, not the resolved value: the sweep validates its + * objects concurrently (`Promise.all`), so the first reader for a datasource + * starts the read and every concurrent sibling awaits the same in-flight + * promise. A rejected read is shared the same way — M objects on one + * unreachable datasource produce M failure rows from ONE connection attempt. + */ + private sweepScopedIntrospect(): (datasource: string) => Promise { + const memo = new Map>(); + return (datasource) => { + let read = memo.get(datasource); + if (!read) { + read = this.config.introspect(datasource); + memo.set(datasource, read); + } + return read; + }; + } + /** * Validate a chosen set of objects, one report. * @@ -691,11 +738,15 @@ export class ExternalDatasourceService implements IExternalDatasourceService { * message rather than rejecting the whole report: one unreachable remote (or * one object whose definition vanished mid-sweep) must not erase the verdicts * of the objects that did validate. + * + * [#10962] All objects in one call share one live schema read per datasource + * (see {@link sweepScopedIntrospect}); the memo dies with this call. */ private async validateEach(objects: ObjectLike[]): Promise { + const readSchema = this.sweepScopedIntrospect(); const results = await Promise.all( objects.map((o) => - this.validateObject(o.name).catch((err): SchemaValidationResult => { + this.validateObjectUsing(o.name, readSchema).catch((err): SchemaValidationResult => { this.logger?.warn(`validateObject('${o.name}') failed`, err); return { ok: false, diff --git a/packages/spec/src/contracts/external-datasource-service.ts b/packages/spec/src/contracts/external-datasource-service.ts index ee2f5d50fd..f91c6c597c 100644 --- a/packages/spec/src/contracts/external-datasource-service.ts +++ b/packages/spec/src/contracts/external-datasource-service.ts @@ -154,6 +154,6 @@ export interface IExternalDatasourceService { /** Validate one federated object against the live remote table. */ validateObject(objectName: string): Promise; - /** Validate every federated object, parallelised per datasource. */ + /** Validate every federated object in parallel; each datasource's live schema is read once per call. */ validateAll(): Promise; }