Coverage for object_streams/models.py: 92%
44 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-02 17:07 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-02 17:07 +0000
1"""Database models for replayable object stream events."""
3from __future__ import annotations
5from django.contrib.contenttypes.models import ContentType
6from django.db import models
8from object_streams.events import ObjectRef
9from object_streams.events import SourceRef
10from object_streams.events import StreamEvent
13__all__ = ("ObjectStreamEvent", "ObjectStreamOutboxState")
16def _content_type_model_label(content_type: ContentType) -> str:
17 model_class = content_type.model_class()
18 if model_class is not None: 18 ↛ 20line 18 didn't jump to line 20 because the condition on line 18 was always true
19 return model_class._meta.label
20 return f"{content_type.app_label}.{content_type.model}"
23class ObjectStreamEvent(models.Model):
24 """Outbox row that gives object streams a replayable global cursor."""
26 cursor = models.BigIntegerField(null=True, blank=True, unique=True)
27 subject_content_type = models.ForeignKey(ContentType, on_delete=models.CASCADE, related_name="+")
28 subject_object_id = models.TextField()
29 source_content_type = models.ForeignKey(
30 ContentType,
31 null=True,
32 blank=True,
33 on_delete=models.SET_NULL,
34 related_name="+",
35 )
36 source_object_id = models.TextField(blank=True, default="")
37 source_history_content_type = models.ForeignKey(
38 ContentType,
39 null=True,
40 blank=True,
41 on_delete=models.SET_NULL,
42 related_name="+",
43 )
44 source_history_id = models.TextField(blank=True, default="")
45 facet = models.CharField(max_length=64, db_index=True)
46 op = models.CharField(max_length=32)
47 changed_fields = models.JSONField(default=list, blank=True)
48 before = models.JSONField(null=True, blank=True)
49 after = models.JSONField(null=True, blank=True)
50 metadata = models.JSONField(default=dict, blank=True)
51 created_at = models.DateTimeField(auto_now_add=True, db_index=True)
53 class Meta:
54 ordering = ["cursor", "id"]
55 indexes = [
56 models.Index(
57 fields=["subject_content_type", "subject_object_id", "cursor"],
58 name="object_stre_subject_8d65d0_idx",
59 ),
60 models.Index(fields=["facet", "cursor"], name="object_stre_facet_407b73_idx"),
61 ]
63 def __str__(self):
64 return f"{self.subject_content_type}:{self.subject_object_id}:{self.id}"
66 def to_stream_event(self) -> StreamEvent:
67 source = None
68 if self.source_content_type_id or self.source_history_content_type_id:
69 source = SourceRef(
70 model=_content_type_model_label(self.source_content_type) if self.source_content_type_id else None,
71 pk=self.source_object_id or None,
72 history_model=(
73 _content_type_model_label(self.source_history_content_type)
74 if self.source_history_content_type_id
75 else None
76 ),
77 history_id=self.source_history_id or None,
78 )
80 return StreamEvent(
81 cursor=self.cursor,
82 subject=ObjectRef(
83 model=_content_type_model_label(self.subject_content_type),
84 pk=self.subject_object_id,
85 ),
86 facet=self.facet,
87 op=self.op,
88 changed_fields=tuple(self.changed_fields or ()),
89 source=source,
90 before=self.before,
91 after=self.after,
92 metadata=self.metadata or {},
93 )
96class ObjectStreamOutboxState(models.Model):
97 """Singleton row recording how far the outbox has been pruned.
99 Replay needs to tell a range that was pruned from a range that never
100 existed, because primary key sequences skip ids for rolled back
101 transactions. The watermark makes that exact instead of inferred.
102 """
104 id = models.PositiveSmallIntegerField(primary_key=True, default=1, editable=False)
105 pruned_through = models.BigIntegerField(default=0)
106 next_cursor = models.BigIntegerField(default=1)
107 broadcasted_through = models.BigIntegerField(default=0)
109 def __str__(self):
110 return f"broadcast through {self.broadcasted_through}, pruned through {self.pruned_through}"