/** * Log10x Retriever REST client + S3 results poller. * * The Retriever's query API is two-phase: * * 1. POST /streamer/query with a body shaped like QueryRequest.java accepts. * The body carries a client-generated `id` field which the engine uses * as the canonical queryId and echoes back in the response. * * 2. The engine dispatches stream workers that buffer matched events to * disk and upload them to S3 under {basePath}/tenx/{target}/qr/{queryId}/ * as JSONL files. The client polls that prefix via the AWS CLI until * every expected worker has written its marker (at the sibling * {basePath}/tenx/{target}/q/{queryId}/ backstop prefix) and the results * prefix is stable. * * Configure with: * - __SAVE_LOG10X_RETRIEVER_URL__: base URL of the query handler (e.g., the NLB). * - __SAVE_LOG10X_RETRIEVER_BUCKET__: S3 bucket holding the retriever index/results. * - __SAVE_LOG10X_RETRIEVER_TARGET__: default target app prefix (e.g., "app" for * the otek demo env). Overridable per query. * - LOG10X_RETRIEVER_POLL_MS: poll interval, default 1500 ms. * - LOG10X_RETRIEVER_TIMEOUT_MS: total poll budget, default 180_000 ms. * The 90s default was too small for archives where the query-handler * spawns dozens of stream workers — observed live on the otel-demo * env taking >60s to write the first marker. 180s covers single * forensic queries; per-sub-window calls in retriever_series get the * same budget but typically finish in seconds because each query is * scoped to a small sub-window. * * Authentication piggybacks on the same X-10X-Auth header the Prometheus * gateway uses (apiKey/envId). Override via LOG10X_RETRIEVER_AUTH_HEADER / * LOG10X_RETRIEVER_AUTH_VALUE if the deployment uses a different scheme. */ import type { EnvConfig } from './environments.js'; import { type RetrieverQueryDiagnostics } from './retriever-diagnostics.js'; import { type QueryDiagnosis } from './query-funnel.js'; export declare class RetrieverNotConfiguredError extends Error { constructor(); } export interface RetrieverQueryRequest { /** Absolute ISO8601, epoch millis, or relative `now("-1h")`. */ from: string; /** Absolute ISO8601, epoch millis, or relative `now()`. */ to: string; /** Optional Bloom-filter search expression (TenX subset: ==, ||, &&, includes). */ search?: string; /** Optional in-memory JS filters applied after the Bloom-scoped fetch. */ filters?: string[]; /** Target app prefix to scope the index scan. Defaults to __SAVE_LOG10X_RETRIEVER_TARGET__. */ target?: string; /** * Tier-1 result-sink redirect: bare-token prefix under which query OUTPUT * (qr/ events, qrs/ summaries, q/ markers, _DONE.json) is written, in the * SAME index bucket. Maps to the engine `queryResultTarget` option. When * set, the MCP polls `tenx/{resultTarget}/qr/{queryId}/` for results; when * omitted, results stay under `target` (legacy). Must be `[A-Za-z0-9_-]+`. */ resultTarget?: string; /** Logical query name (appears in PERF metrics). */ name?: string; /** Hard cap on total events returned after merging per-worker results. */ limit?: number; /** Max milliseconds the engine has to produce results. */ processingTimeMs?: number; /** Max bytes the engine will ship before terminating. */ resultSizeBytes?: number; /** * Internal: set on the auto-localization probe (a narrowed re-submit of a * DISPATCHED_BLIND query that forces local dispatch). Prevents the probe from * recursively probing itself. Not part of the public tool surface. */ __probe?: boolean; /** @deprecated use `search` instead */ pattern?: string; /** Per-query CloudWatch log-level escalation. Maps to the engine's `logLevels` * REST body field; set to escalate one query to DEBUG for 0-result triage. */ logLevels?: string; /** @deprecated format is now a client-side rollup */ format?: 'events' | 'count' | 'aggregated'; /** @deprecated bucket_size is now a client-side rollup */ bucketSize?: string; /** * Gate the events writer (qr/ prefix). Defaults to true so legacy * call sites keep getting raw event JSONL. Set false alongside * `writeSummaries: true` for count/trend tools that don't need * per-event payload — order-of-magnitude bandwidth saving on * high-volume patterns. */ writeResults?: boolean; /** * Gate the per-slice TenXSummary writer (qrs/ prefix). Defaults to * false. Each TenXSummary record carries `summaryVolume` (event * count for that group) + `summaryBytes` (utf8 byte sum) + the * grouping field values; the slice's time bounds are encoded in * the S3 key path (`qrs/{queryId}/{sliceFrom}_{sliceTo}/`) and the * client parses them back into `sliceFromMs`/`sliceToMs` on each * record. */ writeSummaries?: boolean; /** * Caps the per-worker qr/ event DOWNLOAD (not the engine write — the engine * always writes the full set; results_location carries it). Only takes * effect when qrs/ summaries are present (they supply the whole-match count * + rollups, so the client doesn't need the bulk): * - 0 → download NO qr/ events (count/aggregate served by summaries) * - N (> 0) → download worker files until ~N events accumulate, then stop * - undefined → download everything (legacy) * With NO summaries the full set is always downloaded regardless, so the * true count + rollups are never lost. */ maxDownloadEvents?: number; } /** * Result of a `fetchExistingResults` call. */ export interface ExistingResultsResponse { /** Events recovered from the qr/{queryId}/ prefix, sorted by timestamp. */ events: RetrieverEvent[]; /** True when at least one worker uploaded a `.truncated` sibling. */ truncated: boolean; /** Number of `.jsonl` worker files read. */ jsonlObjectCount: number; /** Worker files that failed to download after retries (events incomplete). */ failedWorkerFiles?: number; /** True when the `_DONE.json` marker exists. False when the query is mid-flight. */ done: boolean; /** Resolved target prefix the result paths used. */ target: string; /** Funnel verdict (OK / EMPTY_RANGE / BLOOM_REJECTED_ALL / MATCHED_NO_EVENTS / * DISPATCHED_BLIND / NO_MARKER / INCONCLUSIVE) derived from the _DONE marker * + the events that came back. Lets the agent debug a zero-result fetch. */ diagnosis: QueryDiagnosis; } /** * Fetch results for a previously-submitted retriever query by `queryId`. * * Closes the stranded-queryId dead-end. When `runRetrieverQuery` returns * `partialResults: true` (its MCP-side poll budget exceeded the engine's * actual completion time), the engine still finishes the scan and uploads * results to S3 under `qr/{queryId}/`. Without this helper, those results * are unreachable: re-running `runRetrieverQuery` submits a new queryId. * `log10x_retriever_query_status` calls this helper when invoked with * `fetch_results: true` so the agent can recover the stranded events. * * Reads only — no engine submission. Returns `done: false` when the * `_DONE.json` marker is missing, in which case the caller should poll * status diagnostics again before re-fetching. */ /** * The S3 location where a query's matched events land as JSONL objects: * `{indexSubpath}/tenx/{target}/qr/{queryId}/`. This is the canonical * "list of result objects" a caller reads to get the full match set * beyond the in-context preview, and the handoff point for the * customer's own S3 -> SIEM path. The engine writes one `*.jsonl` per * stream worker plus a `_DONE.json` marker here. */ export declare function retrieverResultsLocation(target: string, queryId: string): Promise<{ bucket: string; prefix: string; uri: string; }>; export declare function fetchExistingResults(queryId: string, options?: { target?: string; }): Promise; /** Marker prefix for the unresolvable-pattern-name error (caller input). */ export declare const UNRESOLVABLE_PATTERN_PREFIX = "Pattern name not resolvable against the archive:"; /** * A Reporter-named pattern CANNOT be matched against the offload archive, and * every previous attempt to do so failed silently. * * Do NOT emit `tenx_user_pattern == ""`. **That field does not * exist.** It appears zero times across the engine, the modules tree and * pipeline-extensions. It is a plausible-looking member of a real family * (`tenx_user_service`, `tenx_user_process` do exist) that the product never * stamps, so a name-scoped archive query built on it returns * BLOOM_REJECTED_ALL with 0 of 40 blobs matched — a confident empty answer. * * Nor is there a working field to substitute. The Bloom index is built over * tokens literally present in the archived bytes, plus template hashes. On * the demo environment, one window: * * includes(text, "FU1__vh8hbY") -> 608 in the text * tenx_hash == "KJrTvZAaIhI" -> 1025 a template hash * severity_level == "INFO" -> 1025 "info" is in the text * message_pattern == "info_cart_..._GetCartAsync..." -> 0 BLOOM_REJECTED_ALL * tenx_user_pattern == "info_cart_..." -> 0 field does not exist * * `message_pattern` is the real `symbolMessageField` (the worker logs * "Enriching TenXObjects with message field: 'message_pattern'"), and it still * cannot be matched: a Symbol Message is a DERIVED label, computed from the * event, never a token in the bytes. No derived label is Bloom-matchable. * * So the name route is not a naming bug with a correct field waiting to be * found. It is unsatisfiable by construction. * * The identity that DOES work is the metrics-side `pattern_hash`, matched as a * text token via {@link buildArchiveHashSearch}. Callers always have it wherever * they have the name: top_patterns returns `pattern_hash` and `symbol_message` * on the same row, and event_lookup resolves a name to a hash. Pass the hash. * * This throws rather than returning a predicate. Erroring with a remedy beats * the previous behaviour of a silent, authoritative-looking zero, which is what * made reviewers conclude the retriever was broken. Affects three call sites: * retriever_query, retriever_series and backfill_metric. * * The inverse — parsing a pattern back out of a search expression — lives in * retriever-fidelity.ts as `extractPatternFromSearch`. */ export declare function buildPatternSearch(pattern: string): string; /** * Build an archive search for a pattern_hash taken from the METRICS surface. * * The metrics surface and the archive do not share an identity space, and that is * the single reason "retrieve the offloaded events for this pattern" came back * empty through the documented chain. * * `tenx_hash` on the archive read path is RE-DERIVED from the grouped event, not * read from the stored envelope (initialize/message/message-template.js stamps it * from the group's symbolSequence). Multi-line grouping collapses many * receiver-stamped lines into ONE read-time identity, so the field value differs * from the per-line hash the Reporter published. On the demo environment, * one window, one archive: * * metrics (top_patterns): FU1__vh8hbY 57.8% NiDD7PpZw48 28.6% +2 more * archive (read-time): KJrTvZAaIhI <- ALL 1,025 events carry this one value * * tenx_hash == "FU1__vh8hbY" -> 0 events (bloomMatched 18, 1.01 MB read) * tenx_hash == "KJrTvZAaIhI" -> 1025 events * includes(text, "FU1__vh8hbY") -> 608 events (59.3% vs 57.8% expected) * includes(text, "NiDD7PpZw48") -> 290 events (28.3% vs 28.6% expected) * * The stamped hash survives in the raw event text, which is also why the Bloom * filter matched blobs for a predicate that could never match an event: the index * is built over that text. Matching it as a TOKEN recovers the correct per-pattern * cohort, and the two independent cohorts reconcile against their metrics shares * to within 1.5 points. * * A field equality cannot be used, and a caller cannot discover the read-time * value without first reading a result file, which is not an API. * * Substring rather than equality is a deliberate tradeoff: an 11-char base64url * hash occurring incidentally elsewhere in an event would be a false positive. * That is preferable to silently returning nothing. The clean fix is engine-side, * having the read path preserve the stored per-line `tenx_hash` (the archive * objects already carry it as a JSON field) rather than re-deriving it. This makes * the documented chain work in the meantime. */ export declare function buildArchiveHashSearch(hash: string): string; export interface RetrieverEvent { timestamp?: string; text?: string; severity_level?: string; tenx_user_service?: string; k8s_namespace?: string; k8s_pod?: string; k8s_container?: string; http_code?: string; service?: string; severity?: string; /** Engine-INTERNAL field-set fingerprint that joins encoded events to * templates.json. NOT the agent-facing identity — that is tenx_hash / * pattern_hash. Do not surface as the event's stable ID. */ templateHash?: string; enrichedFields?: Record; values?: string[]; /** Any additional fields returned by the engine. */ [key: string]: unknown; } export interface RetrieverBucket { timestamp: string; count: number; labels?: Record; } /** * One TenXSummary record uploaded by the engine's per-slice summaries * writer (qrs/ prefix). Aggregated by the pipeline's grouping fields * (typically pattern + service + pod + severity); `summaryVolume` is * the event count contributed to this row by THIS worker. The slice * bounds come from the S3 key path, NOT the record itself. * * The engine emits enrichment fields as NAMED top-level keys (same * pattern the events writer uses, via `$=yield TenXEnv.get("enrichmentFields")`) * so consumers can look up `record.severity_level` / `record.k8s_pod` * directly. This makes the schema self-describing and survives * customer customization of `enrichmentFields`. */ export interface RetrieverSummary { /** Lower bound (epoch ms) of the slice this summary belongs to. */ sliceFromMs: number; /** Upper bound (epoch ms, exclusive) of the slice. */ sliceToMs: number; /** Event count this summary row aggregates. */ summaryVolume: number; /** UTF-8 byte sum of the events in this row. */ summaryBytes: number; /** Hash over the grouping fields (deterministic group key). */ summaryValuesHash?: string; /** * Named enrichment fields (severity_level, tenx_user_service, * k8s_pod, message_pattern, etc.). Order/presence depends on the * deployment's `enrichmentFields` config — consumers should look up * by name, never by position. */ [field: string]: unknown; } export interface RetrieverQueryResponse { /** * How the fan-out was judged complete. * * `marker_count_confirmed` — the coordinator published `expectedMarkers` and * we counted exactly that many. The result set is complete. * * `quiet_window_inferred` — the coordinator could not know the count (the * Lambda flavor dispatches scan tasks to SQS and exits; stream markers are * written by other processes whose number depends on bloom matches), so we * waited for the marker set to go quiet. That is a GUESS. It cannot tell * "every worker finished" from "the next worker is slow", and a short read * here is silent. * * Never render an inferred result as verified-complete. Absent this flag a * forensic query returned 330 of 1,217 events while reporting * `partial_results: false` and "no worker reported partial results". */ completenessBasis?: 'marker_count_confirmed' | 'quiet_window_inferred'; queryId: string; target: string; from: string; to: string; execution: { wallTimeMs: number; eventsMatched: number; workerFiles: number; truncated: boolean; /** * Worker JSONL files that failed to download after retries. > 0 means * the event set (and event-derived rollups) are INCOMPLETE — surface a * partial caveat; the full set is still intact in S3. */ failedWorkerFiles?: number; /** True when the qr/ download was capped (summaries served the rollups). */ downloadCapped?: boolean; /** Total qr/ worker files that exist for this query (vs downloaded). */ totalWorkerFiles?: number; /** Per-slice summary record count when summaries were requested. */ summariesMatched?: number; /** Distinct slice subdirs observed. */ slicesObserved?: number; }; events: RetrieverEvent[]; /** * Populated only when the request set `writeSummaries: true`. Each * record carries its slice bounds from the S3 key path so callers * can bucket by slice (typically 1 minute on the demo cluster's * `queryScanFunctionParallelTimeslice` default). */ summaries?: RetrieverSummary[]; /** * Structured execution diagnostics built by polling CloudWatch Logs for * the query's per-stream-worker streams. Populated when * `LOG10X_RETRIEVER_LOG_GROUP` is set and the CW SDK can reach the log * group. `pollingError` is set in place of a silent undefined when CW * is unreachable — callers should check and degrade explicitly. */ diagnostics?: RetrieverQueryDiagnostics; format?: 'events' | 'count' | 'aggregated'; buckets?: RetrieverBucket[]; countSummary?: { total: number; byService?: Record; bySeverity?: Record; byDay?: Record; }; /** * Which ingress path delivered the query to the retriever engine. * `"http"` — normal path: POST /streamer/query to the query-handler URL. * `"sqs"` — fallback path: SendMessage to the Quarkus ingress queue * (LOG10X_RETRIEVER_QUERY_QUEUE_URL), used when the HTTP URL is a * ClusterIP address unreachable from outside the cluster. */ transport?: 'http' | 'sqs'; /** * Wall time from SQS SendMessage until S3 polling completed, in ms. * Populated only when transport="sqs". */ sqsLatencyMs?: number; } export type RetrieverDetectionPath = 'explicit_env' | 'aws_s3_bucket_pattern' | 'kubectl_service' | 'terraform_state' | 'helm_release_probe'; export interface RetrieverResolution { url?: string; bucket?: string; target?: string; detectionPath?: RetrieverDetectionPath; trace: Array<{ path: RetrieverDetectionPath; status: 'matched' | 'skipped' | 'failed'; reason: string; }>; } /** * Resolve the Retriever URL + bucket from the ambient environment. * * Detection order (first hit wins): * 1. explicit __SAVE_LOG10X_RETRIEVER_URL__ + __SAVE_LOG10X_RETRIEVER_BUCKET__ * 2. AWS conventional bucket naming (log10x-retriever-*, *-log10x-archive) * combined with a kubectl-discovered query-handler service URL * 3. kubectl service probe alone (url resolved, bucket from annotation) * 4. Terraform state file (~/.log10x/retriever.tfstate or LOG10X_TERRAFORM_STATE) */ export declare function resolveRetriever(): Promise; /** * Fast-path synchronous check (explicit env vars only). * Kept for back-compat with any callers that cannot await. * For the full kubectl-probe cascade use isRetrieverConfigured(). */ export declare function isRetrieverConfiguredSync(): boolean; /** * Fix 83a/83b — async gate that consults resolveRetriever so a * helm-probe-discovered install (no env vars set) is treated as configured. * * Resolution order (matches resolveRetriever()): * 1. Env-var fast-path — if both __SAVE_LOG10X_RETRIEVER_URL__ and * __SAVE_LOG10X_RETRIEVER_BUCKET__ are set, return true immediately. * 2. resolveRetriever() — full cascade (terraform state, AWS bucket * pattern, kubectl probe, helm_release_probe). Returns true only when * BOTH url AND bucket are resolved, which is the same precondition * that runRetrieverQuery requires. This closes the 83a gap where * getRetrieverState() returned url-only (helm probe, no bucket) and * isRetrieverConfigured() returned true while the inner tool call still * threw RetrieverNotConfiguredError on the missing bucket. */ export declare function isRetrieverConfigured(): Promise; export declare function formatRetrieverTrace(trace: RetrieverResolution['trace']): string; /** Reset the resolver cache — only for tests. */ export declare function clearRetrieverResolutionCacheForTest(): void; /** * Convert the MCP-level `from`/`to` expressions to a form the retriever engine * reliably parses. * * Empirical behavior of the retriever server: * - `now("-1h")` / `now()` → accepted, matches events * - epoch millis as a string → accepted, matches events * - ISO8601 like `2026-04-15T11:00:00Z` → accepted (HTTP 200), runs the * query, returns ZERO events even when the wall-clock range should match * — the server-side TenXDate parser mishandles ISO8601 and silently * produces a non-matching range. * * To make ISO8601 inputs work reliably, this function converts them to epoch * millis strings before they leave the MCP. `now(...)` expressions are * preserved verbatim (they require server-side evaluation). * * Filed upstream as a retriever engine bug. The client-side conversion below * is a workaround. */ export declare function normalizeTimeExpression(expr: string): string; /** Expose for validation at the tool layer. */ export declare function parseTimeExpression(expr: string): string; /** * Run `fn` over `items` with at most `concurrency` calls in flight, preserving * input order in the returned array. The first rejection propagates and an * abort flag stops idle workers from pulling new items — so one failed S3 read * aborts the whole fetch (matching the old serial loop's throw-on-error * semantics) instead of silently returning a partial result set. * * Why bounded (not Promise.all over everything): the Retriever fans out to * dozens of stream workers, each writing its own JSONL. An unbounded fan-in * would spawn an `aws s3 cp` subprocess per file simultaneously (each buffering * up to 64 MB) and can trip S3 503 SlowDown on a hot prefix. A small pool * captures the bandwidth win without the blowup. */ export declare function mapWithConcurrency(items: readonly T[], concurrency: number, fn: (item: T, index: number) => Promise): Promise; /** * Like mapWithConcurrency but per-item failures do NOT abort the run: * failed slots return undefined and are reported in `failures`. Used by the * result-download paths so one persistently-failing worker file (e.g. S3 * 503 SlowDown that outlives retries) degrades the fetch to partial instead * of losing the whole query. */ export declare function mapWithConcurrencySettled(items: readonly T[], concurrency: number, fn: (item: T, index: number) => Promise): Promise<{ results: Array; failures: Array<{ index: number; error: Error; }>; }>; /** * Download per-worker JSONL in concurrency-bounded BATCHES, stopping once * `budget` events have accumulated. Order-preserving over the files actually * pulled. Used to cap the qr/ download to a preview-sized sample when qrs/ * summaries already supply the whole-match count + rollups, so a huge match * never materializes client-side. */ export declare function downloadEventsUntilBudget(keys: readonly string[], budget: number, concurrency: number, fetchParse: (key: string) => Promise): Promise<{ events: RetrieverEvent[]; failures: number; filesDownloaded: number; stoppedEarly: boolean; }>; /** * Normalize a retriever event's timestamp to epoch-ms. * * The retriever encodes timestamps as `[]` (array wrapping a single * value). The scalar's unit varies by upstream pipeline configuration — * fluent-bit ships seconds, the engine's TenXObject ships millis or nanos * depending on the input adapter. Magnitude-based detection is the only * portable approach. * * Ranges (a present-day epoch ≈ 1.77 × 10^X): * - seconds: ~1.77e9 (10 digits) * - millis: ~1.77e12 (13 digits) * - micros: ~1.77e15 (16 digits) * - nanos: ~1.77e18 (19 digits) * * Boundaries are placed at 10^10, 10^13, 10^16 — one decade below the * current epoch in each unit so 13-digit millis values like 1776851170107 * are correctly classified as millis instead of falsely * matching the looser `>1e12` micros boundary and dividing by 1000 to * land in 1970. Without the decade-below boundary, the entire bucket * histogram aliases to 1970 because 13-digit millis falsely match the * looser micros boundary and get divided by 1000. */ export declare function eventTimestampMs(ev: RetrieverEvent): number; export declare function runRetrieverQuery(env: EnvConfig, req: RetrieverQueryRequest, options?: { timeoutMs?: number; pollIntervalMs?: number; }): Promise;