返回 DeepSeek-Reasonix
session_events.go
根目录 / internal / control / session_events.go
1 package control
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "os"
10 "path/filepath"
11 "reflect"
12 "runtime"
13 "strings"
14 "sync"
15
16 "reasonix/internal/agent"
17 "reasonix/internal/event"
18 "reasonix/internal/provider"
19 "reasonix/internal/session"
20 "reasonix/internal/turnevent"
21 )
22
23 func sessionDirectory(sessionPath string) string {
24 sessionPath = filepath.Clean(strings.TrimSpace(sessionPath))
25 if sessionPath == "." || sessionPath == "" {
26 return ""
27 }
28 // Windows spellings of one path must retain one shared writer identity.
29 if runtime.GOOS == "windows" {
30 sessionPath = strings.ToLower(sessionPath)
31 }
32 parent := filepath.Dir(sessionPath)
33 root := filepath.Join(parent, "sessions-v4")
34 if filepath.Base(parent) == "sessions" {
35 root = filepath.Join(filepath.Dir(parent), "sessions-v4")
36 }
37 id := agent.BranchID(sessionPath)
38 if id == "" {
39 return ""
40 }
41 return filepath.Join(root, id)
42 }
43
44 type sharedSessionEventStore struct {
45 store *session.Session
46 refs int
47 }
48
49 var processSessionEventStores = struct {
50 sync.Mutex
51 stores map[string]*sharedSessionEventStore
52 }{stores: map[string]*sharedSessionEventStore{}}
53
54 func acquireSessionEventStore(dir, id string) (*session.Session, func(context.Context) error, bool, error) {
55 processSessionEventStores.Lock()
56 defer processSessionEventStores.Unlock()
57 if entry := processSessionEventStores.stores[dir]; entry != nil {
58 entry.refs++
59 return entry.store, releaseSessionEventStore(dir, entry), true, nil
60 }
61 store, err := session.Open(dir, id)
62 if errors.Is(err, session.ErrSessionNotFound) {
63 store, err = session.CreateStore(dir, id)
64 }
65 if err != nil {
66 return nil, nil, false, err
67 }
68 entry := &sharedSessionEventStore{store: store, refs: 1}
69 processSessionEventStores.stores[dir] = entry
70 return store, releaseSessionEventStore(dir, entry), false, nil
71 }
72
73 func releaseSessionEventStore(dir string, entry *sharedSessionEventStore) func(context.Context) error {
74 var once sync.Once
75 var releaseErr error
76 return func(ctx context.Context) error {
77 once.Do(func() {
78 // Every ownership handoff is a semantic checkpoint even when another
79 // in-process controller already retains the physical writer.
80 _, releaseErr = entry.store.Flush(ctx)
81 processSessionEventStores.Lock()
82 current := processSessionEventStores.stores[dir]
83 if current != entry {
84 processSessionEventStores.Unlock()
85 return
86 }
87 entry.refs--
88 last := entry.refs == 0
89 if last {
90 delete(processSessionEventStores.stores, dir)
91 }
92 processSessionEventStores.Unlock()
93 if last {
94 releaseErr = errors.Join(releaseErr, entry.store.Close(ctx))
95 }
96 })
97 return releaseErr
98 }
99 }
100
101 func (c *Controller) openSessionEventStore(sessionPath string) (*session.Session, func(context.Context) error, error) {
102 if c.NativeLegacySession() {
103 // The original transcript remains the only message authority. Runtime
104 // envelopes use the legacy ledger; opening must not seed a v4 copy.
105 return nil, nil, nil
106 }
107 if service, runtime, exclusive := c.v3Binding(); runtime != nil {
108 return runtime.Session(), nil, nil
109 } else if exclusive && service != nil {
110 // A service-backed controller needs a published Runtime before admission.
111 // Creating a path-derived sidecar here would reintroduce a second producer
112 // outside the v3.1 ownership boundary.
113 return nil, nil, nil
114 }
115 dir := sessionDirectory(sessionPath)
116 if dir == "" {
117 return nil, nil, nil
118 }
119 id := agent.BranchID(sessionPath)
120 store, release, shared, err := acquireSessionEventStore(dir, id)
121 if err != nil {
122 return nil, nil, err
123 }
124 if !shared {
125 if _, _, err := store.RecoverInterrupted(context.Background()); err != nil {
126 _ = release(context.Background())
127 return nil, nil, err
128 }
129 }
130 snapshot := store.StateSnapshot()
131 if snapshot.EventSequence != 0 {
132 return store, release, nil
133 }
134 return store, release, nil
135 }
136
137 // releaseLegacyEventStoreForImport stops this controller's retired path-bound
138 // event producer before the importer takes a shared freeze lock. Other
139 // controllers keep their own reference; in that case the freeze correctly
140 // refuses to race an active producer instead of copying a moving prefix.
141 func (c *Controller) releaseLegacyEventStoreForImport(ctx context.Context) (func(), error) {
142 if c == nil {
143 return func() {}, nil
144 }
145 path := c.SessionPath()
146 c.turnEvents.mu.Lock()
147 if c.turnEvents.v3 == nil || strings.HasPrefix(c.turnEvents.v3Path, "session:") {
148 c.turnEvents.mu.Unlock()
149 return func() {}, nil
150 }
151 store := c.turnEvents.v3
152 release := c.turnEvents.v3Release
153 c.turnEvents.v3 = nil
154 c.turnEvents.v3Path = ""
155 c.turnEvents.v3Runtime = nil
156 c.turnEvents.v3Release = nil
157 c.turnEvents.mu.Unlock()
158 var err error
159 if release != nil {
160 err = release(ctx)
161 } else {
162 err = store.Close(ctx)
163 }
164 restore := func() {
165 if _, runtime, _ := c.v3Binding(); runtime == nil {
166 c.rebindTurnEvents(path)
167 }
168 }
169 return restore, err
170 }
171
172 // SuspendLegacyEventStoreForImport closes this controller's path-derived
173 // compatibility producer while an idle host prepares a replacement runtime.
174 // The caller must serialize turn admission for the controller. Calling the
175 // returned function restores the producer when candidate preparation fails;
176 // after a successful publication the old controller is retired instead.
177 func (c *Controller) SuspendLegacyEventStoreForImport(ctx context.Context) (func(), error) {
178 return c.releaseLegacyEventStoreForImport(ctx)
179 }
180
181 func (c *Controller) seedSessionEventsFromExecutor(reason string) error {
182 if !c.sessionEventCommitAllowed() {
183 return nil
184 }
185 if snapshot, ok := c.sessionEventSnapshot(); !ok || snapshot.EventSequence != 0 {
186 return nil
187 }
188 if c == nil || c.executor == nil {
189 return nil
190 }
191 // Existing legacy transcripts are loaded by Resume after Controller
192 // construction. Seeding the constructor's system-only placeholder here
193 // would make it look authoritative and discard the loaded history.
194 if reason == "session-open" {
195 if info, err := os.Stat(c.SessionPath()); err == nil && info.Size() > 0 {
196 return nil
197 }
198 }
199 messages := c.executor.Session().Snapshot()
200 for i := range messages {
201 if strings.TrimSpace(messages[i].ID) == "" {
202 messages[i].ID = agent.NewMessageID()
203 }
204 }
205 if err := c.RecordSessionMessages(context.Background(), reason, messages); err != nil {
206 return err
207 }
208 return nil
209 }
210
211 // importLegacyResumeOverPlaceholder repairs the only safe pre-migration
212 // overlap: an older build may have created a system-only v3 sidecar before it
213 // loaded an existing legacy transcript. No executed turn can be present in
214 // that projection, so replacing it with the frozen loaded history is lossless.
215 func (c *Controller) importLegacyResumeOverPlaceholder(incoming *agent.Session) error {
216 if c == nil || incoming == nil {
217 return nil
218 }
219 if !c.sessionEventCommitAllowed() {
220 return nil
221 }
222 snapshot, ok := c.sessionEventSnapshot()
223 if !ok || len(snapshot.Projection.ModelMessages) != 1 || len(incoming.Snapshot()) <= 1 {
224 return nil
225 }
226 messages := incoming.Snapshot()
227 for i := range messages {
228 if strings.TrimSpace(messages[i].ID) == "" {
229 messages[i].ID = agent.NewMessageID()
230 }
231 }
232 return c.replaceSessionEventProjection(context.Background(), "legacy-resume-placeholder", messages)
233 }
234
235 func (c *Controller) restoreExecutorFromSessionEvents() {
236 if c == nil || c.executor == nil {
237 return
238 }
239 snapshot, ok := c.sessionEventSnapshot()
240 if !ok || len(snapshot.Projection.ModelMessages) == 0 {
241 return
242 }
243 if reflect.DeepEqual(c.executor.Session().Snapshot(), snapshot.Projection.ModelMessages) {
244 return
245 }
246 // Called only at construction or explicit Resume. Path rewrites and recovery
247 // rebinds must keep their prepared candidate session instead of adopting the
248 // source log again.
249 c.executor.Session().Replace(append([]provider.Message(nil), snapshot.Projection.ModelMessages...))
250 }
251
252 // replaceSessionModelContext records an Agent/configuration rebuild without
253 // rewriting the UI transcript. The exact serialized model context becomes a
254 // typed event, while the original messages remain available for history.
255 func (c *Controller) replaceSessionModelContext(ctx context.Context, messages []provider.Message, reason string) error {
256 store := c.sessionEventStore()
257 if store == nil {
258 return errors.New("exclusive v3 controller has no session store")
259 }
260 snapshot := store.ExecutionSnapshot()
261 modelMessages := provider.ModelMessages(messages)
262 if modelMessages == nil {
263 modelMessages = []provider.Message{}
264 }
265 payload, err := json.Marshal(map[string]any{
266 "messages": modelMessages,
267 "reason": reason,
268 "sourceSequences": []uint64{snapshot.EventSequence},
269 })
270 if err != nil {
271 return err
272 }
273 events := []session.Event{{Kind: "model/context-replace", Payload: payload}}
274 if strings.TrimSpace(c.ModelRef()) != "" {
275 configPayload, err := json.Marshal(map[string]string{"modelRef": c.ModelRef(), "modelIdentity": c.ModelSelectionIdentity()})
276 if err != nil {
277 return err
278 }
279 events = append(events, session.Event{Kind: "session/config", Payload: configPayload})
280 }
281 digest := sha256.Sum256(payload)
282 c.turnEvents.commitMu.Lock()
283 batch := session.Batch{
284 OperationID: fmt.Sprintf("model-context:%x", digest[:16]),
285 TurnID: snapshot.Projection.TurnID,
286 Events: events,
287 }
288 _, runtime, exclusive := c.v3Binding()
289 if exclusive && runtime != nil && runtime.Session() == store && !runtime.OwnsExecution(c.ExecutionGeneration()) {
290 prepared, prepareErr := store.PrepareBatchContext(ctx, batch.OperationID, batch)
291 if prepareErr == nil {
292 if previous := c.turnEvents.pendingExecutionCommit; previous != nil {
293 previous.Release()
294 }
295 c.turnEvents.pendingExecutionCommit = &prepared
296 }
297 err = prepareErr
298 } else {
299 _, err = c.appendSessionBatch(ctx, store, batch)
300 }
301 c.turnEvents.commitMu.Unlock()
302 if err != nil {
303 return err
304 }
305 if c.executor != nil {
306 c.executor.Session().Replace(append([]provider.Message(nil), modelMessages...))
307 }
308 return nil
309 }
310
311 func (c *Controller) adoptResumeSystemPrompt(incoming *agent.Session) error {
312 if !c.sessionEventCommitAllowed() {
313 return nil
314 }
315 store := c.sessionEventStore()
316 if store == nil || incoming == nil {
317 return nil
318 }
319 snapshot := store.ExecutionSnapshot()
320 persisted, current := snapshot.Projection.ModelMessages, incoming.Snapshot()
321 if len(persisted) == 0 || len(current) == 0 || persisted[0].Role != provider.RoleSystem || current[0].Role != provider.RoleSystem || reflect.DeepEqual(persisted[0], current[0]) {
322 return nil
323 }
324 replaced := append([]provider.Message(nil), persisted...)
325 replaced[0] = current[0]
326 payload, err := json.Marshal(map[string]any{"messages": replaced, "reason": "system-prompt-refresh"})
327 if err != nil {
328 return err
329 }
330 c.turnEvents.commitMu.Lock()
331 defer c.turnEvents.commitMu.Unlock()
332 _, err = c.appendSessionBatch(context.Background(), store, session.Batch{
333 OperationID: fmt.Sprintf("system-prompt-refresh:%d", snapshot.EventSequence+1),
334 TurnID: snapshot.Projection.TurnID,
335 Events: []session.Event{{Kind: "history/replace", Payload: payload}},
336 })
337 return err
338 }
339
340 func (c *Controller) sessionEventStore() *session.Session {
341 if c == nil {
342 return nil
343 }
344 c.turnEvents.mu.RLock()
345 defer c.turnEvents.mu.RUnlock()
346 return c.turnEvents.v3
347 }
348
349 // sessionEventCommitAllowed fences unpublished Desktop replacement runtimes.
350 // Those candidates intentionally share the active process store so they can
351 // inspect the latest in-memory prefix, but they do not own mutation authority
352 // until the tab's final lease handoff succeeds.
353 func (c *Controller) sessionEventCommitAllowed() bool {
354 if _, runtime, _ := c.v3Binding(); runtime != nil {
355 return runtime.OwnsExecution(c.ExecutionGeneration())
356 }
357 if c == nil || !c.managedSessionEvents.Load() {
358 return true
359 }
360 if c.executor == nil || c.executor.Session() == nil {
361 return false
362 }
363 path := c.SessionPath()
364 if path == "" {
365 return true
366 }
367 auth := c.executor.Session().WriteAuthority()
368 return auth != nil && auth.Covers(path)
369 }
370
371 func (c *Controller) sessionEventSnapshot() (session.Snapshot, bool) {
372 store := c.sessionEventStore()
373 if store == nil {
374 return session.Snapshot{}, false
375 }
376 return store.ExecutionSnapshot(), true
377 }
378
379 func (c *Controller) sessionStateSnapshot() (session.Snapshot, bool) {
380 store := c.sessionEventStore()
381 if store == nil {
382 return session.Snapshot{}, false
383 }
384 return store.StateSnapshot(), true
385 }
386
387 // appendSessionEventLocked mirrors the existing event.Sink lifecycle into the
388 // single typed business log. turnEvents.commitMu must be held, which preserves
389 // the exact order assigned by the compatibility envelope adapter.
390 func (c *Controller) appendSessionEventLocked(ctx context.Context, e event.Event) error {
391 if !c.sessionEventCommitAllowed() {
392 return nil
393 }
394 store := c.sessionEventStore()
395 if store == nil {
396 if c.sessionEngineEnabled() {
397 return session.ErrSessionNotRunning
398 }
399 return nil
400 }
401 snapshot := store.ExecutionSnapshot()
402 projection := snapshot.Projection
403 if projection.Recovery != nil && projection.Recovery.State == "recovery_required" && e.Kind != event.TurnDone {
404 // Recovery has sealed business-state mutation. The watchdog observes
405 // the uncooperative worker; late semantic output must never reactivate
406 // tools, interactions, Goal, or Todo.
407 return nil
408 }
409 if e.TurnID == "" {
410 if _, turnID, active := c.currentTurnToken(); active {
411 e.TurnID = turnID
412 }
413 if e.TurnID == "" {
414 if ledger := c.turnEventLedger(); ledger != nil {
415 e.TurnID = ledger.ActiveTurnID()
416 }
417 }
418 if e.TurnID == "" {
419 e.TurnID = projection.TurnID
420 }
421 }
422 events, err := c.sessionEventsFor(e, projection)
423 if err != nil || len(events) == 0 {
424 return err
425 }
426 if e.Kind == event.TurnDone {
427 if err := c.appendTerminationLocked(ctx, e, store, events); err != nil {
428 return fmt.Errorf("%w: %w", turnevent.ErrTurnLedgerUnavailable, err)
429 }
430 return nil
431 }
432 op := fmt.Sprintf("runtime:%s:%d:%d", e.TurnID, snapshot.EventSequence+1, e.Kind)
433 _, err = c.appendSessionBatch(ctx, store, session.Batch{OperationID: op, TurnID: e.TurnID, Events: events})
434 if err != nil {
435 return fmt.Errorf("%w: %w", turnevent.ErrTurnLedgerUnavailable, err)
436 }
437 c.noteCommittedMessagesLocked(events)
438 return nil
439 }
440
441 func (c *Controller) sessionEventsFor(e event.Event, projection session.Projection) ([]session.Event, error) {
442 if e.Kind == event.Notice {
443 return mcpDisplayNoticeEvents(e)
444 }
445 if e.Kind == event.SessionOperation {
446 return sessionMaintenanceEvents(e)
447 }
448 return c.v3EventsFor(e, projection)
449 }
450
451 func (c *Controller) v3EventsFor(e event.Event, projection session.Projection) ([]session.Event, error) {
452 makePayload := func(value any) (json.RawMessage, error) { return json.Marshal(value) }
453 var out []session.Event
454 switch e.Kind {
455 case event.TurnStarted:
456 payload, err := makePayload(map[string]any{"status": event.TurnInProgress})
457 if err != nil {
458 return nil, err
459 }
460 out = append(out, session.Event{Kind: "turn/start", Payload: payload})
461 out = c.appendSubmissionEvent(out, e.TurnID)
462 if e.DomainKind != "" {
463 if e.DomainKind != "goal/state" || len(e.DomainPayload) == 0 {
464 return nil, fmt.Errorf("unsupported turn admission domain event %q", e.DomainKind)
465 }
466 out = append(out, session.Event{Kind: e.DomainKind, Payload: append(json.RawMessage(nil), e.DomainPayload...)})
467 }
468 case event.ToolDispatch:
469 payload, err := makePayload(v3ToolPayload(e.Tool, false))
470 if err != nil {
471 return nil, err
472 }
473 out = append(out, session.Event{Kind: "tool/call", Payload: payload})
474 case event.ToolStarted:
475 payload, err := makePayload(map[string]any{"id": e.Tool.ID, "name": e.Tool.Name})
476 if err != nil {
477 return nil, err
478 }
479 out = append(out, session.Event{Kind: "tool/start", Payload: payload})
480 case event.ToolResult:
481 // Tool results are emitted before Agent mutates its derived conversation
482 // cache. Commit the exact message and structured result together so a
483 // recovered prefix cannot contain only one half.
484 if e.CommittedMessage != nil {
485 messagePayload, err := makePayload(map[string]any{"message": *e.CommittedMessage})
486 if err != nil {
487 return nil, err
488 }
489 out = append(out, session.Event{Kind: "message/complete", Payload: messagePayload})
490 }
491 payload, err := makePayload(v3ToolPayload(e.Tool, true))
492 if err != nil {
493 return nil, err
494 }
495 if e.Tool.TodoWritten {
496 todoPayload, err := makePayload(map[string]any{"todos": e.Tool.Todos})
497 if err != nil {
498 return nil, err
499 }
500 out = append(out, session.Event{Kind: "todo/write", Payload: todoPayload})
501 }
502 out = append(out, session.Event{Kind: "tool/result", Payload: payload})
503 case event.StreamAttempt:
504 payload, err := makePayload(map[string]any{
505 "id": e.StreamAttempt.ID, "messageId": e.MessageID,
506 "action": e.StreamAttempt.Action, "attempt": e.StreamAttempt.Attempt,
507 "max": e.StreamAttempt.Max, "reason": e.StreamAttempt.Reason,
508 })
509 if err != nil {
510 return nil, err
511 }
512 out = append(out, session.Event{Kind: "assistant/attempt", Payload: payload})
513 case event.AskRequest, event.ApprovalRequest, event.MCPInteractionRequest:
514 kind := strings.TrimSpace(e.PromptKind)
515 if kind == "" {
516 switch e.Kind {
517 case event.AskRequest:
518 kind = "ask"
519 case event.MCPInteractionRequest:
520 kind = "mcp"
521 default:
522 kind = "approval"
523 }
524 }
525 payload, err := makePayload(map[string]any{
526 "id": e.ItemID, "toolCallId": e.ItemID, "kind": kind, "state": "pending",
527 "sessionId": e.SessionID, "headId": agent.BranchID(c.SessionPath()),
528 "turnId": e.TurnID, "runtimeEpoch": e.RuntimeEpoch,
529 })
530 if err != nil {
531 return nil, err
532 }
533 out = append(out, session.Event{Kind: "interaction/created", Payload: payload})
534 case event.PromptAnswered:
535 state := strings.TrimSpace(e.InteractionState)
536 if state == "" {
537 state = "answered"
538 }
539 payload, err := makePayload(map[string]any{"id": e.ItemID, "state": state})
540 if err != nil {
541 return nil, err
542 }
543 out = append(out, session.Event{Kind: "interaction/resolved", Payload: payload})
544 if e.DomainKind != "" {
545 out = append(out, session.Event{Kind: e.DomainKind, Payload: append(json.RawMessage(nil), e.DomainPayload...)})
546 }
547 case event.TurnStatusChanged:
548 if e.Status == event.TurnRecoveryRequired && e.Recovery != nil {
549 payload, err := makePayload(e.Recovery)
550 if err != nil {
551 return nil, err
552 }
553 out = append(out, session.Event{Kind: "runtime/recovery", Payload: payload})
554 }
555 case event.CompactionDone:
556 // Context-maintenance commits persist their exact model projection before
557 // this notification is emitted. CompactionDone is presentation-only.
558 case event.TurnDone:
559 return c.terminalSessionEvents(e, projection)
560 }
561 return out, nil
562 }
563
564 // v3ToolPayload retains the structured execution and presentation metadata
565 // emitted by the tool runtime. UI truncation may choose a smaller view later;
566 // the authoritative event must not collapse a result to name plus text.
567 func v3ToolPayload(tool event.Tool, result bool) map[string]any {
568 payload := map[string]any{
569 "id": tool.ID, "name": tool.Name, "args": tool.Args,
570 "runState": tool.RunState, "diagnostic": tool.Diagnostic,
571 "resolvedName": tool.ResolvedName, "capabilityId": tool.CapabilityID,
572 "readOnly": tool.ReadOnly, "truncated": tool.Truncated,
573 "durationMs": tool.DurationMs, "startedAt": tool.StartedAt, "endedAt": tool.EndedAt,
574 "partial": tool.Partial, "argChars": tool.ArgChars, "refreshed": tool.Refreshed,
575 "parentId": tool.ParentID, "attemptId": tool.AttemptID,
576 "subagentRef": tool.SubagentRef, "subagentStatus": tool.SubagentStatus,
577 "subagentErrorCode": tool.SubagentErrorCode, "subagentRetryable": tool.SubagentRetryable,
578 "diff": tool.Diff, "added": tool.Added, "removed": tool.Removed,
579 "profile": tool.Profile, "execution": tool.Execution,
580 "presentedFiles": tool.PresentedFiles,
581 "workspaceMutation": tool.WorkspaceMutation, "workspacePaths": tool.WorkspacePaths,
582 "workspaceAllPaths": tool.WorkspaceAllPaths,
583 }
584 if result {
585 payload["output"] = tool.Output
586 payload["error"] = tool.Err
587 payload["todos"] = tool.Todos
588 payload["todoWritten"] = tool.TodoWritten
589 }
590 return payload
591 }
592
593 // RecordSessionMessages implements agent.SessionEventRecorder. Runtime message
594 // commits arrive here directly; this code never diffs the mutable legacy
595 // transcript to infer missing business events.
596 func (c *Controller) RecordSessionMessages(ctx context.Context, reason string, messages []provider.Message) error {
597 if !c.sessionEventCommitAllowed() {
598 return nil
599 }
600 store := c.sessionEventStore()
601 if store == nil || len(messages) == 0 {
602 return nil
603 }
604 if recovery := store.ExecutionSnapshot().Projection.Recovery; recovery != nil && recovery.State == "recovery_required" {
605 return nil
606 }
607 c.turnEvents.commitMu.Lock()
608 defer c.turnEvents.commitMu.Unlock()
609 if !c.messageCommitAllowedLocked(ctx, store) {
610 return nil
611 }
612 events := make([]session.Event, 0, len(messages))
613 for _, message := range messages {
614 if strings.TrimSpace(message.ID) == "" {
615 return errors.New("record session message: missing stable message id")
616 }
617 payload, err := json.Marshal(map[string]any{"message": message})
618 if err != nil {
619 return err
620 }
621 events = append(events, session.Event{Kind: "message/complete", Payload: payload})
622 }
623 snapshot := store.ExecutionSnapshot()
624 op := fmt.Sprintf("messages:%s:%d", reason, snapshot.EventSequence+1)
625 _, err := c.appendSessionBatch(ctx, store, session.Batch{OperationID: op, TurnID: snapshot.Projection.TurnID, Events: events})
626 if err == nil {
627 c.noteCommittedMessagesLocked(events)
628 c.settleOpenStreamLocked(messages)
629 }
630 return err
631 }
632
633 // RecordSessionModelContext implements agent.SessionModelContextRecorder. It
634 // persists the exact provider-visible projection without mutating the Agent or
635 // the canonical history. A successful return means the accepted commit is on
636 // stable storage through its final sequence.
637 func (c *Controller) RecordSessionModelContext(ctx context.Context, request agent.SessionModelContextCommit) (agent.SessionModelContextCommitResult, error) {
638 var result agent.SessionModelContextCommitResult
639 if !c.sessionEventCommitAllowed() {
640 return result, session.ErrStaleExecution
641 }
642 store := c.sessionEventStore()
643 if store == nil {
644 return result, nil
645 }
646 operationID := strings.TrimSpace(request.OperationID)
647 if operationID == "" {
648 return result, errors.New("record session model context: missing operation id")
649 }
650 modelMessages := provider.ModelMessages(request.Messages)
651 if modelMessages == nil {
652 modelMessages = []provider.Message{}
653 }
654 payload, err := json.Marshal(map[string]any{
655 "messages": modelMessages,
656 "reason": strings.TrimSpace(request.Reason),
657 })
658 if err != nil {
659 return result, err
660 }
661
662 c.turnEvents.commitMu.Lock()
663 if !c.messageCommitAllowedLocked(ctx, store) {
664 c.turnEvents.commitMu.Unlock()
665 return result, session.ErrStaleExecution
666 }
667 commit, err := c.appendSessionBatch(ctx, store, session.Batch{
668 OperationID: "model-context-maintenance:" + operationID,
669 Events: []session.Event{{Kind: "model/context-replace", Payload: payload}},
670 })
671 c.turnEvents.commitMu.Unlock()
672 if err != nil {
673 return result, err
674 }
675 result.Accepted = true
676 receipt, err := store.Flush(ctx)
677 if err != nil {
678 return result, err
679 }
680 if receipt.DurableSequence < commit.LastSequence() {
681 return result, fmt.Errorf("record session model context: durable sequence %d is before commit %d", receipt.DurableSequence, commit.LastSequence())
682 }
683 result.Durable = true
684 return result, nil
685 }
686
687 // RecordSessionMessageUpsert records one explicit stable-message mutation.
688 // This is used for local metadata such as protocol recovery receipts, whose
689 // provider-visible bytes may not change but whose durable UI/business fact must
690 // survive without falling back to a full legacy transcript rewrite.
691 func (c *Controller) RecordSessionMessageUpsert(ctx context.Context, reason string, message provider.Message) error {
692 if !c.sessionEventCommitAllowed() {
693 return nil
694 }
695 store := c.sessionEventStore()
696 if store == nil {
697 return nil
698 }
699 if strings.TrimSpace(message.ID) == "" {
700 return errors.New("record session message upsert: missing stable message id")
701 }
702 payload, err := json.Marshal(map[string]any{"message": message})
703 if err != nil {
704 return err
705 }
706 c.turnEvents.commitMu.Lock()
707 defer c.turnEvents.commitMu.Unlock()
708 if !c.messageCommitAllowedLocked(ctx, store) {
709 return nil
710 }
711 snapshot := store.ExecutionSnapshot()
712 op := fmt.Sprintf("message-upsert:%s:%s:%d", reason, message.ID, snapshot.EventSequence+1)
713 _, err = c.appendSessionBatch(ctx, store, session.Batch{OperationID: op, TurnID: snapshot.Projection.TurnID, Events: []session.Event{{Kind: "message/upsert", Payload: payload}}})
714 return err
715 }
716
717 // replaceSessionEventProjection records an intentional transcript rewrite as
718 // an explicit context event. Recovery, cancellation and host prompt changes use
719 // this path so the typed log remains authoritative without diffing the legacy
720 // message cache at a later checkpoint.
721 func (c *Controller) replaceSessionEventProjection(ctx context.Context, reason string, messages []provider.Message) error {
722 if !c.sessionEventCommitAllowed() {
723 return nil
724 }
725 store := c.sessionEventStore()
726 if store == nil {
727 return nil
728 }
729 payload, err := json.Marshal(map[string]any{
730 "messages": append([]provider.Message(nil), messages...),
731 "reason": reason,
732 })
733 if err != nil {
734 return err
735 }
736 c.turnEvents.commitMu.Lock()
737 defer c.turnEvents.commitMu.Unlock()
738 snapshot := store.ExecutionSnapshot()
739 op := fmt.Sprintf("context-replace:%s:%d", reason, snapshot.EventSequence+1)
740 _, err = c.appendSessionBatch(ctx, store, session.Batch{
741 OperationID: op,
742 TurnID: snapshot.Projection.TurnID,
743 Events: []session.Event{{Kind: "history/replace", Payload: payload}},
744 })
745 return err
746 }
747
748 func (c *Controller) flushSessionEvents(ctx context.Context) (session.DurableReceipt, error) {
749 store := c.sessionEventStore()
750 if store == nil {
751 return session.DurableReceipt{}, nil
752 }
753 return store.Flush(ctx)
754 }
755
756 func (c *Controller) appendDomainState(kind string, payload json.RawMessage, reason string) error {
757 if !c.sessionEventCommitAllowed() {
758 return nil
759 }
760 store := c.sessionEventStore()
761 if store == nil || len(payload) == 0 {
762 return nil
763 }
764 if recovery := store.ExecutionSnapshot().Projection.Recovery; recovery != nil && recovery.State == "recovery_required" {
765 return nil
766 }
767 c.turnEvents.commitMu.Lock()
768 defer c.turnEvents.commitMu.Unlock()
769 if !c.messageCommitAllowedLocked(context.Background(), store) {
770 return nil
771 }
772 snapshot := store.ExecutionSnapshot()
773 op := fmt.Sprintf("domain:%s:%s:%d", kind, reason, snapshot.EventSequence+1)
774 _, err := c.appendSessionBatch(context.Background(), store, session.Batch{OperationID: op, TurnID: snapshot.Projection.TurnID, Events: []session.Event{{Kind: kind, Payload: append(json.RawMessage(nil), payload...)}}})
775 if err != nil {
776 return fmt.Errorf("%w: %w", turnevent.ErrTurnLedgerUnavailable, err)
777 }
778 return nil
779 }
780
780 lines GO