SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
6 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
12.2 KB · 181 lines typescript
Raw Blame History
1/**2 * OpenStreetMap via the Overpass API — features tagged telecom=data_center or building=data_center.3 * License: ODbL 1.0, © OpenStreetMap contributors (attribution required).4 *5 * The world is split into non-overlapping bounding boxes (one Overpass query each, `out center tags;`).6 * Queries run sequentially, spaced by `params.minDelayMs`, with a resumable cursor in connector state so a7 * run interrupted by load shedding (429/504) continues where it stopped. `params.maxBboxesPerRun` bounds8 * one run. Primary endpoint + fallback endpoint; both governed by the Overpass commons usage guidance9 * (≤ ~10 000 requests/day, ≤ ~1 GB/day per user) — we send ~13 requests per full pass.10 */11import type { ConnectorContext, DiscoveredUrl, ExtractedRecord, RawDocument } from "@dci/connectors";12import type { Fact, FacilityStatus, NormalizedEntity, NormalizedFacility, Provenance } from "@dci/core";13import { inferFacilityType, parseMw, parsePartialDate, validLatLng } from "@dci/core";14import { DatasetConnector, httpUrl, iso2, num, sleep, str, uniqNames } from "./shared.js";1516/** [south, west, north, east] — a partition of the inhabited world (lat −60…85). Edges are shared, not overlapped. */17export const WORLD_BBOXES: Array<{ id: string; bbox: [number, number, number, number] }> = [18  { id: "na-west", bbox: [14, -180, 85, -100] },19  { id: "na-east", bbox: [14, -100, 85, -50] },20  { id: "atlantic-north", bbox: [35, -50, 85, -12] },21  { id: "latam", bbox: [-60, -180, 14, -30] },22  { id: "eu-west", bbox: [35, -12, 85, 15] },23  { id: "eu-east", bbox: [35, 15, 85, 45] },24  { id: "africa-west", bbox: [-40, -30, 35, 35] },25  { id: "africa-east", bbox: [-40, 35, 12, 60] },26  { id: "mena", bbox: [12, 35, 35, 60] },27  { id: "central-asia", bbox: [35, 45, 85, 100] },28  { id: "south-asia", bbox: [-10, 60, 35, 100] },29  { id: "east-asia", bbox: [20, 100, 85, 180] },30  { id: "sea-oceania", bbox: [-60, 100, 20, 180] },31];3233interface OsmElement { type: "node" | "way" | "relation"; id: number; lat?: number; lon?: number; center?: { lat: number; lon: number }; tags?: Record<string, string> }34interface OverpassResponse { elements: OsmElement[]; remark?: string; osm3s?: { timestamp_osm_base?: string; copyright?: string } }3536const STATE_CURSOR = "osm-overpass:cursor";37/** Countries writing house number before street; used only to assemble a readable address string. */38const NUMBER_FIRST = new Set(["US", "CA", "GB", "IE", "AU", "NZ", "FR", "IN", "ZA", "SG", "HK", "MY", "PH"]);3940export class OsmOverpassConnector extends DatasetConnector {41  private endpoints(): string[] {42    const primary = String(this.param("endpoint", "https://overpass-api.de/api/interpreter"));43    const fallback = this.param<string | null>("fallbackEndpoint", "https://overpass.private.coffee/api/interpreter");44    return fallback && fallback !== primary ? [primary, fallback] : [primary];45  }46  private query(b: [number, number, number, number]): string {47    const bbox = b.join(",");48    const timeout = this.param<number>("timeoutSec", 180);49    const maxsize = this.param<number>("maxsizeBytes", 536_870_912);50    return `[out:json][timeout:${timeout}][maxsize:${maxsize}];(nwr["telecom"="data_center"](${bbox});nwr["building"="data_center"](${bbox}););out center tags;`;51  }52  private urlFor(endpoint: string, q: string): string { return `${endpoint}?data=${encodeURIComponent(q)}`; }5354  /** Next `maxBboxesPerRun` tiles from the persisted cursor (wraps around after a full pass). */55  async discover(ctx: ConnectorContext): Promise<DiscoveredUrl[]> {56    const tiles = WORLD_BBOXES;57    const max = Math.max(1, Math.min(tiles.length, this.param<number>("maxBboxesPerRun", tiles.length)));58    const cursor = ((await ctx.getState<number>(STATE_CURSOR)) ?? 0) % tiles.length;59    const [primary] = this.endpoints();60    const out: DiscoveredUrl[] = [];61    for (let i = 0; i < max; i++) {62      const idx = (cursor + i) % tiles.length;63      const t = tiles[idx]!;64      out.push({ url: this.urlFor(primary!, this.query(t.bbox)), group: "dataset", pageType: "dataset", priority: 100 - i, meta: { tile: t.id, bbox: t.bbox, index: idx } });65    }66    ctx.log("info", `Overpass: ${out.length}/${tiles.length} tiles this run, starting at #${cursor} (${tiles[cursor]!.id})`);67    return out;68  }6970  /** Sequential, spaced queries; on load shedding (429/504/timeout) back off then try the fallback endpoint. */71  async fetch(ctx: ConnectorContext, u: DiscoveredUrl): Promise<RawDocument> {72    const q = new URL(u.url).searchParams.get("data") ?? "";73    const minDelay = this.param<number>("minDelayMs", 10_000);74    const backoff = this.param<number>("backoffMs", 30_000);75    const timeoutMs = (this.param<number>("timeoutSec", 180) + 30) * 1000;76    const lastAt = this.lastQueryAt;77    if (lastAt) await sleep(minDelay - (Date.now() - lastAt));78    // attempt sequence: primary → fallback → primary … (alternating), each retry after a backoff; load shedding only79    const eps = this.endpoints();80    const maxAttempts = Math.max(1, this.param<number>("maxAttempts", 3));81    let last: RawDocument | null = null;82    for (let attempt = 0; attempt < maxAttempts; attempt++) {83      const ep = eps[attempt % eps.length]!;84      const url = this.urlFor(ep, q);85      this.lastQueryAt = Date.now();86      const doc = await ctx.fetch(url, { group: u.group, accept: "application/json", timeoutMs, maxBytes: 200 * 1024 * 1024 });87      doc.group = u.group; doc.pageType = "dataset"; doc.meta = { ...(doc.meta ?? {}), ...(u.meta ?? {}), endpoint: ep, attempt: attempt + 1 };88      if (!doc.error && doc.status === 200) return doc;89      last = doc;90      const shed = doc.error?.code === "timeout" || doc.error?.code === "connection" || doc.status === 429 || doc.status === 504 || doc.status === 503 || doc.status === 502;91      const why = doc.error ? `${doc.error.code}: ${doc.error.message}` : `HTTP ${doc.status}`;92      if (!shed || attempt === maxAttempts - 1) { ctx.log("warn", `Overpass tile ${tileLabel(u)} failed at ${ep} (${why}); giving up for this run — cursor not advanced`); break; }93      const wait = backoff * (attempt + 1);94      ctx.log("warn", `Overpass tile ${tileLabel(u)}: ${why} at ${ep}; backing off ${wait} ms then trying ${eps[(attempt + 1) % eps.length]}`);95      await sleep(wait);96    }97    return last!;98  }99  private lastQueryAt = 0;100101  async extract(ctx: ConnectorContext, doc: RawDocument): Promise<ExtractedRecord[]> {102    if (doc.error || doc.notModified || doc.status >= 400 || !doc.body.length) return [];103    let json: OverpassResponse | null = null;104    try { json = JSON.parse(doc.text) as OverpassResponse; } catch { ctx.log("warn", `Overpass tile ${String(doc.meta?.tile)}: invalid JSON (${doc.body.length} bytes)`); return []; }105    if (json.remark) ctx.log("warn", `Overpass remark on tile ${String(doc.meta?.tile)}: ${json.remark}`);106    if (json.remark && /runtime error|timed out|out of memory/i.test(json.remark) && !json.elements?.length) return []; // partial/failed result: do not advance cursor107    const els = Array.isArray(json.elements) ? json.elements : [];108    let skipped = 0;109    const out: ExtractedRecord[] = [];110    for (const el of els) {111      const tags = el.tags ?? {};112      if (!str(tags.name) && !str(tags.operator)) { skipped++; continue; }113      const lat = el.type === "node" ? el.lat : el.center?.lat;114      const lon = el.type === "node" ? el.lon : el.center?.lon;115      out.push({ kind: "facility", key: `osm:${el.type}/${el.id}`, data: { type: el.type, id: el.id, lat, lon, tags, osmBase: json.osm3s?.timestamp_osm_base ?? null }, methods: { geo: el.type === "node" ? "osm:node" : "osm:center", _all: "osm-tags" }, certainty: 1, url: doc.finalUrl, pageType: "dataset" });116    }117    ctx.log("info", `Overpass tile ${String(doc.meta?.tile)}: ${els.length} features, ${out.length} kept, ${skipped} without name/operator skipped`);118    // advance the resumable cursor only after a successful tile119    const idx = num(doc.meta?.index);120    if (idx != null && !ctx.dryRun) await ctx.setState(STATE_CURSOR, (idx + 1) % WORLD_BBOXES.length);121    return out;122  }123124  async normalize(ctx: ConnectorContext, records: ExtractedRecord[]): Promise<NormalizedEntity[]> {125    const out: NormalizedEntity[] = [];126    const seen = new Set<string>();127    for (const r of records) {128      if (seen.has(r.key)) continue; seen.add(r.key);129      const d = r.data as { type: string; id: number; lat?: number; lon?: number; tags: Record<string, string> };130      const t = d.tags;131      const prov = (method: string, extra: Partial<Provenance> = {}): Provenance => ctx.provenance(r.url, { method, extractorVersion: this.parserVersion, note: `osm ${d.type}/${d.id}`, ...extra });132      const operator = str(t.operator) ?? str(t.brand);133      const ref = str(t.ref);134      const rawName = str(t.name) ?? str(t["name:en"]);135      const name = rawName ?? (operator ? (ref ? `${operator} ${ref}` : operator) : null);136      if (!name) continue;137      const country = iso2(t["addr:country"]);138      const street = str(t["addr:street"]), house = str(t["addr:housenumber"]);139      const address = street && house ? (!country || NUMBER_FIRST.has(country) ? `${house} ${street}` : `${street} ${house}`) : street ?? str(t["addr:housename"]);140      const lat = num(d.lat), lng = num(d.lon);141      // opening date: start_date / opening_date only; years before 1950 are far more likely the building's than the data center's → dropped142      const openedRaw = str(t.start_date) ?? str(t.opening_date);143      let openedOn = openedRaw ? parsePartialDate(openedRaw) : null;144      if (openedOn && Number(openedOn.slice(0, 4)) < 1950) { ctx.log("debug", `osm ${d.type}/${d.id}: ignoring start_date ${openedRaw} (< 1950)`); openedOn = null; }145      // power: only the explicit data_center:power tag and only when it states MW146      const powerRaw = str(t["data_center:power"]);147      const totalPowerMw = powerRaw && /\b(mw|megawatts?)\b/i.test(powerRaw) ? parseMw(powerRaw) : null;148      const facts: Fact[] = [];149      if (totalPowerMw != null) facts.push({ field: "totalPowerMw", value: totalPowerMw, provenance: prov("osm-tag:data_center:power") });150      if (openedOn) facts.push({ field: "openedOn", value: openedOn, provenance: prov(`osm-tag:${str(t.start_date) ? "start_date" : "opening_date"}`) });151      let status: FacilityStatus | null = "operational";152      if (t.construction || t.building === "construction" || t["telecom"] === "construction") status = "under_construction";153      if (t.disused === "yes" || t.abandoned === "yes" || t["disused:telecom"] || t["abandoned:telecom"]) status = "closed";154      const description = str(t.description) ?? str(t["description:en"]);155      const externalIds: Record<string, string | number> = { osm: `${d.type}/${d.id}` };156      if (/^Q\d+$/.test(t.wikidata ?? "")) externalIds.wikidata = t.wikidata!;157      if (/^Q\d+$/.test(t["operator:wikidata"] ?? "")) externalIds.operator_wikidata = t["operator:wikidata"]!;158      if (ref) externalIds.osm_ref = ref;159      const ent: NormalizedFacility = {160        entityType: "facility", key: r.key, name,161        aliases: uniqNames([rawName ? t["name:en"] : null, t.alt_name, t.short_name, t.official_name, t.old_name, t.brand !== operator ? t.brand : null, t["operator:short"]], [name, operator]),162        operatorName: operator, ownerName: str(t.owner),163        address: address ?? null, city: str(t["addr:city"]), regionName: str(t["addr:state"]) ?? str(t["addr:province"]), countryIso2: country, postalCode: str(t["addr:postcode"]),164        geo: validLatLng(lat, lng) ? { lat: lat!, lng: lng!, precision: "exact", source: "community:osm" } : null,165        status, facilityType: inferFacilityType(`${name} ${description ?? ""} ${t.operator ?? ""}`),166        totalPowerMw, openedOn, website: httpUrl(t.website) ?? httpUrl(t["contact:website"]) ?? httpUrl(t.url),167        externalIds, description, facts, provenance: prov(d.type === "node" ? "osm-tags+node" : "osm-tags+center"),168      };169      if (!ent.geo) { ctx.log("debug", `osm ${d.type}/${d.id}: no usable coordinates`); }170      out.push(ent);171    }172    return out;173  }174}175176/** Tile label for logs — `meta` is not persisted by the document registry, so fall back to the bbox in the URL. */177function tileLabel(u: DiscoveredUrl): string {178  if (u.meta?.tile) return String(u.meta.tile);179  try { return new URL(u.url).searchParams.get("data")?.match(/\(([-\d.,]+)\)/)?.[1] ?? "?"; } catch { return "?"; }180}181