fix(signals): stop pganalyze re-emitting open issues every sync - #108417
posthog[bot] wants to merge 6 commits into
Conversation
pganalyze issues carry no timestamps, so the warehouse source stamps synced_at on every row at every sync. The synced_at cursor then returned every open issue again. A bespoke fetcher now scans the current issue ids, drops ids already in SignalEmissionRecord, and applies max_records after the dedupe, so a large backlog drains over the next syncs. Generated-By: PostHog Desktop Task-Id: 109f4e0b-3184-4081-aab2-910acc7d55d3
|
❌ This pull request was removed from the merge queue because it failed tests. PR #108548 was used for testing. See more details here.
After your PR is submitted to the merge queue, this comment will be automatically updated with its status. If the PR fails, failure details will also be posted here |
🦔 PostHog Review reviewed this pull requestFound 1 must fix, 0 should fix, 0 consider. Published 1 finding (view the review). Stopped resolving comments at 1/3: this pull request was submitted to the merge queue A fix commit would change what was submitted, or remove it from the queue, so the open threads stay with you. |
Generated-By: PostHog Desktop Task-Id: 109f4e0b-3184-4081-aab2-910acc7d55d3
|
Navigate logical layers of code changes, visualize relationships, and explore their blast radius. 📝 WalkthroughWalkthroughThe pganalyze issue source now uses a custom paginated fetcher. It selects the latest row for each issue ID, skips IDs already recorded for the team and source, and returns up to the configured record limit. The signal pipeline can record successfully emitted outputs and outputs removed by actionability filtering. Tests cover repeated syncs, newly discovered issues, backlog paging, latest-row selection, emission failures, and non-actionable issues. Priority: ➖ Normal Severity of issue fixed: Medium Merge Risk: 🟡 Moderate · up to Repeated failures can keep newer issues out of the selected batch, while an issue skipped by dispatch safeguards can be marked processed and never retried. Resolve these paths before merging. Security Architecture ReviewSecurity architecture risk: 🟡 Moderate · up to Moving issue deduplication until after processing improves retries, but a persistent failure among the earliest issues can prevent later issues from being considered. The impact is limited to the affected team’s pganalyze signals; no new security vulnerability was verified. Retained concerns
Security review detailsSecurity Blast Radius
Trust Boundaries and Controls
Resilience and Maintainability Implications
Hardening Proposals
🚥 Pre-merge checks | ✅ 1✅ Passed checks (1 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
🤖 CI report
|
| File | Comment lines | Added lines |
|---|---|---|
products/signals/backend/emission/pganalyze_issues.py |
2 | 63 |
products/signals/backend/emission/pipeline.py |
1 | 26 |
products/signals/backend/emission/registry.py |
1 | 2 |
This check does not block merging. It updates on every push and clears when the share drops.
⚠️ Backend coverage — 34.0% of changed backend lines covered — 27 uncovered
🧪 Backend test coverage
Patch coverage — changed backend lines (products + core): ███████░░░░░░░░░░░░░ 34.0% (14 / 41)
| File | Patch | Uncovered changed lines |
|---|---|---|
products/signals/backend/emission/pganalyze_issues.py |
31.2% | 141–144, 146, 148, 156, 162–164, 173–179, 187–191 |
products/signals/backend/emission/pipeline.py |
37.5% | 472, 531–532, 632–633 |
🤖 Agents: add a test covering the lines above, or note why under "How did you test this code?". Machine-readable gap list: the patch-coverage artifact on this run (gh run download 36591866898 -n patch-coverage), or the coverage-data block at the end of this comment.
Per-product line coverage (touched products)
| Product | Coverage | Lines |
|---|---|---|
platform_features |
██░░░░░░░░░░░░░░░░░░ 12.1% |
7 / 58 |
warehouse_sources_queue |
██████░░░░░░░░░░░░░░ 29.1% |
92 / 316 |
demo |
███████████░░░░░░░░░ 52.8% |
1,411 / 2,673 |
data_tools |
████████████░░░░░░░░ 61.2% |
90 / 147 |
ai_gateway |
███████████████░░░░░ 75.0% |
9 / 12 |
aeo |
███████████████░░░░░ 76.3% |
617 / 809 |
batch_exports |
████████████████░░░░ 81.2% |
21,527 / 26,502 |
apm |
█████████████████░░░ 84.1% |
1,306 / 1,553 |
cdp |
██████████████████░░ 88.2% |
4,545 / 5,155 |
ml_inference |
██████████████████░░ 88.4% |
509 / 576 |
mcp_analytics |
██████████████████░░ 89.2% |
5,038 / 5,651 |
product_tours |
██████████████████░░ 89.3% |
1,331 / 1,491 |
dashboards |
██████████████████░░ 89.5% |
6,904 / 7,714 |
data_warehouse |
██████████████████░░ 90.0% |
14,027 / 15,589 |
signals |
██████████████████░░ 90.2% |
57,829 / 64,109 |
notebooks |
██████████████████░░ 90.2% |
15,287 / 16,945 |
cohorts |
██████████████████░░ 90.4% |
8,420 / 9,316 |
streamlit_apps |
██████████████████░░ 90.6% |
2,623 / 2,895 |
managed_warehouse |
██████████████████░░ 91.0% |
10,252 / 11,263 |
tasks |
██████████████████░░ 91.1% |
75,146 / 82,503 |
data_modeling |
██████████████████░░ 91.4% |
10,554 / 11,543 |
exports |
██████████████████░░ 91.6% |
9,680 / 10,562 |
engineering_analytics |
██████████████████░░ 91.7% |
11,032 / 12,030 |
ai_training |
██████████████████░░ 92.2% |
356 / 386 |
business_knowledge |
██████████████████░░ 92.2% |
7,684 / 8,330 |
conversations |
███████████████████░ 92.5% |
28,734 / 31,062 |
early_access_features |
███████████████████░ 92.6% |
1,332 / 1,439 |
managed_migrations |
███████████████████░ 92.7% |
1,581 / 1,705 |
visual_review |
███████████████████░ 92.8% |
9,247 / 9,969 |
stamphog |
███████████████████░ 92.8% |
8,109 / 8,742 |
canvas |
███████████████████░ 92.8% |
6,873 / 7,405 |
approvals |
███████████████████░ 93.0% |
3,974 / 4,271 |
mcp_registry |
███████████████████░ 93.1% |
1,670 / 1,794 |
notifications |
███████████████████░ 93.2% |
1,145 / 1,229 |
slack_app |
███████████████████░ 93.2% |
13,733 / 14,735 |
error_tracking |
███████████████████░ 93.2% |
16,359 / 17,547 |
surveys |
███████████████████░ 93.3% |
6,571 / 7,040 |
autoresearch |
███████████████████░ 93.6% |
8,481 / 9,061 |
context_layer |
███████████████████░ 93.9% |
3,415 / 3,638 |
web_analytics |
███████████████████░ 94.0% |
21,815 / 23,218 |
alerts |
███████████████████░ 94.0% |
8,569 / 9,114 |
billing_alerts |
███████████████████░ 94.1% |
2,094 / 2,226 |
mcp_store |
███████████████████░ 94.4% |
8,940 / 9,472 |
ai_observability |
███████████████████░ 94.5% |
24,473 / 25,896 |
workflows |
███████████████████░ 94.6% |
15,222 / 16,093 |
wizard |
███████████████████░ 94.7% |
6,150 / 6,496 |
reminders |
███████████████████░ 94.8% |
760 / 802 |
review_hog |
███████████████████░ 95.0% |
11,507 / 12,119 |
endpoints |
███████████████████░ 95.1% |
9,206 / 9,681 |
annotations |
███████████████████░ 95.1% |
817 / 859 |
customer_analytics |
███████████████████░ 95.2% |
25,081 / 26,352 |
legal_documents |
███████████████████░ 95.2% |
2,311 / 2,427 |
marketing_analytics |
███████████████████░ 95.3% |
19,566 / 20,528 |
posthog_ai |
███████████████████░ 95.3% |
2,488 / 2,610 |
experiments |
███████████████████░ 95.4% |
32,457 / 34,020 |
actions |
███████████████████░ 95.5% |
756 / 792 |
logs |
███████████████████░ 95.5% |
15,399 / 16,130 |
data_catalog |
███████████████████░ 95.5% |
4,401 / 4,606 |
tracing |
███████████████████░ 95.6% |
3,536 / 3,699 |
replay_vision |
███████████████████░ 95.6% |
27,675 / 28,939 |
growth |
███████████████████░ 95.7% |
11,228 / 11,734 |
messaging |
███████████████████░ 95.8% |
3,798 / 3,963 |
skills |
███████████████████░ 95.8% |
6,972 / 7,274 |
product_analytics |
███████████████████░ 96.0% |
28,470 / 29,647 |
access_control |
███████████████████░ 96.3% |
7,113 / 7,386 |
revenue_analytics |
███████████████████░ 96.4% |
1,876 / 1,946 |
user_interviews |
███████████████████░ 96.5% |
2,867 / 2,971 |
feature_flags |
███████████████████░ 96.5% |
25,588 / 26,509 |
warehouse_sources |
███████████████████░ 97.2% |
459,581 / 472,600 |
data_quality |
████████████████████ 97.6% |
7,587 / 7,774 |
links |
████████████████████ 97.9% |
234 / 239 |
security |
████████████████████ 98.0% |
1,286 / 1,312 |
metrics |
████████████████████ 98.0% |
4,084 / 4,166 |
analytics_platform |
████████████████████ 98.3% |
2,784 / 2,833 |
pulse |
████████████████████ 98.5% |
2,046 / 2,078 |
live_debugger |
████████████████████ 99.2% |
626 / 631 |
field_notes |
████████████████████ 99.4% |
172 / 173 |
Report-only. Patch coverage = changed backend lines covered vs origin/master. Sorted lowest first.
Known gaps: lines covered only by Temporal tests show as uncovered; core line numbers may drift if master changed the same file.
There was a problem hiding this comment.
Actionable comments posted: 2
ℹ️ Review info
⚙️ Run configuration
Configuration used: Repository: PostHog/posthog/.coderabbit.yaml
Review profile: QUIET
Plan: Enterprise
Run ID: b0dac313-3024-43c0-b0bb-007001f455ec
📒 Files selected for processing (2)
products/signals/backend/emission/pganalyze_issues.pyproducts/signals/backend/emission/tests/test_pganalyze_issues.py
Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 1 remain after this review.
| def pganalyze_issue_record_fetcher( | ||
| team: Team, | ||
| config: SignalSourceTableConfig, | ||
| context: dict[str, Any], | ||
| ) -> list[dict[str, Any]]: | ||
| """Fetch pganalyze issues that were not emitted before. | ||
|
|
||
| pganalyze's getIssues returns no timestamps, so the warehouse source stamps each row with the | ||
| sync time and merges on `id`. Thus the `synced_at` cursor returns every open issue again on each | ||
| sync. The dedupe runs before `max_records` applies, so a backlog larger than one batch drains | ||
| over the next syncs. | ||
| """ | ||
| id_config = config.model_copy(update={"fields": ("id",), "max_records": ISSUE_ID_SCAN_LIMIT}) | ||
| candidate_ids = sorted( | ||
| {str(row["id"]) for row in data_warehouse_record_fetcher(team, id_config, context) if row.get("id")} | ||
| ) | ||
| if not candidate_ids: | ||
| return [] | ||
|
|
||
| already_emitted = set( | ||
| SignalEmissionRecord.objects.filter( | ||
| team=team, | ||
| source_product=config.source_product, | ||
| source_type=config.source_type, | ||
| source_id__in=candidate_ids, | ||
| ).values_list("source_id", flat=True) | ||
| ) | ||
| new_ids = [issue_id for issue_id in candidate_ids if issue_id not in already_emitted][: config.max_records] | ||
| if not new_ids: | ||
| return [] | ||
|
|
There was a problem hiding this comment.
🎯 Functional Correctness | 🟠 Major | 🏗️ Heavy lift
🔎 Supported by static analysis
🏁 Script executed:
sed -n '135,225p' products/signals/backend/emission/pganalyze_issues.py
rg -n 'def data_warehouse_record_fetcher|def .*record_fetcher|ISSUE_ID_SCAN_LIMIT' products/signals/backend/emissionRepository: PostHog/posthog
Length of output: 4036
🏁 Script executed:
sed -n '1,150p' products/signals/backend/emission/fetchers/data_warehouse.py
sed -n '1,140p' products/signals/backend/emission/pganalyze_issues.py
git diff --unified=30 b6aa2611afc7b08d8b30b730529dfd1f7cc39fb6 c9bf7d96a21fc3e239d7c562116943ef7d5a08a6 -- products/signals/backend/emission/pganalyze_issues.py products/signals/backend/emission/fetchers/data_warehouse.py
rg -n "data_warehouse_record_fetcher\(|max_records|partition_field|synced_at|SignalEmissionRecord" products/signals/backend/emission/fetchers products/signals/backend/emission/facade products/signals/backend/emission -g '*.py' | head -240Repository: PostHog/posthog
Length of output: 42043
🏁 Script executed:
sed -n '1,220p' products/signals/backend/emission/tests/test_pganalyze_issues.py
sed -n '1,130p' products/signals/backend/emission/emit_signals.py
rg -n "pganalyze|synced_at|getIssues|issue.*sync|last_synced_at" products integrations dags posthog -g '*.py' -g '*.ts' -g '*.sql' 2>/dev/null | head -260Repository: PostHog/posthog
Length of output: 40827
Paginate the ID scan before applying emission deduplication.
data_warehouse_record_fetcher applies LIMIT 10,000 before SignalEmissionRecord filtering. It has no ordering or advancing cursor. If more than 10,000 open issues remain eligible on each sync, the warehouse can return the same page repeatedly. After that page is recorded, new_ids is empty and issues outside the page never reach _fetch_issues_by_id.
Use an advancing, ordered pagination strategy and apply emission deduplication to each page before stopping at config.max_records. This is separate from the existing failure-after-fetch concern: this path can lose issues even when the selected batch and pipeline complete successfully.
| now = timezone.now() | ||
| SignalEmissionRecord.objects.bulk_create( | ||
| [ | ||
| SignalEmissionRecord( | ||
| team=team, | ||
| source_product=config.source_product, | ||
| source_type=config.source_type, | ||
| source_id=str(row["id"]), | ||
| emitted_at=now, | ||
| ) | ||
| for row in rows | ||
| ], | ||
| ignore_conflicts=True, | ||
| ) | ||
| return rows |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift
🔎 Supported by static analysis
🏁 Script executed:
#!/bin/bash
# Check whether the pipeline records emissions itself, and how activity retries work.
rg -nP -C4 'SignalEmissionRecord' --type=py -g '!**/tests/**'
rg -nP -C4 '\brecord_fetcher\s*\(' --type=py products/signals
rg -nP -C3 'RetryPolicy|maximum_attempts' --type=py products/signalsRepository: PostHog/posthog
Length of output: 41099
🏁 Script executed:
set -e
printf '%s\n' '--- pganalyze module ---'
cat -n products/signals/backend/emission/pganalyze_issues.py | sed -n '1,260p'
printf '%s\n' '--- record_fetcher call chain ---'
rg -n -P -C10 'record_fetcher|SignalEmitterOutput|SignalSourceTableConfig' products/signals/backend/emission --type=py
printf '%s\n' '--- pganalyze workflow/activity references ---'
rg -n -P -C10 'pganalyze|sync.*signal|signal.*sync|execute_activity' products/signals/backend --type=py
printf '%s\n' '--- changed-file diff ---'
git diff b6aa2611afc7b08d8b30b730529dfd1f7cc39fb6 c9bf7d96a21fc3e239d7c562116943ef7d5a08a6 -- products/signals/backend/emission/pganalyze_issues.py products/signals/backend/emissionRepository: PostHog/posthog
Length of output: 45652
🏁 Script executed:
set -e
printf '%s\n' '--- exact fetcher consumers ---'
rg -n -C12 'config\.record_fetcher|record_fetcher\(|\.record_fetcher' products/signals/backend --glob '*.py' --glob '!**/tests/**'
printf '%s\n' '--- pipeline processing and sync entrypoints ---'
rg -n -C10 'actionability_prompt|summarization_prompt|description_summarization|emitter\(' products/signals/backend/emission --glob '*.py' --glob '!**/tests/**'
printf '%s\n' '--- pganalyze sync/retry references ---'
rg -n -C12 'PGANALYZE_ISSUES_CONFIG|pganalyze.*sync|sync.*pganalyze|RetryPolicy|maximum_attempts' products/signals/backend --glob '*.py' --glob '!**/tests/**'Repository: PostHog/posthog
Length of output: 45667
Record pganalyze issues after successful emission.
emit_signals.py calls pganalyze_issue_record_fetcher before run_signal_pipeline. The fetcher records every issue before that pipeline runs. If emission, validation, or summarization fails, the next sync excludes the issue through SignalEmissionRecord, so the issue can be lost. Move the ledger write into the shared pipeline after successful signal creation, or write it after emission completes. Add a regression test for downstream failure and retry.
|
PostHog Review alpha 🦔 If you find any issues helpful - please reply "valid", "invalid", etc., for evaluation purposes 🙏 |
|
/trunk merge |
The warehouse merges pganalyze issues on id inside one weekly synced_at partition, so an issue open across a week boundary has one row per week. The by-id fetch now orders by synced_at and keeps one row per id, so a stale row cannot emit or push a current id out of the batch. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Generated-By: PostHog Desktop Task-Id: b4b321b2-7810-46a3-8fdb-3951d818388f
A new stamphog review started for this PR — the fresh verdict replaces this approval.
…lyze-re-emitting-7b9d7c Generated-By: PostHog Desktop Task-Id: 868dbc0c-8a28-4a69-bc79-f7f6604ef301
A new stamphog review started for this PR — the fresh verdict replaces this approval.
|
/trunk cancel |
Generated-By: PostHog Desktop Task-Id: 868dbc0c-8a28-4a69-bc79-f7f6604ef301
A new stamphog review started for this PR — the fresh verdict replaces this approval.
|
/trunk merge |
There was a problem hiding this comment.
Actionable comments posted: 2
ℹ️ Review info
⚙️ Run configuration
Configuration used: Repository: PostHog/posthog/.coderabbit.yaml
Review profile: QUIET
Plan: Enterprise
Run ID: c7da7fd0-2df3-4498-ac96-4364e5135be6
📒 Files selected for processing (4)
products/signals/backend/emission/pganalyze_issues.pyproducts/signals/backend/emission/pipeline.pyproducts/signals/backend/emission/registry.pyproducts/signals/backend/emission/tests/test_pganalyze_issues.py
Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 0 remain after this review.
| 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) |
There was a problem hiding this comment.
🎯 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.
|
/trunk cancel |
Generated-By: PostHog Desktop Task-Id: 868dbc0c-8a28-4a69-bc79-f7f6604ef301
A new stamphog review started for this PR — the fresh verdict replaces this approval.
|
/trunk merge |
There was a problem hiding this comment.
Approved.
Contained fix confined to the signals emission code for one source. It adds no migrations or API changes, and the new behavior is opt-in per source and covered by tests. A human teammate approved an earlier commit, and the later commit addressed the idempotency concern. The remaining open bot thread is an edge case that can be fixed forward.
- Open CodeRabbit thread: if the first max_records issues fail emission on every sync, they are refetched each time and could starve later issues. This is an edge case and does not block approval; consider a follow-up.
Gate mechanics and policy version
| Gate | Result | |
|---|---|---|
| prerequisites | ✓ | all clear |
| deny-list | ✓ | no deny categories matched |
| size | ✓ | 104L, 3F substantive, 259L/5F incl. docs/generated/snapshots — within ceiling |
| tier | ✓ | T1-agent / T1c-medium (259L, 5F, single-area, fix) |
| stamphog 2.3.1 | .stamphog/policy.yml @ c685668 · reviewed head c685668 |
There was a problem hiding this comment.
Actionable comments posted: 1
ℹ️ Review info
⚙️ Run configuration
Configuration used: Repository: PostHog/posthog/.coderabbit.yaml
Review profile: QUIET
Plan: Enterprise
Run ID: 956c071b-e220-41f3-92d5-9b4191fa3145
📒 Files selected for processing (3)
products/signals/backend/emission/pipeline.pyproducts/signals/backend/emission/tests/test_emit_signals.pyproducts/signals/backend/emission/tests/test_pganalyze_issues.py
Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 2 remain after this review.
| if record_processed_outputs: | ||
| await _record_processed_outputs(team, [output]) |
There was a problem hiding this comment.
🗄️ 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.pyRepository: PostHog/posthog
Length of output: 13500
🏁 Script executed:
sed -n '487,653p' products/signals/backend/emission/pipeline.pyRepository: PostHog/posthog
Length of output: 7888
🏁 Script executed:
rg -n -C 8 'SignalEmissionRecord|pganalyze|processed_outputs' products/signalsRepository: 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.
|
/trunk cancel |
|
/trunk merge --no-batch |
Problem
Teams using pganalyze receive duplicate signals for open issues on every warehouse sync.
Why: Snapshot timestamps advance each sync, so a time cursor alone repeats the same issues and their processing.
Closes #108338
Origin
453fbddChanges
Before:
flowchart LR A[Sync] --> B[Fetch open issues] --> C[Emit again] classDef phYellow fill:#f9bd2b,stroke:#f9bd2b,color:#000; classDef phBlue fill:#1d4aff,stroke:#1d4aff,color:#fff; class A,C phYellow; class B phBlue;After:
flowchart LR A[Sync] --> B[Page latest unprocessed issues] --> C[Process] --> D[Record successful outcomes] classDef phYellow fill:#f9bd2b,stroke:#f9bd2b,color:#000; classDef phBlue fill:#1d4aff,stroke:#1d4aff,color:#fff; class A,D phYellow; class B,C phBlue;How did you test this code?
Release status
Automatic notifications
Docs update
None. No existing scoped document covers this fetcher.
🤖 Agent context
Autonomy: Human-driven (agent-assisted), following the initial autonomous implementation.
Agent: PostHog Desktop: Claude Code, claude-opus-5-5 (initial implementation); Codex, GPT-6 (CI and review fixes).
Skills: investigating-ci-failures, writing-tests, writing-code-comments, writing-pr-descriptions, running-ci-preflight, merging-prs, announcing-behavior-changes. Duplicate search found #108009, which covers pganalyze access errors. Test data is synthetic; no new event-property reads.
Created with PostHog Desktop from an inbox report