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
72 changes: 70 additions & 2 deletions products/signals/backend/emission/pganalyze_issues.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,18 @@
import json
from datetime import datetime
from typing import Any

from structlog import get_logger

from products.signals.backend.emission.fetchers.data_warehouse import data_warehouse_record_fetcher
from posthog.hogql import ast
from posthog.hogql.parser import parse_select
from posthog.hogql.query import execute_hogql_query

from posthog.models import Team

from products.signals.backend.emission.fetchers.data_warehouse import escape_table_name
from products.signals.backend.emission.registry import SignalEmitterOutput, SignalSourceTableConfig
from products.signals.backend.models import SignalEmissionRecord

logger = get_logger(__name__)

Expand Down Expand Up @@ -58,6 +66,8 @@
"synced_at",
)

ISSUE_PAGE_SIZE = 1_000


def _parse_references(record: dict[str, Any]) -> list[dict[str, Any]]:
raw_refs = record.get("references")
Expand Down Expand Up @@ -125,11 +135,69 @@ def _build_extra(record: dict[str, Any], references: list[dict[str, Any]]) -> di
return extra


def _fetch_issue_page(
team: Team, config: SignalSourceTableConfig, context: dict[str, Any], after_id: str
) -> list[dict[str, Any]]:
placeholders: dict[str, Any] = {"after_id": ast.Constant(value=after_id)}
if context.get("last_synced_at") is not None:
window = "parseDateTimeBestEffort(synced_at) > {last_synced_at}"
placeholders["last_synced_at"] = ast.Constant(value=datetime.fromisoformat(context["last_synced_at"]))
else:
window = f"parseDateTimeBestEffort(synced_at) > now() - interval {config.first_sync_lookback_days} day"
# Weekly warehouse partitions can retain older versions of the same issue.
query = f"""
SELECT {", ".join(config.fields)}
FROM {escape_table_name(context["table_name"])}
WHERE {window} AND id > {{after_id}}
ORDER BY id ASC, parseDateTimeBestEffort(synced_at) DESC
LIMIT 1 BY id
LIMIT {ISSUE_PAGE_SIZE}
"""
result = execute_hogql_query(
query=parse_select(query, placeholders=placeholders),
team=team,
query_type="EmitSignalsNewRecords",
bypass_warehouse_access_control=True,
)
if not result.results or not result.columns:
return []
return [dict(zip(result.columns, row)) for row in result.results]


def pganalyze_issue_record_fetcher(
team: Team,
config: SignalSourceTableConfig,
context: dict[str, Any],
) -> list[dict[str, Any]]:
"""pganalyze stamps open issues on every sync, so the time cursor alone cannot deduplicate them."""
records: list[dict[str, Any]] = []
after_id = ""
while len(records) < config.max_records:
rows = _fetch_issue_page(team, config, context, after_id)
if not rows:
break
already_emitted = set(
SignalEmissionRecord.objects.filter(
team=team,
source_product=config.source_product,
source_type=config.source_type,
source_id__in=[str(row["id"]) for row in rows],
).values_list("source_id", flat=True)
)
records.extend(row for row in rows if str(row["id"]) not in already_emitted)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🟠 Major | 🏗️ Heavy lift

Prevent persistent failures from blocking later issue IDs.

If the first config.max_records IDs fail emission on every sync, none enters SignalEmissionRecord. The next sync starts at the same IDs and returns them again, so later valid issues never reach the pipeline. Preserve retries for failed IDs, but rotate or separately queue retries so a permanently failing batch cannot block the backlog.

after_id = str(rows[-1]["id"])
if len(rows) < ISSUE_PAGE_SIZE:
break
return records[: config.max_records]


PGANALYZE_ISSUES_CONFIG = SignalSourceTableConfig(
source_product="pganalyze",
source_type="issue",
emitter=pganalyze_issue_emitter,
record_fetcher=data_warehouse_record_fetcher,
record_fetcher=pganalyze_issue_record_fetcher,
record_processed_outputs=True,
# The fetcher reads only the rows of the latest sync, which are the issues that are open now.
partition_field="synced_at",
partition_field_is_datetime_string=True,
fields=REQUIRED_FIELDS + EXTRA_FIELDS,
Expand Down
30 changes: 29 additions & 1 deletion products/signals/backend/emission/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@
from datetime import datetime
from typing import Any

from django.utils import timezone

import structlog
import posthoganalytics
from anthropic import AsyncAnthropic
Expand All @@ -25,6 +27,7 @@
)
from products.signals.backend.emission.steering import apply_steering, steering_from_config
from products.signals.backend.facade.api import emit_signal
from products.signals.backend.models import SignalEmissionRecord
from products.signals.backend.temporal import metrics
from products.signals.backend.temporal.drop_telemetry import summarize_drop_error
from products.signals.backend.temporal.llm import effort_kwargs
Expand Down Expand Up @@ -465,11 +468,28 @@
)


async def _record_processed_outputs(team: Team, outputs: list[SignalEmitterOutput]) -> None:
await SignalEmissionRecord.objects.abulk_create(
[
SignalEmissionRecord(
team=team,
source_product=output.source_product,
source_type=output.source_type,
source_id=output.source_id,
emitted_at=timezone.now(),
)
for output in outputs
],
ignore_conflicts=True,
)


async def _emit_signals(
team: Team,
organization: Organization,
outputs: list[SignalEmitterOutput],
extra: dict[str, Any],
record_processed_outputs: bool = False,
) -> int:
semaphore = asyncio.Semaphore(EMIT_CONCURRENCY_LIMIT)
_safe_heartbeat()
Expand Down Expand Up @@ -506,10 +526,13 @@
description=output.description,
weight=output.weight,
extra=output.extra,
idempotency_key=output.source_id if record_processed_outputs else None,
)
if record_processed_outputs:
await _record_processed_outputs(team, [output])
Comment on lines +531 to +532

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick win

🔎 Supported by static analysis

🏁 Script executed:

#!/bin/bash
# Inspect the pipeline entrypoint and the facade's pre-dispatch return paths.
rg -n -C 12 'def run_signal_pipeline|_emit_signals\(|is_ai_data_processing_approved|is_source_enabled' \
  products/signals/backend/emission/pipeline.py \
  products/signals/backend/facade/api.py

Repository: PostHog/posthog

Length of output: 13500


🏁 Script executed:

sed -n '487,653p' products/signals/backend/emission/pipeline.py

Repository: PostHog/posthog

Length of output: 7888


🏁 Script executed:

rg -n -C 8 'SignalEmissionRecord|pganalyze|processed_outputs' products/signals

Repository: PostHog/posthog

Length of output: 42502


Do not record outputs that emit_signal silently drops.

When AI data processing is unapproved or the source is disabled, emit_signal returns without dispatching. _bounded_emit treats that return as success and writes SignalEmissionRecord. The pganalyze fetcher then excludes the issue from later syncs. Add both guards before _emit_signals, or record the output only when emit_signal reports a confirmed dispatch.

return True
except Exception as e:
# Fetchers record emission optimistically, so a record lost here is lost for good.
# Sources that record at fetch time cannot retry a record lost here.
# Close the funnel (entered - summarized - filtered - emit_failed = emitted) and
# count the drop, or the loss is invisible outside logs.
error_type, _ = summarize_drop_error(e)
Expand Down Expand Up @@ -539,7 +562,7 @@
return succeeded


async def run_signal_pipeline(

Check warning on line 565 in products/signals/backend/emission/pipeline.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`run_signal_pipeline` has cyclomatic complexity 12 (warn >10)

Check warning on line 565 in products/signals/backend/emission/pipeline.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`run_signal_pipeline` has cyclomatic complexity 12 (warn >10)
team: Team,
config: SignalSourceTableConfig,
records: list[dict[str, Any]],
Expand Down Expand Up @@ -606,6 +629,10 @@
context_fields=config.actionability_context_fields,
)
post_filter_ids = {o.source_id for o in outputs}
if config.record_processed_outputs:
await _record_processed_outputs(
team, [output for source_id, output in pre_filter_by_id.items() if source_id not in post_filter_ids]
)
for source_id, output in pre_filter_by_id.items():
if source_id not in post_filter_ids:
capture_pipeline_stage(
Expand All @@ -621,6 +648,7 @@
organization=organization,
outputs=outputs,
extra=extra,
record_processed_outputs=config.record_processed_outputs,
)
logger.info(f"Emitted {signals_emitted} signals for {source_label}", **extra)
return {"status": "success", "signals_emitted": signals_emitted}
2 changes: 2 additions & 0 deletions products/signals/backend/emission/registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,8 @@ class SignalSourceTableConfig(BaseModel):
emitter: SignalEmitter
# Each source defines how to fetch records — no default, must be explicit
record_fetcher: RecordFetcher
# Snapshot sources deduplicate only after successful emission or an intentional filter.
record_processed_outputs: bool = False
# Field used to filter records by time window (e.g. "created_at")
partition_field: str
# Columns to SELECT — only what the emitter and extra metadata need
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -686,6 +686,7 @@ async def test_passes_correct_args_to_emit_signal(self):
description="bug report",
weight=0.5,
extra={},
idempotency_key=None,
)

@pytest.mark.asyncio
Expand Down
154 changes: 154 additions & 0 deletions products/signals/backend/emission/tests/test_pganalyze_issues.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,27 @@
from datetime import timedelta
from typing import Any

import pytest
from posthog.test.base import BaseTest, ClickhouseTestMixin
from unittest.mock import AsyncMock, patch

from django.utils import timezone

from asgiref.sync import async_to_sync
from parameterized import parameterized

from posthog.hogql import ast
from posthog.hogql.query import execute_hogql_query

from products.signals.backend.emission.pganalyze_issues import (
EXTRA_FIELDS,
PGANALYZE_ISSUES_CONFIG,
pganalyze_issue_emitter,
pganalyze_issue_record_fetcher,
)
from products.signals.backend.emission.pipeline import run_signal_pipeline
from products.signals.backend.emission.tests.conftest import MOCK_PGANALYZE_ISSUE_RECORD
from products.signals.backend.models import SignalEmissionRecord


class TestPgAnalyzeIssueEmitter:
Expand Down Expand Up @@ -105,6 +122,143 @@ def test_has_summarization_prompt(self):
def test_emitter_is_pganalyze_issue_emitter(self):
assert PGANALYZE_ISSUES_CONFIG.emitter is pganalyze_issue_emitter

def test_uses_deduping_fetcher(self):
assert PGANALYZE_ISSUES_CONFIG.record_fetcher is pganalyze_issue_record_fetcher

def test_source_product_and_type(self):
assert PGANALYZE_ISSUES_CONFIG.source_product == "pganalyze"
assert PGANALYZE_ISSUES_CONFIG.source_type == "issue"


class IssuesTable:
def __init__(self, issue_ids: list[str]) -> None:
self.rows = [
{**MOCK_PGANALYZE_ISSUE_RECORD, "id": issue_id, "synced_at": "2026-04-20T07:00:00+00:00"}
for issue_id in issue_ids
]

def execute(self, query: ast.SelectQuery, **kwargs):
query.select_from = ast.JoinExpr(
table=ast.SelectSetQuery.create_from_queries(
[
ast.SelectQuery(
select=[ast.Alias(alias=key, expr=ast.Constant(value=value)) for key, value in row.items()]
)
for row in self.rows
],
"UNION ALL",
)
)
return execute_hogql_query(query=query, **kwargs)


@pytest.mark.django_db
class TestPgAnalyzeIssueRecordFetcher(ClickhouseTestMixin, BaseTest):
context: dict[str, str | None] = {"table_name": "pganalyze.issues", "last_synced_at": "2026-04-20T06:00:00+00:00"}

def _fetch(self, table: IssuesTable, max_records: int = 200) -> list[dict[str, Any]]:
config = PGANALYZE_ISSUES_CONFIG.model_copy(update={"max_records": max_records})
with patch("products.signals.backend.emission.pganalyze_issues.execute_hogql_query", table.execute):
return pganalyze_issue_record_fetcher(self.team, config, self.context)

def _process(self, records: list[dict[str, Any]], emit: AsyncMock | None = None):
config = PGANALYZE_ISSUES_CONFIG.model_copy(update={"actionability_prompt": None})
with patch("products.signals.backend.emission.pipeline.emit_signal", emit or AsyncMock()):
return async_to_sync(run_signal_pipeline)(team=self.team, config=config, records=records, extra={})

def _sync(self, table: IssuesTable, max_records: int = 200) -> list[str]:
records = self._fetch(table, max_records)
self._process(records)
return [row["id"] for row in records]

def test_open_issue_emits_once_across_syncs(self):
table = IssuesTable(["issue_1", "issue_2"])

assert self._sync(table) == ["issue_1", "issue_2"]
assert self._sync(table) == []
assert SignalEmissionRecord.objects.filter(team=self.team, source_product="pganalyze").count() == 2

def test_only_new_issue_emits_next_to_open_ones(self):
self._sync(IssuesTable(["issue_1", "issue_2"]))

assert self._sync(IssuesTable(["issue_1", "issue_2", "issue_3"])) == ["issue_3"]

@parameterized.expand([("single_page", 1000), ("multiple_pages", 2)])
def test_backlog_larger_than_max_records_drains_over_syncs(self, _name, page_size):
table = IssuesTable([f"issue_{i}" for i in reversed(range(5))])
with patch("products.signals.backend.emission.pganalyze_issues.ISSUE_PAGE_SIZE", page_size):
first = self._sync(table, max_records=2)
second = self._sync(table, max_records=2)
third = self._sync(table, max_records=2)

assert len(first) == len(second) == 2
assert len(third) == 1
assert sorted(first + second + third) == [f"issue_{i}" for i in range(5)]
assert self._sync(table, max_records=2) == []

@parameterized.expand([("continuous", False), ("first_sync", True)])
def test_weekly_versions_use_latest_row_before_limit(self, _name, first_sync):
now = timezone.now()
self.context = {
"table_name": "pganalyze.issues",
"last_synced_at": None if first_sync else (now - timedelta(hours=2)).isoformat(),
}
table = IssuesTable(["issue_1", "issue_2"])
for row in table.rows:
row["synced_at"] = (now - timedelta(hours=1)).isoformat()
table.rows = [
{**table.rows[0], "synced_at": (now - timedelta(days=7)).isoformat(), "description": "Old description"},
{**table.rows[0], "synced_at": (now - timedelta(minutes=90)).isoformat(), "severity": "info"},
*table.rows,
]
records = self._fetch(table, max_records=2)

assert [row["id"] for row in records] == ["issue_1", "issue_2"]
assert records[0]["description"] == MOCK_PGANALYZE_ISSUE_RECORD["description"]
assert records[0]["severity"] == MOCK_PGANALYZE_ISSUE_RECORD["severity"]
assert not SignalEmissionRecord.objects.filter(team=self.team).exists()

@parameterized.expand([("all_fail", False), ("partial_failure", True)])
def test_failed_emissions_remain_eligible_for_retry(self, _name, partial_failure):
table = IssuesTable(["issue_1", "issue_2"] if partial_failure else ["issue_1"])
records = self._fetch(table)

async def emit_issue(**kwargs):
if kwargs["source_id"] == "issue_1":
raise RuntimeError("Emission unavailable")

emit = AsyncMock(side_effect=emit_issue)
if partial_failure:
assert self._process(records, emit)["signals_emitted"] == 1
else:
with pytest.raises(RuntimeError, match="All 1 signal emissions failed"):
self._process(records, emit)

assert self._sync(table) == ["issue_1"]
assert self._sync(table) == []

def test_dispatch_keeps_idempotency_key_when_ledger_write_fails(self):
table = IssuesTable(["issue_1"])
records = self._fetch(table)
emit = AsyncMock()
with patch.object(SignalEmissionRecord.objects, "abulk_create", side_effect=RuntimeError("Ledger unavailable")):
with pytest.raises(RuntimeError, match="All 1 signal emissions failed"):
self._process(records, emit)

self._process(self._fetch(table), emit)

assert [call.kwargs["idempotency_key"] for call in emit.await_args_list] == ["issue_1", "issue_1"]
assert self._fetch(table) == []

def test_non_actionable_issues_are_not_reprocessed(self):
from products.signals.backend.emission.tests.test_emit_signals import _make_llm_response

table = IssuesTable(["issue_1"])
records = self._fetch(table)
with patch("products.signals.backend.emission.pipeline.build_async_anthropic_client") as client:
client.return_value.messages.create = AsyncMock(return_value=_make_llm_response("NOT_ACTIONABLE"))
result = async_to_sync(run_signal_pipeline)(
team=self.team, config=PGANALYZE_ISSUES_CONFIG, records=records, extra={}
)
assert result["signals_emitted"] == 0
assert self._fetch(table) == []
Loading