返回 DeepSeek-Reasonix
session_dag_writer.go
根目录 / internal / agent / session_dag_writer.go
1 package agent
2
3 import (
4 "bytes"
5 "encoding/json"
6 "fmt"
7 "io"
8 "log/slog"
9 "os"
10 "path/filepath"
11 "time"
12
13 "reasonix/internal/fileutil"
14 "reasonix/internal/store"
15 )
16
17 // sessionDAGTailRepairMinAge keeps tail repair from truncating a line another
18 // writer is still appending: a torn tail is only cut once the log has been
19 // quiet for this long.
20 const sessionDAGTailRepairMinAge = 2 * time.Second
21
22 type sessionDAGHeader struct {
23 generation int64
24 upgradedFrom int
25 }
26
27 // encodeSessionDAGEntries serializes entries one per line, filling in the
28 // schema version and the default writer/timestamp.
29 func encodeSessionDAGEntries(entries []sessionDAGEntry, now time.Time) ([]byte, error) {
30 var buf bytes.Buffer
31 for i := range entries {
32 e := &entries[i]
33 e.SchemaVersion = sessionDAGSchemaVersion
34 if e.At.IsZero() {
35 e.At = now
36 }
37 if e.Writer == "" {
38 e.Writer = SessionWriterID()
39 }
40 b, err := json.Marshal(e)
41 if err != nil {
42 return nil, fmt.Errorf("encode session entry: %w", err)
43 }
44 buf.Write(b)
45 buf.WriteByte('\n')
46 }
47 return buf.Bytes(), nil
48 }
49
50 // appendSessionDAGEntries writes entries as one contiguous batch and returns
51 // the log size afterwards. Callers hold the session file lock: JSON lines
52 // exceed PIPE_BUF, so the flock, not O_APPEND, is what keeps two writers'
53 // batches from interleaving.
54 func appendSessionDAGEntries(sessionPath string, entries []sessionDAGEntry, sync bool) (int64, error) {
55 path := store.SessionEventLog(sessionPath)
56 if path == "" {
57 return 0, fmt.Errorf("empty session event log path")
58 }
59 if len(entries) == 0 {
60 info, err := os.Stat(path)
61 if err != nil {
62 return 0, err
63 }
64 return info.Size(), nil
65 }
66 fileutil.Crash("dag-append", path)
67 if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
68 return 0, err
69 }
70 data, err := encodeSessionDAGEntries(entries, time.Now().UTC())
71 if err != nil {
72 return 0, err
73 }
74 f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600)
75 if err != nil {
76 return 0, fmt.Errorf("open session event log: %w", err)
77 }
78 if err := f.Chmod(0o600); err != nil {
79 _ = f.Close()
80 return 0, fmt.Errorf("protect session event log: %w", err)
81 }
82 if _, err := f.Write(data); err != nil {
83 _ = f.Close()
84 return 0, fmt.Errorf("append session entries: %w", err)
85 }
86 if sync {
87 if err := f.Sync(); err != nil {
88 _ = f.Close()
89 return 0, err
90 }
91 }
92 info, err := f.Stat()
93 if err != nil {
94 _ = f.Close()
95 return 0, err
96 }
97 return info.Size(), f.Close()
98 }
99
100 // readSessionDAGHeader decodes the leading log entry without reading the rest
101 // of the file, so a writer can notice a rotation (new generation) cheaply.
102 func readSessionDAGHeader(sessionPath string) (sessionDAGHeader, bool, error) {
103 path := store.SessionEventLog(sessionPath)
104 if path == "" {
105 return sessionDAGHeader{}, false, nil
106 }
107 f, err := os.Open(path)
108 if err != nil {
109 if os.IsNotExist(err) {
110 return sessionDAGHeader{}, false, nil
111 }
112 return sessionDAGHeader{}, false, err
113 }
114 defer f.Close()
115 var e sessionDAGEntry
116 if err := json.NewDecoder(io.LimitReader(f, sessionEventProbeMaxBytes)).Decode(&e); err != nil {
117 return sessionDAGHeader{}, false, nil
118 }
119 if e.SchemaVersion != sessionDAGSchemaVersion || e.Type != sessionDAGTypeLog {
120 return sessionDAGHeader{}, false, nil
121 }
122 return sessionDAGHeader{generation: e.Generation, upgradedFrom: e.UpgradedFrom}, true, nil
123 }
124
125 // repairSessionDAGTail truncates a torn tail found by replay once the log has
126 // been quiet long enough that no writer can still be finishing that line. The
127 // discarded bytes go to the .damaged sidecar first. Callers hold the file lock.
128 func repairSessionDAGTail(sessionPath string, st *sessionDAGState, now time.Time) (bool, error) {
129 if st == nil || !st.damaged || st.lastGoodEnd >= st.size {
130 return false, nil
131 }
132 path := store.SessionEventLog(sessionPath)
133 info, err := os.Stat(path)
134 if err != nil {
135 return false, err
136 }
137 if info.Size() != st.size || now.Sub(info.ModTime()) < sessionDAGTailRepairMinAge {
138 return false, nil
139 }
140 if preserveErr := preserveDamagedEventLogTail(sessionPath, path, st.lastGoodEnd, st.size); preserveErr != nil {
141 slog.Warn("session: could not preserve damaged log tail; truncating anyway",
142 "path", path, "from", st.lastGoodEnd, "size", st.size, "err", preserveErr)
143 }
144 if err := os.Truncate(path, st.lastGoodEnd); err != nil {
145 return false, err
146 }
147 if st.lastGoodEnd > 0 {
148 f, err := os.OpenFile(path, os.O_WRONLY|os.O_APPEND, 0o600)
149 if err != nil {
150 return false, err
151 }
152 if _, err := f.Write([]byte{'\n'}); err != nil {
153 _ = f.Close()
154 return false, err
155 }
156 if err := f.Close(); err != nil {
157 return false, err
158 }
159 }
160 st.size = st.lastGoodEnd
161 if st.lastGoodEnd > 0 {
162 st.size++
163 }
164 st.damaged = false
165 return true, nil
166 }
167
167 lines GO