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

1"""Database models for replayable object stream events.""" 

2 

3from __future__ import annotations 

4 

5from django.contrib.contenttypes.models import ContentType 

6from django.db import models 

7 

8from object_streams.events import ObjectRef 

9from object_streams.events import SourceRef 

10from object_streams.events import StreamEvent 

11 

12 

13__all__ = ("ObjectStreamEvent", "ObjectStreamOutboxState") 

14 

15 

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}" 

21 

22 

23class ObjectStreamEvent(models.Model): 

24 """Outbox row that gives object streams a replayable global cursor.""" 

25 

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) 

52 

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 ] 

62 

63 def __str__(self): 

64 return f"{self.subject_content_type}:{self.subject_object_id}:{self.id}" 

65 

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 ) 

79 

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 ) 

94 

95 

96class ObjectStreamOutboxState(models.Model): 

97 """Singleton row recording how far the outbox has been pruned. 

98 

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 """ 

103 

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) 

108 

109 def __str__(self): 

110 return f"broadcast through {self.broadcasted_through}, pruned through {self.pruned_through}"