Streams replies over SSE from any OpenAI-compatible server (oMLX, llama.cpp, Ollama, ...). Single binary that can install itself as an OS service; docker compose bundles llama.cpp + Gemma 4 E2B.
180 lines
4.2 KiB
Go
180 lines
4.2 KiB
Go
// Package chat keeps conversations in memory, one per browser session.
|
|
package chat
|
|
|
|
import (
|
|
"crypto/rand"
|
|
"encoding/hex"
|
|
"sync"
|
|
"time"
|
|
|
|
"git.b0b.be/bdeb/localchat/internal/llm"
|
|
)
|
|
|
|
// State tracks an assistant reply through its lifecycle.
|
|
type State int
|
|
|
|
const (
|
|
Pending State = iota // created, nobody is generating it yet
|
|
Streaming // a stream handler owns it
|
|
Done // finished (successfully or with Err set)
|
|
)
|
|
|
|
// Message is a chat turn shown in the UI.
|
|
type Message struct {
|
|
ID string
|
|
Role string // "user" or "assistant"
|
|
Content string
|
|
State State
|
|
Err string
|
|
}
|
|
|
|
// Conversation is the history of one session.
|
|
type Conversation struct {
|
|
mu sync.Mutex
|
|
messages []*Message
|
|
lastUsed time.Time
|
|
}
|
|
|
|
// Store maps session ids to conversations.
|
|
type Store struct {
|
|
mu sync.Mutex
|
|
convs map[string]*Conversation
|
|
}
|
|
|
|
// NewStore returns an empty store.
|
|
func NewStore() *Store {
|
|
return &Store{convs: make(map[string]*Conversation)}
|
|
}
|
|
|
|
// NewID returns a random hex id, suitable for sessions and messages.
|
|
func NewID() string {
|
|
b := make([]byte, 16)
|
|
rand.Read(b)
|
|
return hex.EncodeToString(b)
|
|
}
|
|
|
|
// Get returns the conversation for a session, creating it if needed.
|
|
func (s *Store) Get(session string) *Conversation {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
c, ok := s.convs[session]
|
|
if !ok {
|
|
c = &Conversation{}
|
|
s.convs[session] = c
|
|
}
|
|
c.mu.Lock()
|
|
c.lastUsed = time.Now()
|
|
c.mu.Unlock()
|
|
return c
|
|
}
|
|
|
|
// Reset forgets a session's conversation.
|
|
func (s *Store) Reset(session string) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
delete(s.convs, session)
|
|
}
|
|
|
|
// Prune drops conversations idle for longer than maxIdle.
|
|
func (s *Store) Prune(maxIdle time.Duration) {
|
|
cutoff := time.Now().Add(-maxIdle)
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
for id, c := range s.convs {
|
|
c.mu.Lock()
|
|
idle := c.lastUsed.Before(cutoff)
|
|
c.mu.Unlock()
|
|
if idle {
|
|
delete(s.convs, id)
|
|
}
|
|
}
|
|
}
|
|
|
|
// Messages returns a snapshot of the conversation.
|
|
func (c *Conversation) Messages() []Message {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
out := make([]Message, len(c.messages))
|
|
for i, m := range c.messages {
|
|
out[i] = *m
|
|
}
|
|
return out
|
|
}
|
|
|
|
// Ask appends a user message plus a pending assistant reply and returns both.
|
|
func (c *Conversation) Ask(text string) (user, reply Message) {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
u := &Message{ID: NewID(), Role: "user", Content: text, State: Done}
|
|
r := &Message{ID: NewID(), Role: "assistant", State: Pending}
|
|
c.messages = append(c.messages, u, r)
|
|
return *u, *r
|
|
}
|
|
|
|
// Claim marks a pending reply as streaming and returns the context to send
|
|
// to the model: the system prompt plus up to maxHistory finished messages
|
|
// before the reply. ok is false if the reply does not exist or someone else
|
|
// already claimed it; msg then holds its current state (if it exists).
|
|
func (c *Conversation) Claim(id, systemPrompt string, maxHistory int) (prompt []llm.Message, msg Message, ok bool) {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
idx := c.index(id)
|
|
if idx < 0 {
|
|
return nil, Message{}, false
|
|
}
|
|
m := c.messages[idx]
|
|
if m.State != Pending {
|
|
return nil, *m, false
|
|
}
|
|
m.State = Streaming
|
|
|
|
var history []llm.Message
|
|
for _, h := range c.messages[:idx] {
|
|
if h.State == Done && h.Err == "" && h.Content != "" {
|
|
history = append(history, llm.Message{Role: h.Role, Content: h.Content})
|
|
}
|
|
}
|
|
if len(history) > maxHistory {
|
|
history = history[len(history)-maxHistory:]
|
|
}
|
|
if systemPrompt != "" {
|
|
prompt = append(prompt, llm.Message{Role: "system", Content: systemPrompt})
|
|
}
|
|
return append(prompt, history...), *m, true
|
|
}
|
|
|
|
// Append adds streamed text to a reply.
|
|
func (c *Conversation) Append(id, delta string) {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
if i := c.index(id); i >= 0 {
|
|
c.messages[i].Content += delta
|
|
}
|
|
}
|
|
|
|
// Finish marks a reply as done, recording err if generation failed, and
|
|
// returns its final state.
|
|
func (c *Conversation) Finish(id string, err error) Message {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
i := c.index(id)
|
|
if i < 0 {
|
|
return Message{ID: id, Role: "assistant", State: Done}
|
|
}
|
|
m := c.messages[i]
|
|
m.State = Done
|
|
if err != nil {
|
|
m.Err = err.Error()
|
|
}
|
|
return *m
|
|
}
|
|
|
|
func (c *Conversation) index(id string) int {
|
|
for i, m := range c.messages {
|
|
if m.ID == id {
|
|
return i
|
|
}
|
|
}
|
|
return -1
|
|
}
|