Skip to content

Commit 41d00e5

Browse files
committed
Merge remote-tracking branch 'review-pr/30' into review/integration
2 parents ee5aeaf + dcc28e1 commit 41d00e5

2 files changed

Lines changed: 72 additions & 9 deletions

File tree

‎apps/api/services/workbook/refresh.py‎

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -52,12 +52,12 @@ def _interval_minutes(policy: dict) -> Optional[int]:
5252
return INTERVAL_MINUTES.get(str(interval).lower())
5353

5454

55-
def _stale_lead_ids(db, workbook_id: str, enrichment_cols: List[dict], ttl_map: dict) -> List[int]:
56-
"""Lead ids with at least one enrichment cell missing or older than its TTL."""
55+
def _stale_row_ids(db, workbook_id: str, enrichment_cols: List[dict], ttl_map: dict) -> List[int]:
56+
"""Workbook row ids with an enrichment cell missing or older than its TTL."""
5757
rows = db.query(WorkbookRow).filter(WorkbookRow.workbook_id == workbook_id).all()
5858
if not rows:
5959
return []
60-
# Index existing enrichments by (lead_id, col)
60+
# Legacy overlay keys follow the execution key (linked lead id or row id).
6161
overlays = db.query(WorkbookEnrichment).filter(
6262
WorkbookEnrichment.workbook_id == workbook_id
6363
).all()
@@ -66,21 +66,22 @@ def _stale_lead_ids(db, workbook_id: str, enrichment_cols: List[dict], ttl_map:
6666
stale = set()
6767
now = _now()
6868
for row in rows:
69-
lead_id = row.lead_id or row.id
69+
row_id = row.id
70+
lead_id = row.lead_id or row_id
7071
for col in enrichment_cols:
7172
cid = col.get("id")
7273
ttl_days = ttl_map.get(col.get("target_field") or cid, DEFAULT_STALENESS_DAYS)
7374
o = seen.get((lead_id, cid))
7475
if o is None or o.status != "complete":
75-
stale.add(lead_id)
76+
stale.add(row_id)
7677
break
7778
upd = o.updated_at
7879
if upd is None:
79-
stale.add(lead_id); break
80+
stale.add(row_id); break
8081
if upd.tzinfo is None:
8182
upd = upd.replace(tzinfo=timezone.utc)
8283
if upd < now - timedelta(days=ttl_days):
83-
stale.add(lead_id); break
84+
stale.add(row_id); break
8485
return list(stale)
8586

8687

@@ -123,9 +124,9 @@ async def _refresh_workbook_impl(workbook_id: str, reason: str, workspace_id: st
123124
# 2) Re-enrich stale rows only
124125
reenriched = 0
125126
with SessionLocal() as db:
126-
stale = _stale_lead_ids(db, workbook_id, enrichment_cols, ttl_map) if enrichment_cols else []
127+
stale = _stale_row_ids(db, workbook_id, enrichment_cols, ttl_map) if enrichment_cols else []
127128
if stale:
128-
result = await run_workbook_enrichment(workbook_id, lead_ids=stale, workspace_id=workspace_id)
129+
result = await run_workbook_enrichment(workbook_id, row_ids=stale, workspace_id=workspace_id)
129130
reenriched = result.get("completed", 0)
130131

131132
with SessionLocal() as db:
Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,62 @@
1+
"""Standing refresh uses persisted workbook row identities, not optional lead links."""
2+
import asyncio
3+
from datetime import datetime, timezone
4+
5+
import pytest
6+
from sqlalchemy import create_engine, event
7+
from sqlalchemy.orm import sessionmaker
8+
9+
from apps.api.database import Base
10+
from apps.api.services.workbook import enrichment, refresh
11+
from apps.api.services.workbook.models import Workbook, WorkbookEnrichment, WorkbookRow
12+
from apps.api.services.workbook.activity_models import WorkbookActivity
13+
14+
15+
@pytest.mark.parametrize("linked", [False, True])
16+
def test_refresh_executes_stale_row_and_preserves_fresh_row(tmp_path, monkeypatch, linked):
17+
connection = create_engine(f"sqlite:///{tmp_path / 'refresh.db'}")
18+
@event.listens_for(connection, "connect")
19+
def foreign_keys(dbapi, _):
20+
dbapi.execute("PRAGMA foreign_keys=ON")
21+
Base.metadata.create_all(connection)
22+
sessions = sessionmaker(bind=connection, autoflush=False)
23+
with sessions() as db:
24+
wb = Workbook(name="Refresh fixture", workspace_id="refresh-test",
25+
columns_config=[{"id": "summary", "type": "ai_formula", "prompt": "Summarize {email}"}])
26+
db.add(wb)
27+
db.flush()
28+
wid = wb.id
29+
stale = WorkbookRow(id=101, workbook_id=wid, workspace_id="refresh-test", position=0,
30+
lead_id=901 if linked else None, data={"email": "stale@example.test"})
31+
fresh = WorkbookRow(id=202, workbook_id=wid, workspace_id="refresh-test", position=1,
32+
lead_id=902 if linked else None, data={"email": "fresh@example.test"},
33+
enrichments={"summary": {"value": "KEEP", "status": "complete"}})
34+
db.add_all([stale, fresh])
35+
db.flush()
36+
db.add(WorkbookEnrichment(workbook_id=wid, workspace_id="refresh-test", lead_id=fresh.lead_id or fresh.id,
37+
column_id="summary", status="complete", value="KEEP",
38+
updated_at=datetime.now(timezone.utc)))
39+
db.commit()
40+
calls = []
41+
async def provider(**kwargs):
42+
calls.append(kwargs["row_cells"]["email"]["value"])
43+
return {"value": "UPDATED", "error": None}
44+
monkeypatch.setattr(refresh, "SessionLocal", sessions)
45+
monkeypatch.setattr(enrichment, "SessionLocal", sessions)
46+
monkeypatch.setattr(enrichment, "execute_ai_column", provider)
47+
monkeypatch.setattr(enrichment, "_make_redis", lambda: None)
48+
monkeypatch.setattr(enrichment, "flush_row_change_emits", lambda *args: None)
49+
try:
50+
result = asyncio.run(refresh.refresh_workbook(wid, workspace_id="refresh-test"))
51+
assert result == {"sourced": 0, "reenriched": 1, "stale_rows": 1}
52+
assert calls == ["stale@example.test"]
53+
with sessions() as db:
54+
assert db.get(WorkbookRow, 101).enrichments["summary"]["value"] == "UPDATED"
55+
assert db.get(WorkbookRow, 202).enrichments["summary"] == {"value": "KEEP", "status": "complete"}
56+
assert {(o.lead_id, o.value) for o in db.query(WorkbookEnrichment)} == {(901 if linked else 101, "UPDATED"), (902 if linked else 202, "KEEP")}
57+
assert db.query(WorkbookActivity).one().kind == "refresh"
58+
# Fresh receipts must make a second scheduled cycle a no-op.
59+
assert asyncio.run(refresh.refresh_workbook(wid, workspace_id="refresh-test"))["reenriched"] == 0
60+
assert len(calls) == 1
61+
finally:
62+
connection.dispose()

0 commit comments

Comments
 (0)