Source

services/InsightsService.class.js

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