9p: goroutine-per-request; *wait files use connection context
Serve now dispatches each 9P request to its own goroutine with a shared response channel, so blocking reads no longer stall the serve loop. *wait files are handled via SessionFileStore.Wait(connCtx, name), which merges the connection context with the session context via context.AfterFunc — the wait unblocks on either connection close or session kill, preventing goroutine leaks. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
596f658703
commit
a6b6e911dd
|
|
@ -95,8 +95,10 @@ type fid struct {
|
|||
|
||||
// connState tracks all open fids for a single 9P connection.
|
||||
type connState struct {
|
||||
mu sync.RWMutex
|
||||
fids map[uint32]*fid
|
||||
mu sync.RWMutex
|
||||
fids map[uint32]*fid
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
}
|
||||
|
||||
// Server is the 9P server for ollie sessions.
|
||||
|
|
@ -191,14 +193,21 @@ func (s *Server) sessionFileStore(sessID string) (*SessionFileStore, bool) {
|
|||
), true
|
||||
}
|
||||
|
||||
// Serve handles a single 9P connection.
|
||||
// Serve handles a single 9P connection. Each request is dispatched to its own
|
||||
// goroutine so blocking reads (e.g. *wait files) do not stall the serve loop.
|
||||
func (s *Server) Serve(conn net.Conn) {
|
||||
defer conn.Close()
|
||||
cs := &connState{fids: make(map[uint32]*fid)}
|
||||
connCtx, connCancel := context.WithCancel(context.Background())
|
||||
cs := &connState{
|
||||
fids: make(map[uint32]*fid),
|
||||
ctx: connCtx,
|
||||
cancel: connCancel,
|
||||
}
|
||||
s.mu.Lock()
|
||||
s.conns = append(s.conns, cs)
|
||||
s.mu.Unlock()
|
||||
defer func() {
|
||||
connCancel()
|
||||
s.mu.Lock()
|
||||
for i, c := range s.conns {
|
||||
if c == cs {
|
||||
|
|
@ -208,16 +217,34 @@ func (s *Server) Serve(conn net.Conn) {
|
|||
}
|
||||
s.mu.Unlock()
|
||||
}()
|
||||
|
||||
responses := make(chan *plan9.Fcall, 16)
|
||||
var wg sync.WaitGroup
|
||||
|
||||
// writer: serialises responses back onto the connection.
|
||||
go func() {
|
||||
for resp := range responses {
|
||||
plan9.WriteFcall(conn, resp) //nolint:errcheck
|
||||
}
|
||||
}()
|
||||
|
||||
for {
|
||||
fc, err := plan9.ReadFcall(conn)
|
||||
if err != nil {
|
||||
if err != io.EOF {
|
||||
plog.Error("read: %v", err)
|
||||
}
|
||||
return
|
||||
break
|
||||
}
|
||||
plan9.WriteFcall(conn, s.handle(cs, fc)) //nolint:errcheck
|
||||
wg.Add(1)
|
||||
go func(fc *plan9.Fcall) {
|
||||
defer wg.Done()
|
||||
responses <- s.handle(cs, fc)
|
||||
}(fc)
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
close(responses)
|
||||
}
|
||||
|
||||
var plog = olog.New("9p")
|
||||
|
|
@ -722,6 +749,14 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
if parts[2] == "fifo.out" && fc.Offset > 0 {
|
||||
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
|
||||
}
|
||||
// *wait files block until a value changes; use connection context.
|
||||
if strings.HasSuffix(parts[2], "wait") {
|
||||
content, err := store.Wait(cs.ctx, parts[2])
|
||||
if err != nil {
|
||||
return errFcall(fc, err.Error())
|
||||
}
|
||||
return s.readSlice(fc, content)
|
||||
}
|
||||
content, err := store.Get(parts[2])
|
||||
if err != nil {
|
||||
plog.Debug("Rread session file err=%v", err)
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package p9
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"strconv"
|
||||
|
|
@ -97,34 +98,6 @@ func (s *SessionFileStore) Get(name string) ([]byte, error) {
|
|||
return nil, nil
|
||||
}
|
||||
return []byte(item), nil
|
||||
case "statewait":
|
||||
current := s.sess.core.State()
|
||||
v, ok := s.sess.core.WaitChange(s.sess.ctx, agent.WatchState, current)
|
||||
if !ok {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(v + "\n"), nil
|
||||
case "usagewait":
|
||||
current := s.sess.core.Usage()
|
||||
v, ok := s.sess.core.WaitChange(s.sess.ctx, agent.WatchUsage, current)
|
||||
if !ok {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(v + "\n"), nil
|
||||
case "ctxszwait":
|
||||
current := s.sess.core.CtxSz()
|
||||
v, ok := s.sess.core.WaitChange(s.sess.ctx, agent.WatchCtxSz, current)
|
||||
if !ok {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(v + "\n"), nil
|
||||
case "cwdwait":
|
||||
current := s.sess.core.CWD()
|
||||
v, ok := s.sess.core.WaitChange(s.sess.ctx, agent.WatchCWD, current)
|
||||
if !ok {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(v + "\n"), nil
|
||||
default:
|
||||
for _, f := range sessionFileList {
|
||||
if f.name == name {
|
||||
|
|
@ -299,6 +272,36 @@ func (s *SessionFileStore) content(name string) string {
|
|||
return ""
|
||||
}
|
||||
|
||||
// Wait blocks until the named *wait file's underlying value changes, then
|
||||
// returns the new value. Unblocks when connCtx or the session context is
|
||||
// cancelled, returning nil with no error (caller returns EOF to client).
|
||||
func (s *SessionFileStore) Wait(connCtx context.Context, name string) ([]byte, error) {
|
||||
ctx, cancel := context.WithCancel(connCtx)
|
||||
defer cancel()
|
||||
// also unblock when the session itself is killed.
|
||||
context.AfterFunc(s.sess.ctx, cancel)
|
||||
|
||||
var field, current string
|
||||
switch name {
|
||||
case "statewait":
|
||||
field, current = agent.WatchState, s.sess.core.State()
|
||||
case "usagewait":
|
||||
field, current = agent.WatchUsage, s.sess.core.Usage()
|
||||
case "ctxszwait":
|
||||
field, current = agent.WatchCtxSz, s.sess.core.CtxSz()
|
||||
case "cwdwait":
|
||||
field, current = agent.WatchCWD, s.sess.core.CWD()
|
||||
default:
|
||||
return nil, fmt.Errorf("%s: not a wait file", name)
|
||||
}
|
||||
|
||||
v, ok := s.sess.core.WaitChange(ctx, field, current)
|
||||
if !ok {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(v + "\n"), nil
|
||||
}
|
||||
|
||||
func (s *SessionFileStore) makePublish() func(agent.Event) {
|
||||
assistantStarted := false
|
||||
return func(ev agent.Event) {
|
||||
|
|
|
|||
Reference in New Issue