Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion libs/datatug/main/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,9 @@
"@sneat/random": "*",
"@sneat/api": "0.27.6",
"firebase": ">=12.0.0 <13.0.0",
"vitest": "^4.0.18"
"vitest": "^4.0.18",
"ajv": "8.20.0",
"yaml": "2.9.0"
},
"sideEffects": false
}
4 changes: 4 additions & 0 deletions libs/datatug/main/src/lib/models/definition/query-def.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@ import { IParameterDef } from './parameter';
import { IRecordsetDef } from './recordset';
import { HttpMethod } from './command-definition';
import { IWidgetRef } from './widget';
import type { BoundedFederation } from '../../queries/public-data/bounded-federation';
import type { PublicDataScenario } from '../../queries/public-data/public-data-scenario';

export enum QueryType {
HTTP = 'HTTP',
Expand Down Expand Up @@ -39,6 +41,7 @@ export interface IQueryDef extends IQueryItem {
recordsets?: IRecordsetDef[];
federation?: {
readonly ovdbBaseUrl: string;
readonly bounds?: BoundedFederation;
readonly tables: readonly { readonly name: string; readonly database?: string; readonly schema?: string; readonly fields: readonly string[] }[];
readonly lookups?: readonly {
readonly database: string;
Expand All @@ -48,6 +51,7 @@ export interface IQueryDef extends IQueryItem {
readonly concurrency?: number;
}[];
};
publicData?: PublicDataScenario;
widgets?: IWidgetRef[];
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
<ion-select-option value="boards">Boards</ion-select-option>
<ion-select-option value="environments">Environments</ion-select-option>
<ion-select-option value="queries">Queries</ion-select-option>
<ion-select-option value="public-data">Public data lookup</ion-select-option>
@if (enableEmptyShellPages) {
<ion-select-option value="tags">Tags</ion-select-option>
}
Expand Down
31 changes: 30 additions & 1 deletion libs/datatug/main/src/lib/queries/federated-query-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ import { IndexedDbDatabase } from '@dalgo/indexeddb';
import type { RunQueryResponse, TypedValue } from '@sneat/datatug-semantic';
import type { IQueryDef, ITextQueryRequest } from '../models/definition/query-def';
import { deleteQueryDatabase, queryStorageError } from './federated-query-storage';
import { createBoundedFederationFetch } from './public-data/bounded-federation';
import { publicDataExceptions, SAVED_SCENARIO_PUBLICATION_BLOCKER, type PublicDataExceptions } from './public-data/public-data-scenario';

type Data = Record<string, unknown>;
interface OvdbRecord { readonly key: string; readonly data: Data }
Expand Down Expand Up @@ -60,7 +62,7 @@ export interface FederatedQueryObserver {
readonly onSourceLoaded?: (event: FederatedSourceLoaded) => void;
}

export type FederatedQueryResult = RunQueryResponse & { readonly totalRows?: number; readonly hasMore?: boolean };
export type FederatedQueryResult = RunQueryResponse & { readonly totalRows?: number; readonly hasMore?: boolean; readonly publicDataExceptions?: PublicDataExceptions; readonly publicDataBytes?: number };
export type FederatedOutputPage = (rows: readonly (readonly TypedValue[])[]) => Promise<void>;
export type FederatedQueryMode = 'full' | 'visible';

Expand Down Expand Up @@ -93,6 +95,33 @@ function retryDelay(attempt: number, signal?: AbortSignal): Promise<void> {

/** Runs each leaf against OVDB directly and merges/aggregates in this runtime. */
export async function runFederatedQuery(definition: IQueryDef, onProgress?: (progress: FederatedQueryProgress) => void, token = '', onOutputPage?: FederatedOutputPage, signal?: AbortSignal, mode: FederatedQueryMode = 'full', onPageReady?: (result: FederatedQueryResult) => void, waitForNextPage?: () => Promise<void>, observer?: FederatedQueryObserver): Promise<FederatedQueryResult> {
if (definition.publicData) throw new Error(SAVED_SCENARIO_PUBLICATION_BLOCKER);
const bounds = definition.federation?.bounds;
if (!bounds) return runFederatedQueryInternal(definition, onProgress, token, onOutputPage, signal, mode, onPageReady, waitForNextPage, observer);
if (mode !== 'full') throw new Error('Bounded public-data runs calculate one selected page. Browse result pages after it finishes.');
const deadline = new AbortController();
const timer = setTimeout(() => deadline.abort(new Error('The bounded lookup exceeded its deadline.')), bounds.timeoutMs);
const combined = signal ? AbortSignal.any([signal, deadline.signal]) : deadline.signal;
try {
const transport = createBoundedFederationFetch(definition.federation?.ovdbBaseUrl ?? '', bounds, observer?.fetch ?? fetch, combined);
let outputRows = 0;
const collected: TypedValue[][] = [];
const result = await runFederatedQueryInternal(definition, onProgress, token, async (rows) => {
combined.throwIfAborted();
outputRows += rows.length;
if (outputRows > bounds.resultRows) throw new Error('The lookup exceeds the result-row bound; ambiguous or multiplied matches require a narrower query.');
if (onOutputPage) await onOutputPage(rows);
else collected.push(...rows.map((row) => [...row]));
}, combined, 'full', undefined, undefined, { ...observer, fetch: transport.fetch });
combined.throwIfAborted();
if (result.recordset.rows.length > bounds.resultRows) throw new Error('The lookup exceeds the result-row bound.');
// Nested joins use DALgo's batch result rather than the flat join callback.
if (!outputRows && result.recordset.rows.length && onOutputPage) await onOutputPage(result.recordset.rows);
return { ...result, ...(!onOutputPage && outputRows ? { recordset: { ...result.recordset, rows: collected } } : {}), publicDataExceptions: publicDataExceptions(bounds, transport.receipt.sources), publicDataBytes: transport.receipt.bytes(), limitations: [...result.limitations, { rowsFiltered: false, hiddenColumns: [], policy: `Selected user page: at most ${bounds.userRows} rows from offset ${bounds.userOffset}; ${bounds.identifierLimit} identifiers, ${bounds.resultRows} results, ${bounds.bytes} bytes, ${bounds.timeoutMs}ms.` }] };
} finally { clearTimeout(timer); }
}

async function runFederatedQueryInternal(definition: IQueryDef, onProgress?: (progress: FederatedQueryProgress) => void, token = '', onOutputPage?: FederatedOutputPage, signal?: AbortSignal, mode: FederatedQueryMode = 'full', onPageReady?: (result: FederatedQueryResult) => void, waitForNextPage?: () => Promise<void>, observer?: FederatedQueryObserver): Promise<FederatedQueryResult> {
// Late-bound so a stubbed global fetch is honoured; an observer's fetch replaces it for this run only.
const httpFetch: typeof fetch = observer?.fetch ?? ((input, init) => fetch(input, init));
const config = definition.federation;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ describe('federated query worker boundary', () => {
postMessage(message: unknown): void { this.posted.push(message); }
addEventListener(): void { /* not needed: no dispose in these runs */ }
removeEventListener(): void { /* not needed */ }
terminate(): void { /* not needed */ }
terminate = vi.fn();
send(data: unknown): void { this.onmessage?.({ data } as MessageEvent); }
}
const definition = { id: 'query', title: 'Query', request: { queryType: QueryType.DTQL, text: '{}' } } as IQueryDef;
Expand Down Expand Up @@ -63,5 +63,18 @@ describe('federated query worker boundary', () => {
withSource.worker.send({ type: 'result', result });
await withSource.run;
});
it('terminates a bounded worker that never answers and rejects instead of keeping synchronous work alive', async () => {
vi.stubGlobal('Worker', FakeWorker);
const service = new FederatedQueryService();
const bounded = { ...definition, federation: { ovdbBaseUrl: 'https://demodb.dev/ovdb', tables: [], bounds: { userRows: 100, userOffset: 0, identifierKind: 'place' as const, identifierLimit: 100, resultRows: 5000, bytes: 5242880, timeoutMs: 100, sources: [{ database: 'user', name: 'Customer', keyField: 'Country' }, { database: 'geo', name: 'countries', keyField: 'id', parent: { database: 'user', name: 'Customer', field: 'Country' } }] } } };
FakeWorker.last = undefined;
const run = service.run(bounded);
const rejected = expect(run).rejects.toThrow(/exceeded its deadline/);
await vi.waitFor(() => expect(FakeWorker.last?.posted.length).toBe(1));
const worker = FakeWorker.last;
await rejected;
expect(worker?.terminate).toHaveBeenCalledOnce();
await expect(service.dispose()).resolves.toBeUndefined();
});
});
});
18 changes: 15 additions & 3 deletions libs/datatug/main/src/lib/queries/federated-query.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,16 @@ export class FederatedQueryService {
return new Promise<FederatedQueryResult>((resolve, reject) => {
const worker = new Worker(new URL('./federated-query.worker.ts', import.meta.url), { type: 'module' });
this.worker = worker;
// A worker can be busy in synchronous DALgo work: enforce the wall-clock
// bound from the UI thread as well as aborting requests inside the worker.
const deadline = definition.federation?.bounds
? setTimeout(() => {
worker.terminate();
if (this.worker === worker) this.worker = undefined;
reject(new Error('The bounded lookup exceeded its deadline. Temporary storage is cleaned before the next run.'));
}, definition.federation.bounds.timeoutMs)
: undefined;
const clearDeadline = (): void => { if (deadline !== undefined) clearTimeout(deadline); };
worker.onmessage = (event: MessageEvent<
| { type: 'progress'; progress: FederatedQueryProgress }
| { type: 'source'; event: FederatedSourceLoaded }
Expand All @@ -45,24 +55,26 @@ export class FederatedQueryService {
if (message.type === 'source') { extras.onSourceLoaded?.(message.event); return; }
if (message.type === 'finished') { onFinished?.(message.totalRows); return; }
if (message.type === 'cleanup-error') return;
if (message.type === 'cancelled') { reject(new Error('The query was cancelled.')); return; }
if (message.type === 'cancelled') { clearDeadline(); reject(new Error('The query was cancelled.')); return; }
if (message.type === 'page' || message.type === 'page-error') {
const pending = this.pageRequests.get(message.requestId);
this.pageRequests.delete(message.requestId);
if (message.type === 'page') pending?.resolve(message.rows);
else pending?.reject(new Error(message.message));
return;
}
if (message.type === 'closed') { worker.terminate(); if (this.worker === worker) this.worker = undefined; return; }
if (message.type === 'closed') { clearDeadline(); worker.terminate(); if (this.worker === worker) this.worker = undefined; return; }
if (message.type === 'result') {
clearDeadline();
if (mode === 'full' && message.result.totalRows === undefined) void this.dispose().catch(() => undefined);
resolve(message.result);
} else {
clearDeadline();
void this.dispose().catch(() => undefined);
reject(new Error(message.message));
}
};
worker.onerror = (event) => { void this.dispose().catch(() => undefined); reject(new Error(event.message || 'The query worker failed.')); };
worker.onerror = (event) => { clearDeadline(); void this.dispose().catch(() => undefined); reject(new Error(event.message || 'The query worker failed.')); };
worker.postMessage({ type: 'run', definition, token, mode, ...(extras.staticSource ? { staticSource: extras.staticSource } : {}) });
});
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
# Hypothetical real-ROR user input

These eight input rows declare a hypothetical `affiliations.Affiliation.ror_id`
string property. They are not a native field of any DemoDB database. The two
published URL IDs occur in the full ROR release dated 2026-09-22 (ROR release
v2.13, JSON member SHA256
`6d032ff473f771e015da0837fe28b881a8f85822ee7815effefe11dc49c4346c`),
verified against `ingitdb/ror-ingitdb` source output at
`bbbec903248680caea04e68f94b9a957b6efc55b`.

Proposed representation is native ROR URL (`ROR:URL`), UTF8 byte-exact identity,
with no trimming, case folding, URL shortening or organization-name inference.
`a1` and `a3` are two distinct affiliations to one organization. Multiple ROR
locations would remain ordered details of that organization, without multiplying
the affiliation denominator. NULL, empty, malformed case/space and synthetic
name inputs remain separate exceptions, preserving their raw strings.

The exact immutable app commit and the schema path/hash identify this proposed
user mapping for a fresh independent reconciler. These files do not grant
semantic acceptance, canonical registration, native-id runtime support,
production eligibility or live journey acceptance. A committed declaration is
still pending until the dedicated decision and canonical companions land.

Fixture rows and schema are app test material under the app repository licence;
referenced ROR metadata is CC0-1.0. Embedded ROR GeoNames location data, when
displayed, carries CC-BY-4.0 attribution separately.
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
[
{ "affiliation_id": "a1", "ror_id": "https://ror.org/000025p04" },
{ "affiliation_id": "a2", "ror_id": "https://ror.org/00002d369" },
{ "affiliation_id": "a3", "ror_id": "https://ror.org/000025p04" },
{ "affiliation_id": "null", "ror_id": null },
{ "affiliation_id": "empty", "ror_id": "" },
{ "affiliation_id": "space", "ror_id": " https://ror.org/000025p04" },
{ "affiliation_id": "case", "ror_id": "https://ROR.org/000025p04" },
{ "affiliation_id": "synthetic", "ror_id": "Synthetic organization" }
]
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
{
"modelspec": "1.0-draft",
"module": { "name": "affiliations", "id": "github.com/datatug/datatug-apps/affiliations", "version": "0.1.0" },
"entities": {
"Affiliation": {
"properties": {
"affiliation_id": { "type": "string", "required": true },
"ror_id": { "type": "string", "required": false }
},
"key": ["affiliation_id"]
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
# Synthetic consumer compatibility fixture

Copied byte-for-byte from `openvaultdb/ovdb` commit `aec07d17be1e074210acefddb858363ce7777704`, `publisher/representation/testdata`, after independent review of the schema1 helper. These example repositories and labels are hypothetical test inputs. They do not register a provider or admit a semantic decision. Tests mutate copies and retain the original pinned bytes here.

The consumer remains closed pending real Directory/provider companions, independent scoped semantic acceptance and runtime publication. Direct native identifier mode is a separate, unfrozen schema proposal; this fixture only exercises the reviewed label bridge shape.

`directory.pinned.json` and `models.pinned.json` are exact canonical index bytes from Directory commit `02db362144d7924c6081dd6768cc0d7187ac9bc3` and ModelSpec registry commit `48b30b250a61385d94d46a70968677870750ea35`. Their SHA256 values match INITIAL_CANONICAL_PINS. They test connection/provider/entity/property coordinates and never grant semantic eligibility. The configured adapter test uses an explicit synthetic representation namespace; its source schema hash is the actual Chinook model JSON at Directory provider commit `26e852cca00101f53a84ef8ee1f1ae389067f5cf`.
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
{
"table": "CustomerCountries",
"rows": [
{
"raw_label": "USA",
"target_key": "US"
},
{
"raw_label": "United Kingdom",
"target_key": "GB"
}
]
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
{
"format": "ovdb-representation-contract/1",
"contracts": [
{
"source": {
"schema": {
"path": "source.modelspec.json",
"sha256": "7043520817a957c3c3e175fc09877238c320eae236f1ebbb5e80957024551905",
"repository": "https://github.com/example/source",
"revision": "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
},
"module": "sample",
"entity": "Customer",
"property": "Country",
"datatype": "string",
"namespace": "sample-country-labels"
},
"target": {
"snapshot": {
"path": "snapshot.json",
"sha256": "6f4ec14c1ba7adad094ec6d95d889c1387a570eb110d86f15c005af7ad8cf447"
},
"keys": {
"path": "keys.json",
"sha256": "5aa000ef1e2c824aa4dd6315256e5e181ba0a334ed94b43aca252d17434f222e"
},
"model": {
"path": "target.modelspec.json",
"sha256": "aa5b9ea22bf18fb6f81b3ab537431c05a99a54384f1773edc4a890ff17bd410d"
},
"module": "geo",
"entity": "Countries",
"property": "iso",
"datatype": "string",
"namespace": "iso-3166-1-alpha-2",
"binding": {
"document": {
"path": "target.meaning.json",
"sha256": "a447619f70c4de82848091c7eed913b54fd6bbc0264ff65c15568a9a27d1d954"
},
"concept": "native-country",
"role": "identifier",
"meaning": {
"document": {
"path": "core.meaning.json",
"sha256": "af890560ed133e6a4aaa40d22c82027fc329e7e9f9cf8da2886b5415650292a2",
"repository": "https://github.com/example/source",
"revision": "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
},
"concept": "country"
}
}
},
"bridge": {
"artifact": {
"path": "bridge.json",
"sha256": "fd225cf26ca051471e983c770e4af49edd8bdce8d429729e53b4722d677a5bc9"
},
"table": "CustomerCountries",
"raw_label_column": "raw_label",
"target_key_column": "target_key",
"serving_identity_column": "serving_id"
},
"policy": {
"transform": "identity",
"equality": "utf8-byte-exact",
"cardinality": "zero-or-one",
"unmatched": "exception",
"collision": "ineligible"
},
"decision": {
"document": {
"path": "decision.md",
"sha256": "ad838e4f5321fdc424a2ae16a4f3288a1e901e11a2968a697c539df2e2ac0f66",
"repository": "https://github.com/example/source",
"revision": "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
},
"scope": "fixture-only/Customer.Country"
}
}
]
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
{
"format": "meaning/draft-1",
"concepts": [
{
"id": "country"
}
]
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
# Synthetic acceptance fixture

Fixture-only Customer.Country bridge scope; never production admission.
Loading
Loading