返回 DeepSeek-Reasonix
session_event_decode_context.go
根目录 / internal / agent / session_event_decode_context.go
1 package agent
2
3 import (
4 "bytes"
5 "context"
6 "encoding/json"
7 "fmt"
8 "io"
9 "reasonix/internal/historywork"
10 )
11
12 func preflightSessionEventMessages(ctx context.Context, path string, raw []byte, existingMessages, existingCollectionItems int, limits sessionReplayLimits) (messageCount, collectionItems int, err error) {
13 dec := json.NewDecoder(&contextReader{ctx: ctx, reader: bytes.NewReader(raw)})
14 tok, err := dec.Token()
15 if err != nil {
16 return 0, existingCollectionItems, err
17 }
18 if delim, ok := tok.(json.Delim); !ok || delim != '[' {
19 return 0, existingCollectionItems, fmt.Errorf("messages must be an array")
20 }
21 collectionItems = existingCollectionItems
22 for dec.More() {
23 if err := ctx.Err(); err != nil {
24 return 0, existingCollectionItems, err
25 }
26 if existingMessages+messageCount >= limits.maxMessages {
27 return 0, existingCollectionItems, sessionReplayLimitError(
28 path, "messages", int64(existingMessages+messageCount+1), int64(limits.maxMessages),
29 )
30 }
31 messageCount++
32 if err := preflightSessionEventValue(ctx, path, dec, &collectionItems, limits.maxCollectionItems); err != nil {
33 return 0, existingCollectionItems, err
34 }
35 }
36 if _, err := dec.Token(); err != nil {
37 return 0, existingCollectionItems, err
38 }
39 return messageCount, collectionItems, nil
40 }
41
42 // preflightSessionEventValue walks JSON without materializing nested values.
43 func preflightSessionEventValue(ctx context.Context, path string, dec *json.Decoder, collectionItems *int, maxCollectionItems int) error {
44 if err := ctx.Err(); err != nil {
45 return err
46 }
47 tok, err := dec.Token()
48 if err != nil {
49 return err
50 }
51 delim, ok := tok.(json.Delim)
52 if !ok {
53 return nil
54 }
55 switch delim {
56 case '{':
57 for dec.More() {
58 key, err := dec.Token()
59 if err != nil {
60 return err
61 }
62 if _, ok := key.(string); !ok {
63 return fmt.Errorf("object key must be a string")
64 }
65 if err := preflightSessionEventValue(ctx, path, dec, collectionItems, maxCollectionItems); err != nil {
66 return err
67 }
68 }
69 end, err := dec.Token()
70 if err != nil {
71 return err
72 }
73 if end != json.Delim('}') {
74 return fmt.Errorf("object is not terminated")
75 }
76 return nil
77 case '[':
78 for dec.More() {
79 if err := ctx.Err(); err != nil {
80 return err
81 }
82 if *collectionItems >= maxCollectionItems {
83 return sessionReplayLimitError(path, "message_collection_items", int64(*collectionItems+1), int64(maxCollectionItems))
84 }
85 (*collectionItems)++
86 if err := preflightSessionEventValue(ctx, path, dec, collectionItems, maxCollectionItems); err != nil {
87 return err
88 }
89 }
90 end, err := dec.Token()
91 if err != nil {
92 return err
93 }
94 if end != json.Delim(']') {
95 return fmt.Errorf("array is not terminated")
96 }
97 return nil
98 default:
99 return fmt.Errorf("unexpected JSON delimiter %q", delim)
100 }
101 }
102
103 type contextReader struct {
104 ctx context.Context
105 reader io.Reader
106 }
107
108 func (r *contextReader) Read(p []byte) (int, error) {
109 if err := r.ctx.Err(); err != nil {
110 return 0, err
111 }
112 if len(p) > historywork.ReadChunk {
113 p = p[:historywork.ReadChunk]
114 }
115 return r.reader.Read(p)
116 }
117
117 lines GO