session: subscribe to event bus for chat log, drop handler from Submit
- NewSession subscribes to bus "event" topic for chat log writing - Remove makePublish/handler threading from all Submit call sites - Adapt tests to new bus-based event delivery
This commit is contained in:
parent
518a31a1e2
commit
afb531c936
1
go.mod
1
go.mod
|
|
@ -9,6 +9,7 @@ require (
|
|||
)
|
||||
|
||||
require (
|
||||
github.com/simonfxr/pubsub v0.0.5 // indirect
|
||||
golang.org/x/sys v0.28.0 // indirect
|
||||
gopkg.in/yaml.v3 v3.0.1 // indirect
|
||||
)
|
||||
|
|
|
|||
7
go.sum
7
go.sum
|
|
@ -2,6 +2,7 @@
|
|||
9fans.net/go v0.0.7/go.mod h1:Rxvbbc1e+1TyGMjAvLthGTyO97t+6JMQ6ly+Lcs9Uf0=
|
||||
dmitri.shuralyov.com/gpu/mtl v0.0.0-20201218220906-28db891af037/go.mod h1:H6x//7gZCb22OMCxBHrMx7a5I7Hp++hsVxbQ4BYO7hU=
|
||||
github.com/BurntSushi/xgb v0.0.0-20160522181843-27f122750802/go.mod h1:IVnqGOEym/WlBOVXweHU+Q+/VP0lqqI8lqeDx9IjBqo=
|
||||
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/go-gl/glfw/v3.3/glfw v0.0.0-20200222043503-6f7a984d4dc4/go.mod h1:tQ2UAYgL5IevRw8kRxooKSPJfGvJ9fJQFa0TUsXzTg8=
|
||||
github.com/hanwen/go-fuse/v2 v2.10.1 h1:QAqZuc9+aBtTou+OPruU/hkYQYCkgPtQd2QaepHkTTs=
|
||||
github.com/hanwen/go-fuse/v2 v2.10.1/go.mod h1:aU7NkGYZUmuJrZapoI3mEcNve7PZTySUOLBuch/vR6U=
|
||||
|
|
@ -9,6 +10,11 @@ github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0
|
|||
github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw=
|
||||
github.com/moby/sys/mountinfo v0.7.2 h1:1shs6aH5s4o5H2zQLn796ADW1wMrIwHsyJ2v9KouLrg=
|
||||
github.com/moby/sys/mountinfo v0.7.2/go.mod h1:1YOa8w8Ih7uW0wALDUgT1dTTSBrZ+HiBLGws92L2RU4=
|
||||
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
|
||||
github.com/simonfxr/pubsub v0.0.5 h1:DJfvFoglqGvwJriIOC5NI5um34n2YX8KAU5+7jv768w=
|
||||
github.com/simonfxr/pubsub v0.0.5/go.mod h1:bQ+B2NEEHZ08VY/0xVFGaGL1KMJ7cl3yLaXe8xvyC24=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
|
||||
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
|
||||
golang.org/x/crypto v0.0.0-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
|
||||
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
|
||||
|
|
@ -42,5 +48,6 @@ golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8T
|
|||
golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
|
||||
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
|
|
|
|||
30
main_test.go
30
main_test.go
|
|
@ -10,6 +10,7 @@ import (
|
|||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/simonfxr/pubsub"
|
||||
"ollie/pkg/agent"
|
||||
"ollie/pkg/backend"
|
||||
olog "ollie/pkg/log"
|
||||
|
|
@ -39,12 +40,13 @@ type stubCore struct {
|
|||
interrupted bool
|
||||
setSessionIDErr error
|
||||
submitCh chan struct{}
|
||||
bus_ *pubsub.Bus
|
||||
}
|
||||
|
||||
func (c *stubCore) Submit(_ context.Context, input string, handler agent.EventHandler) {
|
||||
func (c *stubCore) Submit(_ context.Context, input string) {
|
||||
c.submitted = append(c.submitted, input)
|
||||
if handler != nil && c.reply != "" {
|
||||
handler(agent.Event{Role: "assistant", Content: c.reply})
|
||||
if c.reply != "" {
|
||||
c.Bus().Publish("event", agent.Event{Role: "assistant", Content: c.reply})
|
||||
}
|
||||
if c.submitCh != nil {
|
||||
close(c.submitCh)
|
||||
|
|
@ -87,20 +89,26 @@ func (c *stubCore) WaitChange(ctx context.Context, _, _ string) (string, bool) {
|
|||
}
|
||||
func (c *stubCore) Close() { c.closed = true }
|
||||
func (c *stubCore) ToolCallCount() int64 { return 0 }
|
||||
func (c *stubCore) Bus() *pubsub.Bus {
|
||||
if c.bus_ == nil {
|
||||
c.bus_ = pubsub.NewBus()
|
||||
}
|
||||
return c.bus_
|
||||
}
|
||||
|
||||
type publishCore struct{ *stubCore }
|
||||
func (c *publishCore) Submit(_ context.Context, input string, handler agent.EventHandler) {
|
||||
func (c *publishCore) Submit(_ context.Context, input string) {
|
||||
c.submitted = append(c.submitted, input)
|
||||
if handler == nil { return }
|
||||
handler(agent.Event{Role: "user", Content: input})
|
||||
handler(agent.Event{Role: "assistant", Content: "thinking..."})
|
||||
handler(agent.Event{Role: "call", Name: "fn", Content: "arg1"})
|
||||
handler(agent.Event{Role: "tool", Content: "result"})
|
||||
handler(agent.Event{Role: "assistant", Content: "done"})
|
||||
bus := c.Bus()
|
||||
bus.Publish("event", agent.Event{Role: "user", Content: input})
|
||||
bus.Publish("event", agent.Event{Role: "assistant", Content: "thinking..."})
|
||||
bus.Publish("event", agent.Event{Role: "call", Name: "fn", Content: "arg1"})
|
||||
bus.Publish("event", agent.Event{Role: "tool", Content: "result"})
|
||||
bus.Publish("event", agent.Event{Role: "assistant", Content: "done"})
|
||||
}
|
||||
|
||||
type blockingCore struct{ *stubCore }
|
||||
func (c *blockingCore) Submit(ctx context.Context, input string, handler agent.EventHandler) {
|
||||
func (c *blockingCore) Submit(ctx context.Context, input string) {
|
||||
c.submitted = append(c.submitted, input)
|
||||
<-ctx.Done()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -36,7 +36,9 @@ type Session struct {
|
|||
}
|
||||
|
||||
func NewSession(id string, core agent.Core, ctx context.Context, cancel context.CancelFunc) *Session {
|
||||
return &Session{id: id, Core: core, Ctx: ctx, cancel: cancel}
|
||||
sess := &Session{id: id, Core: core, Ctx: ctx, cancel: cancel}
|
||||
sess.startEventLog()
|
||||
return sess
|
||||
}
|
||||
|
||||
func (sess *Session) RunnableID() string { return sess.id }
|
||||
|
|
@ -71,6 +73,38 @@ func (sess *Session) EnsureTrailingNewline() {
|
|||
sess.mu.Unlock()
|
||||
}
|
||||
|
||||
// startEventLog subscribes to the agent's "event" bus topic and writes
|
||||
// all events to the session chat log.
|
||||
func (sess *Session) startEventLog() {
|
||||
assistantStarted := false
|
||||
sess.Core.Bus().Subscribe("event", func(ev agent.Event) {
|
||||
if ev.Role == "user" {
|
||||
if assistantStarted {
|
||||
sess.AppendLog([]byte("\n"))
|
||||
assistantStarted = false
|
||||
}
|
||||
sess.AppendLog(FormatEvent(ev))
|
||||
return
|
||||
}
|
||||
switch ev.Role {
|
||||
case "assistant":
|
||||
if !assistantStarted {
|
||||
sess.AppendLog([]byte("assistant: "))
|
||||
sess.mu.Lock()
|
||||
sess.ChatOffset = len(sess.log)
|
||||
sess.mu.Unlock()
|
||||
assistantStarted = true
|
||||
}
|
||||
default:
|
||||
if assistantStarted {
|
||||
sess.AppendLog([]byte("\n"))
|
||||
assistantStarted = false
|
||||
}
|
||||
}
|
||||
sess.AppendLog(FormatEvent(ev))
|
||||
})
|
||||
}
|
||||
|
||||
// LogInfo returns the current log length and version atomically.
|
||||
func (sess *Session) LogInfo() (length int, vers uint32) {
|
||||
sess.mu.RLock()
|
||||
|
|
|
|||
|
|
@ -145,9 +145,8 @@ func (h *sessionHelper) fileSpec(name string, mode os.FileMode) fs.FileSpec {
|
|||
h.sess.mu.Lock()
|
||||
h.sess.prevPrompt = []byte(input)
|
||||
h.sess.mu.Unlock()
|
||||
pub := h.makePublish()
|
||||
go func() {
|
||||
h.sess.Core.Submit(h.sess.Ctx, input, pub)
|
||||
h.sess.Core.Submit(h.sess.Ctx, input)
|
||||
h.sess.EnsureTrailingNewline()
|
||||
}()
|
||||
return nil
|
||||
|
|
@ -360,7 +359,7 @@ func (h *sessionHelper) handleCfg(input string) error {
|
|||
if h.sess.Core.IsRunning() {
|
||||
return fmt.Errorf("cannot switch %s while agent is running", k)
|
||||
}
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/agent "+v, h.makePublish())
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/agent "+v)
|
||||
case "backend", "model":
|
||||
if v == "" {
|
||||
continue
|
||||
|
|
@ -488,10 +487,10 @@ func (h *sessionHelper) handleCfg(input string) error {
|
|||
}
|
||||
// Apply backend/model after agent= so manual overrides win.
|
||||
if deferredBackend != "" {
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/backend "+deferredBackend, h.makePublish())
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/backend "+deferredBackend)
|
||||
}
|
||||
if deferredModel != "" {
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/model "+deferredModel, h.makePublish())
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/model "+deferredModel)
|
||||
}
|
||||
if hasParams {
|
||||
if h.sess.Core.IsRunning() {
|
||||
|
|
@ -503,37 +502,6 @@ func (h *sessionHelper) handleCfg(input string) error {
|
|||
}
|
||||
|
||||
|
||||
func (h *sessionHelper) makePublish() func(agent.Event) {
|
||||
assistantStarted := false
|
||||
return func(ev agent.Event) {
|
||||
if ev.Role == "user" {
|
||||
if assistantStarted {
|
||||
h.sess.AppendLog([]byte("\n"))
|
||||
assistantStarted = false
|
||||
}
|
||||
h.sess.AppendLog(FormatEvent(ev))
|
||||
return
|
||||
} else {
|
||||
switch ev.Role {
|
||||
case "assistant":
|
||||
if !assistantStarted {
|
||||
h.sess.AppendLog([]byte("assistant: "))
|
||||
h.sess.mu.Lock()
|
||||
h.sess.ChatOffset = len(h.sess.log)
|
||||
h.sess.mu.Unlock()
|
||||
assistantStarted = true
|
||||
}
|
||||
default:
|
||||
if assistantStarted {
|
||||
h.sess.AppendLog([]byte("\n"))
|
||||
assistantStarted = false
|
||||
}
|
||||
}
|
||||
}
|
||||
h.sess.AppendLog(FormatEvent(ev))
|
||||
}
|
||||
}
|
||||
|
||||
func (h *sessionHelper) handleCtl(input string) error {
|
||||
cmd := strings.Fields(input)
|
||||
if len(cmd) == 0 {
|
||||
|
|
@ -560,7 +528,7 @@ func (h *sessionHelper) handleCtl(input string) error {
|
|||
"agents", "agent", "sessions", "cwd", "skills",
|
||||
"tools", "context", "usage", "cost", "history",
|
||||
"irw", "help":
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/"+input, h.makePublish())
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/"+input)
|
||||
default:
|
||||
return fmt.Errorf("unknown ctl command: %s", cmd[0])
|
||||
}
|
||||
|
|
|
|||
Reference in New Issue