Queue a per-project segment membership count refresh. Pass `delay_until` to schedule the refresh for a future time, e.g. to let an Edge CDC tombstone land in ClickHouse before recounting. No-op when the org has the feature off, or when a refresh for the project is already pending.
(
project: Project,
*,
delay_until: datetime | None = None,
)
| 39 | |
| 40 | |
| 41 | def enqueue_membership_refresh( |
| 42 | project: Project, |
| 43 | *, |
| 44 | delay_until: datetime | None = None, |
| 45 | ) -> None: |
| 46 | """Queue a per-project segment membership count refresh. |
| 47 | |
| 48 | Pass `delay_until` to schedule the refresh for a future time, e.g. to let an |
| 49 | Edge CDC tombstone land in ClickHouse before recounting. |
| 50 | |
| 51 | No-op when the org has the feature off, or when a refresh for the project is |
| 52 | already pending. A pending refresh scheduled sooner than `delay_until` is |
| 53 | pushed out to it. |
| 54 | """ |
| 55 | if not is_membership_enabled(project.organisation): |
| 56 | return |
| 57 | |
| 58 | from segment_membership.tasks import refresh_project_segment_counts |
| 59 | |
| 60 | pending = Task.objects.filter( |
| 61 | task_identifier=refresh_project_segment_counts.task_identifier, |
| 62 | completed=False, |
| 63 | num_failures__lt=3, |
| 64 | serialized_args=Task.serialize_data((project.id,)), |
| 65 | ) |
| 66 | if pending.exists(): |
| 67 | if delay_until is not None: |
| 68 | pending.filter( |
| 69 | Q(scheduled_for__isnull=True) | Q(scheduled_for__lt=delay_until), |
| 70 | ).update(scheduled_for=delay_until) |
| 71 | return |
| 72 | |
| 73 | refresh_project_segment_counts.delay( |
| 74 | args=(project.id,), |
| 75 | delay_until=delay_until, |
| 76 | ) |
| 77 | |
| 78 | |
| 79 | @contextmanager |
searching dependent graphs…