// // ChatController.swift // Zyquo Local // // Author: Simon-Pierre Boucher // Mail: contact@spboucher.ai // import Foundation import Observation /// Drives one streaming generation at a time for the selected conversation: /// send/stop/regenerate/edit-resend, live `` parsing, throttled /// tokens-per-second ticker (≤4 Hz), stats capture, and auto-titling. @MainActor @Observable final class ChatController { private(set) var isGenerating = false /// Visible (non-thinking) streamed text of the in-flight response. private(set) var streamingText = "" /// Streamed `` content of the in-flight response. private(set) var streamingThinking = "" /// True while the stream is inside a `` block. private(set) var isThinking = false /// Live generation speed, updated at most 4 Hz. private(set) var liveTokensPerSecond: Double = 0 /// Conversation currently streaming (may differ from the selection). private(set) var streamingConversationID: UUID? private weak var app: AppModel? private var generationTask: Task? func bind(to app: AppModel) { self.app = app } // MARK: - Actions /// Sends `prompt` in the given conversation, streaming the response. func send(prompt: String, in conversation: Conversation) { guard let app, !isGenerating else { return } var conversation = conversation let userMessage = Message(role: .user, content: prompt) conversation.messages.append(userMessage) app.update(conversation) stream(prompt: prompt, conversation: conversation) } /// Regenerates the last assistant response (optionally after the user /// switched models). func regenerate(in conversation: Conversation) { guard !isGenerating else { return } var conversation = conversation guard let lastUser = conversation.messages.last(where: { $0.role == .user }) else { return } // Drop trailing assistant message(s) after the last user turn. while let last = conversation.messages.last, last.role == .assistant { conversation.messages.removeLast() } app?.update(conversation) stream(prompt: lastUser.content, conversation: conversation, replayingLastUser: true) } /// Edits a previous user message and resends from that point. func editAndResend(messageID: UUID, newText: String, in conversation: Conversation) { guard !isGenerating else { return } var conversation = conversation guard let index = conversation.messages.firstIndex(where: { $0.id == messageID }) else { return } conversation.messages[index].content = newText conversation.messages.removeSubrange((index + 1)...) app?.update(conversation) stream(prompt: newText, conversation: conversation, replayingLastUser: true) } func stop() { generationTask?.cancel() Task { await app?.engine.stopGeneration() } } // MARK: - Streaming core /// `replayingLastUser`: the prompt is already the last user message in /// `conversation.messages`; the engine session must be rebuilt so its /// history excludes it (it is re-sent as the new turn). private func stream(prompt: String, conversation: Conversation, replayingLastUser: Bool = false) { guard let app else { return } let conversationID = conversation.id streamingText = "" streamingThinking = "" isThinking = false liveTokensPerSecond = 0 isGenerating = true streamingConversationID = conversationID generationTask = Task { var parser = ThinkTagParser() var stats: GenerationStats? var finish: GenerationFinishReason = .stop let started = Date() var tokenCount = 0 var lastTick = Date.distantPast do { // The engine session's history must exclude the new prompt: // strip the trailing user message before (re)building. var sessionConversation = conversation if let last = sessionConversation.messages.last, last.role == .user { sessionConversation.messages.removeLast() } if replayingLastUser { try await app.engine.startSession(conversation: sessionConversation) } else { try await app.engine.ensureSession(conversation: sessionConversation) } let events = try await app.engine.generate(prompt: prompt, params: conversation.params) for try await event in events { switch event { case .token(let text): tokenCount += 1 let (visible, thinking, inThink) = parser.consume(text) if !visible.isEmpty { streamingText += visible } if !thinking.isEmpty { streamingThinking += thinking } isThinking = inThink // ≤4 Hz ticker to avoid flicker. let now = Date() if now.timeIntervalSince(lastTick) >= 0.25 { lastTick = now let elapsed = now.timeIntervalSince(started) if elapsed > 0.5 { liveTokensPerSecond = Double(tokenCount) / elapsed } } case .stats(let s): stats = s case .finished(let reason): finish = reason } } } catch { app.lastError = error.localizedDescription } finalize(conversationID: conversationID, stats: stats, finish: finish) } } private func finalize(conversationID: UUID, stats: GenerationStats?, finish: GenerationFinishReason) { defer { isGenerating = false streamingConversationID = nil streamingText = "" streamingThinking = "" isThinking = false liveTokensPerSecond = 0 generationTask = nil } guard let app, var conversation = app.conversations.first(where: { $0.id == conversationID }) else { return } let content = streamingText.trimmingCharacters(in: .whitespacesAndNewlines) let thinking = streamingThinking.trimmingCharacters(in: .whitespacesAndNewlines) guard !content.isEmpty || !thinking.isEmpty else { return } let message = Message( role: .assistant, content: content, thinking: thinking.isEmpty ? nil : thinking, stats: stats.map { MessageStats( timeToFirstToken: $0.timeToFirstToken, tokensPerSecond: $0.tokensPerSecond, promptTokenCount: $0.promptTokenCount, generationTokenCount: $0.generationTokenCount, peakMemoryBytes: $0.peakMemoryBytes ) } ) conversation.messages.append(message) app.update(conversation) _ = finish // reason currently not surfaced beyond stats if conversation.title == "New Chat" { autoTitle(conversation: conversation) } } /// Short, cheap title generation after the first exchange, using the /// loaded model itself. Falls back to a truncated first prompt. private func autoTitle(conversation: Conversation) { guard let app else { return } guard let firstUser = conversation.messages.first(where: { $0.role == .user }) else { return } let fallback = String(firstUser.content.prefix(48)) Task { var title = fallback do { let prompt = """ Reply with a title of at most 5 words for a conversation that starts with \ this message, and nothing else — no quotes, no punctuation at the end: \(String(firstUser.content.prefix(500))) """ let temp = Conversation(params: GenerationParams(temperature: 0.1, maxTokens: 24)) try await app.engine.startSession(conversation: temp) var generated = "" let events = try await app.engine.generate(prompt: prompt, params: temp.params) for try await event in events { if case .token(let t) = event { generated += t } } var parser = ThinkTagParser() let (visible, _, _) = parser.consume(generated) let cleaned = visible .trimmingCharacters(in: .whitespacesAndNewlines) .trimmingCharacters(in: CharacterSet(charactersIn: "\"'.“”")) if !cleaned.isEmpty { title = String(cleaned.prefix(60)) } // Rebind the engine session to the real conversation. try? await app.engine.startSession(conversation: conversation) } catch { // Fallback title already set. } if var c = app.conversations.first(where: { $0.id == conversation.id }) { c.title = title app.update(c, touch: false) } } } } /// Incremental parser splitting a token stream into visible text and /// `` content, robust to tags split across chunks. struct ThinkTagParser { private var inThink = false private var pending = "" /// Returns (visibleDelta, thinkingDelta, isInsideThink). mutating func consume(_ chunk: String) -> (String, String, Bool) { pending += chunk var visible = "" var thinking = "" while true { if inThink { if let range = pending.range(of: "") { thinking += pending[..") thinking += pending.prefix(safe) pending = String(pending.dropFirst(safe)) break } } else { if let range = pending.range(of: "") { visible += pending[..") visible += pending.prefix(safe) pending = String(pending.dropFirst(safe)) break } } } return (visible, thinking, inThink) } /// Length of `text` that can be emitted without cutting a partial `tag` /// suffix that might complete in the next chunk. private func safeEmitLength(of text: String, partial tag: String) -> Int { let maxKeep = min(tag.count - 1, text.count) for keep in stride(from: maxKeep, through: 1, by: -1) { if text.hasSuffix(String(tag.prefix(keep))) { return text.count - keep } } return text.count } }