import ScenidCloudService from './ScenidCloudService.class.js'
import { decodeJwtPayload } from '../jwt.js'
import truncateValue from '../truncate.js'
const _sourceCache = new Map() // serviceUrl → { byName, byId, fetchedAt, promise }
const CACHE_TTL = 5 * 60 * 1000
const CACHE_FRESH_GUARD = 10 * 1000
const isCacheStale = entry => !entry || (Date.now() - entry.fetchedAt) >= CACHE_TTL
const isCacheFresh = entry => entry && (Date.now() - entry.fetchedAt) < CACHE_FRESH_GUARD
/**
* Client for the Scenid Insights Service.
*
* Provides access to data sources and their documents via a factory pattern:
* call `Source(nameOrId)` to get a scoped handle for a specific source, then
* use the returned object's methods to query or mutate documents.
*
* Source names (e.g. `'Orders'`) are resolved to their sourceId automatically
* using a per-service-URL cache backed by `GET /api/v2/source`. The cache is
* populated on first use and refreshed after 5 minutes or on a cache miss.
*
* Obtained via `sdk.asAdmin().Insights()` or `sdk.asUser(idToken).Insights()`.
*
* @extends ScenidCloudService
*
* @example
* const insights = sdk.asAdmin().Insights()
*
* // Read source metadata by name
* const { result } = await insights.Source('Orders').get()
*
* // Create a document
* await insights.Source('Orders').Doc.create({ temperature: 21.5, unit: 'C' })
*/
class InsightsService extends ScenidCloudService {
/**
* @param {string|function(): Promise.<string>} tokenSource
* @param {string} serviceUrl
* @param {Object} [scope]
* @param {boolean} [scope.scoped] - When `true` (set by `sdk.asScript()`), this client is
* hard-locked to whatever `sourceId`/`permissions` the current token's own claims say —
* never a separately-supplied value, so scoping can't drift from what the token actually
* grants. `Source(nameOrId)` ignores name resolution entirely and throws if called with any
* id other than the token's `sourceId`; `Sources(...)` (cross-source read) is disabled; write
* methods (`push`, `Doc.create/update/delete/batchInsert`) throw unless the token's
* `permissions` claim includes `'write'` — unless the token also carries a `dryRun` claim
* (minted for test runs), in which case those calls are mocked instead of throwing: no
* request is sent, and the returned `result` is `{ dryRun: true, wouldSend: <payload> }`
* (payload capped by `../truncate.js`) rather than a real server response.
* @param {boolean} [scope.throwOnError] - Forwarded to `ScenidCloudService` — see its
* constructor docs. `sdk.asScript()` sets this, so a rejected call (e.g. a permission
* or auth mismatch the client-side checks above didn't already catch) throws instead
* of silently resolving as `{ status, result }` with nothing written.
*/
constructor (tokenSource, serviceUrl, { scoped, throwOnError } = {}) {
super(tokenSource, serviceUrl, { throwOnError })
this._scoped = Boolean(scoped)
/**
* Returns a scoped handle for the given source, identified by name or sourceId.
*
* The name is resolved to a sourceId via a local cache on first call. If the
* name is not found in the cache, the cache is refreshed once before falling
* back to using the value as-is.
*
* @param {string} nameOrId - The source's human-readable name (e.g. `'Orders'`) or raw sourceId.
* @returns {SourceHandle}
*
* @example
* const src = insights.Source('Orders')
* const { result } = await src.get(true) // include formulator schema
*/
this.Source = nameOrId => ({
/**
* Fetches metadata for the source, optionally including its formulator validation schema.
*
* @param {boolean} [schema=false] - When `true`, the response includes a `schema`
* object (`{ validation, render, translations }`) that can be passed directly to
* `FormulatorForm`.
* @returns {Promise<ServiceResponse>}
*
* @example
* const { result } = await insights.Source('Orders').get(true)
* // result.schema.validation, result.schema.render, result.schema.translations
*/
get: async (schema = false) => {
await this._guardRead()
const sourceId = await this._resolveSource(nameOrId)
const query = schema ? '?schema=true' : ''
return this.get(`/api/v2/source/${sourceId}${query}`)
},
/**
* Runs an aggregation query against the source's data.
*
* Supports filtering, sorting, field selection, pagination via cursor (`pageKey`),
* and Elasticsearch-style aggregations. Pass `raw: true` to query the raw index
* (unprocessed field values as ingested).
*
* `filter` and `sort` are arrays of real Elasticsearch clause objects, not a flat
* `{field: value}` shorthand — they're passed straight through to the query body.
* `sort` entries use ES's shorthand form (`{field: 'asc'|'desc'}`); `filter` entries
* are full query-DSL clauses (`term`, `range`, `bool`, ...), ANDed together. A
* `text`/`string`-formatted field needs its `.keyword` sub-field for an exact-match
* `term` filter (e.g. `{ term: { 'status.keyword': 'active' } }`), not the bare field.
*
* @param {Object} params - Aggregation parameters.
* @param {number} [params.size] - Number of rows to return.
* @param {string[]} [params.fields] - Fields to include in each row.
* @param {Object[]} [params.sort] - Sort clauses, e.g. `[{ createdAt: 'desc' }]`.
* @param {Object[]} [params.filter] - ES filter clauses, e.g. `[{ term: { 'status.keyword': 'active' } }]`.
* @param {Object} [params.aggs] - Elasticsearch aggregation definitions.
* @param {string} [params.pageKey] - Cursor for the next page (from a previous response's `nextPageKey`).
* @param {boolean} [params.raw] - When `true`, queries the raw index.
* @returns {Promise<ServiceResponse>}
* @throws {Error} `sdk/aggregate-filter-must-be-array` – `filter` was given but isn't an array.
* @throws {Error} `sdk/aggregate-sort-must-be-array` – `sort` was given but isn't an array.
*
* @example
* const { result } = await insights.Source('Orders').aggregate({
* size: 100,
* filter: [{ term: { 'status.keyword': 'active' } }],
* sort: [{ createdAt: 'desc' }]
* })
*/
aggregate: async params => {
await this._guardRead()
if (params?.filter !== undefined && !Array.isArray(params.filter)) throw new Error('sdk/aggregate-filter-must-be-array')
if (params?.sort !== undefined && !Array.isArray(params.sort)) throw new Error('sdk/aggregate-sort-must-be-array')
const sourceId = await this._resolveSource(nameOrId)
return this.post(`/api/v2/source/${sourceId}/aggregate`, params)
},
push: async docs => {
if (await this._guardWrite()) return this._mockWrite(docs)
const sourceId = await this._resolveSource(nameOrId)
return this.post(`/api/v2/source/${sourceId}/push`, docs)
},
/**
* Document-level operations scoped to this source.
* @namespace
*/
Doc: {
/**
* Fetches a single document by its ID from the main (transformed) index.
*
* @param {string} docId - The document's identifier.
* @returns {Promise<ServiceResponse>}
*/
get: async docId => {
await this._guardRead()
const sourceId = await this._resolveSource(nameOrId)
return this.get(`/api/v2/source/${sourceId}/doc/${docId}`)
},
/**
* Creates a new document in the source.
*
* The document is run through the source's field type transformations
* and stored in both the raw index (original values) and the main index
* (transformed values). The document ID is derived from a hash of the
* raw field values, making creates idempotent for identical payloads.
*
* @param {Object} body - Flat document object (no nested keys, no `_id` or `id` fields).
* @param {Object} [options]
* @param {boolean} [options.return] - When `true`, returns `{ id, doc }` instead of `{ id, result: 'created' }`.
* @returns {Promise<ServiceResponse>}
*
* @example
* const { result } = await insights.Source('Orders').Doc.create(
* { temperature: 21.5, location: 'Berlin' },
* { return: true }
* )
* // result.id, result.doc
*/
create: async (body, options = {}) => {
if (await this._guardWrite()) return this._mockWrite({ data: body, ...options })
const sourceId = await this._resolveSource(nameOrId)
return this.post(`/api/v2/source/${sourceId}/doc`, { data: body, ...options })
},
/**
* Updates an existing document by its ID.
*
* By default this is a full replacement: the incoming `body` replaces all
* stored fields. Pass `update: true` for a merge update — the service fetches
* the current raw document, shallow-merges your `body` on top, then re-indexes.
*
* @param {string} docId - The document's identifier.
* @param {Object} body - New field values. Must be flat (no nesting, no reserved fields).
* @param {Object} [options]
* @param {boolean} [options.update] - When `true`, merges `body` into the existing document rather than replacing it.
* @param {boolean} [options.return] - When `true`, returns the resulting document after indexing.
* @returns {Promise<ServiceResponse>}
*
* @example
* // Full replacement
* await insights.Source('Orders').Doc.update('doc-abc', { temperature: 22.0, location: 'Berlin' })
*
* // Partial update — only changes temperature, keeps other fields
* await insights.Source('Orders').Doc.update('doc-abc', { temperature: 22.0 }, { update: true, return: true })
*/
update: async (docId, body, options = {}) => {
if (await this._guardWrite()) return this._mockWrite({ docId, data: body, ...options })
const sourceId = await this._resolveSource(nameOrId)
return this.put(`/api/v2/source/${sourceId}/doc/${docId}`, { data: body, ...options })
},
/**
* Deletes a document from both the raw and main indices.
*
* @param {string} docId - The document's identifier.
* @param {Object} [options]
* @param {boolean} [options.return] - When `true`, returns `{ id, doc }` with the deleted document instead of 204.
* @returns {Promise<ServiceResponse>}
*
* @example
* const { result } = await insights.Source('Orders').Doc.delete('doc-abc', { return: true })
* // result.id, result.doc
*/
delete: async (docId, options = {}) => {
if (await this._guardWrite()) return this._mockWrite({ docId, ...options })
const sourceId = await this._resolveSource(nameOrId)
const query = options.return ? '?return=true' : ''
return this.delete(`/api/v2/source/${sourceId}/doc/${docId}${query}`)
},
/**
* Inserts multiple documents in a single bulk request.
*
* Each document is independently transformed and hashed. The response includes
* counts of created and updated documents and a per-document error list for any
* that failed ES indexing. Pass `return: true` to include the computed document
* IDs in the response.
*
* The batch endpoint also accepts `Content-Type: application/x-ndjson` for
* streaming ingestion from tools like Vector — in that case send raw JSONL lines
* directly without the `{ data }` envelope.
*
* @param {Object[]} docs - Array of flat document objects.
* @param {Object} [options]
* @param {boolean} [options.return] - When `true`, includes `ids` (array of doc IDs) in the response.
* @returns {Promise<ServiceResponse>}
* @throws {Error} `scenid/empty-batch` – `docs` is empty or not an array.
*
* @example
* const { result } = await insights.Source('Orders').Doc.batchInsert([
* { temperature: 21.5, location: 'Berlin' },
* { temperature: 19.0, location: 'Hamburg' }
* ], { return: true })
* // result.created, result.updated, result.failed, result.ids
*/
batchInsert: async (docs, options = {}) => {
if (!Array.isArray(docs) || !docs.length) throw new Error('scenid/empty-batch')
if (await this._guardWrite()) return this._mockWrite({ data: docs, ...options })
const sourceId = await this._resolveSource(nameOrId)
return this.post(`/api/v2/source/${sourceId}/batch`, { data: docs, ...options })
}
}
})
/**
* Returns a handle for fetching documents across multiple sources at once.
*
* Source names (e.g. `['Orders', 'Inventory']`) are resolved to sourceIds via
* the local cache. Results are grouped by source name in the response.
*
* @param {string[]} namesOrIds - List of source names or sourceIds.
* @returns {Object}
*
* @example
* const { result } = await insights.Sources(['Orders', 'Inventory'])
* .Docs.find({ humanId: 'HT-4d44d' })
* // result.total, result.Orders, result.Inventory
*/
this.Sources = namesOrIds => ({
Docs: {
/**
* Fetches documents from all listed sources matching the given filter.
*
* Each key-value pair in `filterObj` is treated as an equality condition.
* Results are grouped by source name.
*
* @param {Object} [filterObj={}] - `{ field: value }` pairs, all eq-matched.
* @param {Object} [options]
* @param {number} [options.size=10] - Max docs to return (server caps at 1000).
* @returns {Promise<ServiceResponse>} `{ total, [sourceName]: docs[] }`
*/
find: async (filterObj = {}, { size = 10 } = {}) => {
if (this._scoped) throw new Error('sdk/script-source-scope-violation')
const resolved = await Promise.all(namesOrIds.map(n => this._resolveSource(n)))
const filter = Object.entries(filterObj).map(([field, value]) => ({ field, operator: 'eq', value }))
const params = new URLSearchParams({ size: String(size) })
if (filter.length) params.set('filter', JSON.stringify(filter))
return this.get(`/api/v2/sources/${resolved.join(',')}/docs?${params}`)
}
}
})
}
async _refreshSourceCache () {
const existing = _sourceCache.get(this.serviceUrl)
if (existing?.promise) {
await existing.promise
return
}
let resolve
const promise = new Promise(r => { resolve = r })
_sourceCache.set(this.serviceUrl, { ...(existing ?? {}), promise })
try {
const { result } = await this.get('/api/v2/source')
const sources = Array.isArray(result) ? result : []
const byName = new Map()
const byId = new Map()
sources.forEach(s => {
byName.set(s.name, s.sourceId)
byId.set(s.sourceId, s.sourceId)
})
_sourceCache.set(this.serviceUrl, { byName, byId, fetchedAt: Date.now(), promise: undefined })
resolve()
} catch (e) {
_sourceCache.delete(this.serviceUrl)
resolve()
throw e
}
}
// Re-decodes on every call rather than caching at construction time — the
// token can rotate under a long-lived daemon script (a fresh line pushed
// down the harness's stdin channel on a timer), so scope must always
// reflect whatever token is currently in effect, not whatever was current
// when this client was first built.
async _getScope () {
const token = await this._getToken()
const payload = decodeJwtPayload(token)
if (payload.token_type !== 'insights_functions') throw new Error('sdk/not-a-script-token')
return { sourceId: payload.sourceId, permissions: payload.permissions || [], dryRun: Boolean(payload.dryRun) }
}
// Returns true when a write call should be mocked instead of performed —
// only for a dry-run test token (see the constructor's `scope` docs) that
// lacks write permission. A regular read-only-configured token still
// throws: that's a real misconfiguration (a scheduled function trying to
// write without write permission), not something to silently paper over.
async _guardWrite () {
if (!this._scoped) return false
const { permissions, dryRun } = await this._getScope()
if (permissions.includes('write')) return false
if (dryRun) return true
throw new Error('sdk/script-read-only')
}
// Read-side mirror of _guardWrite() — a write-only-scoped token (a
// producer that should never be able to pull its own data back out)
// throws before ever reaching the network. No dry-run/mock case here:
// unlike a blocked write, there's nothing meaningful to fake for a read.
async _guardRead () {
if (!this._scoped) return
const { permissions } = await this._getScope()
if (permissions.includes('read')) return
throw new Error('sdk/script-write-only')
}
// Mirrors the real write methods' ServiceResponse shape ({ status, result })
// so a script destructuring `const { result } = await ...` behaves the same
// whether mocked or real. The mock is clearly tagged (`result.dryRun`)
// rather than trying to fake the real endpoint's exact response shape,
// which would risk a script depending on fields that don't mean the same
// thing for a mock as they would for a real write.
_mockWrite (payload) {
return { status: 200, result: { dryRun: true, wouldSend: truncateValue(payload) } }
}
async _resolveSource (key) {
if (this._scoped) {
const { sourceId } = await this._getScope()
if (key !== undefined && key !== sourceId) throw new Error('sdk/script-source-scope-violation')
return sourceId
}
const cache = _sourceCache.get(this.serviceUrl)
if (!isCacheStale(cache)) {
if (cache.byName.has(key)) return cache.byName.get(key)
if (cache.byId.has(key)) return key
if (!isCacheFresh(cache)) {
await this._refreshSourceCache()
const fresh = _sourceCache.get(this.serviceUrl)
if (fresh?.byName.has(key)) return fresh.byName.get(key)
}
throw new Error('insights/source-not-found')
}
await this._refreshSourceCache()
const fresh = _sourceCache.get(this.serviceUrl)
if (fresh?.byName.has(key)) return fresh.byName.get(key)
throw new Error('insights/source-not-found')
}
}
export default InsightsService
Source