返回 DeepSeek-Reasonix
message_identities.go
根目录 / internal / session / message_identities.go
1 package session
2
3 import (
4 "encoding/json"
5 "fmt"
6 "maps"
7 "slices"
8 )
9
10 // messageIdentities is every stable message id the log holds live. The writer
11 // needs it because a Service-owned projection drops durable message bodies, so
12 // the projection alone cannot see an id claimed before its last checkpoint.
13 type messageIdentities map[string]struct{}
14
15 func identitiesOf(ids []string) messageIdentities {
16 set := make(messageIdentities, len(ids))
17 for _, id := range ids {
18 set[id] = struct{}{}
19 }
20 return set
21 }
22
23 func (ids messageIdentities) holds(id string) bool {
24 _, ok := ids[id]
25 return ok
26 }
27
28 func (ids messageIdentities) list() []string {
29 return slices.Collect(maps.Keys(ids))
30 }
31
32 // identityChange is one commit's effect on the live set, kept apart so a
33 // refused commit leaves the set untouched.
34 type identityChange struct {
35 reset bool
36 live map[string]bool
37 // duplicate is the first message/complete id the commit repeats.
38 duplicate string
39 }
40
41 func (c *identityChange) has(ids messageIdentities, id string) bool {
42 if live, ok := c.live[id]; ok {
43 return live
44 }
45 if c.reset {
46 return false
47 }
48 _, ok := ids[id]
49 return ok
50 }
51
52 // changeFor mirrors the projection's message semantics. A payload it cannot
53 // decode is left for the projection to refuse.
54 func (ids messageIdentities) changeFor(commit Commit) identityChange {
55 change := identityChange{live: map[string]bool{}}
56 for _, event := range commit.Events {
57 switch event.Kind {
58 case "message/complete", "message/upsert":
59 id := eventMessageID(event)
60 if id == "" {
61 continue
62 }
63 if event.Kind == "message/complete" && change.has(ids, id) {
64 if change.duplicate == "" {
65 change.duplicate = id
66 }
67 continue
68 }
69 change.live[id] = true
70 case "message/retract":
71 retracted, err := retractedMessageIDs(event, event.Payload)
72 if err != nil {
73 continue
74 }
75 for _, id := range retracted {
76 change.live[id] = false
77 }
78 case "history/replace", "legacy/import":
79 messages, err := replacementEventMessages(event, event.Payload)
80 if err != nil {
81 continue
82 }
83 change.reset, change.live = true, map[string]bool{}
84 for _, message := range messages {
85 if message.ID != "" {
86 change.live[message.ID] = true
87 }
88 }
89 }
90 }
91 return change
92 }
93
94 func (ids *messageIdentities) apply(change identityChange) {
95 if change.reset || *ids == nil {
96 *ids = messageIdentities{}
97 }
98 for id, live := range change.live {
99 if live {
100 (*ids)[id] = struct{}{}
101 } else {
102 delete(*ids, id)
103 }
104 }
105 }
106
107 // admit records a commit a reader replays; a repeated message/complete keeps
108 // the id's first occurrence, as the projection does.
109 func (ids *messageIdentities) admit(commit Commit) {
110 ids.apply(ids.changeFor(commit))
111 }
112
113 // admitEvent reports whether a replayed event is kept, recording in repeated
114 // the sequence of a message/complete whose id is already live.
115 func (ids *messageIdentities) admitEvent(event Event, repeated map[uint64]bool) bool {
116 if event.Kind == "message/complete" && ids.holds(eventMessageID(event)) {
117 repeated[event.Sequence] = true
118 return false
119 }
120 ids.admit(Commit{Events: []Event{event}})
121 return true
122 }
123
124 func eventMessageID(event Event) string {
125 var body struct {
126 Message struct {
127 ID string `json:"id"`
128 } `json:"message"`
129 }
130 if json.Unmarshal(event.Payload, &body) != nil {
131 return ""
132 }
133 return body.Message.ID
134 }
135
136 func duplicateMessageError(id string) error {
137 return fmt.Errorf("%w: message/complete %q", ErrDuplicateMessageID, id)
138 }
139
139 lines GO