返回 DeepSeek-Reasonix
store.go
1 package workspacestate
2
3 import (
4 "bytes"
5 "context"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "os"
10 "path/filepath"
11 "slices"
12 "strings"
13 "sync"
14 "sync/atomic"
15 "time"
16
17 "reasonix/internal/fileutil"
18 filelock "reasonix/internal/identitylock"
19 )
20
21 const (
22 SchemaVersion = 3
23 GlobalWorkspaceID = "global"
24 )
25
26 var (
27 ErrUnsupportedVersion = errors.New("workspace state version is unsupported")
28 ErrWorkspaceNotFound = errors.New("workspace is not registered")
29 ErrSessionNotFound = errors.New("session is not registered")
30 ErrMutationConflict = errors.New("workspace mutation conflicts with persisted state")
31 )
32
33 type Workspace struct {
34 Organization *Organization `json:"organization,omitempty"`
35 ID string `json:"id"`
36 Root string `json:"root"`
37 FormerRoots []string `json:"formerRoots,omitempty"` // global only: roots it was rebound away from
38 Title string `json:"title"`
39 SessionIDs []string `json:"sessionIds"`
40 Visible bool `json:"visible"`
41 CreatedAt time.Time `json:"createdAt"`
42 UpdatedAt time.Time `json:"updatedAt"`
43 extra map[string]json.RawMessage
44 }
45
46 type PendingCreate struct {
47 ParentSessionID string `json:"parentSessionId,omitempty"`
48 Presentation *Presentation `json:"presentation,omitempty"`
49 OperationID string `json:"operationId"`
50 WorkspaceID string `json:"workspaceId"`
51 SessionID string `json:"sessionId"`
52 CreatedAt time.Time `json:"createdAt"`
53 ArchiveSource string `json:"archiveSource,omitempty"`
54 extra map[string]json.RawMessage
55 }
56
57 type State struct {
58 Version int `json:"version"`
59 Generation uint64 `json:"generation"`
60 Initialized bool `json:"initialized"`
61 WorkspaceIDs []string `json:"workspaceIds"`
62 Workspaces map[string]Workspace `json:"workspaces"`
63 ArchivedSessionIDs []string `json:"archivedSessionIds"`
64 PendingCreates map[string]PendingCreate `json:"pendingCreates"`
65 SessionStates map[string]SessionState `json:"sessionStates"`
66 SourceMappings map[string]SourceMapping `json:"sourceMappings"`
67 PendingOperations map[string]Operation `json:"pendingOperations"`
68 RecoveryEntries map[string]RecoveryEntry `json:"recoveryEntries"`
69 Presentation map[string]Presentation `json:"presentation"`
70 TopicRemovals map[string]TopicRemoval `json:"topicRemovals,omitempty"`
71 // Immutable derived index for display copies; never serialized.
72 adoptedTopics map[string]map[string]bool
73 sourceIdentities *sourceIdentityIndex
74 extra map[string]json.RawMessage
75 }
76
77 type Store struct {
78 path string
79 mu sync.Mutex
80 beforeUpgrade func(context.Context) error
81 readBody []byte
82 readSnapshot atomic.Pointer[ReadSnapshot]
83 verificationMu sync.Mutex
84 verification *snapshotVerification
85 }
86
87 func NewStore(path string, beforeUpgrade ...func(context.Context) error) *Store {
88 store := &Store{path: filepath.Clean(path)}
89 if len(beforeUpgrade) > 0 {
90 store.beforeUpgrade = beforeUpgrade[0]
91 }
92 return store
93 }
94
95 func (s *Store) Path() string {
96 if s == nil {
97 return ""
98 }
99 return s.path
100 }
101
102 func (s *Store) Load(ctx context.Context) (State, error) {
103 return s.loadSnapshot(ctx, false)
104 }
105
106 // LoadProjection returns display/ownership metadata only. Recovery and command
107 // callers must use Load: operation journals are intentionally absent here.
108 func (s *Store) LoadProjection(ctx context.Context) (State, error) {
109 return s.loadSnapshot(ctx, true)
110 }
111
112 // LoadProjectionWithVersions returns a mutable display copy and the immutable
113 // invalidation index from the exact same verification. Callers must not infer
114 // that relationship from generation numbers, which legacy writers can retain.
115 func (s *Store) LoadProjectionWithVersions(ctx context.Context) (State, *ReadVersions, error) {
116 snapshot, err := s.VerifySnapshot(ctx)
117 if err != nil {
118 return State{}, nil, err
119 }
120 return snapshot.cloneState(true), snapshot.versions, nil
121 }
122
123 func (s *Store) loadSnapshot(ctx context.Context, projection bool) (State, error) {
124 snapshot, err := s.VerifySnapshot(ctx)
125 if err != nil {
126 return State{}, err
127 }
128 return snapshot.cloneState(projection), nil
129 }
130
131 func (r *ReadSnapshot) cloneState(projection bool) State {
132 state := r.state
133 if projection {
134 state = State{Version: state.Version, Generation: state.Generation, Initialized: state.Initialized,
135 WorkspaceIDs: state.WorkspaceIDs, Workspaces: state.Workspaces,
136 SessionStates: state.SessionStates, SourceMappings: state.SourceMappings, Presentation: state.Presentation,
137 adoptedTopics: r.adoptedTopics}
138 }
139 state.sourceIdentities = r.state.sourceIdentities
140 return cloneSnapshot(state)
141 }
142
143 func (s *Store) verifySnapshotLocked(ctx context.Context) error {
144 if err := ctx.Err(); err != nil {
145 return err
146 }
147 // Compare actual bytes, not timestamps or generation: another supported
148 // writer may replace a file while preserving either of those values.
149 body, err := os.ReadFile(s.path)
150 if errors.Is(err, os.ErrNotExist) {
151 if s.readSnapshot.Load() == nil || s.readBody != nil {
152 s.publishSnapshotLocked(nil, newState())
153 }
154 return nil
155 }
156 if err != nil {
157 return err
158 }
159 if s.readBody == nil || !bytes.Equal(body, s.readBody) {
160 state, err := decodeState(body)
161 if err != nil {
162 return err
163 }
164 s.publishSnapshotLocked(body, state)
165 }
166 return ctx.Err()
167 }
168
169 func (s *Store) RenameWorkspace(ctx context.Context, workspaceID, title string) error {
170 return s.mutate(ctx, func(state *State) error {
171 workspace, ok := state.Workspaces[strings.TrimSpace(workspaceID)]
172 if !ok {
173 return ErrWorkspaceNotFound
174 }
175 workspace.Title = strings.TrimSpace(title)
176 workspace.UpdatedAt = time.Now().UTC()
177 state.Workspaces[workspace.ID] = workspace
178 return nil
179 })
180 }
181
182 func (s *Store) SetWorkspaceVisible(ctx context.Context, workspaceID string, visible bool) error {
183 return s.mutate(ctx, func(state *State) error {
184 workspace, ok := state.Workspaces[strings.TrimSpace(workspaceID)]
185 if !ok {
186 return ErrWorkspaceNotFound
187 }
188 workspace.Visible = visible
189 workspace.UpdatedAt = time.Now().UTC()
190 state.Workspaces[workspace.ID] = workspace
191 return nil
192 })
193 }
194
195 func (s *Store) MoveWorkspace(ctx context.Context, workspaceID, beforeWorkspaceID string) error {
196 workspaceID = strings.TrimSpace(workspaceID)
197 beforeWorkspaceID = strings.TrimSpace(beforeWorkspaceID)
198 return s.mutate(ctx, func(state *State) error {
199 if _, ok := state.Workspaces[workspaceID]; !ok {
200 return ErrWorkspaceNotFound
201 }
202 if beforeWorkspaceID != "" {
203 if _, ok := state.Workspaces[beforeWorkspaceID]; !ok {
204 return ErrWorkspaceNotFound
205 }
206 }
207 state.WorkspaceIDs = insertBefore(remove(state.WorkspaceIDs, workspaceID), workspaceID, beforeWorkspaceID)
208 return nil
209 })
210 }
211
212 // CommitRotation atomically publishes a prepared replacement into its
213 // workspace and, for Clear, archives the source without removing its stable
214 // position from the registry.
215 func (s *Store) CommitRotation(ctx context.Context, operationID, workspaceID, sessionID, beforeSessionID, archiveSessionID string) error {
216 operationID, workspaceID, sessionID = strings.TrimSpace(operationID), strings.TrimSpace(workspaceID), strings.TrimSpace(sessionID)
217 archiveSessionID = strings.TrimSpace(archiveSessionID)
218 if operationID == "" || workspaceID == "" || sessionID == "" {
219 return errors.New("rotation commit requires operation, workspace, and session ids")
220 }
221 return s.mutate(ctx, func(state *State) error {
222 workspace, ok := state.Workspaces[workspaceID]
223 if !ok {
224 return ErrWorkspaceNotFound
225 }
226 pending, ok := state.PendingCreates[sessionID]
227 if !ok || pending.OperationID != operationID || pending.WorkspaceID != workspaceID {
228 if owner, attached := sessionOwner(*state, sessionID); !attached || owner != workspaceID {
229 return ErrMutationConflict
230 }
231 } else if owner, attached := sessionOwner(*state, sessionID); attached && owner != workspaceID {
232 return ErrMutationConflict
233 } else if !attached {
234 workspace.SessionIDs = insertBefore(workspace.SessionIDs, sessionID, beforeSessionID)
235 attachOrganizationSession(&workspace, sessionID, "")
236 mirrorOrganizationOrder(&workspace)
237 workspace.UpdatedAt = time.Now().UTC()
238 state.Workspaces[workspaceID] = workspace
239 }
240 delete(state.PendingCreates, sessionID)
241 if archiveSessionID != "" {
242 if _, exists := sessionOwner(*state, archiveSessionID); !exists {
243 return ErrSessionNotFound
244 }
245 setLifecycle(state, archiveSessionID, Archived)
246 }
247 return nil
248 })
249 }
250
251 func (s *Store) AbortCreate(ctx context.Context, sessionID string) error {
252 return s.mutate(ctx, func(state *State) error {
253 delete(state.PendingCreates, strings.TrimSpace(sessionID))
254 return nil
255 })
256 }
257
258 // AbortCreateIfOperation removes only the caller's reservation. A late cleanup
259 // from an older draft operation must never erase a newer operation's claim.
260 func (s *Store) AbortCreateIfOperation(ctx context.Context, sessionID, operationID string) error {
261 sessionID, operationID = strings.TrimSpace(sessionID), strings.TrimSpace(operationID)
262 return s.mutate(ctx, func(state *State) error {
263 pending, ok := state.PendingCreates[sessionID]
264 if ok && pending.OperationID == operationID {
265 delete(state.PendingCreates, sessionID)
266 }
267 return nil
268 })
269 }
270
271 func (s *Store) MoveSession(ctx context.Context, workspaceID, sessionID, beforeSessionID string) error {
272 return s.moveSession(ctx, workspaceID, sessionID, beforeSessionID, nil)
273 }
274
275 // MoveSessionIfUnchanged reorders one active session only while the caller's
276 // resolved lifecycle generation and workspace owner remain current.
277 func (s *Store) MoveSessionIfUnchanged(ctx context.Context, workspaceID, sessionID, beforeSessionID string, generation uint64) error {
278 return s.moveSession(ctx, workspaceID, sessionID, beforeSessionID, &generation)
279 }
280
281 func (s *Store) moveSession(ctx context.Context, workspaceID, sessionID, beforeSessionID string, generation *uint64) error {
282 return s.mutate(ctx, func(state *State) error {
283 workspace, ok := state.Workspaces[strings.TrimSpace(workspaceID)]
284 if !ok {
285 return ErrWorkspaceNotFound
286 }
287 if !contains(workspace.SessionIDs, sessionID) {
288 return ErrSessionNotFound
289 }
290 status := state.SessionStates[sessionID]
291 if status.Lifecycle != Active {
292 return ErrSessionNotFound
293 }
294 if generation != nil && status.Generation != *generation {
295 return ErrMutationConflict
296 }
297 workspace.SessionIDs = insertBefore(remove(workspace.SessionIDs, sessionID), sessionID, beforeSessionID)
298 if o := workspace.Organization; o != nil {
299 key, before := SessionKey(sessionID), ""
300 if beforeSessionID != "" {
301 before = SessionKey(beforeSessionID)
302 }
303 o.Order = insertBefore(remove(o.Order, key), key, before)
304 o.ManualOrderEnabled = true
305 o.Revision++
306 mirrorOrganizationOrder(&workspace)
307 }
308 workspace.UpdatedAt = time.Now().UTC()
309 state.Workspaces[workspace.ID] = workspace
310 status.Generation++
311 state.SessionStates[sessionID] = status
312 return nil
313 })
314 }
315
316 func (s *Store) ArchiveSession(ctx context.Context, sessionID string) error {
317 return s.SetLifecycle(ctx, []string{sessionID}, Archived)
318 }
319
320 func (s *Store) RestoreSession(ctx context.Context, sessionID string) error {
321 return s.SetLifecycle(ctx, []string{sessionID}, Active)
322 }
323
324 func (s *Store) Contains(ctx context.Context, sessionID string) (bool, error) {
325 snapshot, err := s.VerifySnapshot(ctx)
326 if err != nil {
327 return false, err
328 }
329 _, ok := snapshot.owners[strings.TrimSpace(sessionID)]
330 return ok, nil
331 }
332
333 // WithSessionUnchanged serializes a durable metadata commit with lifecycle and
334 // workspace changes, including writers in other processes. The callback must
335 // not call the registry; it may only commit session content metadata.
336 func (s *Store) WithSessionUnchanged(ctx context.Context, id, workspaceID string, generation uint64, commit func() error) error {
337 s.mu.Lock()
338 defer s.mu.Unlock()
339 release, err := filelock.Acquire(ctx, s.path+".lock")
340 if err != nil {
341 return err
342 }
343 defer release()
344 state, err := load(s.path)
345 if err != nil {
346 return err
347 }
348 owner, ok := sessionOwner(state, id)
349 status := state.SessionStates[id]
350 if !ok || status.Lifecycle != Active {
351 return ErrSessionNotFound
352 }
353 if owner != workspaceID || status.Generation != generation {
354 return ErrMutationConflict
355 }
356 return commit()
357 }
358
359 // WithStateLocked holds the registry's process and file locks while commit
360 // validates a read-only state snapshot and performs a related external write.
361 // The callback must not call this Store.
362 func (s *Store) WithStateLocked(ctx context.Context, commit func(State) error) error {
363 if s == nil || strings.TrimSpace(s.path) == "" || s.path == "." {
364 return errors.New("workspace state path is required")
365 }
366 s.mu.Lock()
367 defer s.mu.Unlock()
368 release, err := filelock.Acquire(ctx, s.path+".lock")
369 if err != nil {
370 return err
371 }
372 defer release()
373 state, err := load(s.path)
374 if err != nil {
375 return err
376 }
377 if commit == nil {
378 return nil
379 }
380 return commit(state)
381 }
382
383 func (s *Store) mutate(ctx context.Context, change func(*State) error) error {
384 if s == nil || strings.TrimSpace(s.path) == "" || s.path == "." {
385 return errors.New("workspace state path is required")
386 }
387 s.mu.Lock()
388 defer s.mu.Unlock()
389 if err := os.MkdirAll(filepath.Dir(s.path), 0o700); err != nil {
390 return err
391 }
392 release, err := filelock.Acquire(ctx, s.path+".lock")
393 if err != nil {
394 return err
395 }
396 defer release()
397 upgrading := false
398 if body, err := os.ReadFile(s.path); err == nil {
399 var header struct {
400 Version int `json:"version"`
401 }
402 upgrading = json.Unmarshal(body, &header) == nil && header.Version < SchemaVersion
403 }
404 if s.beforeUpgrade != nil {
405 body, readErr := os.ReadFile(s.path)
406 if readErr != nil && !os.IsNotExist(readErr) {
407 return readErr
408 }
409 var header struct {
410 Version int `json:"version"`
411 }
412 if len(body) > 0 && json.Unmarshal(body, &header) == nil && header.Version == 1 {
413 if err := s.beforeUpgrade(ctx); err != nil {
414 return err
415 }
416 }
417 }
418 if err := backupV1(s.path); err != nil {
419 return err
420 }
421 if err := backupV2(s.path); err != nil {
422 return err
423 }
424 state, err := load(s.path)
425 if err != nil {
426 return err
427 }
428 before, err := json.Marshal(state)
429 if err != nil {
430 return err
431 }
432 if err := change(&state); err != nil {
433 return err
434 }
435 after, err := json.Marshal(state)
436 if err != nil {
437 return err
438 }
439 if bytes.Equal(before, after) && !upgrading {
440 return nil
441 }
442 state.Generation++
443 state.Initialized = true
444 normalize(&state)
445 if err := validate(state); err != nil {
446 return err
447 }
448 body, err := json.Marshal(state)
449 if err != nil {
450 return err
451 }
452 body = append(body, '\n')
453 // Decode a private copy before committing: mutation inputs may retain slices
454 // or maps, and must never be able to change a published immutable snapshot.
455 published, err := decodeState(body)
456 if err != nil {
457 return err
458 }
459 if err := fileutil.AtomicWriteFileStrict(s.path, body, 0o600); err != nil {
460 return err
461 }
462 s.publishSnapshotLocked(body, published)
463 return nil
464 }
465
466 func load(path string) (State, error) {
467 body, err := os.ReadFile(path)
468 if errors.Is(err, os.ErrNotExist) {
469 return newState(), nil
470 }
471 if err != nil {
472 return State{}, err
473 }
474 return decodeState(body)
475 }
476
477 func decodeState(body []byte) (State, error) {
478 var state State
479 if err := json.Unmarshal(body, &state); err != nil {
480 return State{}, fmt.Errorf("decode workspace state: %w", err)
481 }
482 if state.Version != 1 && state.Version != 2 && state.Version != SchemaVersion {
483 return State{}, fmt.Errorf("%w: %d", ErrUnsupportedVersion, state.Version)
484 }
485 if state.Version == 1 {
486 state.SessionStates = map[string]SessionState{}
487 for _, id := range state.ArchivedSessionIDs {
488 state.SessionStates[id] = SessionState{Lifecycle: Archived, Generation: state.Generation}
489 }
490 state.Version = SchemaVersion
491 } else {
492 if state.WorkspaceIDs == nil || state.Workspaces == nil || state.SessionStates == nil {
493 return State{}, fmt.Errorf("%w: missing required registry fields", ErrUnsupportedVersion)
494 }
495 for _, workspace := range state.Workspaces {
496 for _, id := range workspace.SessionIDs {
497 if _, ok := state.SessionStates[id]; !ok {
498 return State{}, fmt.Errorf("%w: missing session lifecycle", ErrUnsupportedVersion)
499 }
500 }
501 }
502 }
503 state.Version = SchemaVersion
504 normalize(&state)
505 if err := validate(state); err != nil {
506 return State{}, err
507 }
508 return state, nil
509 }
510
511 func newState() State {
512 state := State{Version: SchemaVersion, WorkspaceIDs: []string{}, Workspaces: map[string]Workspace{}, ArchivedSessionIDs: []string{}, PendingCreates: map[string]PendingCreate{}}
513 normalize(&state)
514 return state
515 }
516
517 func normalize(state *State) {
518 if state.TopicRemovals == nil {
519 state.TopicRemovals = map[string]TopicRemoval{}
520 }
521 if state.SessionStates == nil {
522 state.SessionStates = map[string]SessionState{}
523 }
524 if state.SourceMappings == nil {
525 state.SourceMappings = map[string]SourceMapping{}
526 }
527 if state.PendingOperations == nil {
528 state.PendingOperations = map[string]Operation{}
529 }
530 if state.RecoveryEntries == nil {
531 state.RecoveryEntries = map[string]RecoveryEntry{}
532 }
533 if state.Presentation == nil {
534 state.Presentation = map[string]Presentation{}
535 }
536 if state.WorkspaceIDs == nil {
537 state.WorkspaceIDs = []string{}
538 }
539 if state.Workspaces == nil {
540 state.Workspaces = map[string]Workspace{}
541 }
542 if state.ArchivedSessionIDs == nil {
543 state.ArchivedSessionIDs = []string{}
544 }
545 if state.PendingCreates == nil {
546 state.PendingCreates = map[string]PendingCreate{}
547 }
548 for id, workspace := range state.Workspaces {
549 if workspace.SessionIDs == nil {
550 workspace.SessionIDs = []string{}
551 }
552 state.Workspaces[id] = workspace
553 for _, sessionID := range workspace.SessionIDs {
554 if _, ok := state.SessionStates[sessionID]; !ok {
555 state.SessionStates[sessionID] = SessionState{Lifecycle: Active, Generation: state.Generation}
556 }
557 }
558 }
559 state.ArchivedSessionIDs = []string{}
560 for id, status := range state.SessionStates {
561 if status.Lifecycle == Archived {
562 state.ArchivedSessionIDs = append(state.ArchivedSessionIDs, id)
563 }
564 }
565 slices.Sort(state.ArchivedSessionIDs)
566 }
567
568 func validate(state State) error {
569 if err := validateLifecycleState(state); err != nil {
570 return err
571 }
572 seen := map[string]struct{}{}
573 for _, id := range state.WorkspaceIDs {
574 if _, duplicate := seen[id]; duplicate {
575 return fmt.Errorf("workspace state has duplicate workspace %q", id)
576 }
577 seen[id] = struct{}{}
578 workspace, ok := state.Workspaces[id]
579 if !ok || workspace.ID != id {
580 return fmt.Errorf("workspace state has invalid workspace %q", id)
581 }
582 }
583 owners := map[string]string{}
584 for id, workspace := range state.Workspaces {
585 for _, sessionID := range workspace.SessionIDs {
586 if owner, duplicate := owners[sessionID]; duplicate {
587 return fmt.Errorf("session %q belongs to both %q and %q", sessionID, owner, id)
588 }
589 owners[sessionID] = id
590 }
591 }
592 return nil
593 }
594
595 func sessionOwner(state State, sessionID string) (string, bool) {
596 for id, workspace := range state.Workspaces {
597 if contains(workspace.SessionIDs, sessionID) {
598 return id, true
599 }
600 }
601 return "", false
602 }
603
604 func insertBefore(ids []string, id, before string) []string {
605 ids = remove(ids, id)
606 if before != "" {
607 for i, current := range ids {
608 if current == before {
609 return append(append(append([]string{}, ids[:i]...), id), ids[i:]...)
610 }
611 }
612 }
613 return append(ids, id)
614 }
615
616 func remove(ids []string, target string) []string {
617 result := make([]string, 0, len(ids))
618 for _, id := range ids {
619 if id != target {
620 result = append(result, id)
621 }
622 }
623 return result
624 }
625
626 func contains(ids []string, target string) bool {
627 return slices.Contains(ids, target)
628 }
629
630 func (s *State) UnmarshalJSON(body []byte) error {
631 type plain State
632 var decoded plain
633 if err := json.Unmarshal(body, &decoded); err != nil {
634 return err
635 }
636 var fields map[string]json.RawMessage
637 if err := json.Unmarshal(body, &fields); err != nil {
638 return err
639 }
640 for _, key := range []string{"version", "generation", "initialized", "workspaceIds", "workspaces", "archivedSessionIds", "pendingCreates", "sessionStates", "sourceMappings", "pendingOperations", "recoveryEntries", "presentation", "topicRemovals"} {
641 delete(fields, key)
642 }
643 *s = State(decoded)
644 s.extra = fields
645 return nil
646 }
647
648 func (s State) MarshalJSON() ([]byte, error) {
649 type plain State
650 body, err := json.Marshal(plain(s))
651 if err != nil {
652 return nil, err
653 }
654 return mergeUnknown(body, s.extra)
655 }
656
657 func (w *Workspace) UnmarshalJSON(body []byte) error {
658 type plain Workspace
659 var decoded plain
660 if err := json.Unmarshal(body, &decoded); err != nil {
661 return err
662 }
663 var fields map[string]json.RawMessage
664 if err := json.Unmarshal(body, &fields); err != nil {
665 return err
666 }
667 for _, key := range []string{"id", "root", "title", "sessionIds", "visible", "createdAt", "updatedAt", "organization", "formerRoots"} {
668 delete(fields, key)
669 }
670 *w = Workspace(decoded)
671 w.extra = fields
672 return nil
673 }
674
675 func (w Workspace) MarshalJSON() ([]byte, error) {
676 type plain Workspace
677 body, err := json.Marshal(plain(w))
678 if err != nil {
679 return nil, err
680 }
681 return mergeUnknown(body, w.extra)
682 }
683
684 func mergeUnknown(known []byte, extra map[string]json.RawMessage) ([]byte, error) {
685 if len(extra) == 0 {
686 return known, nil
687 }
688 var fields map[string]json.RawMessage
689 if err := json.Unmarshal(known, &fields); err != nil {
690 return nil, err
691 }
692 for key, value := range extra {
693 if _, exists := fields[key]; !exists {
694 fields[key] = bytes.Clone(value)
695 }
696 }
697 return json.Marshal(fields)
698 }
699
699 lines GO