olliesrv/fs: embed *session.Session in SessionNode, kill delegation
SessionNode now embeds *session.Session instead of wrapping it. All delegation methods (ID(), Name(), Ctx(), Pause(), etc.) are removed. Handlers access session fields/methods directly through promotion. Moved pause/resume event publishing into session.Session itself since the session package already owns PublishEvent. Eliminated all s.Session.X stutter from handler code.
This commit is contained in:
parent
1474eeaaed
commit
d71ff80993
|
|
@ -22,7 +22,7 @@ func requestAgentNew(s *SessionNode, data []byte) ([]byte, error) {
|
|||
// Create AgentLog for this agent and wire up the event handler
|
||||
al := NewAgentLog("")
|
||||
s.SetAgentLog(ag.ID(), al)
|
||||
wireAgentEvents(s.ID(), ag, al)
|
||||
wireAgentEvents(s.ID, ag, al)
|
||||
return []byte(ag.ID() + "\n"), nil
|
||||
}
|
||||
|
||||
|
|
@ -36,7 +36,7 @@ func writeAgentPrompt(s *SessionNode, ag *agent.Agent, al *AgentLog, data []byte
|
|||
return nil
|
||||
}
|
||||
go func() {
|
||||
ag.Submit(s.Ctx(), input)
|
||||
ag.Submit(s.Ctx, input)
|
||||
al.EnsureTrailingNewline()
|
||||
}()
|
||||
return nil
|
||||
|
|
@ -78,13 +78,13 @@ func readAgentChatText(al *AgentLog) ([]byte, error) {
|
|||
}
|
||||
|
||||
func streamAgentChat(s *SessionNode, al *AgentLog, cctx context.Context, base string) ([]byte, string, error) {
|
||||
merged, cancel := mergeCtx(s.Ctx(), cctx)
|
||||
merged, cancel := mergeCtx(s.Ctx, cctx)
|
||||
defer cancel()
|
||||
return streamChat(al, merged, base)
|
||||
}
|
||||
|
||||
func streamAgentChatText(s *SessionNode, al *AgentLog, cctx context.Context, base string) ([]byte, string, error) {
|
||||
merged, cancel := mergeCtx(s.Ctx(), cctx)
|
||||
merged, cancel := mergeCtx(s.Ctx, cctx)
|
||||
defer cancel()
|
||||
for {
|
||||
data, nextBase, err := streamChat(al, merged, base)
|
||||
|
|
@ -137,7 +137,7 @@ func stripMarkers(data []byte) []byte {
|
|||
}
|
||||
|
||||
func blockAgentStateWait(s *SessionNode, ag *agent.Agent, cctx context.Context, base string) ([]byte, string, error) {
|
||||
merged, cancel := mergeCtx(s.Ctx(), cctx)
|
||||
merged, cancel := mergeCtx(s.Ctx, cctx)
|
||||
defer cancel()
|
||||
if base == "" {
|
||||
base = ag.State()
|
||||
|
|
@ -185,7 +185,7 @@ func readAgentCfg(s *SessionNode, ag *agent.Agent, al *AgentLog) ([]byte, error)
|
|||
defer al.mu.RUnlock()
|
||||
p := ag.GenParams()
|
||||
var sb strings.Builder
|
||||
fmt.Fprintf(&sb, "name=%s\n", s.ID())
|
||||
fmt.Fprintf(&sb, "name=%s\n", s.ID)
|
||||
fmt.Fprintf(&sb, "backend=%s\n", ag.BackendName())
|
||||
fmt.Fprintf(&sb, "model=%s\n", ag.ModelName())
|
||||
fmt.Fprintf(&sb, "id=%s\n", ag.ID())
|
||||
|
|
@ -232,10 +232,8 @@ func writeAgentCfg(ag *agent.Agent, data []byte) error {
|
|||
func agentCtlHandler(s *SessionNode, ag *agent.Agent) func([]byte) ([]byte, error) {
|
||||
handlers := map[string]func([]string) ([]byte, error){
|
||||
"kill": func(_ []string) ([]byte, error) {
|
||||
if s.Session != nil {
|
||||
s.Session.RemoveAgent(ag.ID())
|
||||
session.PublishEvent("session."+s.ID()+".agent."+ag.ID()+".kill", "")
|
||||
}
|
||||
s.RemoveAgent(ag.ID())
|
||||
session.PublishEvent("session."+s.ID+".agent."+ag.ID()+".kill", "")
|
||||
return []byte("ok\n"), nil
|
||||
},
|
||||
"stop": func(_ []string) ([]byte, error) {
|
||||
|
|
@ -243,7 +241,7 @@ func agentCtlHandler(s *SessionNode, ag *agent.Agent) func([]byte) ([]byte, erro
|
|||
return []byte("ok\n"), nil
|
||||
},
|
||||
"compact": func(_ []string) ([]byte, error) {
|
||||
if err := ag.Compact(s.Ctx()); err != nil {
|
||||
if err := ag.Compact(s.Ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return []byte("ok\n"), nil
|
||||
|
|
|
|||
|
|
@ -40,7 +40,7 @@ func requestSessionNew(data []byte) ([]byte, error) {
|
|||
|
||||
func readSessionEnv(s *SessionNode) ([]byte, error) {
|
||||
var sb strings.Builder
|
||||
fmt.Fprintf(&sb, "OLLIE_SESSION_ID=%s\n", s.ID())
|
||||
fmt.Fprintf(&sb, "OLLIE_SESSION_ID=%s\n", s.ID)
|
||||
for _, e := range os.Environ() {
|
||||
if strings.HasPrefix(e, "OLLIE_") && !strings.HasPrefix(e, "OLLIE_SESSION_ID=") {
|
||||
sb.WriteString(e)
|
||||
|
|
@ -75,9 +75,7 @@ func sessionCtlHandler(s *SessionNode, removeFn func(), renameFn func(string) er
|
|||
return []byte("ok\n"), nil
|
||||
},
|
||||
"save": func(_ []string) ([]byte, error) {
|
||||
if s.Session != nil {
|
||||
s.Session.Save()
|
||||
}
|
||||
s.Save()
|
||||
return []byte("ok\n"), nil
|
||||
},
|
||||
"invalidate": func(_ []string) ([]byte, error) {
|
||||
|
|
@ -119,14 +117,14 @@ func writeSessionName(s *SessionNode, renameFn func(string) error, data []byte)
|
|||
}
|
||||
|
||||
func readSessionID(s *SessionNode) ([]byte, error) {
|
||||
return []byte(s.ID() + "\n"), nil
|
||||
return []byte(s.ID + "\n"), nil
|
||||
}
|
||||
|
||||
func readSessionBypass(b *bypass.Broker, s *SessionNode) ([]byte, error) {
|
||||
if b == nil {
|
||||
return nil, fmt.Errorf("bypass broker not available")
|
||||
}
|
||||
p := b.SessionPolicy(s.ID())
|
||||
p := b.SessionPolicy(s.ID)
|
||||
return p.Marshal()
|
||||
}
|
||||
|
||||
|
|
@ -138,6 +136,6 @@ func writeSessionBypass(b *bypass.Broker, s *SessionNode, data []byte) error {
|
|||
if err := bypass.ParsePolicy(data, &p); err != nil {
|
||||
return err
|
||||
}
|
||||
b.SetSessionPolicy(s.ID(), p)
|
||||
b.SetSessionPolicy(s.ID, p)
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,20 +1,17 @@
|
|||
package fs
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"ollie/cmd/olliesrv/internal/agent"
|
||||
"ollie/cmd/olliesrv/internal/session"
|
||||
"ollie/toolsrv"
|
||||
)
|
||||
|
||||
// SessionNode wraps session.Session for 9P handlers.
|
||||
// It adds AgentLogs (9P chat streaming state) and models cache.
|
||||
// SessionNode holds per-session 9P-layer state: agent chat logs and a cached
|
||||
// model list. It embeds *session.Session for direct access to session operations.
|
||||
type SessionNode struct {
|
||||
Session *session.Session
|
||||
*session.Session
|
||||
|
||||
// AgentLogs for chat streaming, keyed by agent ID.
|
||||
mu sync.RWMutex
|
||||
|
|
@ -31,48 +28,6 @@ func NewSessionNode(core *session.Session) *SessionNode {
|
|||
return &SessionNode{Session: core, AgentLogs: make(map[string]*AgentLog)}
|
||||
}
|
||||
|
||||
// --- Delegation to Core ---
|
||||
|
||||
func (s *SessionNode) ID() string { return s.Session.ID }
|
||||
func (s *SessionNode) Name() string { return s.Session.Name() }
|
||||
func (s *SessionNode) SetName(n string) { s.Session.SetName(n) }
|
||||
func (s *SessionNode) Uname() string { return s.Session.Uname }
|
||||
func (s *SessionNode) Cancel() { if s.Session.Cancel != nil { s.Session.Cancel() } }
|
||||
func (s *SessionNode) IsPaused() bool { return s.Session.IsPaused() }
|
||||
func (s *SessionNode) IsConnected() bool { return s.Session.IsConnected() }
|
||||
|
||||
func (s *SessionNode) Pause() error {
|
||||
err := s.Session.Pause()
|
||||
if err == nil {
|
||||
session.PublishEvent("session."+s.Session.ID+".pause", "")
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *SessionNode) Resume() error {
|
||||
err := s.Session.Resume()
|
||||
if err == nil {
|
||||
session.PublishEvent("session."+s.Session.ID+".resume", "")
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *SessionNode) LoadTool(name string, ag *agent.Agent) error {
|
||||
return s.Session.LoadTool(name, ag)
|
||||
}
|
||||
|
||||
func (s *SessionNode) Ctx() context.Context {
|
||||
return s.Session.Ctx
|
||||
}
|
||||
|
||||
func (s *SessionNode) ToolsConn() *toolsrv.Conn {
|
||||
return s.Session.ToolsConn()
|
||||
}
|
||||
|
||||
func (s *SessionNode) DialToolServer() *toolsrv.Conn {
|
||||
return s.Session.DialToolServer()
|
||||
}
|
||||
|
||||
// --- AgentLog management ---
|
||||
|
||||
func (s *SessionNode) AgentLog() *AgentLog {
|
||||
|
|
@ -112,7 +67,7 @@ func (s *SessionNode) CachedListModels() string {
|
|||
}
|
||||
s.modelsMu.Unlock()
|
||||
|
||||
agents := s.Session.Agents()
|
||||
agents := s.Agents()
|
||||
if len(agents) == 0 {
|
||||
return ""
|
||||
}
|
||||
|
|
|
|||
|
|
@ -197,24 +197,24 @@ func buildSessionChildren(
|
|||
|
||||
// buildAgentBindings produces bindings for each agent in a session.
|
||||
func buildAgentBindings(s *SessionNode, rs *RootState) ([]fsedsl.Binding, error) {
|
||||
agents := s.Session.Agents()
|
||||
agents := s.Agents()
|
||||
var out []fsedsl.Binding
|
||||
|
||||
for _, ag := range agents {
|
||||
a := ag
|
||||
al := s.AgentLogFor(a.ID())
|
||||
if al == nil {
|
||||
al = NewAgentLog(s.Session.Remote)
|
||||
al = NewAgentLog(s.Remote)
|
||||
s.SetAgentLog(a.ID(), al)
|
||||
wireAgentEvents(s.Session.ID, a, al)
|
||||
wireAgentEvents(s.ID, a, al)
|
||||
}
|
||||
|
||||
out = append(out, fsedsl.Binding{
|
||||
Name: a.Name(),
|
||||
Aliases: []string{a.ID()},
|
||||
Remove: func() {
|
||||
s.Session.RemoveAgent(a.ID())
|
||||
session.PublishEvent("session."+s.ID()+".agent."+a.ID()+".kill", "")
|
||||
s.RemoveAgent(a.ID())
|
||||
session.PublishEvent("session."+s.ID+".agent."+a.ID()+".kill", "")
|
||||
},
|
||||
Children: buildAgentChildren(a, al, s),
|
||||
})
|
||||
|
|
|
|||
|
|
@ -229,6 +229,7 @@ func (s *Session) Pause() error {
|
|||
}
|
||||
s.paused = true
|
||||
go s.saveSession()
|
||||
PublishEvent("session."+s.ID+".pause", "")
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -282,6 +283,7 @@ func (s *Session) Resume() error {
|
|||
|
||||
s.paused = false
|
||||
go s.saveSession()
|
||||
PublishEvent("session."+s.ID+".resume", "")
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue