Skip to content

Commit cb6c813

Browse files
authored
Merge pull request #149 from TheLinuxGuy-ssh/main
Space-Track to PostgreSQL Satellite Synchronization Pipeline
2 parents 9a44e5a + 778232f commit cb6c813

2 files changed

Lines changed: 117 additions & 35 deletions

File tree

backend/models/db_models.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -418,7 +418,7 @@ class SpaceWeather(Base):
418418
k_index: Mapped[Optional[int]] = mapped_column(Integer)
419419
description: Mapped[Optional[str]] = mapped_column(Text)
420420
recorded_at: Mapped[datetime.datetime] = mapped_column(
421-
DateTime(timezone=True), server_default=func.now(), index=True
421+
DateTime(timezone=True), server_default=func.now()
422422
)
423423

424424
@property

backend/orbital/spacetrack.py

Lines changed: 116 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -103,22 +103,28 @@ def providers(self, chain: "ProviderChain") -> None:
103103

104104
def authenticate(self) -> bool:
105105
if self._authenticated:
106+
logger.debug("[SpaceTrack] Using cached session (already authenticated)")
106107
return True
107108
if not self.username or not self.password:
108-
logger.warning("SPACETRACK_USERNAME / SPACETRACK_PASSWORD not set.")
109+
logger.warning("[SpaceTrack] SPACETRACK_USERNAME / SPACETRACK_PASSWORD not set — Space-Track unavailable")
109110
return False
110111
try:
112+
logger.info("[SpaceTrack] Authenticating with Space-Track.org ...")
111113
resp = self.client.post(
112114
f"{self.base_url}/ajaxauth/login",
113115
data={"identity": self.username, "password": self.password},
114116
)
115117
if resp.status_code == 200 and "spacetrack_session" in self.client.cookies:
116118
self._authenticated = True
119+
logger.info("[SpaceTrack] Authentication successful")
117120
return True
118121
logger.error(f"[SpaceTrack] Auth failed — HTTP {resp.status_code}")
119122
return False
123+
except httpx.RequestError as exc:
124+
logger.error(f"[SpaceTrack] Network error during auth (Space-Track may be down): {exc}")
125+
return False
120126
except Exception as exc:
121-
logger.error(f"[SpaceTrack] Auth exception: {exc}")
127+
logger.error(f"[SpaceTrack] Auth exception: {exc}", exc_info=True)
122128
return False
123129

124130
def _reset_auth(self):
@@ -145,24 +151,39 @@ def _build_group_path(self, group: str, limit: int) -> Optional[str]:
145151
@_retry(HTTP_MAX_ATTEMPTS, HTTP_BASE_DELAY, HTTP_MAX_DELAY,
146152
retry_on=(httpx.RequestError, httpx.HTTPStatusError))
147153
def _send_request(self, url: str) -> List[Dict[str, Any]]:
154+
logger.debug(f"[SpaceTrack] GET {url}")
148155
resp = self.client.get(url)
149156
resp.raise_for_status()
150157
data = resp.json()
151158
if not isinstance(data, list):
152159
self._reset_auth()
160+
logger.warning(f"[SpaceTrack] Unexpected response type {type(data)} — resetting auth")
153161
raise ValueError(f"Unexpected Space-Track response type: {type(data)}")
162+
logger.debug(f"[SpaceTrack] Response: {len(data)} records")
154163
return data
155164

156165
def fetch_group_json(self, group: str, limit: int = 500) -> List[Dict[str, Any]]:
157-
if not self.authenticate():
158-
return []
159166
path = self._build_group_path(group, limit)
160167
if not path:
168+
logger.warning(f"[SpaceTrack] Unknown group '{group}' — no query path defined")
161169
return []
170+
171+
if not self.authenticate():
172+
logger.warning(f"[SpaceTrack] Skipping fetch for group '{group}' — not authenticated")
173+
return []
174+
162175
try:
163-
return self._send_request(f"{self.base_url}{path}")
176+
data = self._send_request(f"{self.base_url}{path}")
177+
logger.info(f"[SpaceTrack] Fetched {len(data)} records for group '{group}'")
178+
return data
179+
except httpx.HTTPStatusError as exc:
180+
logger.error(f"[SpaceTrack] HTTP {exc.response.status_code} for group '{group}': {exc}")
181+
return []
182+
except httpx.RequestError as exc:
183+
logger.error(f"[SpaceTrack] Network error for group '{group}': {exc}")
184+
return []
164185
except Exception as exc:
165-
logger.error(f"[SpaceTrack] Fetch failed for group '{group}': {exc}")
186+
logger.error(f"[SpaceTrack] Fetch failed for group '{group}': {exc}", exc_info=True)
166187
return []
167188

168189
def fetch_by_catalog_json(self, catalog_number: str) -> Optional[Dict[str, Any]]:
@@ -172,7 +193,12 @@ def fetch_by_catalog_json(self, catalog_number: str) -> Optional[Dict[str, Any]]
172193
f"{catalog_number}/format/json")
173194
try:
174195
data = self._send_request(url)
175-
return data[0] if data else None
196+
result = data[0] if data else None
197+
if result:
198+
logger.info(f"[SpaceTrack] fetch_by_catalog({catalog_number}): found record")
199+
else:
200+
logger.warning(f"[SpaceTrack] fetch_by_catalog({catalog_number}): no record found")
201+
return result
176202
except Exception as exc:
177203
logger.error(f"[SpaceTrack] fetch_by_catalog({catalog_number}) failed: {exc}")
178204
return None
@@ -199,7 +225,10 @@ def _gp_to_doc(
199225
epoch = _parse_epoch(rec.get("EPOCH", ""))
200226
tle1 = rec.get("TLE_LINE1") or rec.get("LINE1")
201227
tle2 = rec.get("TLE_LINE2") or rec.get("LINE2")
202-
now = datetime.datetime.utcnow().isoformat()
228+
now = datetime.datetime.now(datetime.timezone.utc)
229+
230+
if not norad_id:
231+
logger.warning(f"[SpaceTrack] Record missing NORAD_CAT_ID: {rec.get('OBJECT_NAME', '?')}")
203232

204233
return {
205234
"noradId": norad_id,
@@ -228,23 +257,31 @@ def _gp_to_doc(
228257

229258
def _bulk_upsert(
230259
self, db: Session, is_debris: bool, docs: List[Dict[str, Any]]
231-
) -> Tuple[int, List[str]]:
260+
) -> Tuple[int, int, List[str], List[str]]:
232261
"""
233262
Upsert using PostgreSQL INSERT ... ON CONFLICT (noradId) DO UPDATE SET ...
234263
Falls back to individual merge for non-PostgreSQL dialects (e.g. SQLite in tests).
264+
265+
Returns (inserted, updated, failed_ids, skipped_ids).
266+
Skipped ids are records that already exist and were updated in place.
235267
"""
236268
if not docs:
237-
return 0, []
269+
return 0, 0, [], []
238270

239271
model = Debris if is_debris else Satellite
272+
model_cols = {c.name for c in model.__table__.columns}
240273
failed: List[str] = []
241-
written = 0
274+
inserted = 0
275+
updated = 0
242276

243277
dialect = db.bind.dialect.name if db.bind else "postgresql"
244278

245279
if dialect == "postgresql":
280+
filtered_docs = []
281+
for doc in docs:
282+
filtered_docs.append({k: v for k, v in doc.items() if k in model_cols})
246283
try:
247-
stmt = pg_insert(model).values(docs)
284+
stmt = pg_insert(model).values(filtered_docs)
248285
update_cols = {
249286
c.name: c
250287
for c in stmt.excluded
@@ -254,38 +291,45 @@ def _bulk_upsert(
254291
index_elements=["noradId"],
255292
set_=update_cols,
256293
)
257-
db.execute(stmt)
294+
result = db.execute(stmt)
258295
db.commit()
259-
written = len(docs)
296+
inserted = result.rowcount
260297
except Exception as exc:
261298
db.rollback()
262-
logger.error(f"[SpaceTrack] Bulk upsert failed: {exc}")
299+
logger.error(f"[SpaceTrack] Bulk upsert failed: {exc}", exc_info=True)
263300
failed = [d.get("noradId", "?") for d in docs]
264301
else:
265-
# SQLite / other: row-by-row merge
266302
for doc in docs:
267303
try:
304+
safe_doc = {k: v for k, v in doc.items() if k in model_cols}
268305
existing = db.query(model).filter(model.noradId == doc["noradId"]).first()
269306
if existing:
270-
for k, v in doc.items():
307+
for k, v in safe_doc.items():
271308
if k not in ("id", "noradId", "createdAt") and hasattr(existing, k):
272309
setattr(existing, k, v)
310+
updated += 1
273311
else:
274-
db.add(model(**doc))
275-
written += 1
312+
db.add(model(**safe_doc))
313+
inserted += 1
276314
except Exception as exc:
277315
logger.warning(f"[SpaceTrack] Row upsert failed for {doc.get('noradId')}: {exc}")
278316
failed.append(doc.get("noradId", "?"))
279317
db.commit()
280318

281-
return written, failed
319+
logger.info(
320+
f"[SpaceTrack] {model.__tablename__} upsert: "
321+
f"{inserted} inserted, {updated} updated, "
322+
f"{len(failed)} failed, {len(docs) - inserted - updated - len(failed)} skipped"
323+
)
324+
return inserted, updated, failed, []
282325

283326
def _ensure_db_connection(self, db: Session) -> bool:
284327
try:
285328
db.execute(__import__("sqlalchemy").text("SELECT 1"))
329+
logger.debug("[SpaceTrack] Database connection OK")
286330
return True
287331
except Exception as exc:
288-
logger.error(f"[SpaceTrack] DB unreachable: {exc}")
332+
logger.error(f"[SpaceTrack] Database unreachable: {exc}", exc_info=True)
289333
return False
290334

291335
# ------------------------------------------------------------------
@@ -299,16 +343,28 @@ def sync_group(
299343
object_type_override: Optional[str] = None,
300344
limit: Optional[int] = None,
301345
) -> Dict[str, Any]:
346+
logger.info(f"[SpaceTrack] Starting sync for group '{group}' (limit={limit or 500})")
347+
302348
if not self._ensure_db_connection(db):
349+
logger.error(f"[SpaceTrack] DB unreachable for group '{group}'")
303350
return {"group": group, "fetched": 0, "parsed": 0, "upserted": 0,
351+
"inserted": 0, "updated": 0,
304352
"failed": 0, "source": None, "errors": ["db_unreachable"]}
305353

354+
# Check authentication
355+
if not self.authenticate():
356+
logger.warning(f"[SpaceTrack] Not authenticated for group '{group}' — falling through provider chain")
357+
306358
try:
307359
records, source = self.providers.fetch_group(group, limit or 500, db=db)
360+
logger.info(f"[SpaceTrack] Group '{group}': fetched {len(records)} records from '{source}'")
308361
except AllProvidersFailedError as exc:
362+
logger.error(f"[SpaceTrack] All providers failed for group '{group}': {exc.failures}")
309363
return {
310-
"group": group, "fetched": 0, "parsed": 0, "upserted": 0, "failed": 0,
311-
"source": None, "errors": ["all_providers_failed"], "provider_failures": exc.failures,
364+
"group": group, "fetched": 0, "parsed": 0, "upserted": 0,
365+
"inserted": 0, "updated": 0,
366+
"failed": 0, "source": None,
367+
"errors": ["all_providers_failed"], "provider_failures": exc.failures,
312368
}
313369

314370
is_debris = (object_type_override == "DEBRIS" or group in ("analyst", "debris"))
@@ -326,34 +382,53 @@ def sync_group(
326382
parse_errors.append(f"{rec.get('NORAD_CAT_ID', '?')}:{exc}")
327383

328384
try:
329-
upserted, failed_ids = self._bulk_upsert(db, is_debris, docs)
385+
inserted, updated, failed_ids, _ = self._bulk_upsert(db, is_debris, docs)
386+
total_upserted = inserted + updated
387+
total_failed = len(failed_ids) + len(parse_errors)
388+
logger.info(
389+
f"[SpaceTrack] Group '{group}' — fetched={len(records)}, "
390+
f"parsed={len(docs)}, inserted={inserted}, updated={updated}, "
391+
f"failed={total_failed}, source={source}"
392+
)
330393
return {
331394
"group": group, "fetched": len(records), "parsed": len(docs),
332-
"upserted": upserted, "failed": len(failed_ids) + len(parse_errors),
333-
"source": source, "errors": failed_ids + parse_errors,
395+
"upserted": total_upserted, "inserted": inserted, "updated": updated,
396+
"failed": total_failed, "source": source,
397+
"errors": failed_ids + parse_errors,
334398
}
335399
except Exception as exc:
336-
logger.error(f"[Ingest] Upsert failure for '{group}': {exc}")
400+
logger.error(f"[Ingest] Upsert failure for '{group}': {exc}", exc_info=True)
337401
return {
338402
"group": group, "fetched": len(records), "parsed": len(docs),
339-
"upserted": 0, "failed": len(docs) + len(parse_errors),
403+
"upserted": 0, "inserted": 0, "updated": 0,
404+
"failed": len(docs) + len(parse_errors),
340405
"source": source, "errors": [str(exc)] + parse_errors,
341406
}
342407

343408
def sync_all_groups(self, db: Session, limit_per_group: int = 500) -> Dict[str, Any]:
344409
per_group: Dict[str, Any] = {}
345-
total_fetched = total_upserted = total_failed = 0
410+
total_fetched = total_upserted = total_failed = total_inserted = total_updated = 0
346411

347412
for group, type_override, _ in SYNC_GROUPS:
348413
status = self.sync_group(db, group, type_override, limit=limit_per_group)
349414
per_group[group] = status
350-
total_fetched += status["fetched"]
351-
total_upserted += status["upserted"]
352-
total_failed += status["failed"]
415+
total_fetched += status["fetched"]
416+
total_upserted += status.get("upserted", 0)
417+
total_inserted += status.get("inserted", 0)
418+
total_updated += status.get("updated", 0)
419+
total_failed += status["failed"]
420+
421+
logger.info(
422+
f"[SpaceTrack] Sync complete — "
423+
f"fetched={total_fetched}, inserted={total_inserted}, "
424+
f"updated={total_updated}, failed={total_failed}"
425+
)
353426

354427
return {
355428
"total_fetched": total_fetched,
356429
"total_upserted": total_upserted,
430+
"total_inserted": total_inserted,
431+
"total_updated": total_updated,
357432
"total_failed": total_failed,
358433
"groups": per_group,
359434
}
@@ -366,9 +441,16 @@ def sync_by_catalog(self, db: Session, catalog_number: str) -> Optional[Dict]:
366441
is_debris = obj_type == "DEBRIS"
367442
doc = self._gp_to_doc(rec, obj_type)
368443
try:
369-
self._bulk_upsert(db, is_debris, [doc])
444+
inserted, updated, failed_ids, _ = self._bulk_upsert(db, is_debris, [doc])
445+
if failed_ids:
446+
logger.error(f"[SpaceTrack] Single upsert failed for {catalog_number}: {failed_ids}")
447+
else:
448+
logger.info(
449+
f"[SpaceTrack] sync_by_catalog({catalog_number}): "
450+
f"{'inserted' if inserted else 'updated'}"
451+
)
370452
except Exception as exc:
371-
logger.error(f"[SpaceTrack] Single upsert failed for {catalog_number}: {exc}")
453+
logger.error(f"[SpaceTrack] Single upsert failed for {catalog_number}: {exc}", exc_info=True)
372454
model = Debris if is_debris else Satellite
373455
return db.query(model).filter(model.noradId == catalog_number).first()
374456

0 commit comments

Comments
 (0)