Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions apps/api/services/workbook/column_deps.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,25 @@ def strings(value):
if v is not None:
for value in strings(v):
refs.update(m.strip() for m in _REF_RE.findall(value))
# Output adapters also read source columns from mapping keys and Sheets'
# ordered column list. These references need no {placeholder} wrapper.
if col.get("type") == "output":
destination = str(col.get("destination") or "").lower()
config = col.get("destination_config") or {}
if isinstance(config, dict):
sources = []
if destination in {"airtable", "crm", "instantly", "smartlead"}:
field_map = config.get("field_map")
if isinstance(field_map, dict):
sources = field_map.keys()
elif destination == "sheets" and isinstance(config.get("columns"), list):
sources = config["columns"]
for source in sources:
if isinstance(source, str) and source.strip():
if destination in {"instantly", "smartlead"} and "{" in source:
refs.update(m.strip() for m in _REF_RE.findall(source))
else:
refs.add(source.strip())
return refs


Expand Down
135 changes: 135 additions & 0 deletions tests/test_workbook_output_dependencies.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
"""Mapped output columns consume freshly computed values before dispatch."""
import asyncio

import pytest
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker

from apps.api.database import Base
from apps.api.services.workbook import enrichment
from apps.api.services.workbook.column_deps import (
_refs_in, cycle_blocked_columns, downstream_columns, independent_columns,
referencing_columns, topo_sort_columns,
)
from apps.api.services.workbook.models import Workbook, WorkbookRow


@pytest.fixture
def workbook_runner(tmp_path, monkeypatch):
engine = create_engine(f"sqlite:///{tmp_path / 'outputs.db'}")
Base.metadata.create_all(engine)
sessions = sessionmaker(bind=engine, autoflush=False)
monkeypatch.setattr(enrichment, "SessionLocal", sessions)
monkeypatch.setattr(enrichment, "_make_redis", lambda: None)
monkeypatch.setattr(enrichment, "flush_row_change_emits", lambda *args: None)
monkeypatch.setattr(enrichment, "BATCH_ENABLED", False)
sent = []

async def append_row(spreadsheet, values, range_name, workspace_id):
sent.append(values)
return {"success": True, "range": range_name}

async def push_record(fields, base, table, *, workspace_id):
sent.append(fields)
return {"success": True, "record_id": "record-fixture"}

monkeypatch.setattr("apps.api.services.integrations.sheets.append_row", append_row)
monkeypatch.setattr("apps.api.services.integrations.airtable.push_record", push_record)
yield sessions, sent
engine.dispose()


def create_workbook(sessions, destination, *, old_value=None, formula='upper({input})'):
config = ({"columns": ["summary"]} if destination == "sheets"
else {"field_map": {"summary": "Summary"}})
columns = [
{"id": "send", "type": "output", "destination": destination,
"destination_config": config},
{"id": "summary", "name": "Summary", "type": "formula", "formula": formula},
{"id": "input", "type": "input"},
]
with sessions() as db:
wb = Workbook(name="Mapped output", workspace_id="output-deps", columns_config=columns)
db.add(wb)
db.flush()
row = WorkbookRow(workbook_id=wb.id, workspace_id="output-deps", position=0,
data={"input": "new"},
enrichments={} if old_value is None else {
"summary": {"value": old_value, "status": "complete"}})
db.add(row)
db.commit()
return wb.id


@pytest.mark.parametrize("destination", ["sheets", "airtable"])
@pytest.mark.parametrize("old_value", [None, "STALE"])
def test_runner_dispatches_new_producer_value(workbook_runner, destination, old_value):
sessions, sent = workbook_runner
wid = create_workbook(sessions, destination, old_value=old_value)
result = asyncio.run(enrichment.run_workbook_enrichment(
wid, force=True, retry_passes=0, workspace_id="output-deps"))
assert result["completed"] == 2 and result["errors"] == 0
expected = ["NEW"] if destination == "sheets" else {"Summary": "NEW"}
assert sent == [expected]


@pytest.mark.parametrize("destination", ["sheets", "airtable"])
def test_partial_run_without_producer_value_does_not_send(workbook_runner, destination):
sessions, sent = workbook_runner
wid = create_workbook(sessions, destination)
result = asyncio.run(enrichment.run_workbook_enrichment(
wid, column_ids=["send"], force=True, retry_passes=0, workspace_id="output-deps"))
assert sent == []
assert result["completed"] == 0 and result["errors"] == 1


@pytest.mark.parametrize("destination", ["sheets", "airtable"])
def test_failed_producer_does_not_send_old_value(workbook_runner, destination):
sessions, sent = workbook_runner
wid = create_workbook(sessions, destination, old_value="STALE", formula="1 / 0")
result = asyncio.run(enrichment.run_workbook_enrichment(
wid, force=True, retry_passes=0, workspace_id="output-deps"))
assert sent == []
assert result["completed"] == 0 and result["errors"] == 2


def test_mapping_references_join_existing_dependency_and_cycle_guards():
producer = {"id": "summary", "name": "Summary", "type": "formula", "formula": '"ready"'}
send = {"id": "send", "type": "output", "destination": "instantly",
"destination_config": {"field_map": {"Summary": "email", "{first} {last}": "name"}}}
columns = [send, producer]
assert _refs_in(send) == {"Summary", "first", "last"}
assert [col["id"] for col in topo_sort_columns(columns)] == ["summary", "send"]
assert independent_columns(columns) == []
assert referencing_columns(columns, producer) == [send]
assert downstream_columns(columns, {"summary"}) == [send]
assert cycle_blocked_columns([send, {**producer, "formula": "{send}"}]) == {"summary", "send"}


def test_destination_labels_and_unrelated_configs_are_not_source_references():
send = {"id": "send", "type": "output", "destination": "airtable",
"destination_config": {"field_map": {"summary": "email"}, "table": "input"}}
producer = {"id": "email", "type": "formula", "formula": '"independent"'}
assert _refs_in(send) == {"summary"}
assert topo_sort_columns([send, producer]) == [send, producer]
assert _refs_in({"type": "http", "destination_config": {"field_map": {"summary": "email"}}}) == set()


@pytest.mark.parametrize("destination", ["sheets", "airtable"])
def test_partial_run_can_send_existing_producer_value(workbook_runner, destination):
sessions, sent = workbook_runner
wid = create_workbook(sessions, destination, old_value="KEEP")
result = asyncio.run(enrichment.run_workbook_enrichment(
wid, column_ids=["send"], force=True, retry_passes=0, workspace_id="output-deps"))
expected = ["KEEP"] if destination == "sheets" else {"Summary": "KEEP"}
assert sent == [expected]
assert result["completed"] == 1 and result["errors"] == 0


def test_template_output_dependencies_keep_existing_behavior():
send = {"id": "send", "type": "output", "destination": "webhook",
"destination_config": {"body": {"summary": "{summary}"}}}
producer = {"id": "summary", "type": "formula", "formula": '"ready"'}
assert _refs_in(send) == {"summary"}
assert topo_sort_columns([send, producer]) == [producer, send]
assert _refs_in({"type": "output", "destination": "sheets", "destination_config": {}}) == set()
Loading