diff --git a/packages/comunica/.npmignore b/packages/comunica/.npmignore new file mode 100644 index 0000000..db25063 --- /dev/null +++ b/packages/comunica/.npmignore @@ -0,0 +1,4 @@ +*_test.ts +*_bench.ts +_memory_test.ts +*.map diff --git a/packages/comunica/README.md b/packages/comunica/README.md new file mode 100644 index 0000000..4340bcc --- /dev/null +++ b/packages/comunica/README.md @@ -0,0 +1,22 @@ +# `@okikio/comunica` + +Adapter from a caller-owned Comunica `QueryEngine` to `@okikio/sparql`'s result-mode-specific `Queryable` contract. + +```ts +import { QueryEngine } from '@comunica/query-sparql' +import { createClient } from '@okikio/comunica' + +const engine = new QueryEngine() +const client = createClient(engine, { + context: () => ({ sources: [/* caller-owned sources */] }), +}) + +const rows = await client.queryBindings(query) +await client.update(update) +``` + +The package does not create an engine or choose data sources at import time. The supplied engine remains caller-owned. + +Streaming SELECT/graph results destroy the upstream result stream when the consumer returns early or the operation signal aborts. Boolean/update operations are promise-based upstream, so their deeper cancellation behavior depends on the context and capabilities supplied to Comunica. + +The public generic method is `update()`. The adapter translates it to Comunica's upstream `queryVoid()` method internally. diff --git a/packages/comunica/deno.json b/packages/comunica/deno.json new file mode 100644 index 0000000..093b641 --- /dev/null +++ b/packages/comunica/deno.json @@ -0,0 +1,6 @@ +{ + "name": "@okikio/comunica", + "version": "0.1.0", + "license": "MIT", + "exports": { ".": "./mod.ts" } +} diff --git a/packages/comunica/mod.ts b/packages/comunica/mod.ts new file mode 100644 index 0000000..ebd2c17 --- /dev/null +++ b/packages/comunica/mod.ts @@ -0,0 +1,138 @@ +/** Comunica QueryEngine integration for the engine-neutral SPARQL contract. @module */ + +import { fromQuad, fromTerm, type Quad, type Term } from '@okikio/rdf' +import { getQueryText, getUpdateText, type Queryable, type QueryOptionsType } from '@okikio/sparql' +import type { BindingType } from '@okikio/sparql' + +/** Async result stream shape returned by Comunica query methods. */ +export interface ResultStreamType extends AsyncIterable { + /** Node-style streams expose `destroy`; the adapter uses it for early cancellation when available. */ + destroy?(error?: Error): void +} + +/** Minimal Comunica QueryEngine surface used by this adapter. */ +export interface QueryEngineType { + queryBindings(query: string, context?: unknown): Promise> + queryQuads(query: string, context?: unknown): Promise> + queryBoolean(query: string, context?: unknown): Promise + queryVoid(query: string, context?: unknown): Promise +} + +/** Comunica integration configuration. */ +export interface ClientOptionsType { + /** Creates the engine-specific query context for each operation. */ + readonly context?: (options: QueryOptionsType) => unknown +} + +/** Comunica-backed query interface. The supplied QueryEngine remains caller-owned. */ +export interface Client extends Queryable { + readonly engine: QueryEngineType +} + +/** + * Wraps an already-created Comunica QueryEngine without taking engine ownership. + * + * Stream queries destroy their upstream result stream when the consumer returns + * early or the supplied signal aborts. Boolean/void cancellation still depends + * on the query context provided to Comunica because those methods return one + * promise rather than a cancellable stream. + */ +export function createClient(engine: QueryEngineType, options: ClientOptionsType = {}): Client { + return { + engine, + /** Query bindings through the wrapped engine without transferring engine ownership. */ + async queryBindings(query, queryOptions = {}) { + abort(queryOptions.signal) + const stream = await engine.queryBindings(getQueryText(query), options.context?.(queryOptions)) + return mapStream(stream, queryOptions.signal, readBinding) + }, + /** Query quads through the wrapped engine without transferring engine ownership. */ + async queryQuads(query, queryOptions = {}) { + abort(queryOptions.signal) + const stream = await engine.queryQuads(getQueryText(query), options.context?.(queryOptions)) + return mapStream(stream, queryOptions.signal, readQuad) + }, + /** Query boolean through the wrapped engine without transferring engine ownership. */ + async queryBoolean(query, queryOptions = {}) { + abort(queryOptions.signal) + const result = await engine.queryBoolean(getQueryText(query), options.context?.(queryOptions)) + abort(queryOptions.signal) + return result + }, + /** Submits one complete SPARQL Update document through the wrapped engine. */ + async update(update, queryOptions = {}) { + abort(queryOptions.signal) + await engine.queryVoid(getUpdateText(update), options.context?.(queryOptions)) + abort(queryOptions.signal) + }, + } +} + +/** Maps one upstream engine stream and destroys unfinished work when consumption stops early. */ +async function* mapStream( + stream: ResultStreamType, + signal: AbortSignal | undefined, + map: (value: Input) => Output, +): AsyncGenerator { + let complete = false + const onAbort = (): void => stream.destroy?.(abortError(signal)) + signal?.addEventListener('abort', onAbort, { once: true }) + try { + for await (const value of stream) { + abort(signal) + yield map(value) + } + complete = true + } finally { + signal?.removeEventListener('abort', onAbort) + if (!complete) stream.destroy?.() + } +} + +/** Converts one engine-specific binding row into the engine-neutral RDF binding map. */ +function readBinding(value: unknown): BindingType { + if (typeof value !== 'object' || value === null || !('entries' in value) || typeof value.entries !== 'function') { + throw new TypeError('Comunica binding row does not expose entries().') + } + const result = new Map>() + for (const entry of value.entries() as Iterable) { + const [variable, term] = entry + const name = variableName(variable) + if (!isTerm(term)) throw new TypeError(`Comunica binding '${name}' is not an RDF term.`) + result.set(name, fromTerm(term)) + } + return result +} + +/** Converts one engine-specific graph result into a native RDF quad. */ +function readQuad(value: unknown): Quad { + if (!isTerm(value) || value.termType !== 'Quad' || !('subject' in value) || !('predicate' in value) || !('object' in value) || !('graph' in value)) { + throw new TypeError('Comunica graph result contains a non-quad value.') + } + return fromQuad(value as Quad) +} + +/** Normalizes an engine binding key to the SPARQL variable name without its sigil. */ +function variableName(value: unknown): string { + if (typeof value === 'string') return value.replace(/^[?$]/, '') + if (isTerm(value) && value.termType === 'Variable') return value.value + throw new TypeError('Comunica binding key is not a variable or string.') +} + +/** Returns whether the supplied value satisfies the term contract. */ +function isTerm(value: unknown): value is Term { + if (typeof value !== 'object' || value === null) return false + const term = value as Partial + return typeof term.termType === 'string' && typeof term.value === 'string' && typeof term.equals === 'function' +} + +/** Throws the caller supplied abort reason when cancellation has been requested. */ +function abort(signal?: AbortSignal): void { + if (signal?.aborted) throw signal.reason ?? new DOMException('Aborted', 'AbortError') +} + +/** Converts an abort reason into the Error shape required by upstream stream destruction. */ +function abortError(signal: AbortSignal | undefined): Error { + const reason = signal?.reason + return reason instanceof Error ? reason : new DOMException('Aborted', 'AbortError') +} diff --git a/packages/comunica/mod_test.ts b/packages/comunica/mod_test.ts new file mode 100644 index 0000000..9aa7395 --- /dev/null +++ b/packages/comunica/mod_test.ts @@ -0,0 +1,49 @@ +import { describe, it } from 'node:test' +import { expect } from '@std/expect' +import { literal } from '@okikio/rdf' +import { select, triple, update } from '@okikio/sparql' +import { createClient, type ResultStreamType } from './mod.ts' + +describe('@okikio/comunica', () => { + it('forwards caller-owned query context and structured documents', async () => { + let seen: unknown + let queryText = '' + let updateText = '' + const client = createClient({ + queryBoolean: async (query: string, context?: unknown) => { queryText = query; seen = context; return true }, + queryVoid: async (query: string) => { updateText = query }, + queryBindings: async () => ({ [Symbol.asyncIterator]: async function* () {} }), + queryQuads: async () => ({ [Symbol.asyncIterator]: async function* () {} }), + }, { context: () => ({ source: 'memory' }) }) + + const query = select('*').where(triple('?s', '?p', '?o')) + expect(await client.queryBoolean(query)).toBe(true) + await client.update(update().deleteWhere(triple('?s', '?p', '?o'))) + expect(seen).toEqual({ source: 'memory' }) + expect(queryText).toBe(query.build().value) + expect(updateText).toBe('DELETE WHERE { ?s ?p ?o . }') + }) + + it('destroys caller-owned result work when the consumer returns early', async () => { + let destroyed = false + const stream: ResultStreamType>> = { + async *[Symbol.asyncIterator]() { + yield new Map([['name', literal('Alice')]]) + yield new Map([['name', literal('Bob')]]) + }, + destroy() { destroyed = true }, + } + const client = createClient({ + queryBindings: async () => stream, + queryQuads: async () => ({ [Symbol.asyncIterator]: async function* () {} }), + queryBoolean: async () => true, + queryVoid: async () => undefined, + }) + + for await (const row of await client.queryBindings('SELECT ?name WHERE {}')) { + expect(row.get('name')?.value).toBe('Alice') + break + } + expect(destroyed).toBe(true) + }) +}) diff --git a/packages/comunica/package.json b/packages/comunica/package.json new file mode 100644 index 0000000..957a924 --- /dev/null +++ b/packages/comunica/package.json @@ -0,0 +1,23 @@ +{ + "name": "@okikio/comunica", + "version": "0.1.0", + "type": "module", + "sideEffects": false, + "exports": { + ".": "./mod.ts" + }, + "description": "Comunica adapter for the @okikio/sparql query contract.", + "license": "MIT", + "repository": { + "type": "git", + "url": "git+https://github.com/okikio/sparql-client.git", + "directory": "packages/comunica" + }, + "publishConfig": { + "access": "public" + }, + "dependencies": { + "@okikio/rdf": "0.1.0", + "@okikio/sparql": "0.1.0" + } +} diff --git a/packages/oxigraph/.npmignore b/packages/oxigraph/.npmignore new file mode 100644 index 0000000..db25063 --- /dev/null +++ b/packages/oxigraph/.npmignore @@ -0,0 +1,4 @@ +*_test.ts +*_bench.ts +_memory_test.ts +*.map diff --git a/packages/oxigraph/README.md b/packages/oxigraph/README.md new file mode 100644 index 0000000..c5508dd --- /dev/null +++ b/packages/oxigraph/README.md @@ -0,0 +1,18 @@ +# `@okikio/oxigraph` + +Adapter from a caller-owned Oxigraph `Store` to `@okikio/sparql`'s result-mode-specific `Queryable` contract. + +```ts +import { Store } from 'oxigraph' +import { createClient } from '@okikio/oxigraph' + +const store = new Store() +const client = createClient(store) + +const rows = await client.queryBindings(query) +await client.update(update) +``` + +The package does not initialize Wasm or create a store at import time. The supplied store remains caller-owned. + +The current Oxigraph JavaScript Store API is synchronous. `AbortSignal` is checked before an operation begins, but a positive `timeoutMs` is rejected because the adapter cannot truthfully interrupt an already-running synchronous Wasm call. Use an owned Worker/process when interruptibility is required. diff --git a/packages/oxigraph/deno.json b/packages/oxigraph/deno.json new file mode 100644 index 0000000..bebc0af --- /dev/null +++ b/packages/oxigraph/deno.json @@ -0,0 +1,6 @@ +{ + "name": "@okikio/oxigraph", + "version": "0.1.0", + "license": "MIT", + "exports": { ".": "./mod.ts" } +} diff --git a/packages/oxigraph/mod.ts b/packages/oxigraph/mod.ts new file mode 100644 index 0000000..794f569 --- /dev/null +++ b/packages/oxigraph/mod.ts @@ -0,0 +1,110 @@ +/** Oxigraph Store integration for the engine-neutral SPARQL query contract. @module */ + +import { fromQuad, fromTerm, type Quad, type Term } from '@okikio/rdf' +import { getQueryText, getUpdateText, type Queryable, type QueryOptionsType } from '@okikio/sparql' +import type { BindingType } from '@okikio/sparql' + +/** Minimal Oxigraph Store surface used by this adapter. */ +export interface StoreType { + query(query: string, options?: Readonly>): unknown + update(update: string, options?: Readonly>): void +} + +/** Oxigraph adapter options. */ +export interface ClientOptionsType { + /** Static query options forwarded to `Store.query()`. */ + readonly query?: Readonly> + /** Static update options forwarded to `Store.update()`. */ + readonly update?: Readonly> +} + +/** Oxigraph-backed query interface. The supplied store remains caller-owned. */ +export interface Client extends Queryable { + readonly store: StoreType +} + +/** + * Wraps an already-created Oxigraph Store without initializing Wasm or taking ownership. + * + * Oxigraph's current JavaScript Store API is synchronous. Abort signals are + * therefore checked before execution, and timeout requests are rejected rather + * than pretending a synchronous Wasm call can be interrupted. + */ +export function createClient(store: StoreType, options: ClientOptionsType = {}): Client { + return { + store, + /** Query bindings through the wrapped engine without transferring engine ownership. */ + async queryBindings(query, queryOptions = {}) { + prepare(queryOptions) + const result = store.query(getQueryText(query), options.query) + if (!isIterable(result)) throw new TypeError('Oxigraph SELECT query did not return an iterable of bindings.') + return bindings(result) + }, + /** Query quads through the wrapped engine without transferring engine ownership. */ + async queryQuads(query, queryOptions = {}) { + prepare(queryOptions) + const result = store.query(getQueryText(query), options.query) + if (!isIterable(result)) throw new TypeError('Oxigraph graph query did not return an iterable of quads.') + return quads(result) + }, + /** Query boolean through the wrapped engine without transferring engine ownership. */ + async queryBoolean(query, queryOptions = {}) { + prepare(queryOptions) + const result = store.query(getQueryText(query), options.query) + if (typeof result !== 'boolean') throw new TypeError('Oxigraph ASK query did not return a boolean.') + return result + }, + /** Submits one complete SPARQL Update document through the wrapped engine. */ + async update(update, queryOptions = {}) { + prepare(queryOptions) + store.update(getUpdateText(update), options.update) + }, + } +} + +/** Validates operation controls that the synchronous engine can actually honor. */ +function prepare(options: QueryOptionsType): void { + if (options.signal?.aborted) throw options.signal.reason ?? new DOMException('Aborted', 'AbortError') + if (options.timeoutMs !== undefined && options.timeoutMs !== null && options.timeoutMs > 0) { + throw new TypeError('Oxigraph synchronous Store queries cannot honor timeoutMs. Run the store in an owned Worker when interruptibility is required.') + } +} + +/** Adapts synchronous engine binding rows to the asynchronous query contract. */ +async function* bindings(values: Iterable): AsyncGenerator { + for (const value of values) { + if (!(value instanceof Map)) throw new TypeError('Oxigraph SELECT row is not a Map.') + const binding = new Map>() + for (const [name, term] of value) { + if (typeof name !== 'string') throw new TypeError('Oxigraph binding variable name is not a string.') + if (!isTerm(term)) throw new TypeError(`Oxigraph binding '${name}' is not an RDF term.`) + binding.set(name.replace(/^[?$]/, ''), fromTerm(term)) + } + yield binding + } +} + +/** Adapts synchronous engine graph rows to the asynchronous query contract. */ +async function* quads(values: Iterable): AsyncGenerator { + for (const value of values) { + if (!isQuad(value)) throw new TypeError('Oxigraph graph result contains a non-quad value.') + yield fromQuad(value) + } +} + +/** Returns whether the supplied value satisfies the iterable contract. */ +function isIterable(value: unknown): value is Iterable { + return typeof value === 'object' && value !== null && Symbol.iterator in value +} + +/** Returns whether the supplied value satisfies the term contract. */ +function isTerm(value: unknown): value is Term { + if (typeof value !== 'object' || value === null) return false + const record = value as Partial + return typeof record.termType === 'string' && typeof record.value === 'string' && typeof record.equals === 'function' +} + +/** Returns whether the supplied value satisfies the quad contract. */ +function isQuad(value: unknown): value is Quad { + return isTerm(value) && value.termType === 'Quad' && 'subject' in value && 'predicate' in value && 'object' in value && 'graph' in value +} diff --git a/packages/oxigraph/mod_test.ts b/packages/oxigraph/mod_test.ts new file mode 100644 index 0000000..5edc3b7 --- /dev/null +++ b/packages/oxigraph/mod_test.ts @@ -0,0 +1,31 @@ +import { describe, it } from 'node:test' +import { expect } from '@std/expect' +import { select, triple, update } from '@okikio/sparql' +import { createClient } from './mod.ts' + +describe('@okikio/oxigraph', () => { + it('does not pretend a synchronous Store supports timeout cancellation', async () => { + const client = createClient({ query: () => true, update: () => undefined }) + let kind = '' + try { + await client.queryBoolean('ASK {}', { timeoutMs: 1 }) + } catch (error) { + kind = error instanceof Error ? error.name : '' + } + expect(kind).toBe('TypeError') + }) + + it('accepts structured query and update documents without taking store ownership', async () => { + let queryText = '' + let updateText = '' + const client = createClient({ + query(query: string) { queryText = query; return true }, + update(value: string) { updateText = value }, + }) + const query = select('*').where(triple('?s', '?p', '?o')) + expect(await client.queryBoolean(query)).toBe(true) + await client.update(update().deleteWhere(triple('?s', '?p', '?o'))) + expect(queryText).toBe(query.build().value) + expect(updateText).toBe('DELETE WHERE { ?s ?p ?o . }') + }) +}) diff --git a/packages/oxigraph/package.json b/packages/oxigraph/package.json new file mode 100644 index 0000000..c8a9726 --- /dev/null +++ b/packages/oxigraph/package.json @@ -0,0 +1,23 @@ +{ + "name": "@okikio/oxigraph", + "version": "0.1.0", + "type": "module", + "sideEffects": false, + "exports": { + ".": "./mod.ts" + }, + "description": "Oxigraph adapter for the @okikio/sparql query contract.", + "license": "MIT", + "repository": { + "type": "git", + "url": "git+https://github.com/okikio/sparql-client.git", + "directory": "packages/oxigraph" + }, + "publishConfig": { + "access": "public" + }, + "dependencies": { + "@okikio/rdf": "0.1.0", + "@okikio/sparql": "0.1.0" + } +}