返回 DeepSeek-Reasonix
message_event_payload.go
根目录 / internal / session / message_event_payload.go
1 package session
2
3 import (
4 "encoding/json"
5 "errors"
6 "fmt"
7 "strings"
8
9 "reasonix/internal/provider"
10 )
11
12 type messageRetractPayload struct {
13 MessageIDs []string `json:"messageIds"`
14 Reason string `json:"reason,omitempty"`
15 }
16
17 func retractedMessageIDs(event Event, payload json.RawMessage) ([]string, error) {
18 var body messageRetractPayload
19 if err := strictPayload(payload, &body); err != nil {
20 return nil, damagedPayload(event, err)
21 }
22 if len(body.MessageIDs) == 0 {
23 return nil, damagedPayload(event, fmt.Errorf("empty messageIds"))
24 }
25 seen := make(map[string]bool, len(body.MessageIDs))
26 for _, id := range body.MessageIDs {
27 if strings.TrimSpace(id) == "" || strings.TrimSpace(id) != id || seen[id] {
28 return nil, damagedPayload(event, fmt.Errorf("invalid or duplicate message id"))
29 }
30 seen[id] = true
31 }
32 return body.MessageIDs, nil
33 }
34
35 // Execution, recovery, history and search must decode the same durable event
36 // schema. Projection-specific DTOs used to reject valid rewrite/import metadata.
37 type historyReplacePayload struct {
38 Messages []provider.Message `json:"messages"`
39 Reason string `json:"reason,omitempty"`
40 Sources []uint64 `json:"sourceSequences,omitempty"`
41 }
42
43 type legacyImportPayload struct {
44 Source Source `json:"source"`
45 Messages []provider.Message `json:"messages"`
46 Goal json.RawMessage `json:"goal,omitempty"`
47 ModelRef string `json:"modelRef,omitempty"`
48 ModelIdentity string `json:"modelIdentity,omitempty"`
49 }
50
51 // ErrDuplicateMessageID refuses a write that would give two messages one
52 // stable id; every reader of the log keys a message by that id.
53 var ErrDuplicateMessageID = errors.New("session: duplicate stable message id")
54
55 // replacementEventMessages keeps the first occurrence of each id, so a log
56 // written before the writer refused duplicates projects and indexes alike.
57 func replacementEventMessages(event Event, payload json.RawMessage) ([]provider.Message, error) {
58 messages, err := decodeReplacementMessages(event, payload)
59 if err != nil {
60 return nil, err
61 }
62 return firstMessageOccurrences(messages), nil
63 }
64
65 func decodeReplacementMessages(event Event, payload json.RawMessage) ([]provider.Message, error) {
66 var messages []provider.Message
67 var err error
68 if event.Kind == "legacy/import" {
69 var body legacyImportPayload
70 err = strictPayload(payload, &body)
71 messages = body.Messages
72 } else {
73 var body historyReplacePayload
74 err = strictPayload(payload, &body)
75 messages = body.Messages
76 }
77 if err != nil || messages == nil {
78 return nil, damagedPayload(event, err)
79 }
80 return messages, nil
81 }
82
83 func firstMessageOccurrences(messages []provider.Message) []provider.Message {
84 if _, duplicated := duplicateMessageID(messages); !duplicated {
85 return messages
86 }
87 seen := make(map[string]bool, len(messages))
88 unique := make([]provider.Message, 0, len(messages))
89 for _, message := range messages {
90 if message.ID != "" && seen[message.ID] {
91 continue
92 }
93 seen[message.ID] = true
94 unique = append(unique, message)
95 }
96 return unique
97 }
98
99 func duplicateMessageID(messages []provider.Message) (string, bool) {
100 seen := make(map[string]bool, len(messages))
101 for _, message := range messages {
102 if message.ID == "" {
103 continue
104 }
105 if seen[message.ID] {
106 return message.ID, true
107 }
108 seen[message.ID] = true
109 }
110 return "", false
111 }
112
113 // checkReplacementIdentities refuses a new replacement that repeats an id. A
114 // payload that does not decode is the projection's to refuse.
115 func checkReplacementIdentities(event Event) error {
116 if event.Kind != "history/replace" && event.Kind != "legacy/import" {
117 return nil
118 }
119 messages, err := decodeReplacementMessages(event, event.Payload)
120 if err != nil {
121 return nil
122 }
123 if id, duplicated := duplicateMessageID(messages); duplicated {
124 return fmt.Errorf("%w: %s message %q", ErrDuplicateMessageID, event.Kind, id)
125 }
126 return nil
127 }
128
128 lines GO