From 129046c2784583c1cebae02b3f3ade5563886232 Mon Sep 17 00:00:00 2001 From: Kris McGinnes Date: Fri, 25 Sep 2026 17:10:38 -0500 Subject: [PATCH 1/4] Ask the database for distinct edge connections instead of scanning edges --- CONTEXT.md | 6 +- ...-database-for-distinct-edge-connections.md | 95 ++++ docs/guides/troubleshooting.md | 6 +- .../discoveryError.test.ts | 113 ++++ .../fetchEdgeConnections/discoveryError.ts | 105 ++++ .../discoveryPlan.test.ts | 318 +++++++++++ .../fetchEdgeConnections/discoveryPlan.ts | 278 ++++++++++ .../edgeConnectionsTemplate.test.ts | 41 +- .../edgeConnectionsTemplate.ts | 96 +++- .../fetchEdgeConnections/index.test.ts | 504 ++++++++++++++---- .../gremlin/fetchEdgeConnections/index.ts | 237 ++++++-- .../src/connector/gremlin/gremlinExplorer.ts | 11 +- .../src/connector/gremlin/types.ts | 12 + .../connector/queries/edgeConnectionsQuery.ts | 8 +- .../src/connector/useGEFetchTypes.ts | 6 + ...geConnectionDiscoveryStatusButton.test.tsx | 49 +- .../src/utils/createDisplayError.test.ts | 33 ++ .../src/utils/createDisplayError.ts | 23 + .../src/utils/createErrorDetails.test.ts | 29 + .../src/utils/createErrorDetails.ts | 14 +- 20 files changed, 1797 insertions(+), 187 deletions(-) create mode 100644 docs/adr/20260918-ask-the-database-for-distinct-edge-connections.md create mode 100644 packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryError.test.ts create mode 100644 packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryError.ts create mode 100644 packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryPlan.test.ts create mode 100644 packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryPlan.ts diff --git a/CONTEXT.md b/CONTEXT.md index 07cc07337c..0d13375098 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -62,6 +62,10 @@ A React-based web application that lets users visually explore graph databases w **Schema Sync**: The process that queries the database to discover vertex types, edge types, and their attributes. Required before a user can explore a new Connection. _Avoid_: Fetch, load +**Edge Connection Discovery**: The process that finds which Edge Connections exist. Runs after Schema Sync completes and is a separate step: Schema Sync finds the Edge Types, Edge Connection Discovery finds how they link Vertex Types. Only the Schema View depends on it, so a failure degrades that view and leaves the rest of the app working. _Avoid_: Relationship Discovery (follows the Edge Connection entry), edge discovery + +**Complete** / **Sampled**: The two strategies Edge Connection Discovery can take. **Complete** reads every edge, so it finds every Edge Connection. **Sampled** caps the edges it reads per Edge Type, so it is bounded on a large graph but will miss an Edge Connection that occurs rarely. The strategy is chosen up front by comparing predicted cost, and a Complete scan is abandoned for Sampled if it proves too large for the database. Both strategies may split their work across several requests; neither split is a third strategy. _Avoid_: Full, exhaustive, partial, approximate, mode (for the pair) + **Schema**: The discovered structure of a connected graph database — vertex types, edge types, their attributes, and how they connect. Populated by Schema Sync when a Connection is first used; not user-defined. _Avoid_: Model, structure **Exported Connection File**: The on-disk JSON format a user gets when they export a Connection (`saveConfigurationToFile`), and which import consumes. It bundles the connection config with a snapshot of the Schema (`lastUpdate` is an ISO string on disk). A single Zod schema in `parseConnectionFile.ts` is the source of truth: both the writer and the importer target the same inferred type (`ExportedConnectionFile`). The writer assigns a `Date` for `lastUpdate` and `JSON.stringify` serializes it to the ISO string; the parser coerces it back via `z.coerce.date()`. The schema is lenient (every level is a `looseObject`), so unknown and legacy fields — styling, `__inferred`/`__matches` on prefixes, attribute `dataType` — pass through untouched. It is intentionally decoupled from the in-memory configuration and from the IndexedDB storage shape, so the wire format can evolve independently. On import it is split — the connection lands in `configurationAtom`, the schema in `schemaAtom`. _Avoid_: Configuration file (the wire format is not the in-memory or persisted shape) @@ -85,7 +89,7 @@ A React-based web application that lets users visually explore graph databases w - A **Schema** contains **Vertex Types**, **Edge Types**, and **Edge Connections** - A **Vertex** has one or more **Vertex Types** and zero or more **Properties** - An **Edge** connects exactly two **Vertices** (source → target), has one **Edge Type**, and zero or more **Properties** -- An **Edge Connection** links a source **Vertex Type** to a target **Vertex Type** via an **Edge Type** +- An **Edge Connection** links a source **Vertex Type** to a target **Vertex Type** via an **Edge Type**, and is found by **Edge Connection Discovery** - A **Session** belongs to a **Connection** and contains **Vertices** and **Edges** - **Neighbors** are **Vertices** one hop away from a given **Vertex** - **Styles** are scoped per **Vertex Type** (**Vertex Styles**) and **Edge Type** (**Edge Styles**) diff --git a/docs/adr/20260918-ask-the-database-for-distinct-edge-connections.md b/docs/adr/20260918-ask-the-database-for-distinct-edge-connections.md new file mode 100644 index 0000000000..4fa59e7352 --- /dev/null +++ b/docs/adr/20260918-ask-the-database-for-distinct-edge-connections.md @@ -0,0 +1,95 @@ +# ADR — Ask the database for distinct edge connections instead of sampling edges + +- **Status:** Accepted +- **Date:** 2026-09-18 +- **Related:** Fixes [#2141](https://github.com/aws/graph-explorer/issues/2141). Replaces the Gremlin half of [#2100](https://github.com/aws/graph-explorer/pull/2100), which batched edge connection discovery for [#2085](https://github.com/aws/graph-explorer/issues/2085). openCypher and SPARQL keep their `#2100` shapes. + +## Context + +Edge connection discovery finds which vertex types each edge type connects. Since 3.2.2 the Gremlin connector batches it into one query per 100 edge types: + +```gremlin +g.E().hasLabel('A','B','C') + .group().by(label()) + .by(limit(10000).project('sourceType','targetType') + .by(outV().label()).by(inV().label()).dedup().fold()) +``` + +The per-type `limit` caps sampling inside each bucket but not the scan that fills the buckets. `group()` enumerates every edge of the batched types first, and Neptune does not convert it to a native step, so it materializes one item per edge outside the engine. A customer with three edge types over 19.9M edges gets `MemoryLimitExceededException` after about 30 seconds and an error in the Schema View. + +Reproduced on a purpose-built 1M-edge graph: the same failure at 35.5s. The profile shows 57,538 items crossing into the query layer on a 57k-edge graph and `GroupStep` consuming 88% of runtime, with an explicit "not supported natively" warning on both DFE and non-DFE Neptune. + +`#2100`'s doc comment claims a per-type scan cap "is not expressible in one native TinkerPop 3.6.2 request". That is false, and the counter-example already shipped two directories away in `fetchSchema/edgesSchemaTemplate.ts`. + +Three facts, all measured live, shaped the decision: + +- **`groupCount()` keyed by a projection is native and memory-lean.** It returns one item regardless of graph size, because the accumulator is keyed by the answer rather than the input. Verified native and correct on Neptune 1.2.1.0 (our TinkerPop 3.6.2 floor), 1.3.5.0, 1.4.5.1 on an instance with no DFE, and 1.4.7.0, plus reference TinkerPop 3.6.2 and 3.7.3. +- **The ceiling for an unsampled scan is unpredictable from the client.** On DFE it fails on memory between 100,000 and 150,000 edges in under 12 seconds. On a non-DFE instance of the same hardware class it survives 550,000 and fails on time at 1M. Different limits, different failure modes, and neither DFE presence nor instance memory is visible to us. +- **Several plausible query shapes are silently wrong.** See Considered Options. + +## Decision + +**Ask for the distinct `(edgeType, sourceVertexType, targetVertexType)` combinations directly, with two strategies and one query shape for each.** + +Complete: + +```gremlin +g.E()[.hasLabel(...)] + .groupCount().by(project('e','s','t') + .by(label()).by(outV().label().fold()).by(inV().label().fold())) +``` + +Sampled, one limited branch per edge type: + +```gremlin +g.V().limit(1).union(V().outE('A').limit(10000), V().outE('B').limit(10000), ...) + .group().by(label()) + .by(project('s','t').by(outV().label().fold()).by(inV().label().fold()).groupCount()) +``` + +The endpoint labels are folded because engines disagree on what `label()` emits for a multi-label vertex. Neptune 1.2 and 1.4 emit one `::` composite; 1.3.5 emits each label as its own value, so an unfolded `by(outV().label())` kept only the first and silently dropped the vertex's other types. Measured on a fixture of `a::b` and `d::e::f` vertices on all four versions. The parser splits every folded entry on `::`. + +1. **Complete strategy.** No `limit`. Returns every edge connection plus exact edge counts. Split across requests when the graph exceeds the scan budget; chunking is this strategy issued N times, not a third strategy. +2. **Sampled strategy.** Each edge type gets its own limited `union()` branch, `g.V().limit(1).union(V().outE('A').limit(10000), ...)`, 10 edge types per request, grouped by edge type before counting: `.group().by(label()).by(project('s','t').by(outV().label().fold()).by(inV().label().fold()).groupCount())`. Grouping first is for DFE, which could not count one `project()` key across several full branches: on 1.3.5.0 two took 54s and five timed out, where grouping first handled ten in 9.4s. Grouping holds each type's sample until it is counted, so a request holds at most the batch times the limit, 100,000 edges; the complete strategy has no limit and keeps `groupCount()`, because grouping an unlimited 550,000-edge type ran out of memory on 1.4.5.1. Matches the earlier one-request-per-type semantics with a tenth of the requests. On Neptune's default engine each branch is an index lookup by edge label, and it returned the same combinations and counts as one request per type on 1.4.5.1 and 1.4.7.0, including missing the same planted rare combinations. On TinkerGraph it samples in vertex order rather than edge order, so it misses different rare combinations than a per-type request would, not more. The batch is sized for the worst case of every branch reaching its limit, 100,000 edges, because a type's count is unknown until read: 10 full branches took 9.7s on a db.t3.medium, 100 took 116s and left the instance refusing even single-type samples on memory for two minutes. The budget is per request, and requests run four at a time; four concurrent requests of 10 full branches finished in 36s on both db.t3.medium instances, with and without DFE, and left them answering normally. A separate cap of 100 branches keeps the batch under Gremlin Server's Groovy argument limit if the per-type limit ever shrinks. `DiscoveryStrategy` carries a third value, `none`, for a schema with no edge types; that is a guard clause rather than a strategy choice. +3. **Choose by comparing predicted cost.** Complete scales with edge count, sampled with edge type count. Below the budget, take complete. Above it, take the cheaper of `edgeTypes x 1.5s` and `edges x 85us`. Sampled stays charged per edge type after batching, because a branch that reaches its limit costs as much as the separate request it replaced. With no recorded edge total, sample only when `edgeTypes x 1.5s` fits 60 seconds, 40 edge types or fewer, and otherwise run a complete scan in chunks of 100 edge types. +4. **Scan budget is 50,000 edges, and it is not a safety prediction.** It is the highest volume that never failed on any configuration tested. Given the measured spread, no constant can predict the real ceiling. +5. **A second bound caps the query text, at 60,000 characters of edge type names.** The edge total says how much a request reads; it says nothing about how much it carries. A graph with thousands of edge types would otherwise name them all in one filter, so the two bounds are enforced independently and a chunk splits when either is reached. Also the largest size measured rather than the point of failure: 5,008 names of 11 characters completed in 7.1s on 1.4.7.0. +6. **The degrade path is what recovers from a budget we cannot predict.** A complete request that fails because it was too big abandons the complete strategy and redoes the whole discovery as sampled. The trigger is our own request timeout, with Neptune's `MemoryLimitExceededException` and `TimeLimitExceededException` as a fast path. Fail fast on the first such error and do not retry individual chunks. +7. **An automatic complete request is bounded at 20 seconds, separately from the connection's fetch timeout.** That timeout defaults to four minutes, so leaving the degrade path to it would make a graph above the budget wait four minutes for a schema it could have sampled in one. 20 seconds is over three times the cost of reading the whole scan budget at the slowest rate measured, and under the 30 to 35 seconds at which Neptune failed in both reproductions, so we give up before the database does. +8. **Degrading applies to every complete scan, because the strategy is always chosen automatically for now.** A per-connection setting to force either strategy is the next step. A forced complete should report its failure rather than degrade, since silently sampling would contradict the setting. +9. **The `count` the complete strategy returns is not persisted.** The `EdgeConnection.count` field stays unpopulated, because the same field would be capped and misleading on the sampled path. + +The decision is evaluated in the Gremlin connector, by a pure `planDiscovery` that takes the edge types and the edge total, and returns the strategy with the requests it takes. + +`EdgeConnectionsRequest` carries `{ edgeTypes, totalEdges }`, and openCypher and SPARQL ignore the total. + +An abandoned attempt is cancelled rather than merely ignored. `mapWithConcurrency` stops pulling new work once a callback rejects, and every request carries its own `queryId` and abort signal, so aborting the attempt closes the connections still open and the proxy turns each one into a `cancelQuery` at the database. Without that, the chunks already in flight would keep scanning for a result nobody is waiting for, and on the degrade path they would compete with the sampled requests that replace them. + +Giving up produces an `EdgeConnectionDiscoveryError` rather than the database's error, because only we know which recovery path is still open: a failed sample has already exhausted the cheaper options, so what is left is the connection's fetch timeout or the database's own limits. The original error is the `cause`, so the error details dialog still shows what the database said. Failures that are not about size propagate untouched, since the existing display wording already reads a refused connection or a bad URL correctly. + +## Considered options + +- **`groupCount().by(project(...))` (chosen).** Native on every engine tested, one item returned, correct direction, and every label of a multi-label vertex preserved once the endpoint labels are folded. +- **The same shape with unfolded endpoint labels.** Rejected as silently wrong on 1.3.5.0, where it returned `b` for an `a::b` vertex in both the complete and sampled strategies. The shape it replaced, `group().by(label()).by(project(...).dedup().fold())`, returned both labels there, so this was a regression. +- **Keep `group().by(label()).by(limit(...))`.** Rejected: the shape that fails. Not native, and holds one item per edge. +- **`dedup()` on a projected map.** Rejected as silently wrong. Returned 5 of 13 expected combinations on Neptune 1.2.1.0, because `dedup()` does not compare projected maps by content. Same defect in `g.V().outE().project(...).dedup()` and `path().by(label())`. +- **`groupCount()` keyed by `union(label(), outV().label(), inV().label()).fold()`.** Rejected as silently wrong, and the most dangerous of the three: it is native and looks correct, but `union()` does not guarantee order and DFE permuted the key on 1.3.5.0, reporting `contains` as `airport -> continent` when the truth is `continent -> airport`. A named `project()` key is what makes the chosen shape safe. +- **Anchored per-type sub-traversals,** `g.V().limit(1).project(types).by(V().outE(type).limit(n)...)`. Rejected: Neptune cannot run a traversal inside `project().by()` natively, so it costs roughly 0.23s per type. Its miss of an endpoint pair occurring 50 times against 550,000 was measured on 1.3.5.0, where one request per type misses the same pair, so it is not evidence against sampling from `V()`. The `union()` form of the same idea is native and is the chosen sampled strategy. +- **One shared limit per batch,** `g.E().hasLabel(A, B, ...).limit(n)`. Rejected as silently wrong: on 1.4.7.0 `route` filled a limit of anything from 100 to 20,000 on its own, and 99 of 100 edge types in the batch came back empty. +- **Larger sampled batches.** Rejected. Beyond the time and memory cost of full branches above, Gremlin Server compiles the script as Groovy and a call cannot take more than about 250 arguments, so `union()` fails outright above 250 branches on 3.6.2 and 252 on 3.7 and 3.8. Neptune's DFE engine crashed a db.t3.medium on a 200-branch batch. +- **`project()` inside each `union()` branch.** Rejected: faster on DFE, but on the default engine it returned only the first label of a multi-label vertex, and ran out of memory on 1.4.5.1 at 250 branches where the chosen shape did not. +- **`union` with mid-traversal `E()` per arm.** Rejected: requires TinkerPop 3.7, so it drops Neptune 1.2.x and 1.3.x, and it only ties the 3.6.2-compatible options on the case we care about while still scaling linearly in edge type count. +- **Raise the minimum TinkerPop version to 3.7.** Rejected: the profile shows `E().hasLabel()` inside a union arm is not index-backed on Neptune, so a sparse edge type still scans the whole edge store per arm. We would pay a compatibility cost and still need the same strategy split. + +## Consequences + +- **The `#2085` request collapse is preserved and improved for most graphs.** The 10,015-edge-type graph generated to test this goes from 101 requests to 3 at the two budgets, and to 1 on any graph under the scan budget. The character budget costs one request over what the scan budget alone would ask for, and buys a bound that holds however long the edge type names are. +- **Sampled coverage is engine-dependent, and only complete looks at every edge.** On Neptune, sampled found a pair occurring 50 times in 550,050. On reference TinkerPop 3.6.2 and 3.7.3 the same test missed it, because TinkerGraph iterates in insertion order. Any user-facing description of sampling must say this plainly rather than imply it is theoretical. +- **Attempting complete and degrading costs the user time.** On a graph above the budget, the first sync spends 8 to 12 seconds discovering that complete does not fit, or 20 seconds when the database neither answers nor refuses. That is the price of the budget being a guess rather than a prediction. +- **A 20 second bound can degrade a scan that would have finished.** A graph whose edges are concentrated in one edge type gets chunks balanced by edge type count, not by edge count, so a single chunk can read far more than the budget and still have been on its way to succeeding. The cost of getting it wrong is a sampled schema, which can miss a rare connection; the cost of not bounding it is a four minute wait. Worth it in that trade, but it is a trade. +- **The cost constants are Neptune-derived.** 85 microseconds per edge and 1.5s per request come from Neptune measurements and will be wrong for JanusGraph on Cassandra. Tolerable because the crossover is insensitive, with real graphs sitting one to three orders of magnitude away from it, but the numbers are not universal. +- **Sampled's flat per-branch cost depends on an edge type index.** Neptune turns each `V().outE(X)` branch into a label index lookup. TinkerGraph has no such index, so a branch for a sparse type walks every vertex. Batching still wins there, because the fixed cost of a request is paid once: 250 sparse types on a 1.8M-edge TinkerGraph took 4.3s as one request and 209s as 250. +- **Sampled batches are sized for the worst case.** A graph of mostly sparse edge types would finish in far fewer requests at a larger batch, but nothing before the read says which types are sparse. Weighting batches by the per-type edge counts the schema sync sometimes records is left to #2254. +- **Nothing above 1M edges is tested.** The reported graph is 19.9M, and graphs an order of magnitude larger than that exist. Sampled work per request is capped by construction, so scale affects timing rather than correctness, but the budget is extrapolated. +- **A graph that is both very large and has very many edge types has no cheap option.** Sampling it is still one request per 10 edge types, and scanning it takes hours. No reported case has that shape, but our visibility is limited to graphs that generate a ticket, so this is unmeasured rather than ruled out. +- **The `EdgeConnectionDetails` panel keeps showing an overcount.** It displays the edge type's aggregate total on an edge connection panel, which overcounts whenever an edge type has more than one connection. The complete strategy could fix this for free, but wiring it would make a displayed number depend on which strategy ran. Tracked separately. diff --git a/docs/guides/troubleshooting.md b/docs/guides/troubleshooting.md index eb07e8a446..83a69d2d67 100644 --- a/docs/guides/troubleshooting.md +++ b/docs/guides/troubleshooting.md @@ -125,7 +125,11 @@ If a request is cancelled instead, Graph Explorer shows a plain "Request cancell ### Out of Memory -This can happen when your database is very large. Graph Explorer does its best to support larger databases and is always improving. Please [file an issue](https://github.com/aws/graph-explorer/issues/new/choose) if you encounter this situation. +This can happen when your database is very large. Graph Explorer does its best to support larger databases and is always improving. + +For a Gremlin connection, this is often the database running out of memory while discovering edge connections for the Schema view. Graph Explorer samples each edge type instead of scanning every edge when a graph is too large, and if that still fails you can raise the query timeout in the database configuration, such as the DB cluster parameter group for Neptune, or use an instance with more memory. + +Otherwise, please [file an issue](https://github.com/aws/graph-explorer/issues/new/choose) if you encounter this situation. ### Proxy Server Cannot Be Reached diff --git a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryError.test.ts b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryError.test.ts new file mode 100644 index 0000000000..6b4597529a --- /dev/null +++ b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryError.test.ts @@ -0,0 +1,113 @@ +import { DatabaseTimeoutError, FetchTimeoutError, NetworkError } from "@/utils"; + +import { + EdgeConnectionDiscoveryError, + type FailedDiscovery, + isTooBig, +} from "./discoveryError"; + +function attempt(overrides: Partial = {}): FailedDiscovery { + return { + strategy: "sampled", + requests: 1, + totalEdges: 10, + degraded: false, + cause: "database-limit", + ...overrides, + }; +} + +describe("isTooBig", () => { + it("is true for a FetchTimeoutError", () => { + expect(isTooBig(new FetchTimeoutError(1000, new Error("aborted")))).toBe( + true, + ); + }); + + it("is true for a DatabaseTimeoutError", () => { + expect( + isTooBig( + new DatabaseTimeoutError( + "Query cannot be completed", + 500, + {}, + "TimeLimitExceededException", + ), + ), + ).toBe(true); + }); + + it("is true for a NetworkError carrying the memory limit code", () => { + expect( + isTooBig( + new NetworkError("Query cannot be completed", 500, { + code: "MemoryLimitExceededException", + }), + ), + ).toBe(true); + }); + + it("is true when the memory limit code is nested in a cause", () => { + expect( + isTooBig( + new NetworkError("Query cannot be completed", 500, { + cause: { code: "MemoryLimitExceededException" }, + }), + ), + ).toBe(true); + }); + + it("is false for an unrelated NetworkError", () => { + expect( + isTooBig( + new NetworkError("Query cannot be completed", 500, { + code: "MalformedQueryException", + }), + ), + ).toBe(false); + }); + + it("is false for a user cancellation, so it never looks like a size problem", () => { + expect(isTooBig(new DOMException("Aborted", "AbortError"))).toBe(false); + }); + + it("is false for a plain error", () => { + expect(isTooBig(new Error("Network error"))).toBe(false); + }); +}); + +describe("EdgeConnectionDiscoveryError recovery text", () => { + it("points only at the Fetch Timeout, not the parameter group, when a cheaper pass exhausted our own fetch timeout", () => { + const error = new EdgeConnectionDiscoveryError( + attempt({ cause: "fetch-timeout" }), + new Error("cause"), + ); + expect(error.recovery).toContain("Fetch Timeout"); + expect(error.recovery).toContain("advanced options"); + expect(error.recovery).not.toContain("parameter group"); + expect(error.recovery).toContain( + "the Schema view shows node types without the edge connections between them", + ); + }); + + it("points at the database's own query timeout and the parameter group when a cheaper pass exhausted it", () => { + const error = new EdgeConnectionDiscoveryError( + attempt({ cause: "database-limit" }), + new Error("cause"), + ); + expect(error.recovery).toContain("DB cluster parameter group"); + expect(error.recovery).toContain( + "the Schema view shows node types without the edge connections between them", + ); + }); + + it("includes the cause in the structured details, keyed apart from the JS cause", () => { + const jsCause = new Error("Query cannot be completed"); + const error = new EdgeConnectionDiscoveryError( + attempt({ cause: "fetch-timeout" }), + jsCause, + ); + expect(error.details).toMatchObject({ failureCause: "fetch-timeout" }); + expect(error.cause).toBe(jsCause); + }); +}); diff --git a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryError.ts b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryError.ts new file mode 100644 index 0000000000..fb5338041d --- /dev/null +++ b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryError.ts @@ -0,0 +1,105 @@ +import { DatabaseTimeoutError, FetchTimeoutError, NetworkError } from "@/utils"; + +import type { DiscoveryStrategy } from "./discoveryPlan"; + +/** Neptune's code for a query that asked for more memory than the instance had. */ +const MEMORY_LIMIT_ERROR_CODE = "MemoryLimitExceededException"; + +/** Whether the database gave up because one request asked for too much at once. */ +export function isTooBig(error: unknown): boolean { + return ( + error instanceof FetchTimeoutError || + error instanceof DatabaseTimeoutError || + memoryLimitCode(error) !== undefined + ); +} + +/** + * The database's own memory-limit code, from either shape the body arrives in. + * Reading only the top level would miss a nested code and cost the degrade path + * its fast exit, leaving the user to wait out the request bound instead. + */ +function memoryLimitCode(error: unknown): string | undefined { + const data = error instanceof NetworkError ? error.data : undefined; + const code = data?.code ?? data?.cause?.code; + return code === MEMORY_LIMIT_ERROR_CODE ? code : undefined; +} + +/** + * Which side gave up. A fetch timeout is the connection's own bound running out, + * fixed in the connection's settings; a database limit is the database itself + * refusing the request, fixed in the database's configuration. + */ +export type FailureCause = "fetch-timeout" | "database-limit"; + +/** Classifies a size failure that `isTooBig` already confirmed. */ +export function causeOf(error: unknown): FailureCause { + return error instanceof FetchTimeoutError + ? "fetch-timeout" + : "database-limit"; +} + +/** What edge connection discovery had already tried when it gave up. */ +export type FailedDiscovery = { + strategy: DiscoveryStrategy; + requests: number; + totalEdges: number | undefined; + /** A complete scan was already abandoned as too large before this attempt. */ + degraded: boolean; + /** Which side gave up: the connection's fetch timeout, or the database itself. */ + cause: FailureCause; +}; + +/** + * Edge connection discovery has nothing cheaper left to try. + * + * Thrown only for a size failure, so the error reaches the user carrying the one + * thing the generic wording cannot give them: which recovery path is open. Other + * failures propagate untouched, because the existing display branches already + * read a refused connection or a bad URL correctly. + */ +export class EdgeConnectionDiscoveryError extends Error { + readonly attempt: FailedDiscovery; + /** What the user can do about it, in the order worth trying. */ + readonly recovery: string; + + constructor(attempt: FailedDiscovery, cause: unknown) { + super(describeFailure(attempt), { cause }); + // A literal rather than the class name, because the production build + // minifies class names and the error details dialog shows this. + this.name = "EdgeConnectionDiscoveryError"; + this.attempt = attempt; + this.recovery = describeRecovery(attempt); + } + + /** Structured context for the error details dialog. */ + get details() { + return { + strategy: this.attempt.strategy, + requests: this.attempt.requests, + totalEdges: this.attempt.totalEdges, + completeScanAbandoned: this.attempt.degraded, + // Named apart from `cause`, which `createErrorDetails` reserves for the + // serialized JS `Error.cause` and would otherwise overwrite this. + failureCause: this.attempt.cause, + }; + } +} + +/** + * Keyed on whether a complete scan was already abandoned, never on the + * strategy: either way discovery has run out of cheaper options. + */ +function describeFailure({ degraded }: FailedDiscovery): string { + if (degraded) { + return "The database could not discover edge connections either way. Scanning every edge was too large, and sampling each edge type failed as well."; + } + return "The database could not sample the edges of each edge type to discover edge connections."; +} + +function describeRecovery({ cause }: FailedDiscovery): string { + if (cause === "fetch-timeout") { + return "Raise the Fetch Timeout in this connection's advanced options, or clear it, since this request may simply need longer than that allows. Until then, the Schema view shows node types without the edge connections between them."; + } + return "Raise the query timeout in the database configuration, such as the DB cluster parameter group for Neptune, or use an instance with more memory. Until then, the Schema view shows node types without the edge connections between them."; +} diff --git a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryPlan.test.ts b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryPlan.test.ts new file mode 100644 index 0000000000..a84f50ec63 --- /dev/null +++ b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryPlan.test.ts @@ -0,0 +1,318 @@ +import { createEdgeType, type EdgeType } from "@/core"; +import { DEFAULT_SAMPLE_SIZE } from "@/utils"; + +import { + EDGE_TYPES_PER_CHUNK, + EDGE_TYPES_PER_SAMPLE, + MAX_UNION_BRANCHES, + SAMPLE_EDGE_BUDGET, + LABEL_BUDGET_CHARS, + COMPLETE_ATTEMPT_TIMEOUT_MS, + planDiscovery, + planSampling, + SCAN_BUDGET, +} from "./discoveryPlan"; + +function edgeTypes(count: number): EdgeType[] { + return Array.from({ length: count }, (_, i) => createEdgeType(`edge${i}`)); +} + +/** Every edge type appears exactly once across the planned requests. */ +function coveredTypes(requests: { edgeTypes?: EdgeType[] }[]) { + return requests.flatMap(r => r.edgeTypes ?? []); +} + +describe("Gremlin > planDiscovery", () => { + it("should plan nothing when the schema has no edge types", () => { + const plan = planDiscovery({ + edgeTypes: [], + totalEdges: 1_000_000, + }); + + expect(plan).toStrictEqual({ strategy: "none", requests: [] }); + }); + + describe("auto", () => { + it("should scan the whole graph in one unfiltered request when it fits the budget", () => { + const plan = planDiscovery({ + edgeTypes: edgeTypes(12), + totalEdges: 40_000, + }); + + expect(plan).toStrictEqual({ + strategy: "complete", + requests: [{}], + requestTimeoutMs: COMPLETE_ATTEMPT_TIMEOUT_MS, + }); + }); + + it("should treat the budget as inclusive", () => { + const plan = planDiscovery({ + edgeTypes: edgeTypes(3), + totalEdges: SCAN_BUDGET, + }); + + expect(plan.strategy).toBe("complete"); + expect(plan.requests).toStrictEqual([{}]); + }); + + it("should sample when the graph is too large to scan but has few edge types", () => { + const types = edgeTypes(3); + const plan = planDiscovery({ + edgeTypes: types, + totalEdges: 19_928_805, + }); + + expect(plan).toStrictEqual({ + strategy: "sampled", + requests: [{ edgeTypes: types, limitPerType: DEFAULT_SAMPLE_SIZE }], + }); + }); + + it("should chunk a complete scan when there are more edge types than the scan is worth", () => { + const types = edgeTypes(10_015); + const plan = planDiscovery({ + edgeTypes: types, + totalEdges: 68_582, + }); + + expect(plan.strategy).toBe("complete"); + // The edge count alone asks for 2 chunks. Naming 5,008 edge types in one + // filter overruns the character budget, so it splits once more. + expect(plan.requests).toHaveLength(3); + expect(plan.requests.every(r => !("limitPerType" in r))).toBe(true); + expect(coveredTypes(plan.requests)).toStrictEqual(types); + }); + + it("should chunk a complete scan for a graph that is both large and wide", () => { + const types = edgeTypes(2_010); + const plan = planDiscovery({ + edgeTypes: types, + totalEdges: 1_015_639, + }); + + expect(plan.strategy).toBe("complete"); + expect(plan.requests).toHaveLength(21); + expect(coveredTypes(plan.requests)).toStrictEqual(types); + }); + }); + + describe("auto with no recorded edge total", () => { + it("should sample when there are few enough edge types to stay bounded", () => { + const types = edgeTypes(3); + const plan = planDiscovery({ + edgeTypes: types, + totalEdges: undefined, + }); + + expect(plan).toStrictEqual({ + strategy: "sampled", + requests: [{ edgeTypes: types, limitPerType: DEFAULT_SAMPLE_SIZE }], + }); + }); + + it("should chunk a complete scan by edge type when there are too many to sample", () => { + const types = edgeTypes(10_015); + const plan = planDiscovery({ + edgeTypes: types, + totalEdges: undefined, + }); + + expect(plan.strategy).toBe("complete"); + expect(plan.requests).toHaveLength( + Math.ceil(types.length / EDGE_TYPES_PER_CHUNK), + ); + expect(coveredTypes(plan.requests)).toStrictEqual(types); + }); + }); + + describe("sampling after a complete scan was abandoned", () => { + it("should sample every edge type whatever the graph's size", () => { + const types = edgeTypes(4); + const plan = planSampling(types); + + expect(plan).toStrictEqual({ + strategy: "sampled", + requests: [{ edgeTypes: types, limitPerType: DEFAULT_SAMPLE_SIZE }], + }); + }); + }); + + describe("an edge total that is not a usable number", () => { + // `totalEdges` is cast out of the summary API response and copied verbatim + // from an imported connection file, so neither source guarantees a number. + // A request naming zero edge types reads as "no filter" downstream, which is + // the unbounded query this whole module exists to avoid. + it.each([ + ["NaN", NaN], + ["Infinity", Infinity], + ["an object", {} as unknown as number], + ["null", null as unknown as number], + ["a negative count", -1], + ["a string", "19928805" as unknown as number], + ])( + "should never plan an empty edge type filter for %s", + (_label, total) => { + const plan = planDiscovery({ + edgeTypes: edgeTypes(3), + totalEdges: total, + }); + + for (const request of plan.requests) { + expect(request.edgeTypes).not.toStrictEqual([]); + } + }, + ); + + it("should decide as though the total were unrecorded", () => { + const types = edgeTypes(3); + const unusable = planDiscovery({ + edgeTypes: types, + totalEdges: NaN, + }); + const unrecorded = planDiscovery({ + edgeTypes: types, + totalEdges: undefined, + }); + + expect(unusable).toStrictEqual(unrecorded); + }); + }); + + describe("batching a sampled pass", () => { + it("should name several edge types per request, each with its own limit", () => { + const types = edgeTypes(25); + const plan = planSampling(types); + + expect(plan.requests).toHaveLength( + Math.ceil(types.length / EDGE_TYPES_PER_SAMPLE), + ); + expect(coveredTypes(plan.requests)).toStrictEqual(types); + }); + + it("should keep the worst case of every branch reaching its limit inside the sample budget", () => { + // A type's edge count is unknown until it is read, so every branch is + // assumed full. On a small instance 100 full branches ran for 116s. + const plan = planSampling(edgeTypes(1_000)); + + for (const request of plan.requests) { + expect( + (request.edgeTypes ?? []).length * DEFAULT_SAMPLE_SIZE, + ).toBeLessThanOrEqual(SAMPLE_EDGE_BUDGET); + } + }); + + it("should never name more union branches than Gremlin Server can compile", () => { + // Gremlin Server compiles the script as Groovy, where one call takes at + // most about 250 arguments, so a wider union() fails outright. + const plan = planSampling(edgeTypes(1_000)); + + for (const request of plan.requests) { + expect((request.edgeTypes ?? []).length).toBeLessThanOrEqual( + MAX_UNION_BRANCHES, + ); + } + }); + + it("should split a sampled batch early when the edge type names are long", () => { + const types = Array.from({ length: EDGE_TYPES_PER_SAMPLE }, (_, i) => + createEdgeType(String(i).padStart(LABEL_BUDGET_CHARS / 2, "t")), + ); + const plan = planSampling(types); + + expect(plan.requests.length).toBeGreaterThan(1); + expect(coveredTypes(plan.requests)).toStrictEqual(types); + }); + }); + + describe("bounding a complete attempt", () => { + // Without a bound of its own, a complete attempt runs until the + // connection-wide fetch timeout, which defaults to four minutes. The degrade + // path only helps if the attempt gives up in seconds. + it("should bound every automatic complete scan, chunked or not", () => { + const whole = planDiscovery({ + edgeTypes: edgeTypes(3), + totalEdges: 40_000, + }); + const chunked = planDiscovery({ + edgeTypes: edgeTypes(2_010), + totalEdges: 1_015_639, + }); + + expect(whole.requestTimeoutMs).toBe(COMPLETE_ATTEMPT_TIMEOUT_MS); + expect(chunked.requestTimeoutMs).toBe(COMPLETE_ATTEMPT_TIMEOUT_MS); + }); + + it("should not bound a sampled plan, whose work is already capped per request", () => { + const plan = planDiscovery({ + edgeTypes: edgeTypes(3), + totalEdges: 19_928_805, + }); + + expect(plan.requestTimeoutMs).toBeUndefined(); + }); + + it("should bound a scan planned without an edge total, where the volume is a guess", () => { + const plan = planDiscovery({ + edgeTypes: edgeTypes(10_015), + totalEdges: undefined, + }); + + expect(plan.strategy).toBe("complete"); + expect(plan.requestTimeoutMs).toBe(COMPLETE_ATTEMPT_TIMEOUT_MS); + }); + }); + + describe("query size", () => { + function longEdgeTypes(count: number, nameLength: number): EdgeType[] { + return Array.from({ length: count }, (_, i) => + createEdgeType(`${String(i).padStart(nameLength, "t")}`), + ); + } + + it("should keep the reported request count small for a graph with very many edge types", () => { + const types = edgeTypes(10_015); + const plan = planDiscovery({ + edgeTypes: types, + totalEdges: 68_582, + }); + + // 3.2.2 issued 101 requests for this graph. The point of the change is + // that the count comes from the work, not from the edge type count. + expect(plan.requests.length).toBeLessThan(10); + expect(coveredTypes(plan.requests)).toStrictEqual(types); + }); + + it("should split a chunk further when the edge type names are long", () => { + const short = planDiscovery({ + edgeTypes: longEdgeTypes(10_000, 8), + totalEdges: 200_000, + }); + const long = planDiscovery({ + edgeTypes: longEdgeTypes(10_000, 400), + totalEdges: 200_000, + }); + + // Same edge count and same edge type count, so the volume plan is + // identical. Only the rendered query text differs. + expect(long.requests.length).toBeGreaterThan(short.requests.length); + }); + + it("should hold every chunk inside the label budget", () => { + const types = longEdgeTypes(5_000, 300); + const plan = planDiscovery({ + edgeTypes: types, + totalEdges: 500_000, + }); + + for (const request of plan.requests) { + const rendered = (request.edgeTypes ?? []).reduce( + (total, type) => total + type.length, + 0, + ); + expect(rendered).toBeLessThanOrEqual(LABEL_BUDGET_CHARS); + } + expect(coveredTypes(plan.requests)).toStrictEqual(types); + }); + }); +}); diff --git a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryPlan.ts b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryPlan.ts new file mode 100644 index 0000000000..416a78dc47 --- /dev/null +++ b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/discoveryPlan.ts @@ -0,0 +1,278 @@ +import type { EdgeType } from "@/core"; + +import { DEFAULT_SAMPLE_SIZE } from "@/utils"; + +/** + * Edges a single request is allowed to scan. + * + * The highest volume that never failed on any configuration we measured, not a + * prediction of what is safe. The real ceiling moved by more than 4x between two + * instances of the same hardware class, failing on memory at 150,000 edges with + * Neptune's DFE engine and on time at 1,000,000 without it, and neither DFE + * presence nor instance memory is visible from here. The degrade path, not this + * number, is what gets an answer back when the guess is wrong. See the ADR. + */ +export const SCAN_BUDGET = 50_000; + +/** + * Edge types per request when chunking a graph whose edge total is unknown. + * + * Shares the value of `DEFAULT_BATCH_REQUEST_SIZE` and nothing else. That one is + * a fan-out width for batching independent queries; this one is a guess at how + * much of an unmeasured graph one scan can carry, and it moves with the cost + * model in the ADR. Deduplicating them would tie two unrelated decisions + * together. + */ +export const EDGE_TYPES_PER_CHUNK = 100; + +/** + * Characters of edge type names one request may name in its filter. + * + * The edge total says how many edges a request scans, not how much query text it + * carries, so a graph with thousands of edge types needs its own bound. The + * largest filter measured live was about 55,000 characters, 5,008 names of 11, + * which completed in 7.1s on Neptune 1.4.7.0. Like the scan budget this is the + * largest size known to work rather than the point where it breaks. + */ +export const LABEL_BUDGET_CHARS = 60_000; + +/** + * Edges one sampled request may read, counting every edge type it names as + * though it reached the per-type limit. + * + * A branch that reaches the limit costs about a second on a small instance + * whether it shares a request or not, so batching saves round trips, not reads. + * Measured on Neptune 1.4.5.1 (db.t3.medium): 10 full branches took 9.7s, 100 + * took 116s, just under the query timeout, and left the instance refusing even a + * single-type sample on memory for two minutes afterwards. The worst case is + * assumed because a type's edge count is not known before it is read. + * + * It is also the most edges one request holds in memory at once, because the + * sampled query groups by edge type and keeps each type's sample until it is + * counted. The budget is per request, and requests run four at a time: four + * concurrent requests of 10 full branches each finished in 36s on both + * db.t3.medium instances measured, with and without DFE, and a single-type + * sample sent straight afterwards succeeded. + */ +export const SAMPLE_EDGE_BUDGET = 100_000; + +/** + * Most `union()` branches one sampled request may name. Gremlin Server compiles + * the script as Groovy, where a call takes at most about 250 arguments, and DFE + * crashed a db.t3.medium on 200 branches. + */ +export const MAX_UNION_BRANCHES = 100; + +/** Edge types per sampled request, derived so the worst case fits the sample budget. */ +export const EDGE_TYPES_PER_SAMPLE = Math.min( + MAX_UNION_BRANCHES, + Math.max(1, Math.floor(SAMPLE_EDGE_BUDGET / DEFAULT_SAMPLE_SIZE)), +); + +/** Quotes and the separator each name costs on top of its own characters. */ +const LABEL_OVERHEAD_CHARS = 4; + +/** + * Measured cost of sampling one edge type. Range seen across four engines: 0.9s + * to 1.8s. Still charged per type, not per batched request, because a branch + * that reaches the limit costs as much as the separate request it replaced. + */ +const PER_REQUEST_MS = 1_500; + +/** Measured cost of scanning one edge, in microseconds. Range seen: 74 to 115. */ +const PER_EDGE_SCAN_US = 85; + +/** How long we are willing to spend when the edge total is unknown and we must guess. */ +const UNKNOWN_TOTAL_BUDGET_MS = 60_000; + +/** + * How long one request of a complete scan may run before it is abandoned. + * + * Reading the whole scan budget costs 4.3s at the measured per-edge rate, and + * 5.8s at the slowest rate seen, so this only fires on a request reading far more + * than the budget asked for. Chunks are balanced by edge type count rather than + * by edge count, so one chunk holding a dominant edge type is exactly that case. + * + * Below the 30 to 35 seconds at which Neptune failed in both reproductions, on + * purpose: if the database's own error arrives first the bound has done nothing. + */ +export const COMPLETE_ATTEMPT_TIMEOUT_MS = 20_000; + +/** + * Complete reads every edge and reports exact counts; sampled caps the edges it + * reads per edge type and can miss a connection that occurs rarely. Chunking is + * not a third strategy, it is complete split so no single request is too large. + */ +export type DiscoveryStrategy = "none" | "complete" | "sampled"; + +/** Reads every matching edge. */ +export type ScanRequest = { + /** Absent scans every edge type, which is cheaper than naming them all. */ + edgeTypes?: EdgeType[]; +}; + +/** Reads at most `limitPerType` edges of each named edge type, in one request. */ +export type SampleRequest = { + edgeTypes: EdgeType[]; + limitPerType: number; +}; + +export type DiscoveryRequest = ScanRequest | SampleRequest; + +export type DiscoveryPlan = { + strategy: DiscoveryStrategy; + requests: DiscoveryRequest[]; + /** + * How long one request may run before it is abandoned, when the plan can + * predict its cost. Absent means the connection's own fetch timeout is the + * only bound. + */ + requestTimeoutMs?: number; +}; + +/** + * Chooses how to discover edge connections and lays out the requests it takes. + * + * Complete costs grow with the number of edges, sampled costs with the number of + * edge types, so above the scan budget we take whichever quantity is smaller. + * Pure, so the choice is testable without a database. + * + * @param totalEdges The graph's edge count, absent when the cached schema + * predates our recording it. Absence costs extra requests, never correctness. + */ +export function planDiscovery({ + edgeTypes, + totalEdges, +}: { + edgeTypes: EdgeType[]; + totalEdges: number | undefined; +}): DiscoveryPlan { + if (edgeTypes.length === 0) { + return { strategy: "none", requests: [] }; + } + + const edgeTotal = toEdgeTotal(totalEdges); + + if (shouldSample(edgeTypes.length, edgeTotal)) { + return planSampling(edgeTypes); + } + + return { + strategy: "complete", + requests: chunkForCompleteScan(edgeTypes, edgeTotal), + requestTimeoutMs: COMPLETE_ATTEMPT_TIMEOUT_MS, + }; +} + +/** + * Lays out a sampled pass over every edge type, whatever the graph's size. The + * degrade path uses it directly once a complete scan proved too large. + */ +export function planSampling(edgeTypes: EdgeType[]): DiscoveryPlan { + return { + strategy: "sampled", + requests: chunkEdgeTypes(edgeTypes, EDGE_TYPES_PER_SAMPLE).map(chunk => ({ + edgeTypes: chunk, + limitPerType: DEFAULT_SAMPLE_SIZE, + })), + }; +} + +/** + * Accepts an edge total only when it can actually size work, and treats anything + * else as unrecorded. + * + * The declared type says `number | undefined`, but the value is cast out of the + * summary API response and copied verbatim out of an imported connection file, so + * neither source guarantees one. Arithmetic on a non-number yields `NaN`, which + * compares false against every threshold and would send the planner down the + * complete path with a chunk count it cannot use. + */ +export function toEdgeTotal(value: number | undefined): number | undefined { + return typeof value === "number" && Number.isFinite(value) && value >= 0 + ? value + : undefined; +} + +/** Whether sampling is the cheaper way to cover this graph. */ +function shouldSample( + edgeTypeCount: number, + totalEdges: number | undefined, +): boolean { + if (totalEdges === undefined) { + // Decide on edge type count alone, preferring whichever stays bounded + // whatever the graph turns out to be. + return edgeTypeCount * PER_REQUEST_MS <= UNKNOWN_TOTAL_BUDGET_MS; + } + if (totalEdges <= SCAN_BUDGET) { + return false; + } + const sampledMs = edgeTypeCount * PER_REQUEST_MS; + const scanMs = (totalEdges * PER_EDGE_SCAN_US) / 1_000; + return sampledMs <= scanMs; +} + +/** + * Splits a complete scan into requests that stay inside two independent bounds: + * the edges one request may scan, and the characters its filter may name. + * + * A scan that fits in one request drops the filter entirely, so its query text is + * a constant regardless of how many edge types the graph has. Beyond that every + * request has to name its types, and the two bounds constrain different things: + * the edge total says how much a request reads, the character budget says how + * much it carries. Chunking cannot go finer than one edge type per request, so a + * graph far above the scan budget with few edge types still gets chunks larger + * than the budget asks for, which is the case the degrade path exists for. + */ +function chunkForCompleteScan( + edgeTypes: EdgeType[], + totalEdges: number | undefined, +): ScanRequest[] { + const wanted = + totalEdges === undefined + ? Math.ceil(edgeTypes.length / EDGE_TYPES_PER_CHUNK) + : Math.ceil(totalEdges / SCAN_BUDGET); + const chunkCount = Math.min(Math.max(wanted, 1), edgeTypes.length); + + if (chunkCount === 1) { + return [{}]; + } + + return chunkEdgeTypes( + edgeTypes, + Math.ceil(edgeTypes.length / chunkCount), + ).map(chunk => ({ edgeTypes: chunk })); +} + +/** + * Splits edge types into chunks of at most `maxPerChunk`, splitting early when a + * chunk's names would overrun the label budget. + */ +function chunkEdgeTypes( + edgeTypes: EdgeType[], + maxPerChunk: number, +): EdgeType[][] { + const chunks: EdgeType[][] = []; + let chunk: EdgeType[] = []; + let chars = 0; + + for (const edgeType of edgeTypes) { + const cost = edgeType.length + LABEL_OVERHEAD_CHARS; + const full = + chunk.length >= maxPerChunk || chars + cost > LABEL_BUDGET_CHARS; + if (chunk.length > 0 && full) { + chunks.push(chunk); + chunk = []; + chars = 0; + } + chunk.push(edgeType); + chars += cost; + } + // Only ever pushed non-empty, because a request naming zero edge types reads + // downstream as no filter at all, which is the unbounded scan. + if (chunk.length > 0) { + chunks.push(chunk); + } + + return chunks; +} diff --git a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/edgeConnectionsTemplate.test.ts b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/edgeConnectionsTemplate.test.ts index c2cf08f99c..0af44a6189 100644 --- a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/edgeConnectionsTemplate.test.ts +++ b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/edgeConnectionsTemplate.test.ts @@ -4,9 +4,35 @@ import { normalizeWithNoSpace as normalize } from "@/utils/testing"; import edgeConnectionsTemplate from "./edgeConnectionsTemplate"; describe("Gremlin > edgeConnectionsTemplate", () => { + it("should count the distinct edge type and endpoint label combinations over every edge", () => { + const template = edgeConnectionsTemplate({}); + + expect(normalize(template)).toBe( + normalize(` + g.E() + .groupCount() + .by( + project('e', 's', 't') + .by(label()) + .by(outV().label().fold()) + .by(inV().label().fold()) + ) + `), + ); + }); + + it("should filter to the given edge types", () => { + const template = edgeConnectionsTemplate({ + edgeTypes: [createEdgeType("route"), createEdgeType("contains")], + }); + + expect(template).toContain("g.E().hasLabel('route', 'contains')"); + }); + it("should give each sampled edge type its own limited branch", () => { const template = edgeConnectionsTemplate({ - types: [createEdgeType("route"), createEdgeType("contains")], + edgeTypes: [createEdgeType("route"), createEdgeType("contains")], + limitPerType: 10000, }); // A single limit after hasLabel() would be shared, and a dominant type would @@ -29,11 +55,20 @@ describe("Gremlin > edgeConnectionsTemplate", () => { ); }); - it("should escape special characters in the edge type", () => { + it("should escape special characters in a sampled edge type", () => { const template = edgeConnectionsTemplate({ - types: [createEdgeType("edge'with'quotes")], + edgeTypes: [createEdgeType("edge'with'quotes")], + limitPerType: 10000, }); expect(template).toContain("V().outE('edge\\'with\\'quotes')"); }); + + it("should escape special characters in the edge type", () => { + const template = edgeConnectionsTemplate({ + edgeTypes: [createEdgeType("edge'with'quotes")], + }); + + expect(template).toContain("hasLabel('edge\\'with\\'quotes')"); + }); }); diff --git a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/edgeConnectionsTemplate.ts b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/edgeConnectionsTemplate.ts index 8be796f378..f2ece4abea 100644 --- a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/edgeConnectionsTemplate.ts +++ b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/edgeConnectionsTemplate.ts @@ -1,51 +1,95 @@ -import type { EdgeType } from "@/core"; +import { query } from "@/utils"; -import { DEFAULT_SAMPLE_SIZE, query } from "@/utils"; +import type { + DiscoveryRequest, + SampleRequest, + ScanRequest, +} from "./discoveryPlan"; import { fragment } from "../fragments"; /** - * Keys of the projected endpoint labels. The template writes them and the - * response parser reads them, so both derive from this object. + * Keys of the projected triple that identifies one edge connection. + * + * Single letters because they repeat once per distinct combination, and a graph + * with ten thousand edge types returns ten thousand of them. The template writes + * them and the response parser reads them, so both derive from this object. */ export const projectionKeys = { + edgeType: "e", sourceType: "s", targetType: "t", } as const; /** - * Returns a Gremlin query that counts the endpoint label combinations of up to - * `DEFAULT_SAMPLE_SIZE` edges of each given type, grouped by edge type. - * - * One limit after `hasLabel(A, B, ...)` would be shared, so a dominant type fills - * it and the rest come back empty. Each type gets its own `union()` branch - * instead, and on Neptune each branch is an index lookup by edge label. - * Mid-traversal `V()` rather than `E()`, which needs TinkerPop 3.7, and anchored - * on `V().limit(1)` because anchoring on `inject()` is not native on Neptune. + * Returns a Gremlin query that counts the distinct + * `(edge type, source label, target label)` combinations in one request. * - * Grouped by edge type before counting because Neptune's DFE engine cannot count - * one `project()` key across several full branches: two took 54s and five timed - * out, where grouping first handled ten in 9s. + * `groupCount()` keyed by a `project()` is native on every Neptune engine we + * tested and on reference TinkerPop 3.6.2, and its accumulator is keyed by the + * answer rather than the input, so it holds one entry per distinct combination + * instead of one per edge. That is what makes it survive a graph the previous + * `group().by(label())` shape ran out of memory on. * - * Callers send 10 types per request, so one request reads at most 100,000 - * edges. On a db.t3.medium, 100 types in one request took 116s and left the - * instance refusing even a single-type sample for two minutes afterwards. + * The key must be a named `project()`. A `union()` of the three labels is also + * native but does not guarantee order, and Neptune's DFE engine permuted it, + * silently reporting edges in the wrong direction. See the ADR. * * The endpoint labels are folded because engines disagree on what `label()` * emits for a multi-label vertex. Neptune 1.4 emits one `::` composite, but * 1.3.5 emits each label separately, and a bare `by(outV().label())` keeps only * the first, silently dropping the vertex's other types. + * + * A sampled request is shaped differently. See {@link sampledEdgesTemplate}. + */ +export default function edgeConnectionsTemplate(request: DiscoveryRequest) { + return "limitPerType" in request + ? sampledEdgesTemplate(request) + : scannedEdgesTemplate(request); +} + +function scannedEdgesTemplate({ edgeTypes }: ScanRequest) { + const edges = edgeTypes?.length + ? `g.E().hasLabel(${edgeTypes.map(fragment.identifier).join(", ")})` + : "g.E()"; + const keys = Object.values(projectionKeys).map(fragment.identifier); + + return query` + ${edges} + .groupCount() + .by( + project(${keys.join(", ")}) + .by(label()) + .by(outV().label().fold()) + .by(inV().label().fold()) + ) + `; +} + +/** + * Counts the endpoint label combinations of up to `limitPerType` edges of each + * named type, grouped by edge type. + * + * One limit after `hasLabel(A, B, ...)` would be shared, so a dominant type fills + * it and the rest come back empty. Each type gets its own `union()` branch + * instead, and on Neptune each branch is an index lookup by edge label. + * Mid-traversal `V()` rather than `E()`, which needs TinkerPop 3.7, and anchored + * on `V().limit(1)` because anchoring on `inject()` is not native on Neptune. + * + * Grouped by edge type first, unlike the scan, because Neptune's DFE engine + * cannot count one `project()` key across several full branches: two took 54s + * and five timed out, where grouping first handled ten in 9s. Grouping holds + * each type's sample until it is counted, which the per-type limit bounds, and + * is why the scan, which has no limit, cannot use it. */ -export default function edgeConnectionsTemplate({ - types, -}: { - types: EdgeType[]; -}) { - const limit = fragment.number(DEFAULT_SAMPLE_SIZE); - const branches = types.map( +function sampledEdgesTemplate({ edgeTypes, limitPerType }: SampleRequest) { + const limit = fragment.number(limitPerType); + const branches = edgeTypes.map( type => `V().outE(${fragment.identifier(type)}).limit(${limit})`, ); - const keys = Object.values(projectionKeys).map(fragment.identifier); + const keys = [projectionKeys.sourceType, projectionKeys.targetType].map( + fragment.identifier, + ); return query` g.V().limit(1).union(${branches.join(", ")}) diff --git a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/index.test.ts b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/index.test.ts index 6745aaaf6e..a0b1ba61e4 100644 --- a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/index.test.ts +++ b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/index.test.ts @@ -1,6 +1,12 @@ -import { vi } from "vitest"; +import { type Mock, vi } from "vitest"; import { createEdgeType, createVertexType, type EdgeType } from "@/core"; +import { + DatabaseTimeoutError, + FetchTimeoutError, + logger, + NetworkError, +} from "@/utils"; import { createGInt64, createGList, @@ -9,33 +15,43 @@ import { } from "@/utils/testing"; import fetchEdgeConnections from "."; +import { EdgeConnectionDiscoveryError } from "./discoveryError"; +import { EDGE_TYPES_PER_SAMPLE } from "./discoveryPlan"; -/** One distinct `(edge type, source labels, target labels)` combination. */ -type Combination = [ - edgeType: string, - sourceTypes: string[], - targetTypes: string[], -]; - -/** Builds the `group().by(label())` response: edge type to `(s, t)` counts. */ -function sampleResponse(...combinations: Combination[]) { - const byEdgeType = new Map(); - for (const combination of combinations) { - const [edgeType] = combination; - byEdgeType.set(edgeType, [ - ...(byEdgeType.get(edgeType) ?? []), - combination, - ]); +/** One distinct `(edge type, source label, target label)` combination. */ +type Triple = [edgeType: string, sourceType: string, targetType: string]; + +/** One projected key, with each endpoint's labels folded into a list as the template asks. */ +function tripleKey(e: string, s: string[], t: string[]) { + return createGMap({ e, s: createGList(s), t: createGList(t) }); +} + +/** Builds the flat `groupCount().by(project(...))` response: one g:Map of triple to count. */ +function countResponse(...triples: Triple[]) { + return createGremlinResponse( + createGMap( + new Map( + triples.map(([e, s, t]) => [tripleKey(e, [s], [t]), createGInt64(1)]), + ), + ), + ); +} + +/** Builds the `group().by(label())` response a sampled request returns. */ +function sampleResponse(...triples: Triple[]) { + const byEdgeType = new Map(); + for (const triple of triples) { + byEdgeType.set(triple[0], [...(byEdgeType.get(triple[0]) ?? []), triple]); } return createGremlinResponse( createGMap( new Map( - [...byEdgeType].map(([edgeType, ofType]) => [ - edgeType, + [...byEdgeType].map(([e, ofType]) => [ + e, createGMap( new Map( ofType.map(([, s, t]) => [ - createGMap({ s: createGList(s), t: createGList(t) }), + createGMap({ s: createGList([s]), t: createGList([t]) }), createGInt64(1), ]), ), @@ -46,33 +62,78 @@ function sampleResponse(...combinations: Combination[]) { ); } -/** A `group()` over no edges returns an empty map. */ +/** A `groupCount()` over a graph with no matching edges returns an empty map. */ const emptyResponse = createGremlinResponse(createGMap({})); +function memoryLimitError() { + return new NetworkError("Query cannot be completed", 500, { + code: "MemoryLimitExceededException", + }); +} + +function databaseTimeoutError() { + return new DatabaseTimeoutError( + "Query cannot be completed", + 500, + {}, + "TimeLimitExceededException", + ); +} + +function fetchTimeoutError() { + return new FetchTimeoutError( + 240_000, + new DOMException("Aborted", "TimeoutError"), + ); +} + function edgeTypes(count: number): EdgeType[] { return Array.from({ length: count }, (_, i) => createEdgeType(`edge${i}`)); } +/** Every request carries its own id and signal, so one can be cancelled alone. */ +const perRequest = { + queryId: expect.any(String), + signal: expect.any(AbortSignal), +}; + +/** Awaits a discovery that is expected to give up, and returns why. */ +async function discoveryErrorFrom( + discovery: Promise, +): Promise { + const error = await discovery.then( + () => undefined, + (thrown: unknown) => thrown, + ); + expect(error).toBeInstanceOf(EdgeConnectionDiscoveryError); + return error as EdgeConnectionDiscoveryError; +} + +/** The signal handed to the nth request, which the caller aborts when it gives up. */ +function signalOfCall(gremlinFetch: Mock, call: number): AbortSignal { + return gremlinFetch.mock.calls[call][1].signal; +} + describe("Gremlin > fetchEdgeConnections", () => { - it("should sample several edge types in one request and regroup by edge type", async () => { + it("should ask for the distinct combinations in one request when the graph fits the budget", async () => { const gremlinFetch = vi .fn() .mockResolvedValueOnce( - sampleResponse( - ["route", ["airport"], ["airport"]], - ["contains", ["country"], ["airport"]], + countResponse( + ["route", "airport", "airport"], + ["contains", "country", "airport"], ), ); const result = await fetchEdgeConnections(gremlinFetch, { edgeTypes: [createEdgeType("route"), createEdgeType("contains")], + totalEdges: 5_000, }); expect(gremlinFetch).toHaveBeenCalledTimes(1); expect(gremlinFetch).toHaveBeenCalledWith( - expect.stringContaining( - "V().outE('route').limit(10000), V().outE('contains').limit(10000)", - ), + expect.stringContaining("g.E()\n .groupCount()"), + perRequest, ); expect(result).toStrictEqual({ edgeConnections: [ @@ -90,24 +151,32 @@ describe("Gremlin > fetchEdgeConnections", () => { }); }); - it("should send 10 edge types per request", async () => { - const gremlinFetch = vi.fn().mockResolvedValue(emptyResponse); - const types = edgeTypes(25); - - await fetchEdgeConnections(gremlinFetch, { edgeTypes: types }); - - const queries = gremlinFetch.mock.calls.map(call => call[0] as string); - expect( - queries.map(q => (q.match(/outE\('edge\d+'\)/g) ?? []).length), - ).toStrictEqual([10, 10, 5]); - // Every input type is covered across the requests - const all = queries.join("\n"); - for (const type of types) { - expect(all).toContain(`outE('${type}')`); - } + it("should sample several edge types in one request when the graph is too large to scan", async () => { + const gremlinFetch = vi + .fn() + .mockResolvedValueOnce( + sampleResponse( + ["route", "airport", "airport"], + ["contains", "country", "airport"], + ), + ); + + const result = await fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route"), createEdgeType("contains")], + totalEdges: 19_928_805, + }); + + expect(gremlinFetch).toHaveBeenCalledTimes(1); + expect(gremlinFetch).toHaveBeenCalledWith( + expect.stringContaining( + "V().outE('route').limit(10000), V().outE('contains').limit(10000)", + ), + perRequest, + ); + expect(result.edgeConnections).toHaveLength(2); }); - it("should return empty array when no edge types provided", async () => { + it("should return nothing without querying when the schema has no edge types", async () => { const gremlinFetch = vi.fn(); const result = await fetchEdgeConnections(gremlinFetch, { edgeTypes: [] }); @@ -116,11 +185,12 @@ describe("Gremlin > fetchEdgeConnections", () => { expect(result).toStrictEqual({ edgeConnections: [] }); }); - it("should return empty array when no edge connections exist", async () => { + it("should return nothing when the graph has no edge connections", async () => { const gremlinFetch = vi.fn().mockResolvedValue(emptyResponse); const result = await fetchEdgeConnections(gremlinFetch, { edgeTypes: [createEdgeType("route")], + totalEdges: 0, }); expect(result).toStrictEqual({ edgeConnections: [] }); @@ -129,11 +199,12 @@ describe("Gremlin > fetchEdgeConnections", () => { it("should deduplicate combinations returned by more than one request", async () => { const gremlinFetch = vi .fn() - .mockResolvedValue(sampleResponse(["route", ["airport"], ["airport"]])); + .mockResolvedValue(sampleResponse(["route", "airport", "airport"])); const result = await fetchEdgeConnections(gremlinFetch, { - // One more than a request carries, so it takes two. - edgeTypes: [createEdgeType("route"), ...edgeTypes(10)], + // One more than a sampled request carries, so the pass takes two. + edgeTypes: [createEdgeType("route"), ...edgeTypes(EDGE_TYPES_PER_SAMPLE)], + totalEdges: 19_928_805, }); expect(gremlinFetch).toHaveBeenCalledTimes(2); @@ -148,29 +219,16 @@ describe("Gremlin > fetchEdgeConnections", () => { }); }); - it("should propagate errors from fetch", async () => { - const gremlinFetch = vi.fn().mockRejectedValue(new Error("Network error")); - - await expect( - fetchEdgeConnections(gremlinFetch, { - edgeTypes: [createEdgeType("route")], - }), - ).rejects.toThrow(new Error("Network error")); - }); - - it("should expand Neptune multi-label composites on both ends", async () => { + it("should expand Neptune multi-label endpoints on both ends", async () => { const gremlinFetch = vi .fn() .mockResolvedValueOnce( - sampleResponse([ - "worksAt", - ["Person::Employee"], - ["Company::Organization"], - ]), + countResponse(["worksAt", "Person::Employee", "Company::Organization"]), ); const result = await fetchEdgeConnections(gremlinFetch, { edgeTypes: [createEdgeType("worksAt")], + totalEdges: 10, }); expect(result).toStrictEqual({ @@ -206,11 +264,21 @@ describe("Gremlin > fetchEdgeConnections", () => { const gremlinFetch = vi .fn() .mockResolvedValueOnce( - sampleResponse(["worksAt", ["Person", "Employee"], ["Company"]]), + createGremlinResponse( + createGMap( + new Map([ + [ + tripleKey("worksAt", ["Person", "Employee"], ["Company"]), + createGInt64(1), + ], + ]), + ), + ), ); const result = await fetchEdgeConnections(gremlinFetch, { edgeTypes: [createEdgeType("worksAt")], + totalEdges: 10, }); expect(result.edgeConnections).toStrictEqual([ @@ -227,20 +295,19 @@ describe("Gremlin > fetchEdgeConnections", () => { ]); }); - it("should read the projected labels by key, whatever order the keys arrive in", async () => { + it("should keep every label of a multi-label endpoint when sampling", async () => { const gremlinFetch = vi.fn().mockResolvedValueOnce( createGremlinResponse( createGMap( new Map([ [ - "contains", + "worksAt", createGMap( new Map([ [ - // target-then-source key order createGMap({ - t: createGList(["airport"]), - s: createGList(["country"]), + s: createGList(["Person", "Employee"]), + t: createGList(["Company"]), }), createGInt64(1), ], @@ -252,8 +319,67 @@ describe("Gremlin > fetchEdgeConnections", () => { ), ); + const result = await fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("worksAt")], + totalEdges: 19_928_805, + }); + + expect(result.edgeConnections.map(c => c.sourceVertexType)).toStrictEqual([ + createVertexType("Person"), + createVertexType("Employee"), + ]); + }); + + it("should ignore edge types that are not in the schema", async () => { + // An unfiltered scan sees every edge type in the graph, including ones the + // schema does not know about and the app therefore cannot render. + const gremlinFetch = vi + .fn() + .mockResolvedValueOnce( + countResponse( + ["route", "airport", "airport"], + ["undiscovered", "airport", "airport"], + ), + ); + + const result = await fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 10, + }); + + expect(result.edgeConnections).toStrictEqual([ + { + sourceVertexType: createVertexType("airport"), + edgeType: createEdgeType("route"), + targetVertexType: createVertexType("airport"), + }, + ]); + }); + + it("should read the projected triple by key, whatever order the keys arrive in", async () => { + // The whole reason the key is a named project() rather than a union() is that + // Neptune's DFE engine permutes an unnamed key and silently inverts the edge + // direction. Reading by name is what makes the shape safe, so pin it. + const gremlinFetch = vi.fn().mockResolvedValueOnce( + createGremlinResponse( + createGMap( + new Map([ + [ + createGMap({ + t: createGList(["airport"]), + e: "contains", + s: createGList(["country"]), + }), + createGInt64(1), + ], + ]), + ), + ), + ); + const result = await fetchEdgeConnections(gremlinFetch, { edgeTypes: [createEdgeType("contains")], + totalEdges: 10, }); expect(result.edgeConnections).toStrictEqual([ @@ -269,29 +395,15 @@ describe("Gremlin > fetchEdgeConnections", () => { const gremlinFetch = vi.fn().mockResolvedValueOnce( createGremlinResponse( createGMap( - new Map([ + new Map< + ReturnType, + ReturnType + >([ [ - "route", - createGMap( - new Map([ - [ - createGMap({ s: createGList(["airport"]) }), - createGInt64(1), - ], - [ - createGMap({ t: createGList(["airport"]) }), - createGInt64(1), - ], - [ - createGMap({ - s: createGList(["airport"]), - t: createGList(["airport"]), - }), - createGInt64(1), - ], - ]), - ), + createGMap({ e: "route", s: createGList(["airport"]) }), + createGInt64(1), ], + [tripleKey("route", ["airport"], ["airport"]), createGInt64(1)], ]), ), ), @@ -299,6 +411,7 @@ describe("Gremlin > fetchEdgeConnections", () => { const result = await fetchEdgeConnections(gremlinFetch, { edgeTypes: [createEdgeType("route")], + totalEdges: 10, }); expect(result.edgeConnections).toStrictEqual([ @@ -309,4 +422,213 @@ describe("Gremlin > fetchEdgeConnections", () => { }, ]); }); + + describe("degrading a complete scan that was too large", () => { + it.each([ + ["a memory limit", memoryLimitError], + ["a database timeout", databaseTimeoutError], + ["our own fetch timeout", fetchTimeoutError], + ])( + "should redo discovery as sampled after %s", + async (_label, makeError) => { + const gremlinFetch = vi + .fn() + .mockRejectedValueOnce(makeError()) + .mockResolvedValue(sampleResponse(["route", "airport", "airport"])); + + const result = await fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 10, + }); + + expect(gremlinFetch).toHaveBeenCalledTimes(2); + expect(gremlinFetch).toHaveBeenLastCalledWith( + expect.stringContaining("V().outE('route').limit(10000)"), + perRequest, + ); + expect(result.edgeConnections).toHaveLength(1); + }, + ); + + it("should abandon the remaining chunks rather than finish them", async () => { + const gremlinFetch = vi + .fn() + .mockRejectedValueOnce(memoryLimitError()) + .mockResolvedValue(emptyResponse); + + const types = edgeTypes(500); + await fetchEdgeConnections(gremlinFetch, { + edgeTypes: types, + totalEdges: 5_000_000, + }); + + // 100 complete chunks were planned. The failure stops the pool, so only the + // requests already in flight run before the sampled requests. + const sampledRequests = Math.ceil(types.length / EDGE_TYPES_PER_SAMPLE); + expect(gremlinFetch.mock.calls.length).toBeLessThan( + sampledRequests + 100, + ); + expect(gremlinFetch.mock.calls.length).toBeGreaterThanOrEqual( + sampledRequests + 1, + ); + }); + + it("should point at the Fetch Timeout, not the parameter group, when a sampled pass exhausts our own fetch timeout", async () => { + const gremlinFetch = vi.fn().mockRejectedValue(fetchTimeoutError()); + + const error = await discoveryErrorFrom( + fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 19_928_805, + }), + ); + + expect(error.recovery).toContain("Fetch Timeout"); + expect(error.recovery).not.toContain("parameter group"); + expect(error.details).toMatchObject({ failureCause: "fetch-timeout" }); + }); + + it("should point at the database's own query timeout when a sampled pass exhausts it", async () => { + const gremlinFetch = vi.fn().mockRejectedValue(databaseTimeoutError()); + + const error = await discoveryErrorFrom( + fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 19_928_805, + }), + ); + + expect(error.recovery).toContain("DB cluster parameter group"); + expect(error.details).toMatchObject({ failureCause: "database-limit" }); + }); + + it("should report an unusable edge total as unrecorded, like the planner does", async () => { + const gremlinFetch = vi.fn().mockRejectedValue(memoryLimitError()); + + const error = await discoveryErrorFrom( + fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: "19928805" as unknown as number, + }), + ); + + // Showing the raw value would have the error details disagree with the + // plan, which ignored it. + expect(error.details.totalEdges).toBeUndefined(); + }); + + it("should take the fast degrade path when the code arrives nested in a cause", async () => { + const gremlinFetch = vi + .fn() + .mockRejectedValueOnce( + new NetworkError("Query cannot be completed", 500, { + cause: { code: "MemoryLimitExceededException" }, + }), + ) + .mockResolvedValue(sampleResponse(["route", "airport", "airport"])); + + const result = await fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 10, + }); + + expect(result.edgeConnections).toHaveLength(1); + }); + + it("should not degrade a sampled pass, because there is nothing cheaper to try", async () => { + const gremlinFetch = vi.fn().mockRejectedValue(memoryLimitError()); + + await expect( + fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 19_928_805, + }), + ).rejects.toThrow(/could not sample/); + expect(gremlinFetch).toHaveBeenCalledTimes(1); + }); + + it("should say both strategies were tried when sampling fails after degrading", async () => { + const gremlinFetch = vi.fn().mockRejectedValue(memoryLimitError()); + + const error = await discoveryErrorFrom( + fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 10, + }), + ); + + expect(error.details).toMatchObject({ + strategy: "sampled", + completeScanAbandoned: true, + }); + expect(error.message).toContain("sampling each edge type failed as well"); + }); + + it("should record the degrade at warn level, where a user will see it without dev tools", async () => { + const gremlinFetch = vi + .fn() + .mockRejectedValueOnce(memoryLimitError()) + .mockResolvedValue(sampleResponse(["route", "airport", "airport"])); + + await fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 10, + }); + + // Discovery succeeded, so nothing else tells the user the schema they are + // looking at came from a sample rather than a full scan. + expect(logger.warn).toHaveBeenCalledWith( + expect.stringContaining("sampling each edge type instead"), + expect.anything(), + ); + expect(logger.error).not.toHaveBeenCalled(); + }); + + it("should cancel a request still in flight when it abandons the attempt", async () => { + const gremlinFetch: Mock = vi.fn().mockImplementation((query: string) => + query.includes("limit(10000)") || gremlinFetch.mock.calls.length === 1 + ? Promise.reject(memoryLimitError()) + : // Never settles, so this chunk is still in flight when the first + // one fails and the whole complete attempt is abandoned. + new Promise(() => {}), + ); + + await expect( + fetchEdgeConnections(gremlinFetch, { + edgeTypes: edgeTypes(500), + totalEdges: 5_000_000, + }), + ).rejects.toThrow(/could not discover edge connections either way/); + + expect(signalOfCall(gremlinFetch, 1).aborted).toBe(true); + }); + + it("should propagate an error that is not about the query being too large", async () => { + const gremlinFetch = vi + .fn() + .mockRejectedValue(new Error("Network error")); + + await expect( + fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 10, + }), + ).rejects.toThrow("Network error"); + expect(gremlinFetch).toHaveBeenCalledTimes(1); + }); + + it("should propagate a cancellation rather than treating it as a size problem", async () => { + const gremlinFetch = vi + .fn() + .mockRejectedValue(new DOMException("Aborted", "AbortError")); + + await expect( + fetchEdgeConnections(gremlinFetch, { + edgeTypes: [createEdgeType("route")], + totalEdges: 10, + }), + ).rejects.toThrow("Aborted"); + expect(gremlinFetch).toHaveBeenCalledTimes(1); + }); + }); }); diff --git a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/index.ts b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/index.ts index 28f54ec5b6..7c4df3af11 100644 --- a/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/index.ts +++ b/packages/graph-explorer/src/connector/gremlin/fetchEdgeConnections/index.ts @@ -1,4 +1,4 @@ -import { chunk } from "lodash"; +import { v4 } from "uuid"; import type { EdgeConnectionsRequest, @@ -10,66 +10,215 @@ import { createEdgeType, createVertexType, type EdgeConnection, + type EdgeType, } from "@/core"; -import { DEFAULT_CONCURRENT_REQUESTS_LIMIT, mapWithConcurrency } from "@/utils"; +import { + DEFAULT_CONCURRENT_REQUESTS_LIMIT, + logger, + mapWithConcurrency, +} from "@/utils"; import type { GInt64, GMapWithValue, GremlinFetch } from "../types"; +import type { FailedDiscovery } from "./discoveryError"; +import type { DiscoveryPlan } from "./discoveryPlan"; +import { anySignal } from "../../utils/anySignal"; import { parseGMap } from "../mappers/parseGMap"; import { splitLabel } from "../splitLabel"; +import { + causeOf, + EdgeConnectionDiscoveryError, + isTooBig, +} from "./discoveryError"; +import { planDiscovery, planSampling, toEdgeTotal } from "./discoveryPlan"; import edgeConnectionsTemplate, { projectionKeys, } from "./edgeConnectionsTemplate"; -/** Edge types per request. See {@link edgeConnectionsTemplate}. */ -const EDGE_TYPES_PER_SAMPLE = 10; - /** Every label of one endpoint vertex, folded by the template. */ type EndpointLabels = { "@type": "g:List"; "@value": string[] }; /** The projected labels that key one `groupCount()` entry. */ -type ProjectedLabels = GMapWithValue; +type ProjectedLabels = GMapWithValue; /** Endpoint label combinations counted by `groupCount()`. */ type CountedLabels = GMapWithValue; -/** One map from edge type to its `(s, t)` counts. */ -type RawEdgeConnectionsResponse = { - result: { - data: { - "@type": "g:List"; - "@value": Array>; - }; - }; +type RawResponse = { + result: { data: { "@type": "g:List"; "@value": Value[] } }; }; +/** + * A scan returns one map keyed by the `(e, s, t)` projection, or an empty one + * when no edge matched. + */ +type RawScanResponse = RawResponse; + +/** A sample returns one map from edge type to its `(s, t)` counts. */ +type RawSampleResponse = RawResponse>; + /** One endpoint label combination, read by key before it is validated. */ type Combination = { - edgeType: string; - sourceLabels: EndpointLabels | undefined; - targetLabels: EndpointLabels | undefined; + edgeType: string | EndpointLabels | undefined; + sourceLabels: string | EndpointLabels | undefined; + targetLabels: string | EndpointLabels | undefined; }; export default async function fetchEdgeConnections( gremlinFetch: GremlinFetch, req: EdgeConnectionsRequest, ): Promise { - const batches = chunk(req.edgeTypes, EDGE_TYPES_PER_SAMPLE); - const responses = await mapWithConcurrency( - batches, - DEFAULT_CONCURRENT_REQUESTS_LIMIT, - batch => - gremlinFetch( - edgeConnectionsTemplate({ types: batch }), - ), + const plan = planDiscovery({ + edgeTypes: req.edgeTypes, + totalEdges: req.totalEdges, + }); + + logger.log("[Edge connection discovery] Planned", { + strategy: plan.strategy, + requests: plan.requests.length, + requestTimeoutMs: plan.requestTimeoutMs, + edgeTypes: req.edgeTypes.length, + totalEdges: req.totalEdges, + }); + + try { + return await runPlan(gremlinFetch, plan, req.edgeTypes); + } catch (error) { + if (!isTooBig(error)) { + throw error; + } + + // A sampled pass has nothing cheaper to fall back to. + if (plan.strategy !== "complete") { + throw giveUp( + plan, + { totalEdges: req.totalEdges, degraded: false }, + error, + ); + } + + logger.warn( + "[Edge connection discovery] A complete scan was too large for the database, sampling each edge type instead", + error, + ); + + const sampled = planSampling(req.edgeTypes); + + try { + return await runPlan(gremlinFetch, sampled, req.edgeTypes); + } catch (sampledError) { + if (!isTooBig(sampledError)) { + throw sampledError; + } + throw giveUp( + sampled, + { totalEdges: req.totalEdges, degraded: true }, + sampledError, + ); + } + } +} + +/** Reports a size failure with the recovery path that is still open. */ +function giveUp( + plan: DiscoveryPlan, + attempt: Omit, + cause: unknown, +): EdgeConnectionDiscoveryError { + const error = new EdgeConnectionDiscoveryError( + { + ...attempt, + strategy: plan.strategy, + requests: plan.requests.length, + cause: causeOf(cause), + // Through the same guard the planner used, so the error reports the total + // the plan was actually made from. The raw value is cast out of a response + // and may not be a number at all. + totalEdges: toEdgeTotal(attempt.totalEdges), + }, + cause, ); + logger.error(`[Edge connection discovery] Gave up. ${error.recovery}`, error); + return error; +} + +async function runPlan( + gremlinFetch: GremlinFetch, + plan: DiscoveryPlan, + schemaEdgeTypes: EdgeType[], +): Promise { + const startedAt = performance.now(); + const abandon = new AbortController(); + + try { + const responses = await mapWithConcurrency( + plan.requests, + DEFAULT_CONCURRENT_REQUESTS_LIMIT, + async request => { + const response = await gremlinFetch< + RawScanResponse | RawSampleResponse + >(edgeConnectionsTemplate(request), { + // Per request, so the proxy cancels this scan at the database rather + // than whatever else the connection happens to be doing. + queryId: v4(), + signal: anySignal(abandon.signal, requestTimeoutSignal(plan)), + }); + return "limitPerType" in request + ? sampledCombinations(response as RawSampleResponse) + : scannedCombinations(response as RawScanResponse); + }, + ); + + const edgeConnections = parseEdgeConnections( + responses.flat(), + schemaEdgeTypes, + ); + logger.log("[Edge connection discovery] Finished", { + strategy: plan.strategy, + requests: plan.requests.length, + edgeConnections: edgeConnections.length, + elapsedMs: Math.round(performance.now() - startedAt), + }); + return { edgeConnections }; + } finally { + // Whatever is still in flight belongs to an attempt nobody is waiting for + // any more. Aborting closes the connection to the proxy, which turns that + // into a `cancelQuery` for the `queryId` the request carried, so the + // database stops scanning too. + abandon.abort(); + } +} + +function requestTimeoutSignal(plan: DiscoveryPlan): AbortSignal | undefined { + return plan.requestTimeoutMs === undefined + ? undefined + : AbortSignal.timeout(plan.requestTimeoutMs); +} +/** + * Flattens the counted triples into edge connections, expanding Neptune `::` + * composite labels on both endpoints. + * + * The counts are read and discarded. `EdgeConnection.count` stays unpopulated + * because the same field would be capped, and so misleading, whenever the + * sampled strategy produced it. + * + * @param schemaEdgeTypes Edge types the app can render. An unfiltered scan sees + * every edge type in the graph, including ones discovery was not asked about. + */ +function parseEdgeConnections( + combinations: Combination[], + schemaEdgeTypes: EdgeType[], +): EdgeConnection[] { + const knownEdgeTypes = new Set(schemaEdgeTypes); const seen = new Set(); const edgeConnections: EdgeConnection[] = []; - for (const { edgeType, sourceLabels, targetLabels } of responses.flatMap( - sampledCombinations, - )) { + for (const { edgeType, sourceLabels, targetLabels } of combinations) { + if (typeof edgeType !== "string" || !knownEdgeTypes.has(edgeType)) { + continue; + } + for (const sourceType of endpointTypes(sourceLabels)) { for (const targetType of endpointTypes(targetLabels)) { const connection: EdgeConnection = { @@ -90,21 +239,28 @@ export default async function fetchEdgeConnections( } } - return { edgeConnections }; + return edgeConnections; } -/** - * Reads the counted endpoint label combinations by key. The counts are - * discarded, because a sampled count is capped and so would mislead. - */ -function sampledCombinations( - response: RawEdgeConnectionsResponse, -): Combination[] { +function scannedCombinations(response: RawScanResponse): Combination[] { + return response.result.data["@value"].flatMap(counts => + [...parseGMap(counts).keys()].map(key => { + const labels = parseGMap(key); + return { + edgeType: labels.get(projectionKeys.edgeType), + sourceLabels: labels.get(projectionKeys.sourceType), + targetLabels: labels.get(projectionKeys.targetType), + }; + }), + ); +} + +function sampledCombinations(response: RawSampleResponse): Combination[] { return response.result.data["@value"].flatMap(byEdgeType => [...parseGMap(byEdgeType)].flatMap( ([edgeType, counts]) => [...parseGMap(counts).keys()].map(key => { - const labels = parseGMap(key); + const labels = parseGMap(key); return { edgeType, sourceLabels: labels.get(projectionKeys.sourceType), @@ -120,6 +276,9 @@ function sampledCombinations( * multi-label vertex into one `::` composite, 1.3.5 into one entry per label, so * every entry is split. */ -function endpointTypes(labels: EndpointLabels | undefined): string[] { - return labels?.["@value"].flatMap(label => splitLabel(label)) ?? []; +function endpointTypes(labels: string | EndpointLabels | undefined): string[] { + if (labels === undefined || typeof labels === "string") { + return []; + } + return labels["@value"].flatMap(label => splitLabel(label)); } diff --git a/packages/graph-explorer/src/connector/gremlin/gremlinExplorer.ts b/packages/graph-explorer/src/connector/gremlin/gremlinExplorer.ts index 874ac404d7..3257fb165c 100644 --- a/packages/graph-explorer/src/connector/gremlin/gremlinExplorer.ts +++ b/packages/graph-explorer/src/connector/gremlin/gremlinExplorer.ts @@ -9,6 +9,7 @@ import type { Explorer, ExplorerRequestOptions } from "../useGEFetchTypes"; import type { GraphSummary, GremlinFetch } from "./types"; import { fetchDatabaseRequest } from "../fetchDatabaseRequest"; +import { anySignal } from "../utils/anySignal"; import { edgeDetails } from "./edgeDetails"; import fetchEdgeConnections from "./fetchEdgeConnections"; import fetchNeighbors from "./fetchNeighbors"; @@ -24,15 +25,16 @@ function _gremlinFetch( featureFlags: FeatureFlags, options?: ExplorerRequestOptions, ): GremlinFetch { - return async (queryTemplate: string) => { + return async (queryTemplate: string, requestOptions?) => { logger.debug(queryTemplate); const body = JSON.stringify({ query: queryTemplate }); const headers: HeadersInit = { "Content-Type": "application/json", Accept: "application/vnd.gremlin-v3.0+json", }; - if (options?.queryId && connection.proxyConnection === true) { - headers.queryId = options.queryId; + const queryId = requestOptions?.queryId ?? options?.queryId; + if (queryId && connection.proxyConnection === true) { + headers.queryId = queryId; } return fetchDatabaseRequest( @@ -44,6 +46,9 @@ function _gremlinFetch( headers, body, ...options, + // Merged rather than overridden, so a per-request bound cannot outlive + // the cancellation the caller already holds. + signal: anySignal(options?.signal, requestOptions?.signal), }, ); }; diff --git a/packages/graph-explorer/src/connector/gremlin/types.ts b/packages/graph-explorer/src/connector/gremlin/types.ts index 84cb6eaa74..5b945650c4 100644 --- a/packages/graph-explorer/src/connector/gremlin/types.ts +++ b/packages/graph-explorer/src/connector/gremlin/types.ts @@ -137,8 +137,20 @@ export type GAnyValue = | GType | null; +/** + * Overrides for one request out of many sharing a `GremlinFetch`, so a caller + * that issues several can identify and cancel them individually. + */ +export type GremlinRequestOptions = { + /** Lets the proxy cancel this request at the database. Proxy connections only. */ + queryId?: string; + /** Combined with the connection's timeout and the caller's own signal, not replacing them. */ + signal?: AbortSignal; +}; + export type GremlinFetch = ( queryTemplate: string, + requestOptions?: GremlinRequestOptions, ) => Promise; export type GraphSummary = { diff --git a/packages/graph-explorer/src/connector/queries/edgeConnectionsQuery.ts b/packages/graph-explorer/src/connector/queries/edgeConnectionsQuery.ts index 633cdfc31e..ba2599d6ea 100644 --- a/packages/graph-explorer/src/connector/queries/edgeConnectionsQuery.ts +++ b/packages/graph-explorer/src/connector/queries/edgeConnectionsQuery.ts @@ -31,9 +31,13 @@ export function edgeConnectionsQuery( ) { // Sort edge types to keep the order consistent over time to increase the chance of hitting cache const sortedEdgeTypes = activeSchema?.edges.map(e => e.type).toSorted() ?? []; + // The edge total decides how the Gremlin connector discovers, so a graph that + // has grown or shrunk past a threshold must not be served an answer gathered + // under the old size. + const totalEdges = activeSchema?.totalEdges; return queryOptions({ - queryKey: ["schema", "edgeConnections", sortedEdgeTypes], + queryKey: ["schema", "edgeConnections", sortedEdgeTypes, totalEdges], staleTime: Infinity, retryOnMount: false, // A failed fetch marks the query stale, so remounting observers (the @@ -57,7 +61,7 @@ export function edgeConnectionsQuery( try { const results = await explorer.fetchEdgeConnections( - { edgeTypes: sortedEdgeTypes }, + { edgeTypes: sortedEdgeTypes, totalEdges }, { signal }, ); diff --git a/packages/graph-explorer/src/connector/useGEFetchTypes.ts b/packages/graph-explorer/src/connector/useGEFetchTypes.ts index 659c1d48f0..0572708f1e 100644 --- a/packages/graph-explorer/src/connector/useGEFetchTypes.ts +++ b/packages/graph-explorer/src/connector/useGEFetchTypes.ts @@ -207,6 +207,12 @@ export type RawQueryResponse = { export type EdgeConnectionsRequest = { edgeTypes: EdgeType[]; + /** + * Total edges in the graph, which the Gremlin connector uses to size the work. + * Absent when the cached schema predates us recording it, which costs extra + * requests rather than correctness. + */ + totalEdges?: number; }; export type EdgeConnectionsResponse = { diff --git a/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.test.tsx b/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.test.tsx index 998cc6c0e3..3e886a5047 100644 --- a/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.test.tsx +++ b/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.test.tsx @@ -4,10 +4,10 @@ import userEvent from "@testing-library/user-event"; import { beforeEach, describe, expect, test, vi } from "vitest"; import { TooltipProvider } from "@/components"; +import { EdgeConnectionDiscoveryError } from "@/connector/gremlin/fetchEdgeConnections/discoveryError"; import { createEdgeType, getAppStore, schemaAtom } from "@/core"; import { createQueryClient } from "@/core/queryClient"; -import { DatabaseTimeoutError, NetworkError } from "@/utils"; -import { createDisplayError } from "@/utils/createDisplayError"; +import { NetworkError } from "@/utils"; import { DbState, FakeExplorer, @@ -64,12 +64,17 @@ describe("EdgeConnectionDiscoveryStatusButton", () => { return state; } - test("shows a button with the error title, message, Error Details, and Retry in the popover when discovery rejects", async () => { - const error = new DatabaseTimeoutError( - "Query timed out", - 500, - { code: "TimeLimitExceededException" }, - "TimeLimitExceededException", + test("shows the error title, recovery text, Error Details, and Retry in the popover when discovery rejects", async () => { + const cause = new Error("Query timed out"); + const error = new EdgeConnectionDiscoveryError( + { + strategy: "sampled", + requests: 3, + totalEdges: 100, + degraded: false, + cause: "database-limit", + }, + cause, ); vi.spyOn(explorer, "fetchEdgeConnections").mockRejectedValue(error); @@ -83,16 +88,11 @@ describe("EdgeConnectionDiscoveryStatusButton", () => { await user.click(button); await waitFor(() => { - expect(screen.getByText("Database query timed out")).toBeInTheDocument(); + expect( + screen.getByText("Could not discover edge connections"), + ).toBeInTheDocument(); }); - expect( - screen.getByText(createDisplayError(error).message), - ).toBeInTheDocument(); - expect( - screen.getByText( - "Node types are still shown. Retry, or open Error Details to see the full error.", - ), - ).toBeInTheDocument(); + expect(screen.getByText(new RegExp(error.recovery))).toBeInTheDocument(); // Both actions share a size so the popover's buttons line up. expect( screen.getByRole("button", { name: /error details/i }), @@ -104,7 +104,16 @@ describe("EdgeConnectionDiscoveryStatusButton", () => { }); test("hides the button and clears the stored failure flag once retry resolves", async () => { - const error = new Error("boom"); + const error = new EdgeConnectionDiscoveryError( + { + strategy: "sampled", + requests: 1, + totalEdges: 10, + degraded: false, + cause: "database-limit", + }, + new Error("boom"), + ); const spy = vi .spyOn(explorer, "fetchEdgeConnections") .mockRejectedValueOnce(error) @@ -120,7 +129,9 @@ describe("EdgeConnectionDiscoveryStatusButton", () => { const user = userEvent.setup(); await user.click(button); await waitFor(() => { - expect(screen.getByText("Something went wrong")).toBeInTheDocument(); + expect( + screen.getByText("Could not discover edge connections"), + ).toBeInTheDocument(); }); await user.click(screen.getByRole("button", { name: /retry/i })); diff --git a/packages/graph-explorer/src/utils/createDisplayError.test.ts b/packages/graph-explorer/src/utils/createDisplayError.test.ts index d2b83c4334..54147c51a8 100644 --- a/packages/graph-explorer/src/utils/createDisplayError.test.ts +++ b/packages/graph-explorer/src/utils/createDisplayError.test.ts @@ -1,6 +1,7 @@ // @vitest-environment happy-dom import { z } from "zod"; +import { EdgeConnectionDiscoveryError } from "@/connector/gremlin/fetchEdgeConnections/discoveryError"; import { EmptyIdentifierError, QueryValueError, @@ -221,6 +222,38 @@ describe("createDisplayError", () => { }); }); + it("Should handle the database running out of memory", () => { + const result = createDisplayError({ code: "MemoryLimitExceededException" }); + expect(result).toStrictEqual({ + title: "Not enough memory", + message: + "The database ran out of memory answering the query. Try a smaller request, or use an instance with more memory.", + }); + }); + + it("Should give edge connection discovery its own recovery instructions", () => { + const result = createDisplayError( + new EdgeConnectionDiscoveryError( + { + strategy: "sampled", + requests: 1, + totalEdges: 19_928_805, + degraded: true, + cause: "database-limit", + }, + new NetworkError("Query cannot be completed", 500, { + code: "MemoryLimitExceededException", + }), + ), + ); + + expect(result.title).toBe("Could not discover edge connections"); + // The generic memory branch would say "try a smaller request", which is not + // something the user can do here. The database configuration is. + expect(result.message).toContain("DB cluster parameter group"); + expect(result.message).not.toContain("smaller request"); + }); + it("Should handle malformed query", () => { const result = createDisplayError({ code: "MalformedQueryException" }); expect(result).toStrictEqual({ diff --git a/packages/graph-explorer/src/utils/createDisplayError.ts b/packages/graph-explorer/src/utils/createDisplayError.ts index 3a078eeaea..f53baf5b02 100644 --- a/packages/graph-explorer/src/utils/createDisplayError.ts +++ b/packages/graph-explorer/src/utils/createDisplayError.ts @@ -1,5 +1,6 @@ import { ZodError } from "zod"; +import { EdgeConnectionDiscoveryError } from "@/connector/gremlin/fetchEdgeConnections/discoveryError"; import { EmptyIdentifierError, QueryValueError, @@ -39,6 +40,15 @@ const UNREACHABLE_HOST_CODES = new Set(["ENOTFOUND", "EAI_AGAIN"]); * @returns A `DisplayError` that contains a title and message. */ export function createDisplayError(error: any): DisplayError { + // First, because it already knows more about the failure than any code on the + // response it wraps: which strategies were tried and what the user can change. + if (error instanceof EdgeConnectionDiscoveryError) { + return { + title: "Could not discover edge connections", + message: `${error.message} ${error.recovery}`, + }; + } + const data = error instanceof NetworkError ? error.data @@ -88,6 +98,19 @@ export function createDisplayError(error: any): DisplayError { "Please check the database URL in the connection and try again.", }; } + if ( + data.code === "MemoryLimitExceededException" || + data.cause?.code === "MemoryLimitExceededException" + ) { + // The query asked for more memory than the instance had, which is a + // property of the query rather than of the connection. + return { + title: "Not enough memory", + message: + "The database ran out of memory answering the query. Try a smaller request, or use an instance with more memory.", + }; + } + // Malformed query if ( data.code === "MalformedQueryException" || diff --git a/packages/graph-explorer/src/utils/createErrorDetails.test.ts b/packages/graph-explorer/src/utils/createErrorDetails.test.ts index 55c26c8dc9..b6ce93ca97 100644 --- a/packages/graph-explorer/src/utils/createErrorDetails.test.ts +++ b/packages/graph-explorer/src/utils/createErrorDetails.test.ts @@ -1,5 +1,6 @@ import { z } from "zod"; +import { EdgeConnectionDiscoveryError } from "@/connector/gremlin/fetchEdgeConnections/discoveryError"; import { QueryValueError, UnescapableValueError, @@ -358,6 +359,34 @@ describe("createErrorDetails", () => { }); }); + describe("EdgeConnectionDiscoveryError", () => { + it("serializes the attempt alongside the database's own error", () => { + const error = new EdgeConnectionDiscoveryError( + { + strategy: "sampled", + requests: 3, + totalEdges: 19_928_805, + degraded: true, + cause: "database-limit", + }, + new NetworkError("Query cannot be completed", 500, { + code: "MemoryLimitExceededException", + }), + ); + + const details = createErrorDetails(error); + + expect(details.name).toBe("EdgeConnectionDiscoveryError"); + // Our interpretation and the database's error, so support can read both. + expect(JSON.parse(details.data!)).toMatchObject({ + strategy: "sampled", + completeScanAbandoned: true, + failureCause: "database-limit", + cause: { name: "NetworkError", statusCode: 500 }, + }); + }); + }); + describe("DOMException", () => { it("returns name and message from an AbortError", () => { const error = new DOMException("The operation was aborted", "AbortError"); diff --git a/packages/graph-explorer/src/utils/createErrorDetails.ts b/packages/graph-explorer/src/utils/createErrorDetails.ts index de8dbd1831..53eebd5c82 100644 --- a/packages/graph-explorer/src/utils/createErrorDetails.ts +++ b/packages/graph-explorer/src/utils/createErrorDetails.ts @@ -1,5 +1,6 @@ import { ZodError } from "zod"; +import { EdgeConnectionDiscoveryError } from "@/connector/gremlin/fetchEdgeConnections/discoveryError"; import { QueryValueError } from "@/connector/queryValueError"; import { NetworkError } from "./NetworkError"; @@ -42,11 +43,20 @@ export function createErrorDetails(error: unknown): ErrorDetails { data: JSON.stringify(error.issues, null, 2), }; } - if (error instanceof QueryValueError) { + // Errors that carry their own structured context. The cause is serialized with + // it, so the database's own error is still readable under our interpretation. + if ( + error instanceof QueryValueError || + error instanceof EdgeConnectionDiscoveryError + ) { + const data: Record = { ...error.details }; + if (error.cause) { + data.cause = serializeCause(error.cause); + } return { name: error.name, message: error.message, - data: JSON.stringify(error.details, null, 2), + data: JSON.stringify(data, null, 2), }; } if (error instanceof Error) { From 76b09f74f476b874f06a16885949acd3a88f6b58 Mon Sep 17 00:00:00 2001 From: Kris McGinnes Date: Sat, 26 Sep 2026 12:43:16 -0500 Subject: [PATCH 2/4] Pin the edge total in the refetch loop test so the query keeps its key --- .../components/SchemaDiscoveryBoundary.integration.test.tsx | 3 +++ 1 file changed, 3 insertions(+) diff --git a/packages/graph-explorer/src/components/SchemaDiscoveryBoundary.integration.test.tsx b/packages/graph-explorer/src/components/SchemaDiscoveryBoundary.integration.test.tsx index 1e1ed241ff..e9326463cb 100644 --- a/packages/graph-explorer/src/components/SchemaDiscoveryBoundary.integration.test.tsx +++ b/packages/graph-explorer/src/components/SchemaDiscoveryBoundary.integration.test.tsx @@ -106,6 +106,9 @@ describe("SchemaDiscoveryBoundary against the real store", () => { state.activeSchema.edges = []; state.addTestableEdgeToGraph(edge); state.activeSchema.edgeConnections = []; + // Match the total the refresh reports, so the edge connection query keeps + // its key and only a loop could fetch it more than once. + state.activeSchema.totalEdges = 1; const store = getAppStore(); state.applyTo(store); From ac09c2d051486650c80133430e60114093b2a365 Mon Sep 17 00:00:00 2001 From: Kris McGinnes Date: Sat, 26 Sep 2026 13:12:09 -0500 Subject: [PATCH 3/4] Fetch edge connections once, under the refreshed schema's key, on Refresh Schema Refetching through the observer ran against the key from before the refresh, so a changed edge total fetched under the stale key and again under the new one. --- docs/agents/schema.md | 2 + .../src/hooks/useSchemaSync.test.ts | 41 ++++++++++++++++++- .../graph-explorer/src/hooks/useSchemaSync.ts | 13 +++++- 3 files changed, 52 insertions(+), 4 deletions(-) diff --git a/docs/agents/schema.md b/docs/agents/schema.md index 1eff3e6a3c..342c5caf38 100644 --- a/docs/agents/schema.md +++ b/docs/agents/schema.md @@ -27,6 +27,8 @@ If the edge connection query fails, the schema records only that it failed, via `edgeConnectionNotice(schema, error)` (in `src/hooks/edgeConnectionNotice.ts`, next to `useSchemaSync`) checks failure first (the live query error or `lastEdgeConnectionSyncFail`) and only then `edgeConnections == null`, because partial connections added by exploration after a failure must still report the failure. `useEdgeConnectionNotice()` wraps it with `useMaybeActiveSchema()` and `useSchemaSync().edgeDiscoveryQuery` so the toolbar button, the sidebar details, and the connection detail panel share one resolution of the notice. +The edge query's key includes the sorted edge types and the schema's `totalEdges`, so a schema refresh usually moves it to a new key. `refreshSchema` therefore fetches `edgeConnectionsQuery(refreshedSchema)` through the query client rather than calling the observer's `refetch()`, which would still target the pre-refresh key and cost a second request once the observer switched. + `useCancelSchemaSync` cancels the edge query with `revert: false`. Reverting would restore its never-fetched state, and every observer that remounts under `SchemaDiscoveryBoundary` loads a query without data on mount (TanStack ignores `refetchOnMount` then), so the fetch would restart. Settled as a `CancelledError`, `retryOnMount: false` holds it, and `edgeConnectionNotice` reads the cancellation as not discovered rather than failed. ## Incremental Schema Growth diff --git a/packages/graph-explorer/src/hooks/useSchemaSync.test.ts b/packages/graph-explorer/src/hooks/useSchemaSync.test.ts index 4ef2281736..08eb0fceb4 100644 --- a/packages/graph-explorer/src/hooks/useSchemaSync.test.ts +++ b/packages/graph-explorer/src/hooks/useSchemaSync.test.ts @@ -1,3 +1,4 @@ +import { createRandomInteger } from "@shared/utils/testing"; // @vitest-environment happy-dom import { act, waitFor } from "@testing-library/react"; import { beforeEach, describe, expect, it, vi } from "vitest"; @@ -16,6 +17,7 @@ import { createRandomEdgeTypeConfig, createRandomRawConfiguration, createRandomVertexTypeConfig, + createTestableEdge, DbState, FakeExplorer, flushPendingAtomUpdates, @@ -325,8 +327,15 @@ describe("useSchemaSync", () => { }); it("should refetch edge discovery query", async () => { - const edgeType = createEdgeType("worksAt"); - const state = createStateWithSchema([], [edgeType]); + // The refresh must report the same edge types and total, so the edge + // query keeps its key and refetches rather than moving to a new one. + const edge = createTestableEdge(); + explorer.addTestableEdge(edge); + const state = new DbState(explorer); + state.activeSchema.edges = []; + state.addTestableEdgeToGraph(edge); + state.activeSchema.edgeConnections = undefined; + state.activeSchema.totalEdges = 1; const fetchEdgeConnectionsSpy = vi.spyOn( explorer, "fetchEdgeConnections", @@ -344,6 +353,34 @@ describe("useSchemaSync", () => { expect(fetchEdgeConnectionsSpy).toHaveBeenCalledTimes(2); }); + + it("fetches edge connections once, with the refreshed edge total, when the total changes", async () => { + const edge = createTestableEdge(); + explorer.addTestableEdge(edge); + const state = new DbState(explorer); + state.activeSchema.edges = []; + state.addTestableEdgeToGraph(edge); + state.activeSchema.edgeConnections = []; + // The refresh reports the explorer's single edge + state.activeSchema.totalEdges = 1 + createRandomInteger(); + const fetchEdgeConnectionsSpy = vi.spyOn( + explorer, + "fetchEdgeConnections", + ); + + const { result } = renderHookWithState(() => useSchemaSync(), state); + + await act(async () => { + await result.current.refreshSchema(); + }); + await waitFor(() => { + expect(result.current.isFetching).toBe(false); + }); + + expect(fetchEdgeConnectionsSpy.mock.calls).toStrictEqual([ + [{ edgeTypes: [edge.type], totalEdges: 1 }, expect.anything()], + ]); + }); }); describe("staleTime behavior", () => { diff --git a/packages/graph-explorer/src/hooks/useSchemaSync.ts b/packages/graph-explorer/src/hooks/useSchemaSync.ts index df93729c8b..f444a39973 100644 --- a/packages/graph-explorer/src/hooks/useSchemaSync.ts +++ b/packages/graph-explorer/src/hooks/useSchemaSync.ts @@ -50,6 +50,7 @@ export function useSchemaSync() { // The schema and connectionId must update in the same render so the query // options stay consistent when switching connections. const activeSchema = useAtomValue(maybeActiveSchemaAtom); + const queryClient = useQueryClient(); const connectionId = useAtomValue(activeConfigurationAtom); const schemaDiscoveryQuery = useQuery( @@ -68,8 +69,16 @@ export function useSchemaSync() { const refreshSchema = async () => { logger.log("Refreshing schema"); - await schemaDiscoveryQuery.refetch(); - await edgeDiscoveryQuery.refetch(); + const { data: refreshedSchema } = await schemaDiscoveryQuery.refetch(); + // The edge query's key derives from the schema, so refetching through the + // observer would fetch under the pre-refresh key, and the next render would + // fetch again under the new one. Fetch the new key directly instead; the + // observer finds it settled or in flight when it switches. + await queryClient.prefetchQuery({ + ...edgeConnectionsQuery(refreshedSchema), + // An explicit refresh fetches even when cached data exists. + staleTime: 0, + }); }; return { From dcc60ed9b90e5a5f9f0fab44b1de8d2d1df17b1d Mon Sep 17 00:00:00 2001 From: Kris McGinnes Date: Sat, 26 Sep 2026 13:13:41 -0500 Subject: [PATCH 4/4] Let the edge connection discovery error speak for itself in the popover Its display title repeats the popover title and its message already carries the recovery, so drop the lead-in and the generic next step for it. --- ...geConnectionDiscoveryStatusButton.test.tsx | 21 +++++++++++++++---- .../EdgeConnectionDiscoveryStatusButton.tsx | 8 +++++++ 2 files changed, 25 insertions(+), 4 deletions(-) diff --git a/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.test.tsx b/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.test.tsx index 3e886a5047..f929f49a20 100644 --- a/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.test.tsx +++ b/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.test.tsx @@ -64,7 +64,7 @@ describe("EdgeConnectionDiscoveryStatusButton", () => { return state; } - test("shows the error title, recovery text, Error Details, and Retry in the popover when discovery rejects", async () => { + test("shows the failure and its recovery text, without a repeated title or generic next step, plus Error Details and Retry", async () => { const cause = new Error("Query timed out"); const error = new EdgeConnectionDiscoveryError( { @@ -89,10 +89,18 @@ describe("EdgeConnectionDiscoveryStatusButton", () => { await waitFor(() => { expect( - screen.getByText("Could not discover edge connections"), + screen.getByText("Could not discover Relationships"), ).toBeInTheDocument(); }); - expect(screen.getByText(new RegExp(error.recovery))).toBeInTheDocument(); + expect( + screen.getByText(`${error.message} ${error.recovery}`), + ).toBeInTheDocument(); + expect( + screen.queryByText("Could not discover edge connections"), + ).not.toBeInTheDocument(); + expect( + screen.queryByText(/Node types are still shown/), + ).not.toBeInTheDocument(); // Both actions share a size so the popover's buttons line up. expect( screen.getByRole("button", { name: /error details/i }), @@ -130,7 +138,7 @@ describe("EdgeConnectionDiscoveryStatusButton", () => { await user.click(button); await waitFor(() => { expect( - screen.getByText("Could not discover edge connections"), + screen.getByText("Could not discover Relationships"), ).toBeInTheDocument(); }); await user.click(screen.getByRole("button", { name: /retry/i })); @@ -219,6 +227,11 @@ describe("EdgeConnectionDiscoveryStatusButton", () => { await waitFor(() => { expect(screen.getByText("Connection refused")).toBeInTheDocument(); }); + expect( + screen.getByText( + "Node types are still shown. Retry, or open Error Details to see the full error.", + ), + ).toBeInTheDocument(); }); test("shows the not-discovered button with a Synchronize action in the popover before discovery has ever completed", async () => { diff --git a/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.tsx b/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.tsx index c5380774e5..310ca73c98 100644 --- a/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.tsx +++ b/packages/graph-explorer/src/modules/SchemaGraph/EdgeConnectionDiscoveryStatusButton.tsx @@ -16,6 +16,7 @@ import { PopoverTitle, PopoverTrigger, } from "@/components"; +import { EdgeConnectionDiscoveryError } from "@/connector/gremlin/fetchEdgeConnections/discoveryError"; import { useEdgeConnectionNotice } from "@/hooks/useEdgeConnectionNotice"; import useTranslations from "@/hooks/useTranslations"; import { createDisplayError } from "@/utils/createDisplayError"; @@ -96,6 +97,13 @@ function describeFailure(error: Error | null): FailureDescription { } const displayError = createDisplayError(error); + + // Its title restates the popover's, and its message already ends with what + // to change and that node types are still shown. + if (error instanceof EdgeConnectionDiscoveryError) { + return { lead: null, message: displayError.message, nextStep: null }; + } + return { lead: displayError.title, message: displayError.message,