| 85 | c.mu.Unlock() |
| 86 | return false |
| 87 | } |
| 88 | runCtx, cancel := context.WithCancel(context.Background()) |
| 89 | c.currentRunning = true |
| 90 | c.currentCancel = cancel |
| 91 | c.currentRunID++ |
| 92 | runID := c.currentRunID |
| 93 | c.mu.Unlock() |
| 94 | go c.runExtraction(runCtx, payload, runID) |
| 95 | return true |
| 96 | } |
| 97 | |
| 98 | func (c *ExtractionController) incrementalFabricContext(state *AgentState) *extractionContext { |
| 99 | messages, messageIDs, startIndex, endIndex, sessionID, consumerID := c.incrementalFabricMessages(state) |
| 100 | return &extractionContext{ |
| 101 | Messages: messages, MessageIDs: messageIDs, StartIndex: startIndex, EndIndex: endIndex, |
| 102 | SessionID: sessionID, ConsumerID: consumerID, TurnCount: state.TurnCount, |
| 103 | UserTurnCount: state.UserTurnCount, State: state, |
| 104 | } |
| 105 | } |
| 106 | |
| 107 | func (c *ExtractionController) incrementalFabricMessages(state *AgentState) ([]map[string]any, []string, int, int, string, string) { |
| 108 | sessionID := firstNonEmptyString(c.SourceSessionID, state.MemorySessionID) |
| 109 | if sessionID == "" { |
| 110 | sessionID = "runtime-" + stableFabricTextID(c.Config.ProjectRoot()) + "-" + |
| 111 | firstNonEmptyString(c.SourceAgentID, "main") |
| 112 | } |
| 113 | consumerID := "memory-fabric-ingest:" + firstNonEmptyString(c.SourceAgentID, "main") |
| 114 | start := state.MemoryExtractionCursor |
| 115 | if start < 0 { |
| 116 | start = 0 |