| 1 | package session |
| 2 | |
| 3 | import ( |
| 4 | "encoding/json" |
| 5 | "fmt" |
| 6 | "strings" |
| 7 | |
| 8 | "reasonix/internal/provider" |
| 9 | ) |
| 10 | |
| 11 | type messageRetractPayload struct { |
| 12 | MessageIDs []string `json:"messageIds"` |
| 13 | Reason string `json:"reason,omitempty"` |
| 14 | } |
| 15 | |
| 16 | func retractedMessageIDs(event Event, payload json.RawMessage) ([]string, error) { |
| 17 | var body messageRetractPayload |
| 18 | if err := strictPayload(payload, &body); err != nil { |
| 19 | return nil, damagedPayload(event, err) |
| 20 | } |
| 21 | if len(body.MessageIDs) == 0 { |
| 22 | return nil, damagedPayload(event, fmt.Errorf("empty messageIds")) |
| 23 | } |
| 24 | seen := make(map[string]bool, len(body.MessageIDs)) |
| 25 | for _, id := range body.MessageIDs { |
| 26 | if strings.TrimSpace(id) == "" || strings.TrimSpace(id) != id || seen[id] { |
| 27 | return nil, damagedPayload(event, fmt.Errorf("invalid or duplicate message id")) |
| 28 | } |
| 29 | seen[id] = true |
| 30 | } |
| 31 | return body.MessageIDs, nil |
| 32 | } |
| 33 | |
| 34 | // Execution, recovery, history and search must decode the same durable event |
| 35 | // schema. Projection-specific DTOs used to reject valid rewrite/import metadata. |
| 36 | type historyReplacePayload struct { |
| 37 | Messages []provider.Message `json:"messages"` |
| 38 | Reason string `json:"reason,omitempty"` |
| 39 | Sources []uint64 `json:"sourceSequences,omitempty"` |
| 40 | } |
| 41 | |
| 42 | type legacyImportPayload struct { |
| 43 | Source Source `json:"source"` |
| 44 | Messages []provider.Message `json:"messages"` |
| 45 | Goal json.RawMessage `json:"goal,omitempty"` |
| 46 | ModelRef string `json:"modelRef,omitempty"` |
| 47 | ModelIdentity string `json:"modelIdentity,omitempty"` |
| 48 | } |
| 49 | |
| 50 | func replacementEventMessages(event Event, payload json.RawMessage) ([]provider.Message, error) { |
| 51 | var messages []provider.Message |
| 52 | var err error |
| 53 | if event.Kind == "legacy/import" { |
| 54 | var body legacyImportPayload |
| 55 | err = strictPayload(payload, &body) |
| 56 | messages = body.Messages |
| 57 | } else { |
| 58 | var body historyReplacePayload |
| 59 | err = strictPayload(payload, &body) |
| 60 | messages = body.Messages |
| 61 | } |
| 62 | if err != nil || messages == nil { |
| 63 | return nil, damagedPayload(event, err) |
| 64 | } |
| 65 | return messages, nil |
| 66 | } |
| 67 |