Skip to content

trunk-merge/pr-89156/15cdeb55-b6d4-42fc-adbb-2b2605afdcda - #107630

Closed
trunk-io[bot] wants to merge 311 commits into
masterfrom
trunk-merge/pr-89156/15cdeb55-b6d4-42fc-adbb-2b2605afdcda
Closed

trunk-io[bot] wants to merge 311 commits into
masterfrom
trunk-merge/pr-89156/15cdeb55-b6d4-42fc-adbb-2b2605afdcda

Conversation

@trunk-io

@trunk-io trunk-io Bot commented Sep 28, 2026

Copy link
Copy Markdown
Trunk Merge Pull Request Banner

This pull request was created and is being managed by Trunk Merge.

This pull request is based on the master branch at SHA cd7f5b6241b6f6b661536a2ff70aeb958598c8ae.

See more details about each PR in the batch here:

When CI completes, this pull request will be closed automatically.

Pull Requests Being Tested

This pull request is testing a batch with the changes from pull requests 89156, 107255, and 107252 - batching documentation.

Dependencies

This pull request depends on the changes from pull requests 107389, 102336, 107079, and 104978.

dmarchuk and others added 30 commits September 23, 2026 15:44
mypy rejected `**dict[str, object]` against fetch_app_metric_totals' typed keyword parameters. Annotate the parameterised filter mapping as dict[str, Any] so the unpack type-checks.
A double-click in the SQL editor selects one identifier, and the save
dialogs treat any selection as the query to save. The dialog then sends
that fragment to the API, which rejects it as an invalid query.

The dialogs now parse a selected candidate before the person can submit,
and hold the save with a message that names the selection as the cause.
Only a selection is judged, because clearing it is always a way forward.
The whole editor contents are left to the API. If the parser fails to
load the check passes, so a missing parser cannot block a save.

The candidate preview also moves above the name field, so the query
being saved is visible before anything is typed.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The message told the person to clear the selection, which cannot be done
while the dialog is open. It now names closing the dialog first.

Adds a dialog-level test for the view, endpoint and metric dialogs. The
helper tests cover the parse verdict but not the wiring, so a dropped
saveTarget error would have let a selected identifier submit.

Generated-By: PostHog Desktop
Task-Id: 7b0812e1-1261-40dc-bdad-aad6a7172281
Generated-By: PostHog Desktop
Task-Id: 7b0812e1-1261-40dc-bdad-aad6a7172281
In raw mode the text goes to the connection's own engine in its own
dialect, so parsing it as HogQL refused valid SQL. The save dialogs now
pass sendRawQueryEnabled and the check stands down.

Also makes the refusal message name the last step, saving again.

Generated-By: PostHog Desktop
Task-Id: 7b0812e1-1261-40dc-bdad-aad6a7172281
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016F8jY3Du54fcQUdCQkr3ZE
…and validate daily

Review round on the Temporal layer:

- kickoff holds the pipeline row lock through run_training, as start_training does, and a failed launch rolls back the deduction and raises for the activity retry
- discovery pauses a pipeline whose creator lost team access and skips one outside the flag rollout
- validation runs on every sweep; scoring and kickoff run on a cadence day compared by calendar day, and kickoff starts after scoring ends
- scoring and validation re-read the pipeline status
- the queue defaults to the general-purpose fleet, a deploy keeps the schedule's paused state, and child timeouts cover every retry
… once per day

Second review round on the Temporal layer:

- kickoff applies the Desktop access and Tasks usage gates that /train applies, before the row lock
- kickoff launches at most one run per UTC day, so a retry after a lost activity response cannot start a second paid run
- scoring reads the current champion inside its activity, so a retry after a mid-scoring promotion scores with the new one
- discovery passes the loaded organization id to the flag check
- the general-purpose worker redeploys on autoresearch backend changes
…nator run

- inference and validation heartbeat through HeartbeaterSync with a two-minute heartbeat timeout, so a lost worker is retried in minutes
- the schedule gives the coordinator a twelve-hour execution timeout, so a stuck run ends before the next tick instead of making SKIP drop it
- the temporal AGENTS.md says which logic stays in the sweep and which lives in the packages it calls
…unch

Kickoff waits for scoring, hours after discovery checked the flag, so the launch gate checks the autoresearch rollout again.
Discovery leaves out a pipeline whose organization is deactivated (a null is_active included) or pending deletion, and the kickoff gate checks again before a paid launch.
mayteio and others added 25 commits September 28, 2026 11:55
@coderabbitai

coderabbitai Bot commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Review in Change Stack →

Navigate logical layers of code changes, visualize relationships, and explore their blast radius.

📝 Walkthrough

Walkthrough

The changes add HogQL traversal metadata to published catalogs, autoresearch Temporal workflows and scheduling, configurable calendar sync intervals, and HeyReach and MoEngage warehouse sources. They also change Hogland sandbox URL validation and request routing, add flag-based hiding and replacement guidance for MCP loop tools, and update worker deployment coverage.

Priority: ➖ Normal

Merge Risk: 🟡 Moderate · up to a432a

The new HeyReach source can skip earlier pages when a sync is retried after a completed run. The daily autoresearch job can time out as the number of pipelines grows. The loop-to-workflow guidance can send users to steps that fail for active loops or omit earlier run history. These issues should be fixed before merge.

Security Architecture Review

Security architecture risk: 🟡 Moderate · up to a432a

Scheduled inference can repeat externally visible prediction work after an interrupted attempt. The review found controls that limit cross-team execution and sandbox routing, but several integration boundaries remain only partially verified.

Retained concerns

  • Medium · reliability · inferred: The newly scheduled inference activity retries a scorer that creates a run and emits predictions on every invocation. Interruption after emission, or loss of a successful activity response, can repeat a prediction date before recovery completes, affecting tenant prediction-data integrity. Downstream event deduplication was not established.
Security review details

Security Blast Radius

  • inferred — The inference retry exposure is independently reachable for pipelines selected by the daily coordinator, potentially across active teams; the inspected activity reloads each pipeline under its team scope rather than granting cross-team access.

Security Findings and Attack Paths

  • inferred — No cross-tenant or credential-exfiltration path was established. The supported integrity path is interruption or retry after inference side effects, followed by another invocation of the run-creating scorer; event-sink deduplication remains unverified.

Trust Boundaries and Controls

  • observed — The Hogland proxy bypass is restricted to URLs passing the shared exact-origin gate, and redirects are disabled on that path. Non-Hogland, non-credential commands still permit redirects on the standard request path; that behavior was already present before this change.

Resilience and Maintainability Implications

  • observed — Training kickoff locks the pipeline and checks budget, active runs, and same-day runs before deducting budget and launching inside a transaction. These controls limit repeated local training launches; they do not establish idempotency for inference emissions.

Hardening Proposals

  • proposed — Make scheduled inference recoverable by identifying and reconciling work by pipeline, model, and prediction date before repeating run creation or prediction emission.
🚥 Pre-merge checks | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Description check ⚠️ Warning The description only documents Trunk Merge batch metadata, tested pull requests, and dependencies. It does not include the required Problem, Changes, testing, Release status, Automatic notifications, … Replace or supplement the Trunk Merge text with a standalone repository-formatted description. Explain the user problem, visible changes, tests and untested areas, select exactly one release-status option, state changelog and docs actions, …
Full details: Description check

Explanation

The description only documents Trunk Merge batch metadata, tested pull requests, and dependencies. It does not include the required Problem, Changes, testing, Release status, Automatic notifications, Docs update, or Agent context sections.

Resolution

Replace or supplement the Trunk Merge text with a standalone repository-formatted description. Explain the user problem, visible changes, tests and untested areas, select exactly one release-status option, state changelog and docs actions, and complete or remove the Agent context section as applicable.

  • Fix all pre-merge checks with AI
✨ Finishing Touches
📝 Generate docstrings
  • Commit to this branch
  • Create a new PR
🧪 Generate unit tests (beta)
  • Commit to this branch
  • Create a new PR

Comment @coderabbitai help to get the list of available commands.

@coderabbitai coderabbitai Bot left a comment

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.

Actionable comments posted: 9

Note

Quiet mode is enabled, so only the most important comments were posted inline. Other review comments are grouped below.

🟡 Other comments (3)
posthog/hogql/catalog_traversal.py-186-204 (1)

186-204: 🚀 Performance & Scalability | 🟡 Minor | ⚡ Quick win

Cache the negative result in _canonical_name.

The early return on Line 190-191 skips the cache. Line 199-201 also returns without writing to the cache. When serialization raises, each later lookup of the same target runs serialize_fields again and increments canonical_unserializable again. The per-path checks at Line 176 and Line 247 repeat this work, and the omission counts grow on every repeat. Store None in the cache before returning.

Proposed fix
         except Exception:
             self.omissions["canonical_unserializable"] += 1
+            self.canonical_name_cache[id(target)] = None
             return None
products/tasks/backend/presentation/views/api.py-3469-3469 (1)

3469-3469: 🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win

Reject malformed ports as invalid sandbox URLs. If connection.sandbox_url contains a nonnumeric port, the delegated check reads target.port outside its try block and raises ValueError. The command endpoint then returns a server error instead of its invalid-URL response. Catch port-parsing errors in products/tasks/backend/logic/services/agent_command.py and return False.

products/customer_analytics/frontend/scenes/CustomerAnalyticsConfigurationScene/calendar/calendarSyncLogic.ts-263-263 (1)

263-263: 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win

Keep the saved interval visible when status refresh fails. After a successful POST, loadStatuses() starts an asynchronous refresh. Its error handler returns the previous statuses, while finally clears the saving state. If that refresh fails for an idle integration, the selector shows the old interval again until another load occurs. Update the displayed interval after the successful POST, and let a later refresh reconcile it. (keajs.org)

🧹 Nitpick comments (2)
posthog/hogql/catalog_traversal.py (1)

249-249: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Handle a missing chain component explicitly.

When a FieldTraverser names a field that does not exist, current.fields[str(component)] raises KeyError. The broad except Exception in _relation_for_field catches that error, so the catalog still builds. The problem is that this relies on the broad handler. Raise ResolutionError directly so the handler's intent stays clear and a real bug does not get recorded as target_unresolvable.

Proposed fix
-            next_value = current.fields[str(component)]
+            next_value = current.fields.get(str(component))
+            if next_value is None:
+                raise ResolutionError("catalog traversal reached an unknown field")
products/warehouse_sources/backend/temporal/data_imports/sources/heyreach/heyreach.py (1)

109-115: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Add the missing return annotation on _top_level_rows.

The function has no return type. The repo guideline says "Annotate every signature". Use -> Resource, or the type that rest_api_resource returns.

Source: Coding guidelines


ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository: PostHog/posthog/.coderabbit.yaml

Review profile: QUIET

Plan: Enterprise

Run ID: 6de4135d-bf00-4698-b23d-f361e6cdf548

📥 Commits

Reviewing files that changed from the base of the PR and between cd7f5b6 and a432a36.

⛔ Files ignored due to path filters (3)
  • products/customer_analytics/frontend/generated/api.schemas.ts is excluded by !**/generated/**
  • products/customer_analytics/frontend/generated/api.ts is excluded by !**/generated/**
  • products/customer_analytics/frontend/generated/api.zod.ts is excluded by !**/generated/**
📒 Files selected for processing (59)
  • .github/workflows/container-images-cd.yml
  • docs/internal/hogql-language-service.md
  • posthog/api/services/query.py
  • posthog/api/test/test_query_service.py
  • posthog/hogql/catalog_traversal.py
  • posthog/hogql/language_service.py
  • posthog/hogql/test/test_language_service.py
  • posthog/management/commands/start_temporal_worker.py
  • posthog/settings/temporal.py
  • posthog/temporal/schedule.py
  • products/autoresearch/backend/dataset/AGENTS.md
  • products/autoresearch/backend/evaluation/AGENTS.md
  • products/autoresearch/backend/facade/temporal.py
  • products/autoresearch/backend/inference/AGENTS.md
  • products/autoresearch/backend/temporal/AGENTS.md
  • products/autoresearch/backend/temporal/CLAUDE.md
  • products/autoresearch/backend/temporal/__init__.py
  • products/autoresearch/backend/temporal/schedule.py
  • products/autoresearch/backend/temporal/test_workflows.py
  • products/autoresearch/backend/temporal/workflows.py
  • products/autoresearch/backend/training/AGENTS.md
  • products/customer_analytics/backend/facade/api.py
  • products/customer_analytics/backend/facade/contracts.py
  • products/customer_analytics/backend/logic/calendar_sync.py
  • products/customer_analytics/backend/presentation/views/serializers.py
  • products/customer_analytics/backend/presentation/views/views.py
  • products/customer_analytics/backend/temporal/calendar_sync.py
  • products/customer_analytics/backend/test/test_calendar_sync.py
  • products/customer_analytics/backend/test/test_views.py
  • products/customer_analytics/frontend/scenes/CustomerAnalyticsConfigurationScene/calendar/CalendarSyncConfig.tsx
  • products/customer_analytics/frontend/scenes/CustomerAnalyticsConfigurationScene/calendar/calendarSyncLogic.ts
  • products/customer_analytics/mcp/tools.yaml
  • products/tasks/backend/facade/api.py
  • products/tasks/backend/presentation/views/api.py
  • products/tasks/backend/tests/test_api.py
  • products/tasks/backend/tests/test_sandbox_url_validation.py
  • products/tasks/mcp/tools.yaml
  • products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md
  • products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/heyreach.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/moengage.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/heyreach/canonical_descriptions.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/heyreach/heyreach.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/heyreach/settings.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/heyreach/source.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/heyreach/tests/test_heyreach.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/heyreach/tests/test_heyreach_source.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/moengage/canonical_descriptions.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/moengage/moengage.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/moengage/settings.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/moengage/source.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/moengage/tests/test_moengage.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/moengage/tests/test_moengage_source.py
  • services/mcp/schema/generated-tool-definitions.json
  • services/mcp/schema/tool-definitions-all.json
  • services/mcp/schema/tool-definitions.json
  • services/mcp/scripts/generate-tools.ts
  • services/mcp/scripts/yaml-config-schema.ts
  • services/mcp/src/api/generated.ts
  • services/mcp/tests/unit/tool-filtering.test.ts
💤 Files with no reviewable changes (4)
  • products/autoresearch/backend/evaluation/AGENTS.md
  • products/autoresearch/backend/training/AGENTS.md
  • products/autoresearch/backend/inference/AGENTS.md
  • products/autoresearch/backend/dataset/AGENTS.md

Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 4 remain after this review.

Comment on lines +592 to +597
active = await workflow.execute_activity(
activity_load_active_pipelines,
LoadActivePipelinesInput(),
start_to_close_timeout=timedelta(minutes=2),
retry_policy=_COORDINATOR_RETRY,
)

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.

🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

Size the discovery activity for per-pipeline network calls.

activity_load_active_pipelines loops over every live pipeline across all teams. For each pipeline, it runs _resolve_acting_user, which is a DB query. It also runs has_autoresearch_access, which can make a remote flag evaluation through posthog_feature_flag_value or feature_enabled(only_evaluate_locally=False).

The activity has a 2-minute start_to_close_timeout and no heartbeat. When the pipeline count grows or flag latency rises, the activity times out, and each of the 3 retries repeats the whole sweep. If all retries time out, the daily tick fails and no pipeline is validated or scored.

To fix this, either raise the timeout and add HeartbeaterSync() with a heartbeat_timeout, or split discovery into bounded pages.

Proposed fix
         active = await workflow.execute_activity(
             activity_load_active_pipelines,
             LoadActivePipelinesInput(),
-            start_to_close_timeout=timedelta(minutes=2),
+            start_to_close_timeout=timedelta(minutes=15),
+            heartbeat_timeout=_HEARTBEAT_TIMEOUT,
             retry_policy=_COORDINATOR_RETRY,
         )

Also wrap the activity body in with HeartbeaterSync():.

📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
active = await workflow.execute_activity(
activity_load_active_pipelines,
LoadActivePipelinesInput(),
start_to_close_timeout=timedelta(minutes=2),
retry_policy=_COORDINATOR_RETRY,
)
active = await workflow.execute_activity(
activity_load_active_pipelines,
LoadActivePipelinesInput(),
start_to_close_timeout=timedelta(minutes=15),
heartbeat_timeout=_HEARTBEAT_TIMEOUT,
retry_policy=_COORDINATOR_RETRY,
)

counts: CalendarSyncCounts


def get_calendar_sync_interval(config: dict) -> int:

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.

📐 Maintainability & Code Quality | 🟠 Major | ⚡ Quick win

Specify the calendar config types.

The bare dict annotation leaves key and value types implicit in this new signature. Use a typed mapping, such as dict[str, object], and keep the exact-integer check when narrowing the stored value. As per coding guidelines, “Write as if mypy --strict were on. Annotate every signature, avoid Any.”

Source: Coding guidelines

sync_interval_minutes = serializers.IntegerField(help_text="Minutes between scheduled syncs: 5, 15, 30, or 60.")

def validate_sync_interval_minutes(self, value: int) -> int:
if value not in (5, 15, 30, 60):

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

Use one set of allowed intervals.

This tuple duplicates ALLOWED_SYNC_INTERVALS in products/customer_analytics/backend/logic/calendar_sync.py. If the two lists diverge, the API can save an interval that the status getter and scheduler treat as 60 minutes. Move the choices to a shared contract and use them in both places. As per path instructions, code should “say everything once and only once.”

Source: Path instructions

summary="Set Google account sync interval",
)
@action(methods=["POST"], detail=False, url_path="interval")
def interval(self, request: ValidatedRequest, *args, **kwargs) -> Response:

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.

📐 Maintainability & Code Quality | 🟠 Major | ⚡ Quick win

Annotate the new action’s variadic parameters.

interval leaves *args and **kwargs unannotated. Add parameter types to this new signature. As per coding guidelines, “Write as if mypy --strict were on. Annotate every signature, avoid Any.”

Source: Coding guidelines

},
)

def test_sync_interval_defaults_and_is_scoped_to_one_account(self):

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.

📐 Maintainability & Code Quality | 🟠 Major | ⚡ Quick win

Annotate the new test signatures. Add -> None to all three new test methods. As per coding guidelines, “Write as if mypy --strict were on. Annotate every signature.”

Also applies to: 3545-3545, 3562-3562

Source: Coding guidelines

"https://hogland.prod-us.posthog.dev/v1/hogboxes/box-1/proxy/8080/command",
)

def test_command_blocks_hogland_sandbox_url_when_unconfigured(self):

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.

📐 Maintainability & Code Quality | 🟠 Major | ⚡ Quick win

Annotate the new test signatures. The new tests omit parameter or return annotations.

  • products/tasks/backend/tests/test_api.py#L14326-L14326: add -> None.
  • products/tasks/backend/tests/test_api.py#L14299-L14299: annotate both mock parameters as MagicMock and add -> None.
  • products/tasks/backend/tests/test_api.py#L14349-L14349: annotate both mock parameters as MagicMock and add -> None.
  • products/tasks/backend/tests/test_sandbox_url_validation.py#L44-L44: annotate url: str, expected: bool, and -> None.

As per coding guidelines, “Write as if mypy --strict were on. Annotate every signature.”

📍 Affects 2 files
  • products/tasks/backend/tests/test_api.py#L14326-L14326 (this comment)
  • products/tasks/backend/tests/test_api.py#L14299-L14299
  • products/tasks/backend/tests/test_api.py#L14349-L14349
  • products/tasks/backend/tests/test_sandbox_url_validation.py#L44-L44

Source: Coding guidelines

Comment on lines +195 to +196
superseded_by:
- workflows-archive

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 | ⚡ Quick win

Provide the required disable step before archiving an active loop.

If an active loop receives a loops-destroy request, the redirect names workflows-archive. The supplied workflows-archive definition requires the workflow to be draft or disabled first. Directly following this redirect fails for an active loop. Name the disable step and its tool before workflows-archive, or provide a replacement that accepts active loops.

Comment on lines +330 to +338
superseded_by:
- tasks-list
- tasks-runs-retrieve
- workflows-list-invocations
redirect_hint:
'Loops in this project run on workflows: read the building-loops skill, then use the workflows tools with
origin_product "loops". For run history, call tasks-list with the hog_flow_id of the loop and archived set
to all: it lists the tasks the loop created, newest first, each with the status and error of its latest run.
Then call tasks-runs-retrieve on a run for its branch and output, including any PR URL.

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 | ⚡ Quick win

Include tasks-runs-list in the run-history redirect.

If a loop-created task has multiple runs, tasks-list exposes only its latest run. The hint then directs the caller to tasks-runs-retrieve, but gives no way to find earlier run IDs. Add tasks-runs-list to the successors and direct callers to list each task's runs before retrieving run details.

Comment on lines +122 to +124
def save_checkpoint(state: Optional[dict[str, Any]]) -> None:
if state and state.get("offset"):
resumable_source_manager.save_state(HeyReachResumeConfig(offset=int(state["offset"])))

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.

🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

🔎 Supported by static analysis

🏁 Script executed:

#!/bin/bash
rg -nP '\bclear_state\s*\(' products/warehouse_sources/backend/temporal -g '!**/tests/**' -C3 | head -80

Repository: PostHog/posthog

Length of output: 9616


🏁 Script executed:

#!/bin/bash
set -e
printf '%s\n' '--- HeyReach source outline and relevant source ---'
ast-grep outline products/warehouse_sources/backend/temporal/data_imports/sources/heyreach/heyreach.py
sed -n '1,190p' products/warehouse_sources/backend/temporal/data_imports/sources/heyreach/heyreach.py
printf '%s\n' '--- ResumableSourceManager definitions and references ---'
rg -n -C4 'class ResumableSourceManager|def (save_state|clear_state|can_resume)|ResumableSourceManager' products/warehouse_sources/backend/temporal -g '*.py' | head -240
printf '%s\n' '--- SourceResponse consumption and success cleanup candidates ---'
rg -n -C4 'SourceResponse|clear_state\(\)|save_checkpoint' products/warehouse_sources/backend/temporal -g '*.py' | head -300

Repository: PostHog/posthog

Length of output: 41700


🏁 Script executed:

#!/bin/bash
set -e
printf '%s\n' '--- locate relevant definitions ---'
rg -l 'class ResumableSourceManager|class SourceResponse|def _run' products/warehouse_sources/backend/temporal -g '*.py'
printf '%s\n' '--- manager implementation ---'
manager=$(rg -l 'class ResumableSourceManager' products/warehouse_sources/backend/temporal -g '*.py' | head -1)
cat -n "$manager"
printf '%s\n' '--- source response and workflow definitions ---'
for file in $(rg -l 'class SourceResponse|def _run' products/warehouse_sources/backend/temporal -g '*.py'); do
  printf '\n--- %s ---\n' "$file"
  rg -n -C8 'class SourceResponse|def _run' "$file"
done

Repository: PostHog/posthog

Length of output: 42400


🏁 Script executed:

#!/bin/bash
set -e
printf '%s\n' '--- resumable manager ---'
sed -n '1,230p' products/warehouse_sources/backend/temporal/data_imports/sources/common/resumable.py
printf '%s\n' '--- shared import workflow ---'
sed -n '900,1065p' products/warehouse_sources/backend/temporal/data_imports/workflow_activities/import_data_sync.py
printf '%s\n' '--- REST resource and paginator definitions ---'
rg -l 'class OffsetPaginator|def rest_api_resource|resume_hook|initial_paginator_state' products/warehouse_sources/backend/temporal/data_imports/sources/common -g '*.py'
for file in $(rg -l 'class OffsetPaginator|def rest_api_resource|resume_hook|initial_paginator_state' products/warehouse_sources/backend/temporal/data_imports/sources/common -g '*.py'); do
  printf '\n--- %s ---\n' "$file"
  rg -n -C12 'class OffsetPaginator|def rest_api_resource|resume_hook|initial_paginator_state' "$file"
done

Repository: PostHog/posthog

Length of output: 42553


🏁 Script executed:

#!/bin/bash
set -e
printf '%s\n' '--- cleanup calls in workflow and pipelines ---'
rg -n -C5 'clear_state\s*\(' products/warehouse_sources/backend/temporal/data_imports/workflow_activities products/warehouse_sources/backend/temporal/data_imports/pipelines -g '*.py' || true
printf '%s\n' '--- pipeline run definitions and resumable manager handling ---'
rg -n -C12 'class PipelineNonDLT|class PipelineV3|def run\(|resumable_source_manager' products/warehouse_sources/backend/temporal/data_imports/pipelines -g '*.py' | head -260

Repository: PostHog/posthog

Length of output: 35378


Clear the resume state after the final page.

RESTClient.paginate calls save_checkpoint with None after the terminal page. The callback ignores that value. The previous offset can remain in Redis for 24 hours, so a later attempt with the same job ID can skip earlier pages.

🐛 Suggested fix
     def save_checkpoint(state: Optional[dict[str, Any]]) -> None:
+        if state is None:
+            resumable_source_manager.clear_state()
+            return
         if state and state.get("offset"):
             resumable_source_manager.save_state(HeyReachResumeConfig(offset=int(state["offset"])))
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
def save_checkpoint(state: Optional[dict[str, Any]]) -> None:
if state and state.get("offset"):
resumable_source_manager.save_state(HeyReachResumeConfig(offset=int(state["offset"])))
def save_checkpoint(state: Optional[dict[str, Any]]) -> None:
if state is None:
resumable_source_manager.clear_state()
return
if state and state.get("offset"):
resumable_source_manager.save_state(HeyReachResumeConfig(offset=int(state["offset"])))

@trunk-io

trunk-io Bot commented Sep 28, 2026 •

Copy link
Copy Markdown
Author

Static Badge   Static Badge   Static Badge

Failed Test Failure Summary Logs
personalAPIKeysLogic preserves the `*` (all access) scope regardless of flag state The test exceeded the maximum allowed time of 5000 milliseconds and timed out. Logs ↗︎

View Full Report ↗︎ ⋅ Docs

@trunk-io trunk-io Bot closed this Sep 28, 2026
@trunk-io
trunk-io Bot deleted the trunk-merge/pr-89156/15cdeb55-b6d4-42fc-adbb-2b2605afdcda branch September 28, 2026 12:03
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.