Skip to content

Commit 1a7bc42

Browse files
authored
Merge pull request #64 from debpalash/review/secure-integration-20261007
fix: integrate reviewed PR backlog with security and regression checks
2 parents 1210ada + 205bec6 commit 1a7bc42

64 files changed

Lines changed: 2504 additions & 119 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.env.example‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,7 @@
6262
# GOOGLE_API_KEY= # Google Custom Search API key
6363
# GOOGLE_CSE_ID= # Google Custom Search Engine ID
6464
# SERPINGAPI_API_KEY= # Serping API key (Google SERP; keyed search fallback) — https://serpingapi.com
65+
# SERPLY_API_KEY= # Serply API key (Google SERP; keyed search fallback) — https://serply.io
6566

6667
# === Email Outreach (optional, also configurable in Settings) ===
6768
# SMTP_HOST=smtp.gmail.com

‎.github/workflows/ci.yml‎

Lines changed: 28 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -75,5 +75,32 @@ jobs:
7575
bun-version: '1.4.2'
7676
- run: bun install --frozen-lockfile
7777
- run: bun run lint
78-
- run: bun test apps/web/tests/api.lead-create.test.ts
78+
- run: bun test apps/web/tests
7979
- run: bun run build
80+
81+
frontend-browser:
82+
name: Frontend browser regressions
83+
runs-on: ubuntu-latest
84+
timeout-minutes: 15
85+
steps:
86+
- uses: actions/checkout@v7
87+
with:
88+
persist-credentials: false
89+
- uses: astral-sh/setup-uv@v7
90+
with:
91+
python-version: '3.13'
92+
enable-cache: true
93+
- uses: oven-sh/setup-bun@v2
94+
with:
95+
bun-version: '1.4.2'
96+
- run: bun install --frozen-lockfile
97+
- run: uv run --frozen playwright install --with-deps chromium
98+
- name: Run mounted React interaction regressions
99+
shell: bash
100+
run: |
101+
export OPENGTM_BUN="$(command -v bun)"
102+
export OPENGTM_CHROMIUM="$(uv run --frozen python -c 'from playwright.sync_api import sync_playwright; p = sync_playwright().start(); print(p.chromium.executable_path); p.stop()')"
103+
uv run --frozen pytest -q \
104+
tests/test_workbook_socket_browser.py \
105+
tests/test_lead_selection_browser.py \
106+
tests/test_workbook_clipboard_browser.py

‎README.md‎

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44
<p align="center">
55
<a href="https://opengtm.palash.dev">Docs</a> ·
66
<a href="#run-it">Install</a> ·
7+
<a href="#co-maintainer-wanted">Co-maintainer wanted</a> ·
78
<a href="https://github.com/debpalash/OpenGTM/releases">Releases</a> ·
89
<a href="LICENSE">AGPL-3.0</a>
910
</p>
@@ -16,6 +17,17 @@ with cited agents, watches for buying intent, and sends qualified records to
1617
the tools your team already uses. Your data stays on your infrastructure;
1718
third-party calls use keys you choose. A zero-key demo is included.
1819

20+
## Co-maintainer wanted
21+
22+
We're looking for a co-maintainer to help shape OpenGTM, review pull requests,
23+
ship releases, and build a faster enrichment backend. Experience with **Go,
24+
Python, PostgreSQL, or provider integrations** is especially welcome as we
25+
prepare the backend migration below.
26+
27+
Interested? [Join the rewrite issue](https://github.com/debpalash/OpenGTM/issues/33)
28+
with a short introduction and the areas you'd like to own. Focused contributions
29+
are welcome too; start with [CONTRIBUTING.md](CONTRIBUTING.md).
30+
1931
## Run it
2032

2133
You need **Git, Docker with Compose v2, and roughly 4 GB of available RAM**.
@@ -156,6 +168,30 @@ FastAPI ──► PostgreSQL + forced workspace RLS
156168
bridge, and `apps/docs` the documentation. Read the
157169
[architecture guide](docs/architecture.md) for tenancy and failure behavior.
158170

171+
## Go backend migration
172+
173+
An incremental migration to a **Go backend with Python specialist workers** is
174+
on the roadmap, with **optional Rust acceleration** where benchmarks justify it.
175+
The current backend is Python/FastAPI; the migration has not shipped yet.
176+
177+
- **Go:** the primary backend language for APIs, enrichment orchestration,
178+
durable job workers, and scheduling. Priorities include pooled HTTP clients,
179+
bounded concurrency, provider rate limits, cancellation, and batched writes.
180+
- **Python:** AI research, browser automation, and specialized integrations.
181+
- **Rust, optional:** performance-critical parsing, normalization, and
182+
deduplication where profiling and end-to-end benchmarks show a benefit.
183+
- **Supporting stack:** PostgreSQL for durable data and workspace isolation,
184+
Redis for progress updates, and TypeScript/React for the UI.
185+
186+
The goal is hundreds of completed enrichments per second, subject to provider
187+
limits and workload. We'll validate throughput, latency, memory use, and retry
188+
correctness with benchmarks before making performance claims. Small modules,
189+
generated API/database contracts, and existing behavior tests will support fast
190+
AI-assisted development while preserving tenant isolation and billing correctness.
191+
See the [backend rewrite proposal](docs/plans/go-python-backend-rewrite.md) and
192+
[tracking issue](https://github.com/debpalash/OpenGTM/issues/33) for milestones,
193+
correctness requirements, benchmark gates, and rollback.
194+
159195
## Scope and contribution
160196

161197
OpenGTM covers the discover → enrich → segment → act → learn loop. It is not a

‎apps/api/routers/crm.py‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
from fastapi import APIRouter, Depends
1+
from fastapi import APIRouter, Depends, HTTPException
22
from fastapi.responses import StreamingResponse
33
from sqlalchemy.orm import Session
44
from typing import List, Optional, Dict
@@ -135,14 +135,16 @@ async def update_email_data(
135135
):
136136
record = db.query(EmailData).filter(EmailData.id == id).first()
137137
if not record:
138-
return {"status": "not found", "error": f"Record with id {id} not found"}, 404
138+
raise HTTPException(status_code=404, detail=f"Record with id {id} not found")
139139

140140
update_data = data.dict(exclude_unset=True)
141141
for key, value in update_data.items():
142142
setattr(record, key, value)
143143

144144
db.commit()
145-
return {"status": "success", "data": record.__dict__}
145+
return {"status": "success", "data": {
146+
column.name: getattr(record, column.name) for column in EmailData.__table__.columns
147+
}}
146148

147149

148150
@router.delete("/{id}")

‎apps/api/routers/leads.py‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -132,9 +132,8 @@ def list_leads(
132132
limit=limit,
133133
offset=offset,
134134
order_by=order_by,
135+
exclude_dead=status_filter != "dead",
135136
)
136-
if status_filter != "dead":
137-
leads = [l for l in leads if l.status != "dead"]
138137
result = [l.to_dict() for l in leads]
139138
db.close()
140139
return result

‎apps/api/routers/workbooks.py‎

Lines changed: 51 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
import binascii
77
import csv
88
import io
9+
import hashlib
910
import json
1011
import logging
1112
import os
@@ -28,7 +29,7 @@
2829
WorkbookListResponse, WorkbookWithLeadsResponse,
2930
WorkbookLeadRow, EnrichmentOverlay,
3031
RunWorkbookRequest, RunWorkbookResponse, RunCellRequest,
31-
AddColumnRequest, ExportRequest,
32+
AddColumnRequest,
3233
AddRowsRequest, ImportRowsRequest, DeleteRowsRequest, DeleteMatchingRowsRequest, BulkUpdateRowsRequest,
3334
GenerateColumnRequest, GenerateColumnResponse,
3435
WorkbookViewCreate, WorkbookViewUpdate,
@@ -133,7 +134,8 @@ def _workbook_rows_query(db: Session, wb: Workbook, view_id: Optional[str], sear
133134
if expression is None:
134135
raise HTTPException(status_code=409, detail="Saved view references a missing filter column. Repair the view before continuing.")
135136
normalized = sa_func.lower(sa_func.coalesce(expression, ""))
136-
expected = str(rule.get("value") or "").lower()
137+
comparison = rule.get("value")
138+
expected = str("" if comparison is None else comparison).lower()
137139
operation = rule.get("op")
138140
if operation == "equals":
139141
query = query.filter(normalized == expected)
@@ -205,6 +207,46 @@ def _decode_query_cursor(cursor: str, expected_values: int) -> list:
205207
raise HTTPException(status_code=400, detail="Invalid workbook row cursor")
206208

207209

210+
def _workbook_cursor_scope(query, ordering: list, view_id: Optional[str]) -> str:
211+
# Bind row anchors to the owned workbook, effective filters and ordering.
212+
compiled = [expression.compile() for expression in (query.statement, *ordering)]
213+
identity = [view_id, *[(str(item), item.params) for item in compiled]]
214+
return hashlib.sha256(json.dumps(identity, sort_keys=True).encode()).hexdigest()
215+
216+
217+
def _cursor_values_digest(values: list) -> str:
218+
return hashlib.sha256(json.dumps(values, separators=(",", ":")).encode()).hexdigest()
219+
220+
221+
def _encode_workbook_query_cursor(values: list, row_id: int, scope: str) -> str:
222+
cursor = _encode_query_cursor(values)
223+
if len(cursor) <= 512:
224+
return cursor
225+
# Large cell values must not produce an unusable URL. Recover their tuple
226+
# from an owned row on continuation instead of copying the text into it.
227+
payload = json.dumps({"v": 2, "row": row_id, "scope": scope, "digest": _cursor_values_digest(values)}, separators=(",", ":")).encode()
228+
return base64.urlsafe_b64encode(payload).decode().rstrip("=")
229+
230+
231+
def _decode_workbook_query_cursor(cursor: str, query, cursor_terms: list, scope: str) -> list:
232+
try:
233+
raw = base64.urlsafe_b64decode(cursor + "=" * (-len(cursor) % 4))
234+
payload = json.loads(raw)
235+
if isinstance(payload, dict) and payload.get("v") == 2:
236+
row_id = payload.get("row")
237+
if type(row_id) is not int or row_id < 1 or payload.get("scope") != scope:
238+
raise ValueError
239+
values = query.with_entities(*(expression for expression, _ in cursor_terms)).filter(
240+
WorkbookRow.id == row_id
241+
).first()
242+
if values is None or payload.get("digest") != _cursor_values_digest(list(values)):
243+
raise ValueError
244+
return list(values)
245+
except (ValueError, TypeError, json.JSONDecodeError, binascii.Error):
246+
raise HTTPException(status_code=400, detail="Invalid workbook row cursor; restart pagination")
247+
return _decode_query_cursor(cursor, len(cursor_terms))
248+
249+
208250
def _encode_connector_run_cursor(run: ConnectorRun) -> str:
209251
return _encode_query_cursor([run.created_at.isoformat(), run.id])
210252

@@ -552,11 +594,12 @@ async def get_workbook(
552594
query, ordering, cursor_terms = _workbook_rows_query(db, wb, view_id, search)
553595
query_total = query.with_entities(sa_func.count(WorkbookRow.id)).scalar() or 0
554596
using_cursor = cursor_mode or cursor is not None
597+
cursor_scope = _workbook_cursor_scope(query, ordering, view_id) if len(cursor_terms) > 2 else ""
555598
if cursor:
556599
if len(cursor_terms) == 2:
557600
cursor_values = list(_decode_row_cursor(cursor))
558601
else:
559-
cursor_values = _decode_query_cursor(cursor, len(cursor_terms))
602+
cursor_values = _decode_workbook_query_cursor(cursor, query, cursor_terms, cursor_scope)
560603
query = query.filter(_cursor_after_filter(cursor_terms, cursor_values))
561604
if using_cursor:
562605
fetched = query.order_by(*ordering).limit(page_size + 1).all()
@@ -575,7 +618,7 @@ async def get_workbook(
575618
cursor_values = db.query(
576619
*(expression for expression, _ in cursor_terms)
577620
).filter(WorkbookRow.id == last_row.id).one()
578-
next_cursor = _encode_query_cursor(list(cursor_values))
621+
next_cursor = _encode_workbook_query_cursor(list(cursor_values), last_row.id, cursor_scope)
579622

580623
# Older v2 mirrors omitted research metadata. Recover it read-only from
581624
# the matching cell receipt, never from a different workbook/tenant or
@@ -2070,7 +2113,10 @@ def _identity(d: dict) -> str:
20702113
)
20712114
except Exception as _e:
20722115
logger.warning("on_row_added emit (add_rows) failed: %s", _e)
2073-
return {"added": added, "skipped_duplicates": skipped, "total_rows": max_pos + added + 1}
2116+
total_rows = db.query(sa_func.count(WorkbookRow.id)).filter(
2117+
WorkbookRow.workbook_id == wb.id
2118+
).scalar() or 0
2119+
return {"added": added, "skipped_duplicates": skipped, "total_rows": total_rows}
20742120

20752121

20762122
@router.delete("/{workbook_id}/rows")

‎apps/api/services/audiences/scheduler.py‎

Lines changed: 21 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -171,6 +171,7 @@ def bootstrap_audience_schedules() -> int:
171171
with SessionLocal() as db:
172172
due_query = db.query(
173173
AudienceSchedule.audience_id, AudienceSchedule.workspace_id,
174+
AudienceSchedule.next_refresh_at,
174175
).filter(
175176
AudienceSchedule.enabled.is_(True),
176177
or_(
@@ -185,7 +186,7 @@ def bootstrap_audience_schedules() -> int:
185186
).all()
186187
if not identities:
187188
break
188-
for audience_id, workspace_id in identities:
189+
for audience_id, workspace_id, due_at in identities:
189190
with workspace_scope(workspace_id):
190191
with SessionLocal() as db:
191192
audience = db.query(Audience).filter(
@@ -195,8 +196,25 @@ def bootstrap_audience_schedules() -> int:
195196
remove_schedule(db, audience_id)
196197
db.commit()
197198
continue
198-
schedule_next(db, audience, now=now)
199-
enqueued += 1
199+
active = db.query(Job).filter(
200+
Job.type == "audience_refresh",
201+
Job.status.in_(("pending", "processing")),
202+
Job.fire_key.like(f"audience_refresh:{audience_id}:%"),
203+
).first()
204+
if active is not None:
205+
continue
206+
from apps.api.services.job_scheduling import enqueue_job_once
207+
208+
occurrence = _as_utc(due_at).isoformat() if due_at else "bootstrap"
209+
job = enqueue_job_once(
210+
db, job_type="audience_refresh",
211+
payload={"workspace_id": workspace_id, "audience_id": audience_id},
212+
fire_key=f"audience_refresh:{audience_id}:{occurrence}",
213+
next_run_at=now,
214+
)
215+
db.commit()
216+
if job is not None:
217+
enqueued += 1
200218
last_audience_id = identities[-1][0]
201219
logger.info("bootstrap_audience_schedules enqueued %d refresh(es)", enqueued)
202220
return enqueued

‎apps/api/services/automations/engine.py‎

Lines changed: 7 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -93,13 +93,7 @@ def cols_for(wb_id: str):
9393

9494
def _lead_data_for_action(row, columns_config) -> dict:
9595
"""Build the lead_data dict execute_output_column / templates expect."""
96-
data = dict(row.data or {})
97-
# merge enrichment overlay (column_id -> value) so {col} placeholders resolve
98-
for cid, cell in (row.enrichments or {}).items():
99-
if isinstance(cell, dict) and "value" in cell:
100-
data.setdefault(cid, cell.get("value"))
101-
else:
102-
data.setdefault(cid, cell)
96+
data = _row_cells(row)
10397
data.setdefault("id", row.lead_id or row.id)
10498
if row.lead_id:
10599
data.setdefault("lead_id", row.lead_id)
@@ -493,18 +487,21 @@ def bootstrap_schedules():
493487
return 0
494488
from apps.api.core.tenancy import workspace_scope
495489
from apps.api.services.automations.models import ScheduledTrigger, Trigger
490+
from sqlalchemy import or_
496491

497492
now = datetime.now(timezone.utc)
498493
enqueued = 0
499494
# read mirror WITHOUT a workspace scope (mirror is non-RLS by design)
500495
with SessionLocal() as db:
501496
due = (
502497
db.query(ScheduledTrigger)
503-
.filter(ScheduledTrigger.enabled.is_(True))
498+
.filter(
499+
ScheduledTrigger.enabled.is_(True),
500+
or_(ScheduledTrigger.next_run_at.is_(None), ScheduledTrigger.next_run_at <= now),
501+
)
504502
.all()
505503
)
506-
due_ids = [(r.trigger_id, r.workspace_id) for r in due
507-
if r.next_run_at is None or r.next_run_at <= now]
504+
due_ids = [(r.trigger_id, r.workspace_id) for r in due]
508505
for trigger_id, workspace_id in due_ids:
509506
with workspace_scope(workspace_id):
510507
with SessionLocal() as db:

‎apps/api/services/governance/retention.py‎

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -186,7 +186,23 @@ def bootstrap_retention_schedules() -> int:
186186
with SessionLocal() as db:
187187
policy = db.query(RetentionPolicy).filter(RetentionPolicy.workspace_id == workspace_id).first()
188188
if policy and policy.enabled and not policy.legal_hold:
189+
from apps.api.models import Job
189190
from apps.api.services.job_scheduling import enqueue_job_once
190-
enqueue_job_once(db, job_type="retention_enforce", payload={"workspace_id": workspace_id}, fire_key=f"retention:{workspace_id}:{now.date().isoformat()}")
191-
schedule_policy(db, policy, now=now); count += 1
191+
active = db.query(Job).filter(
192+
Job.type == "retention_enforce",
193+
Job.status.in_(("pending", "processing")),
194+
Job.fire_key.like(f"retention:{workspace_id}:%"),
195+
).first()
196+
if active is not None:
197+
continue
198+
occurrence = (comparable or now).date().isoformat()
199+
job = enqueue_job_once(
200+
db, job_type="retention_enforce",
201+
payload={"workspace_id": workspace_id},
202+
fire_key=f"retention:{workspace_id}:{occurrence}",
203+
next_run_at=now,
204+
)
205+
db.commit()
206+
if job is not None:
207+
count += 1
192208
return count

‎apps/api/services/leadgen/config.py‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,7 @@
126126
GOOGLE_CSE_KEY = os.getenv("GOOGLE_CSE_KEY", "")
127127
BRAVE_SEARCH_KEY = os.getenv("BRAVE_SEARCH_KEY", "")
128128
SERPINGAPI_API_KEY = os.getenv("SERPINGAPI_API_KEY", "")
129+
SERPLY_API_KEY = os.getenv("SERPLY_API_KEY", "")
129130

130131
# ── DDG result cache + adaptive backoff (see services/leadgen/search_cache.py)
131132
# In-process only (no DB). Both behaviours are behind flags, default ON; set

0 commit comments

Comments
 (0)