Merge pull request #29 from avalon-vanguard/fix/etl-gaia-floor
Refuse a Gaia answer that came back short, and read the body inside the retry
This commit is contained in:
@@ -49,9 +49,11 @@ jobs:
|
|||||||
key: gaia-dr3-${{ hashFiles('tools/etl/sources/gaia.ts') }}
|
key: gaia-dr3-${{ hashFiles('tools/etl/sources/gaia.ts') }}
|
||||||
|
|
||||||
# Every other source is fetched live on this fresh runner. A failed fetch fails the run by
|
# Every other source is fetched live on this fresh runner. A failed fetch fails the run by
|
||||||
# design — no refresh is better than a partial one. That includes Gaia on a cold cache:
|
# design — no refresh is better than a partial one. That includes Gaia on a cold cache, by
|
||||||
# the ETL skips it when unreachable, and the merge gate in build.ts then refuses a
|
# two different paths: its Hipparcos cross-match is required, so an unreachable archive
|
||||||
# catalogue it contributed nothing to.
|
# fails the run from fetchStars itself, while its main query is skipped when unreachable and
|
||||||
|
# the merge gate in build.ts then refuses a catalogue it contributed nothing to. An archive
|
||||||
|
# that answers short rather than not at all is caught in fetchGaiaStars.
|
||||||
- name: Rebuild the datasets
|
- name: Rebuild the datasets
|
||||||
run: npm run etl
|
run: npm run etl
|
||||||
|
|
||||||
|
|||||||
+12
-2
@@ -66,8 +66,18 @@ function validateStars(stars: StarRecord[]): void {
|
|||||||
* The other failure leaves no close pair at all, because proper motion had already carried the
|
* The other failure leaves no close pair at all, because proper motion had already carried the
|
||||||
* two entries tens of arcseconds apart — the 2026-08-24 refresh, where HYG sat at epoch 2000.0
|
* two entries tens of arcseconds apart — the 2026-08-24 refresh, where HYG sat at epoch 2000.0
|
||||||
* and Gaia at J2016.0. What it does leave is HYG rows that found no counterpart: 36 056 of them
|
* and Gaia at J2016.0. What it does leave is HYG rows that found no counterpart: 36 056 of them
|
||||||
* against the 10 876 Gaia genuinely lacks (bright stars it saturates on, red dwarfs past its
|
* against the 10 886 today, and no counterpart was possible for most of those. Two thirds of them,
|
||||||
* magnitude cut).
|
* 6 835, are the stars Gaia measures but the main query never downloads, because Gaia's parallax
|
||||||
|
* puts them past `ETL_GAIA_DISTANCE_PC` while Hipparcos put them inside `ETL_STAR_DISTANCE_PC`;
|
||||||
|
* they are every star in the published catalogue beyond 250 pc. The rest are what Gaia genuinely
|
||||||
|
* lacks: bright stars it saturates on, red dwarfs past its magnitude cut. So the headroom left to
|
||||||
|
* the ceiling tracks the gap between those two cutoffs as much as Gaia's completeness.
|
||||||
|
*
|
||||||
|
* This bounds a merge that went wrong, and — loosely — a Gaia download that came back short: a
|
||||||
|
* truncated answer leaves the HYG rows whose counterpart it dropped without one, so survivors go
|
||||||
|
* *up*, not down. Measured against the published catalogue: 10 886 today, 11 004 at nine tenths of
|
||||||
|
* the rows, 12 711 at half, 16 258 at a third. So this ceiling only catches a truncation past about
|
||||||
|
* two thirds, and `fetchGaiaStars` catches the shallower ones with its own row floor.
|
||||||
*/
|
*/
|
||||||
const MAX_UNMERGED_TWINS = 100;
|
const MAX_UNMERGED_TWINS = 100;
|
||||||
const MAX_HYG_SURVIVORS = 15_000;
|
const MAX_HYG_SURVIVORS = 15_000;
|
||||||
|
|||||||
@@ -3,7 +3,7 @@ import { writeFileSync } from 'node:fs';
|
|||||||
import { mergeStarCatalogues, placementDistancePc } from '../../src/app/shared/astro/star-merge';
|
import { mergeStarCatalogues, placementDistancePc } from '../../src/app/shared/astro/star-merge';
|
||||||
import { encodeStarCatalog } from '../../src/app/shared/models/star-catalog';
|
import { encodeStarCatalog } from '../../src/app/shared/models/star-catalog';
|
||||||
import { StarRecord, SUN_STAR_ID } from '../../src/app/shared/models/star.model';
|
import { StarRecord, SUN_STAR_ID } from '../../src/app/shared/models/star.model';
|
||||||
import { fetchGaiaDistancesByHip } from './sources/gaia';
|
import { fetchGaiaDistancesByHip, GaiaAnswerError } from './sources/gaia';
|
||||||
import { positionalSources } from './sources/registry';
|
import { positionalSources } from './sources/registry';
|
||||||
import { PARALLAX_PRECISION_MAS } from './sources/star-sources';
|
import { PARALLAX_PRECISION_MAS } from './sources/star-sources';
|
||||||
import { parseCsvObjects, parseOptionalNumber } from './lib/csv';
|
import { parseCsvObjects, parseOptionalNumber } from './lib/csv';
|
||||||
@@ -142,7 +142,9 @@ export async function fetchStars(): Promise<StarRecord[]> {
|
|||||||
*
|
*
|
||||||
* A source that cannot be reached is reported and skipped here rather than thrown, so a run still
|
* A source that cannot be reached is reported and skipped here rather than thrown, so a run still
|
||||||
* gets as far as validation and says what it has. Whether that may be published is decided
|
* gets as far as validation and says what it has. Whether that may be published is decided
|
||||||
* there: `validateMerge` in build.ts refuses a catalogue Gaia contributed nothing to.
|
* there: `validateMerge` in build.ts refuses a catalogue Gaia contributed nothing to. A source
|
||||||
|
* that answered with something unusable ({@link GaiaAnswerError}) is a different matter, and stops
|
||||||
|
* the run where it happened rather than being reported later as an outage.
|
||||||
*/
|
*/
|
||||||
async function mergeWithOtherSources(hygStars: StarRecord[]): Promise<StarRecord[]> {
|
async function mergeWithOtherSources(hygStars: StarRecord[]): Promise<StarRecord[]> {
|
||||||
const others = positionalSources().filter((source) => source.id !== 'hyg');
|
const others = positionalSources().filter((source) => source.id !== 'hyg');
|
||||||
@@ -160,6 +162,11 @@ async function mergeWithOtherSources(hygStars: StarRecord[]): Promise<StarRecord
|
|||||||
stars: await source.fetch!()
|
stars: await source.fetch!()
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
|
// An answer that cannot be worked with is not an outage: skipping it would write a
|
||||||
|
// half-catalogue over the published assets before the merge gate got to say so.
|
||||||
|
if (error instanceof GaiaAnswerError) {
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
console.log(` skipping ${source.name}: ${error instanceof Error ? error.message : error}`);
|
console.log(` skipping ${source.name}: ${error instanceof Error ? error.message : error}`);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -34,9 +34,15 @@ const RETRY_DELAYS_MS = [30_000, 120_000];
|
|||||||
|
|
||||||
async function fetchText(url: string): Promise<string> {
|
async function fetchText(url: string): Promise<string> {
|
||||||
for (let attempt = 0; ; attempt++) {
|
for (let attempt = 0; ; attempt++) {
|
||||||
const response = await fetch(url).catch((error: unknown) => (error instanceof Error ? error : new Error(String(error))));
|
let response = await fetch(url).catch((error: unknown) => (error instanceof Error ? error : new Error(String(error))));
|
||||||
if (!(response instanceof Error) && response.ok) {
|
if (!(response instanceof Error) && response.ok) {
|
||||||
return response.text();
|
// Read inside the loop, because the body is where these downloads fail: the Gaia CSV is
|
||||||
|
// 57 MB, and a connection reset part-way through rejects here, long after the 200.
|
||||||
|
const body = await response.text().catch((error: unknown) => (error instanceof Error ? error : new Error(String(error))));
|
||||||
|
if (typeof body === 'string') {
|
||||||
|
return body;
|
||||||
|
}
|
||||||
|
response = body;
|
||||||
}
|
}
|
||||||
const reason = response instanceof Error ? response.message : `${response.status} ${response.statusText}`;
|
const reason = response instanceof Error ? response.message : `${response.status} ${response.statusText}`;
|
||||||
// A 4xx is the request's own fault, and waiting will not change the answer.
|
// A 4xx is the request's own fault, and waiting will not change the answer.
|
||||||
|
|||||||
@@ -38,9 +38,12 @@ const CATALOGUE_EPOCH = 2000.0;
|
|||||||
* past anything this map draws, so the limit here is a payload decision: the catalogue is baked
|
* past anything this map draws, so the limit here is a payload decision: the catalogue is baked
|
||||||
* into a static asset that a browser downloads before the first frame.
|
* into a static asset that a browser downloads before the first frame.
|
||||||
*/
|
*/
|
||||||
const DISTANCE_CUTOFF_PC = Number(process.env['ETL_GAIA_DISTANCE_PC'] ?? 250);
|
const DEFAULT_DISTANCE_CUTOFF_PC = 250;
|
||||||
const MAGNITUDE_LIMIT = Number(process.env['ETL_GAIA_MAGNITUDE_LIMIT'] ?? 12);
|
const DEFAULT_MAGNITUDE_LIMIT = 12;
|
||||||
const ROW_LIMIT = Number(process.env['ETL_GAIA_ROW_LIMIT'] ?? 500000);
|
const DEFAULT_ROW_LIMIT = 500_000;
|
||||||
|
const DISTANCE_CUTOFF_PC = Number(process.env['ETL_GAIA_DISTANCE_PC'] ?? DEFAULT_DISTANCE_CUTOFF_PC);
|
||||||
|
const MAGNITUDE_LIMIT = Number(process.env['ETL_GAIA_MAGNITUDE_LIMIT'] ?? DEFAULT_MAGNITUDE_LIMIT);
|
||||||
|
const ROW_LIMIT = Number(process.env['ETL_GAIA_ROW_LIMIT'] ?? DEFAULT_ROW_LIMIT);
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Relative parallax error above which a star is dropped: a parallax measured to worse than 20%
|
* Relative parallax error above which a star is dropped: a parallax measured to worse than 20%
|
||||||
@@ -67,6 +70,35 @@ function buildQuery(): string {
|
|||||||
].join(' ');
|
].join(' ');
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* How many rows the query above holds when nothing is overridden: 412 765, and DR3 is a finished
|
||||||
|
* data release, so that number only moves when the query does. It lives here, under the query, so
|
||||||
|
* that an edit to any of its filters is made with the count it invalidates in view.
|
||||||
|
*
|
||||||
|
* Checked because a short answer looks exactly like a complete one. The TAP service truncates on
|
||||||
|
* its own timeout and still serves a well-formed CSV with a 200, and the rows are ordered by
|
||||||
|
* magnitude, so what comes back is the bright half — the half HYG overlaps. The merge gate in
|
||||||
|
* `build.ts` would then see Gaia stars present and a survivor count barely moved, and pass a
|
||||||
|
* catalogue missing two hundred thousand stars, which the weekly job would publish and the runner
|
||||||
|
* would cache for the weeks after it. Same failure, and same guard, as
|
||||||
|
* {@link MIN_USABLE_HIP_DISTANCES} below.
|
||||||
|
*
|
||||||
|
* Only checked when nothing is overridden: the environment overrides exist to fetch a smaller
|
||||||
|
* slice on purpose.
|
||||||
|
*/
|
||||||
|
const DEFAULT_QUERY_ROWS = 412_765;
|
||||||
|
const MIN_ROW_SHARE = 0.95;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* An answer the archive gave that cannot be worked with, as against an archive that gave none.
|
||||||
|
*
|
||||||
|
* `fetchStars` skips a source it cannot reach and leaves the merge gate to judge the result. That
|
||||||
|
* is right for an outage and wrong for a truncated CSV, which would be skipped, cached, and land
|
||||||
|
* as "the archive was unreachable" long after the assets had been overwritten — so these throws
|
||||||
|
* are marked, and rethrown there.
|
||||||
|
*/
|
||||||
|
export class GaiaAnswerError extends Error {}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Gaia publishes no spectral classifications, but `bp_rp` is a colour index on the same footing
|
* Gaia publishes no spectral classifications, but `bp_rp` is a colour index on the same footing
|
||||||
* as HYG's `ci` — so the app's existing colour and spectral-class handling works unchanged, and
|
* as HYG's `ci` — so the app's existing colour and spectral-class handling works unchanged, and
|
||||||
@@ -90,8 +122,22 @@ export async function fetchGaiaStars(): Promise<StarRecord[]> {
|
|||||||
// Keyed by the whole request, so a response cached for other columns, another order, or
|
// Keyed by the whole request, so a response cached for other columns, another order, or
|
||||||
// another endpoint can never be mistaken for this one — the cache records only that some
|
// another endpoint can never be mistaken for this one — the cache records only that some
|
||||||
// response arrived, not what it answered.
|
// response arrived, not what it answered.
|
||||||
const csv = await fetchTextCached(url, `gaia-dr3-${createHash('sha1').update(url).digest('hex').slice(0, 8)}.csv`);
|
const cacheKey = `gaia-dr3-${createHash('sha1').update(url).digest('hex').slice(0, 8)}.csv`;
|
||||||
|
const csv = await fetchTextCached(url, cacheKey);
|
||||||
const rows = parseCsvObjects(csv);
|
const rows = parseCsvObjects(csv);
|
||||||
|
const jobsQuery = DISTANCE_CUTOFF_PC === DEFAULT_DISTANCE_CUTOFF_PC && MAGNITUDE_LIMIT === DEFAULT_MAGNITUDE_LIMIT && ROW_LIMIT === DEFAULT_ROW_LIMIT;
|
||||||
|
if (jobsQuery && rows.length < DEFAULT_QUERY_ROWS * MIN_ROW_SHARE) {
|
||||||
|
throw new GaiaAnswerError(
|
||||||
|
`Gaia returned ${rows.length} rows, not the ~${DEFAULT_QUERY_ROWS} this query holds — the answer was cut short, it was an error page ` +
|
||||||
|
`served with a 200, or the query was edited without updating DEFAULT_QUERY_ROWS; delete tools/etl/.cache/${cacheKey} once the archive answers properly`
|
||||||
|
);
|
||||||
|
}
|
||||||
|
// Not gated on `jobsQuery` like the floor above it: the only ways to reach this cap are the
|
||||||
|
// overrides that *widen* the query, and they are exactly when it is worth saying. What it must
|
||||||
|
// not fire on is a deliberately smaller slice, where filling the limit is the whole point.
|
||||||
|
if (ROW_LIMIT >= DEFAULT_ROW_LIMIT && rows.length >= ROW_LIMIT) {
|
||||||
|
throw new GaiaAnswerError(`Gaia returned the query's own ${ROW_LIMIT}-row limit, so it is the limit deciding what the map holds; raise ETL_GAIA_ROW_LIMIT.`);
|
||||||
|
}
|
||||||
const stars: StarRecord[] = [];
|
const stars: StarRecord[] = [];
|
||||||
|
|
||||||
rows.forEach((row, index) => {
|
rows.forEach((row, index) => {
|
||||||
|
|||||||
Reference in New Issue
Block a user