/** * OpenStreetMap via the Overpass API — features tagged telecom=data_center or building=data_center. * License: ODbL 1.0, © OpenStreetMap contributors (attribution required). * * The world is split into non-overlapping bounding boxes (one Overpass query each, `out center tags;`). * Queries run sequentially, spaced by `params.minDelayMs`, with a resumable cursor in connector state so a * run interrupted by load shedding (429/504) continues where it stopped. `params.maxBboxesPerRun` bounds * one run. Primary endpoint + fallback endpoint; both governed by the Overpass commons usage guidance * (≤ ~10 000 requests/day, ≤ ~1 GB/day per user) — we send ~13 requests per full pass. */ import type { ConnectorContext, DiscoveredUrl, ExtractedRecord, RawDocument } from "@dci/connectors"; import type { Fact, FacilityStatus, NormalizedEntity, NormalizedFacility, Provenance } from "@dci/core"; import { inferFacilityType, parseMw, parsePartialDate, validLatLng } from "@dci/core"; import { DatasetConnector, httpUrl, iso2, num, sleep, str, uniqNames } from "./shared.js"; /** [south, west, north, east] — a partition of the inhabited world (lat −60…85). Edges are shared, not overlapped. */ export const WORLD_BBOXES: Array<{ id: string; bbox: [number, number, number, number] }> = [ { id: "na-west", bbox: [14, -180, 85, -100] }, { id: "na-east", bbox: [14, -100, 85, -50] }, { id: "atlantic-north", bbox: [35, -50, 85, -12] }, { id: "latam", bbox: [-60, -180, 14, -30] }, { id: "eu-west", bbox: [35, -12, 85, 15] }, { id: "eu-east", bbox: [35, 15, 85, 45] }, { id: "africa-west", bbox: [-40, -30, 35, 35] }, { id: "africa-east", bbox: [-40, 35, 12, 60] }, { id: "mena", bbox: [12, 35, 35, 60] }, { id: "central-asia", bbox: [35, 45, 85, 100] }, { id: "south-asia", bbox: [-10, 60, 35, 100] }, { id: "east-asia", bbox: [20, 100, 85, 180] }, { id: "sea-oceania", bbox: [-60, 100, 20, 180] }, ]; interface OsmElement { type: "node" | "way" | "relation"; id: number; lat?: number; lon?: number; center?: { lat: number; lon: number }; tags?: Record } interface OverpassResponse { elements: OsmElement[]; remark?: string; osm3s?: { timestamp_osm_base?: string; copyright?: string } } const STATE_CURSOR = "osm-overpass:cursor"; /** Countries writing house number before street; used only to assemble a readable address string. */ const NUMBER_FIRST = new Set(["US", "CA", "GB", "IE", "AU", "NZ", "FR", "IN", "ZA", "SG", "HK", "MY", "PH"]); export class OsmOverpassConnector extends DatasetConnector { private endpoints(): string[] { const primary = String(this.param("endpoint", "https://overpass-api.de/api/interpreter")); const fallback = this.param("fallbackEndpoint", "https://overpass.private.coffee/api/interpreter"); return fallback && fallback !== primary ? [primary, fallback] : [primary]; } private query(b: [number, number, number, number]): string { const bbox = b.join(","); const timeout = this.param("timeoutSec", 180); const maxsize = this.param("maxsizeBytes", 536_870_912); return `[out:json][timeout:${timeout}][maxsize:${maxsize}];(nwr["telecom"="data_center"](${bbox});nwr["building"="data_center"](${bbox}););out center tags;`; } private urlFor(endpoint: string, q: string): string { return `${endpoint}?data=${encodeURIComponent(q)}`; } /** Next `maxBboxesPerRun` tiles from the persisted cursor (wraps around after a full pass). */ async discover(ctx: ConnectorContext): Promise { const tiles = WORLD_BBOXES; const max = Math.max(1, Math.min(tiles.length, this.param("maxBboxesPerRun", tiles.length))); const cursor = ((await ctx.getState(STATE_CURSOR)) ?? 0) % tiles.length; const [primary] = this.endpoints(); const out: DiscoveredUrl[] = []; for (let i = 0; i < max; i++) { const idx = (cursor + i) % tiles.length; const t = tiles[idx]!; 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 } }); } ctx.log("info", `Overpass: ${out.length}/${tiles.length} tiles this run, starting at #${cursor} (${tiles[cursor]!.id})`); return out; } /** Sequential, spaced queries; on load shedding (429/504/timeout) back off then try the fallback endpoint. */ async fetch(ctx: ConnectorContext, u: DiscoveredUrl): Promise { const q = new URL(u.url).searchParams.get("data") ?? ""; const minDelay = this.param("minDelayMs", 10_000); const backoff = this.param("backoffMs", 30_000); const timeoutMs = (this.param("timeoutSec", 180) + 30) * 1000; const lastAt = this.lastQueryAt; if (lastAt) await sleep(minDelay - (Date.now() - lastAt)); // attempt sequence: primary → fallback → primary … (alternating), each retry after a backoff; load shedding only const eps = this.endpoints(); const maxAttempts = Math.max(1, this.param("maxAttempts", 3)); let last: RawDocument | null = null; for (let attempt = 0; attempt < maxAttempts; attempt++) { const ep = eps[attempt % eps.length]!; const url = this.urlFor(ep, q); this.lastQueryAt = Date.now(); const doc = await ctx.fetch(url, { group: u.group, accept: "application/json", timeoutMs, maxBytes: 200 * 1024 * 1024 }); doc.group = u.group; doc.pageType = "dataset"; doc.meta = { ...(doc.meta ?? {}), ...(u.meta ?? {}), endpoint: ep, attempt: attempt + 1 }; if (!doc.error && doc.status === 200) return doc; last = doc; const shed = doc.error?.code === "timeout" || doc.error?.code === "connection" || doc.status === 429 || doc.status === 504 || doc.status === 503 || doc.status === 502; const why = doc.error ? `${doc.error.code}: ${doc.error.message}` : `HTTP ${doc.status}`; 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; } const wait = backoff * (attempt + 1); ctx.log("warn", `Overpass tile ${tileLabel(u)}: ${why} at ${ep}; backing off ${wait} ms then trying ${eps[(attempt + 1) % eps.length]}`); await sleep(wait); } return last!; } private lastQueryAt = 0; async extract(ctx: ConnectorContext, doc: RawDocument): Promise { if (doc.error || doc.notModified || doc.status >= 400 || !doc.body.length) return []; let json: OverpassResponse | null = null; 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 []; } if (json.remark) ctx.log("warn", `Overpass remark on tile ${String(doc.meta?.tile)}: ${json.remark}`); if (json.remark && /runtime error|timed out|out of memory/i.test(json.remark) && !json.elements?.length) return []; // partial/failed result: do not advance cursor const els = Array.isArray(json.elements) ? json.elements : []; let skipped = 0; const out: ExtractedRecord[] = []; for (const el of els) { const tags = el.tags ?? {}; if (!str(tags.name) && !str(tags.operator)) { skipped++; continue; } const lat = el.type === "node" ? el.lat : el.center?.lat; const lon = el.type === "node" ? el.lon : el.center?.lon; 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" }); } ctx.log("info", `Overpass tile ${String(doc.meta?.tile)}: ${els.length} features, ${out.length} kept, ${skipped} without name/operator skipped`); // advance the resumable cursor only after a successful tile const idx = num(doc.meta?.index); if (idx != null && !ctx.dryRun) await ctx.setState(STATE_CURSOR, (idx + 1) % WORLD_BBOXES.length); return out; } async normalize(ctx: ConnectorContext, records: ExtractedRecord[]): Promise { const out: NormalizedEntity[] = []; const seen = new Set(); for (const r of records) { if (seen.has(r.key)) continue; seen.add(r.key); const d = r.data as { type: string; id: number; lat?: number; lon?: number; tags: Record }; const t = d.tags; const prov = (method: string, extra: Partial = {}): Provenance => ctx.provenance(r.url, { method, extractorVersion: this.parserVersion, note: `osm ${d.type}/${d.id}`, ...extra }); const operator = str(t.operator) ?? str(t.brand); const ref = str(t.ref); const rawName = str(t.name) ?? str(t["name:en"]); const name = rawName ?? (operator ? (ref ? `${operator} ${ref}` : operator) : null); if (!name) continue; const country = iso2(t["addr:country"]); const street = str(t["addr:street"]), house = str(t["addr:housenumber"]); const address = street && house ? (!country || NUMBER_FIRST.has(country) ? `${house} ${street}` : `${street} ${house}`) : street ?? str(t["addr:housename"]); const lat = num(d.lat), lng = num(d.lon); // opening date: start_date / opening_date only; years before 1950 are far more likely the building's than the data center's → dropped const openedRaw = str(t.start_date) ?? str(t.opening_date); let openedOn = openedRaw ? parsePartialDate(openedRaw) : null; if (openedOn && Number(openedOn.slice(0, 4)) < 1950) { ctx.log("debug", `osm ${d.type}/${d.id}: ignoring start_date ${openedRaw} (< 1950)`); openedOn = null; } // power: only the explicit data_center:power tag and only when it states MW const powerRaw = str(t["data_center:power"]); const totalPowerMw = powerRaw && /\b(mw|megawatts?)\b/i.test(powerRaw) ? parseMw(powerRaw) : null; const facts: Fact[] = []; if (totalPowerMw != null) facts.push({ field: "totalPowerMw", value: totalPowerMw, provenance: prov("osm-tag:data_center:power") }); if (openedOn) facts.push({ field: "openedOn", value: openedOn, provenance: prov(`osm-tag:${str(t.start_date) ? "start_date" : "opening_date"}`) }); let status: FacilityStatus | null = "operational"; if (t.construction || t.building === "construction" || t["telecom"] === "construction") status = "under_construction"; if (t.disused === "yes" || t.abandoned === "yes" || t["disused:telecom"] || t["abandoned:telecom"]) status = "closed"; const description = str(t.description) ?? str(t["description:en"]); const externalIds: Record = { osm: `${d.type}/${d.id}` }; if (/^Q\d+$/.test(t.wikidata ?? "")) externalIds.wikidata = t.wikidata!; if (/^Q\d+$/.test(t["operator:wikidata"] ?? "")) externalIds.operator_wikidata = t["operator:wikidata"]!; if (ref) externalIds.osm_ref = ref; const ent: NormalizedFacility = { entityType: "facility", key: r.key, name, 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]), operatorName: operator, ownerName: str(t.owner), address: address ?? null, city: str(t["addr:city"]), regionName: str(t["addr:state"]) ?? str(t["addr:province"]), countryIso2: country, postalCode: str(t["addr:postcode"]), geo: validLatLng(lat, lng) ? { lat: lat!, lng: lng!, precision: "exact", source: "community:osm" } : null, status, facilityType: inferFacilityType(`${name} ${description ?? ""} ${t.operator ?? ""}`), totalPowerMw, openedOn, website: httpUrl(t.website) ?? httpUrl(t["contact:website"]) ?? httpUrl(t.url), externalIds, description, facts, provenance: prov(d.type === "node" ? "osm-tags+node" : "osm-tags+center"), }; if (!ent.geo) { ctx.log("debug", `osm ${d.type}/${d.id}: no usable coordinates`); } out.push(ent); } return out; } } /** Tile label for logs — `meta` is not persisted by the document registry, so fall back to the bbox in the URL. */ function tileLabel(u: DiscoveredUrl): string { if (u.meta?.tile) return String(u.meta.tile); try { return new URL(u.url).searchParams.get("data")?.match(/\(([-\d.,]+)\)/)?.[1] ?? "?"; } catch { return "?"; } }