Skip to content

Commit 28fc2ba

Browse files
committed
fix: eliminate DetachedInstanceError — pass plain IDs across thread boundary
Root cause: enrichment runs in a background thread with its own DB session. The old code expunged ORM objects and merged them into the main session, but merge/refresh on cross-session objects fails when PgBouncer drops the idle connection. Fix: EnrichmentResult now carries property_id and scenario_id as plain UUIDs. The thread sets these before closing its session, then the SSE handler queries fresh objects from the main session using db.get(Model, id). No ORM objects cross the thread boundary. Same pattern applied to the Bricked enrichment thread — removed the __dict__.update() hack and the broken db.refresh() call.
1 parent 0f53e30 commit 28fc2ba

2 files changed

Lines changed: 37 additions & 40 deletions

File tree

‎backend/core/property_data/service.py‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -29,8 +29,10 @@ class ProviderStatus:
2929
@dataclass
3030
class EnrichmentResult:
3131
"""Result of enriching a property from external data sources."""
32-
property: Any = None # Property model instance
33-
scenario: Any = None # AnalysisScenario model instance
32+
property: Any = None # Property model instance (only valid within creating session)
33+
scenario: Any = None # AnalysisScenario model instance (only valid within creating session)
34+
property_id: Optional[Any] = None # Plain UUID — safe to pass across thread boundaries
35+
scenario_id: Optional[Any] = None # Plain UUID — safe to pass across thread boundaries
3436
fields_populated: dict[str, Any] = field(default_factory=dict)
3537
fields_missing: list[str] = field(default_factory=list)
3638
confidence: dict[str, str] = field(default_factory=dict)
@@ -191,6 +193,7 @@ def enrich_property(
191193
if existing:
192194
logger.info("enrich_property: found existing property %s, skipping provider calls", existing.id)
193195
result.property = existing
196+
result.property_id = existing.id
194197
result.is_existing = True
195198
result.status = "existing"
196199

@@ -225,6 +228,7 @@ def enrich_property(
225228
db.add(scenario)
226229
db.flush()
227230
result.scenario = scenario
231+
result.scenario_id = scenario.id
228232
return result
229233

230234
# 3. Create Property record immediately (always created, even if providers fail)
@@ -244,6 +248,7 @@ def enrich_property(
244248
db.add(prop)
245249
db.flush() # get prop.id for DataSourceEvent FK
246250
result.property = prop
251+
result.property_id = prop.id
247252

248253
# 4. Fetch from each provider
249254
logger.info("enrich_property: created property %s, calling providers: %s", prop.id, providers)
@@ -354,6 +359,7 @@ def enrich_property(
354359
db.add(scenario)
355360
db.flush()
356361
result.scenario = scenario
362+
result.scenario_id = scenario.id
357363

358364
logger.info(
359365
"enrich_property: done → status=%s, purchase_price=%s, monthly_rent=%s, arv=%s, repair=%s, fields_missing=%s",

‎backend/routers/analysis.py‎

Lines changed: 29 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -454,10 +454,10 @@ def _enrich_sync():
454454
providers=["rentcast"],
455455
)
456456
thread_db.commit()
457-
if result.property:
458-
thread_db.expunge(result.property)
459-
if result.scenario:
460-
thread_db.expunge(result.scenario)
457+
# Don't pass ORM objects across threads — just IDs.
458+
# property_id and scenario_id are set by enrich_property().
459+
result.property = None
460+
result.scenario = None
461461
return result
462462
finally:
463463
thread_db.close()
@@ -470,23 +470,22 @@ def _enrich_sync():
470470

471471
enrichment = await asyncio.to_thread(_enrich_sync) if _use_thread else _enrich_sync()
472472

473-
# The enrichment call takes 5-30s (RentCast + Bricked). During
474-
# this time the main db session sits idle. Railway's PgBouncer
475-
# or PostgreSQL may drop idle connections. expire_all() marks
476-
# all cached objects as stale so the next attribute access or
477-
# flush triggers a fresh query on a live connection, without
478-
# detaching objects from the session (unlike rollback()).
473+
# After the thread completes, the main db session may have a
474+
# stale connection (Railway PgBouncer drops idle connections).
475+
# expire_all() forces fresh queries on next access.
479476
if _use_thread:
480477
db.expire_all()
481478

482-
# Re-attach detached ORM objects to the main session.
483-
# The thread session expunged them; merge() copies their state
484-
# into the main session so lazy loads and flushes work.
479+
# Load fresh ORM objects in the main session using plain IDs
480+
# that were safely passed across the thread boundary.
481+
from models.properties import Property as _Prop
482+
from models.analysis_scenarios import AnalysisScenario as _AS
483+
485484
if _use_thread:
486-
if enrichment.property:
487-
enrichment.property = db.merge(enrichment.property)
488-
if enrichment.scenario:
489-
enrichment.scenario = db.merge(enrichment.scenario)
485+
if enrichment.property_id:
486+
enrichment.property = db.get(_Prop, enrichment.property_id)
487+
if enrichment.scenario_id:
488+
enrichment.scenario = db.get(_AS, enrichment.scenario_id)
490489

491490
# Persist client-side geocoding data if provided
492491
if enrichment.property and (lat is not None and lng is not None):
@@ -527,25 +526,24 @@ def _enrich_sync():
527526
yield _sse("status", {"stage": "fetching_advanced_data"})
528527
try:
529528
if _use_thread:
530-
prop_id = enrichment.property.id
531-
scenario_id = enrichment.scenario.id
529+
prop_id = enrichment.property_id
530+
scenario_id = enrichment.scenario_id
532531

533532
def _bricked_sync():
534533
bdb = SessionLocal()
535534
try:
536-
bprop = bdb.get(type(enrichment.property), prop_id)
537-
bscen = bdb.get(type(enrichment.scenario), scenario_id)
535+
bprop = bdb.get(_Prop, prop_id)
536+
bscen = bdb.get(_AS, scenario_id)
538537
result = enrich_with_bricked(bprop, bscen, address, bdb)
539538
bdb.commit()
540-
bdb.expunge(bscen)
541-
enrichment.scenario.__dict__.update(bscen.__dict__)
542539
return result
543540
finally:
544541
bdb.close()
545542

546543
bricked_status = await asyncio.to_thread(_bricked_sync)
547-
# Bricked takes 15-30s — expire cached objects
548544
db.expire_all()
545+
# Re-query scenario in main session to pick up bricked data
546+
enrichment.scenario = db.get(_AS, scenario_id)
549547
else:
550548
bricked_status = enrich_with_bricked(
551549
enrichment.property, enrichment.scenario, address, db,
@@ -557,21 +555,11 @@ def _bricked_sync():
557555
)
558556
db.commit()
559557

560-
# Re-query scenario from main session to pick up bricked data.
561-
# The thread session's __dict__.update() bypasses SQLAlchemy
562-
# instrumentation, and refresh() fails on objects that were
563-
# merged from a different session.
564-
if _use_thread and enrichment.scenario:
565-
from models.analysis_scenarios import AnalysisScenario as _AS
566-
refreshed = db.get(_AS, enrichment.scenario.id)
567-
if refreshed:
568-
enrichment.scenario = refreshed
569-
570558
yield _sse("enrichment_update", {
571559
"bricked_status": bricked_status.status,
572560
"bricked_latency_ms": bricked_status.latency_ms,
573-
"has_arv": enrichment.scenario.after_repair_value is not None,
574-
"has_repair": enrichment.scenario.repair_cost is not None,
561+
"has_arv": enrichment.scenario.after_repair_value is not None if enrichment.scenario else False,
562+
"has_repair": enrichment.scenario.repair_cost is not None if enrichment.scenario else False,
575563
})
576564
except Exception:
577565
logger.warning(
@@ -593,8 +581,11 @@ def _bricked_sync():
593581

594582
# Stage 4: AI Narrative — generate for new properties OR existing
595583
# ones that never got a narrative (e.g. first attempt crashed)
596-
needs_narrative = not enrichment.is_existing or not enrichment.scenario.ai_narrative
597-
if enrichment.scenario and needs_narrative:
584+
needs_narrative = (
585+
enrichment.scenario
586+
and (not enrichment.is_existing or not enrichment.scenario.ai_narrative)
587+
)
588+
if needs_narrative:
598589
yield _sse("status", {"stage": "generating_narrative"})
599590

600591
narrative_resp = await _generate_narrative(

0 commit comments

Comments
 (0)