( ctx context.Context, state *bootState, sessions SessionManager, now func() time.Time, )
| 65 | } |
| 66 | |
| 67 | func newDaemonMemoryExtractor( |
| 68 | ctx context.Context, |
| 69 | state *bootState, |
| 70 | sessions SessionManager, |
| 71 | now func() time.Time, |
| 72 | ) (*daemonMemoryExtractor, error) { |
| 73 | if state == nil || state.memoryStore == nil || !state.cfg.Memory.Enabled || !state.cfg.Memory.Extractor.Enabled { |
| 74 | return nil, nil |
| 75 | } |
| 76 | forkSessions, ok := sessions.(memoryExtractorSessionManager) |
| 77 | if !ok { |
| 78 | return nil, errors.New("daemon: session manager does not implement memory extractor spawn surface") |
| 79 | } |
| 80 | if now == nil { |
| 81 | now = func() time.Time { |
| 82 | return time.Now().UTC() |
| 83 | } |
| 84 | } |
| 85 | workspaceRoots := &sync.Map{} |
| 86 | forked := &forkedMemoryExtractor{ |
| 87 | sessions: forkSessions, |
| 88 | defaultAgent: firstNonEmptyString(state.cfg.Defaults.Agent, state.cfg.Memory.Dream.Agent), |
| 89 | model: state.cfg.Memory.Extractor.Model, |
| 90 | deadline: state.cfg.Memory.Extractor.Deadline, |
| 91 | logger: state.logger, |
| 92 | now: now, |
| 93 | workspaceRoots: workspaceRoots, |
| 94 | } |
| 95 | runtime, err := extractorpkg.NewRuntime( |
| 96 | context.WithoutCancel(ctx), |
| 97 | state.globalMemoryDir, |
| 98 | forked, |
| 99 | extractorpkg.WithEventSink(state.memoryStore), |
| 100 | extractorpkg.WithLogger(state.logger), |
| 101 | extractorpkg.WithClock(now), |
| 102 | extractorpkg.WithCoalesceMax(state.cfg.Memory.Extractor.Queue.CoalesceMax), |
| 103 | extractorpkg.WithQueueCapacity(state.cfg.Memory.Extractor.Queue.Capacity), |
| 104 | extractorpkg.WithThrottleTurns(state.cfg.Memory.Extractor.ThrottleTurns), |
| 105 | extractorpkg.WithInboxPath(state.cfg.Memory.Extractor.InboxPath), |
| 106 | ) |
| 107 | if err != nil { |
| 108 | return nil, fmt.Errorf("daemon: create memory extractor runtime: %w", err) |
| 109 | } |
| 110 | sink := &daemonMemoryProposalSink{ |
| 111 | base: state.memoryStore, |
| 112 | workspaceResolver: state.workspaceResolver, |
| 113 | } |
| 114 | consumer, err := extractorpkg.NewInboxConsumer( |
| 115 | state.globalMemoryDir, |
| 116 | sink, |
| 117 | extractorpkg.WithConsumerEventSink(state.memoryStore), |
| 118 | extractorpkg.WithConsumerLogger(state.logger), |
| 119 | extractorpkg.WithConsumerClock(now), |
| 120 | extractorpkg.WithConsumerInboxPath(state.cfg.Memory.Extractor.InboxPath), |
| 121 | extractorpkg.WithConsumerFailurePath(state.cfg.Memory.Extractor.DLQPath), |
| 122 | ) |
| 123 | if err != nil { |
| 124 | return nil, fmt.Errorf("daemon: create memory extractor inbox consumer: %w", err) |
no test coverage detected