Skip to content

Commit 5836e78

Browse files
committed
feat(segments): Refresh membership counts on segment edit
Wire the per-project segment membership count refresh into the segment create, update, clone and LaunchDarkly import paths via a new `enqueue_membership_refresh` service, so edits refresh counts within minutes instead of waiting for the daily backfill. The service debounces per project and no-ops when the feature flag is off. beep boop
1 parent c4670cb commit 5836e78

8 files changed

Lines changed: 266 additions & 2 deletions

File tree

api/integrations/launch_darkly/services.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@
3939
from integrations.launch_darkly.types import Clause
4040
from projects.models import Project
4141
from projects.tags.models import Tag
42+
from segment_membership.services import enqueue_membership_refresh
4243
from segments.models import Condition, Segment, SegmentRule
4344
from users.models import FFAdminUser
4445
from util.db import closing_stale_connections
@@ -1179,3 +1180,6 @@ def process_import_request(
11791180
import_request.status["deprecated_flag_count"] = sum(
11801181
1 for ld_flag in ld_flags if ld_flag["deprecated"]
11811182
)
1183+
1184+
# Refresh membership counts for the segments the import just created.
1185+
enqueue_membership_refresh(import_request.project)

api/segment_membership/services.py

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
from flag_engine.context.types import EvaluationContext
88
from flagsmith_sql_flag_engine import TranslateContext, translate_segment
99
from flagsmith_sql_flag_engine.dialects import ClickHouseDialect
10+
from task_processor.models import Task
1011

1112
from integrations.flagsmith.client import get_openfeature_client
1213
from organisations.models import Organisation
@@ -28,6 +29,29 @@ def is_membership_enabled(organisation: Organisation) -> bool:
2829
)
2930

3031

32+
def enqueue_membership_refresh(project: Project) -> None:
33+
"""Queue a per-project segment membership count refresh after a canonical
34+
segment is created or edited.
35+
36+
No-op when the org has the feature off, or when a refresh for the project
37+
is already pending or running.
38+
"""
39+
if not is_membership_enabled(project.organisation):
40+
return
41+
42+
from segment_membership.tasks import refresh_project_segment_counts
43+
44+
if Task.objects.filter(
45+
task_identifier=refresh_project_segment_counts.task_identifier,
46+
completed=False,
47+
num_failures__lt=3,
48+
serialized_args=Task.serialize_data((project.id,)),
49+
).exists():
50+
return
51+
52+
refresh_project_segment_counts.delay(args=(project.id,))
53+
54+
3155
@contextmanager
3256
def open_clickhouse_cursor(
3357
*, log_comment: str | None = None

api/segments/serializers.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
from metadata.serializers import MetadataSerializer, MetadataSerializerMixin
1212
from projects.models import Project
1313
from segment_membership.models import SegmentMembershipCount
14+
from segment_membership.services import enqueue_membership_refresh
1415
from segments.models import Condition, Segment, SegmentRule
1516

1617
logger = structlog.get_logger(__name__)
@@ -145,6 +146,7 @@ def create(self, validated_data: dict[str, Any]): # type: ignore[no-untyped-def
145146
metadata_data = validated_data.pop("metadata", [])
146147
segment = super().create(validated_data) # type: ignore[no-untyped-call]
147148
self._update_metadata(segment, metadata_data)
149+
enqueue_membership_refresh(segment.project)
148150
return segment
149151

150152
def update(self, segment: Segment, validated_data: dict[str, Any]): # type: ignore[no-untyped-def]
@@ -159,6 +161,7 @@ def update(self, segment: Segment, validated_data: dict[str, Any]): # type: ign
159161
)
160162
segment = super().update(segment, validated_data) # type: ignore[no-untyped-call]
161163
self._update_metadata(segment, metadata)
164+
enqueue_membership_refresh(segment.project)
162165
return segment
163166

164167
def _get_rules_and_conditions_without_deleted(

api/segments/views.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
)
2323
from features.versioning.models import EnvironmentFeatureVersion
2424
from projects.models import Project
25+
from segment_membership.services import enqueue_membership_refresh
2526

2627
from .models import Segment
2728
from .permissions import SegmentPermissions
@@ -180,6 +181,7 @@ def clone(self, request: Request, *args: Any, **kwargs: Any) -> Response:
180181
serializer = CloneSegmentSerializer(data=request.data)
181182
serializer.is_valid(raise_exception=True)
182183
clone = source_segment.clone(name=serializer.validated_data["name"])
184+
enqueue_membership_refresh(clone.project)
183185
return Response(SegmentSerializer(clone).data, status=status.HTTP_201_CREATED)
184186

185187

Lines changed: 96 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,96 @@
1+
import json
2+
3+
from django.urls import reverse
4+
from pytest_mock import MockerFixture
5+
from rest_framework import status
6+
from rest_framework.test import APIClient
7+
8+
from projects.models import Project
9+
10+
11+
def test_update_segment__edit__enqueues_membership_refresh(
12+
admin_client: APIClient,
13+
project: int,
14+
segment: int,
15+
mocker: MockerFixture,
16+
) -> None:
17+
# Given
18+
enqueue_membership_refresh_mock = mocker.patch(
19+
"segments.serializers.enqueue_membership_refresh"
20+
)
21+
url = reverse(
22+
"api-v1:projects:project-segments-detail",
23+
args=[project, segment],
24+
)
25+
data = {
26+
"name": "renamed",
27+
"project": project,
28+
"rules": [{"type": "ALL", "rules": [], "conditions": []}],
29+
}
30+
31+
# When
32+
response = admin_client.put(
33+
url, data=json.dumps(data), content_type="application/json"
34+
)
35+
36+
# Then
37+
assert response.status_code == status.HTTP_200_OK
38+
enqueue_membership_refresh_mock.assert_called_once_with(
39+
Project.objects.get(pk=project)
40+
)
41+
42+
43+
def test_create_segment__new_segment__enqueues_membership_refresh(
44+
admin_client: APIClient,
45+
project: int,
46+
mocker: MockerFixture,
47+
) -> None:
48+
# Given
49+
enqueue_membership_refresh_mock = mocker.patch(
50+
"segments.serializers.enqueue_membership_refresh"
51+
)
52+
url = reverse("api-v1:projects:project-segments-list", args=[project])
53+
data = {
54+
"name": "new-segment",
55+
"project": project,
56+
"rules": [{"type": "ALL", "rules": [], "conditions": []}],
57+
}
58+
59+
# When
60+
response = admin_client.post(
61+
url, data=json.dumps(data), content_type="application/json"
62+
)
63+
64+
# Then
65+
assert response.status_code == status.HTTP_201_CREATED
66+
enqueue_membership_refresh_mock.assert_called_once_with(
67+
Project.objects.get(pk=project)
68+
)
69+
70+
71+
def test_clone_segment__clone__enqueues_membership_refresh(
72+
admin_client: APIClient,
73+
project: int,
74+
segment: int,
75+
mocker: MockerFixture,
76+
) -> None:
77+
# Given
78+
enqueue_membership_refresh_mock = mocker.patch(
79+
"segments.views.enqueue_membership_refresh"
80+
)
81+
url = reverse(
82+
"api-v1:projects:project-segments-clone",
83+
args=[project, segment],
84+
)
85+
data = {"name": "cloned-segment"}
86+
87+
# When
88+
response = admin_client.post(
89+
url, data=json.dumps(data), content_type="application/json"
90+
)
91+
92+
# Then
93+
assert response.status_code == status.HTTP_201_CREATED
94+
enqueue_membership_refresh_mock.assert_called_once_with(
95+
Project.objects.get(pk=project)
96+
)

api/tests/unit/integrations/launch_darkly/test_services.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99
from django.conf import settings
1010
from django.core import signing
1111
from flag_engine.segments import constants as segment_constants
12+
from pytest_mock import MockerFixture
1213
from requests.exceptions import HTTPError, RequestException, Timeout
1314

1415
from environments.identities.models import Identity
@@ -646,3 +647,23 @@ def test_serialize_variation_value__various_types__returns_expected(
646647

647648
# Then
648649
assert result == expected
650+
651+
652+
@pytest.mark.django_db(transaction=True)
653+
def test_process_import_request__import__enqueues_membership_refresh(
654+
import_request: LaunchDarklyImportRequest,
655+
project: Project,
656+
mocker: MockerFixture,
657+
) -> None:
658+
# Given
659+
enqueue_membership_refresh_mock = mocker.patch(
660+
"integrations.launch_darkly.services.enqueue_membership_refresh"
661+
)
662+
663+
# When
664+
# the import (which bulk-creates segments) completes
665+
process_import_request(import_request)
666+
667+
# Then
668+
# it triggers a single membership refresh for the imported project
669+
enqueue_membership_refresh_mock.assert_called_once_with(project)

api/tests/unit/segment_membership/test_unit_segment_membership_services.py

Lines changed: 114 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,20 @@
11
from unittest.mock import MagicMock
22

3+
from common.test_tools import RunTasksFixture
4+
from pytest_django.fixtures import SettingsWrapper
35
from pytest_mock import MockerFixture
46

57
from environments.models import Environment
68
from organisations.models import Organisation
79
from projects.models import Project
810
from segment_membership.services import (
911
compute_segment_counts_for_project,
12+
enqueue_membership_refresh,
1013
get_projects_to_process,
1114
is_membership_enabled,
1215
open_clickhouse_cursor,
1316
)
17+
from segment_membership.tasks import refresh_project_segment_counts
1418
from segments.models import Segment, SegmentRule
1519
from tests.types import EnableFeaturesFixture
1620

@@ -220,3 +224,113 @@ def test_compute_segment_counts_for_project__untranslatable_segment__skips(
220224
# Then
221225
assert result == []
222226
cursor.execute.assert_not_called()
227+
228+
229+
def test_enqueue_membership_refresh__flag_on__enqueues_refresh(
230+
run_tasks: RunTasksFixture,
231+
mocker: MockerFixture,
232+
settings: SettingsWrapper,
233+
project: Project,
234+
segment: Segment,
235+
enable_features: EnableFeaturesFixture,
236+
) -> None:
237+
# Given
238+
# an org with the flag on
239+
enable_features("segment_membership_inspection")
240+
settings.CLICKHOUSE_ENABLED = True
241+
mocker.patch("segment_membership.tasks.open_clickhouse_cursor")
242+
compute_segment_counts_for_project_mock = mocker.patch(
243+
"segment_membership.tasks.compute_segment_counts_for_project",
244+
return_value=[],
245+
)
246+
247+
# When
248+
enqueue_membership_refresh(project)
249+
run_tasks(num_tasks=2)
250+
251+
# Then
252+
# exactly one refresh runs, for the project
253+
compute_segment_counts_for_project_mock.assert_called_once_with(project, mocker.ANY)
254+
255+
256+
def test_enqueue_membership_refresh__flag_off__does_not_enqueue(
257+
run_tasks: RunTasksFixture,
258+
mocker: MockerFixture,
259+
settings: SettingsWrapper,
260+
project: Project,
261+
) -> None:
262+
# Given
263+
# the org's flag is off
264+
settings.CLICKHOUSE_ENABLED = True
265+
mocker.patch("segment_membership.tasks.open_clickhouse_cursor")
266+
compute_segment_counts_for_project_mock = mocker.patch(
267+
"segment_membership.tasks.compute_segment_counts_for_project",
268+
return_value=[],
269+
)
270+
271+
# When
272+
enqueue_membership_refresh(project)
273+
run_tasks(num_tasks=1)
274+
275+
# Then
276+
# no refresh runs
277+
compute_segment_counts_for_project_mock.assert_not_called()
278+
279+
280+
def test_enqueue_membership_refresh__refresh_already_pending__debounces(
281+
run_tasks: RunTasksFixture,
282+
mocker: MockerFixture,
283+
settings: SettingsWrapper,
284+
project: Project,
285+
enable_features: EnableFeaturesFixture,
286+
) -> None:
287+
# Given
288+
# a refresh for the same project is already pending
289+
enable_features("segment_membership_inspection")
290+
settings.CLICKHOUSE_ENABLED = True
291+
mocker.patch("segment_membership.tasks.open_clickhouse_cursor")
292+
compute_segment_counts_for_project_mock = mocker.patch(
293+
"segment_membership.tasks.compute_segment_counts_for_project",
294+
return_value=[],
295+
)
296+
refresh_project_segment_counts.delay(args=(project.id,))
297+
298+
# When
299+
enqueue_membership_refresh(project)
300+
run_tasks(num_tasks=2)
301+
302+
# Then
303+
# the call did not enqueue a second refresh -- only the pending one ran
304+
compute_segment_counts_for_project_mock.assert_called_once_with(project, mocker.ANY)
305+
306+
307+
def test_enqueue_membership_refresh__pending_for_other_project__still_enqueues(
308+
run_tasks: RunTasksFixture,
309+
mocker: MockerFixture,
310+
settings: SettingsWrapper,
311+
project: Project,
312+
project_b: Project,
313+
enable_features: EnableFeaturesFixture,
314+
) -> None:
315+
# Given
316+
# a refresh is pending for a different project only
317+
enable_features("segment_membership_inspection")
318+
settings.CLICKHOUSE_ENABLED = True
319+
mocker.patch("segment_membership.tasks.open_clickhouse_cursor")
320+
compute_segment_counts_for_project_mock = mocker.patch(
321+
"segment_membership.tasks.compute_segment_counts_for_project",
322+
return_value=[],
323+
)
324+
refresh_project_segment_counts.delay(args=(project_b.id,))
325+
326+
# When
327+
enqueue_membership_refresh(project)
328+
run_tasks(num_tasks=3)
329+
330+
# Then
331+
# the debounce is scoped per project: both refreshes run
332+
assert compute_segment_counts_for_project_mock.call_count == 2
333+
compute_segment_counts_for_project_mock.assert_has_calls(
334+
[mocker.call(project, mocker.ANY), mocker.call(project_b, mocker.ANY)],
335+
any_order=True,
336+
)

docs/docs/deployment-self-hosting/observability/_events-catalogue.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -376,7 +376,7 @@ Attributes:
376376
### `segment_membership.compute.segment.skipped`
377377

378378
Logged at `error` from:
379-
- `api/segment_membership/services.py:96`
379+
- `api/segment_membership/services.py:120`
380380

381381
Attributes:
382382
- `project.id`
@@ -414,7 +414,7 @@ Attributes:
414414
### `segments.serializers.segment_revision_created`
415415

416416
Logged at `info` from:
417-
- `api/segments/serializers.py:155`
417+
- `api/segments/serializers.py:157`
418418

419419
Attributes:
420420
- `revision_id`

0 commit comments

Comments
 (0)