9p: remove tailable machinery from non-chat session files
Tailing works well for append-only logs but is fundamentally racy for files that get overwritten: fast state changes collapse into a single truncation signal, so tail reads the wrong (already-stale) value. Remove mutableFile, mutableVers, trackMutable, and the per-fid readBuf entirely. Only chat retains append-only stat reporting (Qid.Vers + Length) for tail -f. All other session files (state, usage, backend, model, agent, cwd) are plain poll-on-read files. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
a876190310
commit
d10bdf6fcc
|
|
@ -42,7 +42,6 @@ import (
|
|||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"9fans.net/go/plan9"
|
||||
|
|
@ -60,23 +59,6 @@ const (
|
|||
QTFile = plan9.QTFILE
|
||||
)
|
||||
|
||||
// mutableFile tracks version and last-known value for a tailable session file.
|
||||
// Qid.Vers is bumped on change so tail -f detects rewrites via stat polling.
|
||||
// On change, truncated is set to 1 so the next stat reports Length 0,
|
||||
// causing tail to detect truncation and re-read from offset 0.
|
||||
type mutableFile struct {
|
||||
vers uint32
|
||||
last string
|
||||
truncated int32 // atomic: 1 after change, cleared by stat
|
||||
}
|
||||
|
||||
func (m *mutableFile) update(val string) {
|
||||
if val != m.last {
|
||||
m.last = val
|
||||
m.vers++
|
||||
atomic.StoreInt32(&m.truncated, 1)
|
||||
}
|
||||
}
|
||||
|
||||
// session holds all state for one ollie agent session.
|
||||
type session struct {
|
||||
|
|
@ -90,10 +72,6 @@ type session struct {
|
|||
chatLog []byte
|
||||
chatVers uint32 // incremented on each append; used as Qid.Vers
|
||||
chatOffset int // byte position of chat EOF immediately before last prompt submit
|
||||
// mutableVers tracks changes to tailable mutable-state files
|
||||
// (state, backend, agent, model, usage, cwd). Bumped on change;
|
||||
// reported via Qid.Vers in stat so tail -f detects truncation/rewrite.
|
||||
mutableVers map[string]*mutableFile
|
||||
}
|
||||
|
||||
|
||||
|
|
@ -108,17 +86,6 @@ func (sess *session) appendChat(data []byte) {
|
|||
sess.mu.Unlock()
|
||||
}
|
||||
|
||||
// trackMutable snapshots state and usage into their version trackers.
|
||||
// Called from the publish callback after each event.
|
||||
func (sess *session) trackMutable() {
|
||||
state := sess.core.State()
|
||||
usage := sess.core.Usage()
|
||||
sess.mu.Lock()
|
||||
sess.mutableVers["state"].update(state)
|
||||
sess.mutableVers["usage"].update(usage)
|
||||
sess.mu.Unlock()
|
||||
}
|
||||
|
||||
// fid tracks per-descriptor state for a single 9P connection.
|
||||
type fid struct {
|
||||
path string
|
||||
|
|
@ -1203,14 +1170,6 @@ func (s *Server) createSession(args []string) error {
|
|||
core: core,
|
||||
ctx: ctx,
|
||||
cancel: cancel,
|
||||
mutableVers: map[string]*mutableFile{
|
||||
"state": {last: core.State()},
|
||||
"backend": {last: core.BackendName()},
|
||||
"agent": {last: core.AgentName()},
|
||||
"model": {last: core.ModelName()},
|
||||
"usage": {last: core.Usage()},
|
||||
"cwd": {last: core.CWD()},
|
||||
},
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
|
|
@ -1544,27 +1503,11 @@ func (s *Server) makeStat(path string) plan9.Dir {
|
|||
sess := s.sessions[parts[1]]
|
||||
s.mu.RUnlock()
|
||||
if sess != nil {
|
||||
switch base {
|
||||
case "chat":
|
||||
if base == "chat" {
|
||||
sess.mu.RLock()
|
||||
dir.Length = uint64(len(sess.chatLog))
|
||||
dir.Qid.Vers = sess.chatVers
|
||||
sess.mu.RUnlock()
|
||||
default:
|
||||
sess.mu.RLock()
|
||||
mf := sess.mutableVers[base]
|
||||
if mf != nil {
|
||||
dir.Qid.Vers = mf.vers
|
||||
}
|
||||
sess.mu.RUnlock()
|
||||
if mf != nil && atomic.CompareAndSwapInt32(&mf.truncated, 1, 0) {
|
||||
// Report size 0 once so tail -f detects truncation.
|
||||
dir.Length = 0
|
||||
} else if store, ok := s.sessionFileStore(parts[1]); ok {
|
||||
if info, err := store.Stat(base); err == nil {
|
||||
dir.Length = uint64(info.Size())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -112,7 +112,6 @@ func (s *SessionFileStore) Put(name string, data []byte) error {
|
|||
|
||||
case "prompt":
|
||||
s.sess.core.Submit(s.sess.ctx, input, s.makePublish())
|
||||
s.sess.trackMutable()
|
||||
|
||||
case "enqueue":
|
||||
s.sess.core.Queue(input)
|
||||
|
|
@ -125,36 +124,23 @@ func (s *SessionFileStore) Put(name string, data []byte) error {
|
|||
return fmt.Errorf("cannot switch backend while agent is running")
|
||||
}
|
||||
s.sess.core.Submit(s.sess.ctx, "/backend "+input, s.makePublish())
|
||||
s.sess.mu.Lock()
|
||||
s.sess.mutableVers["backend"].update(s.sess.core.BackendName())
|
||||
s.sess.mutableVers["model"].update(s.sess.core.ModelName())
|
||||
s.sess.mu.Unlock()
|
||||
|
||||
case "agent":
|
||||
if s.sess.core.IsRunning() {
|
||||
return fmt.Errorf("cannot switch agent while agent is running")
|
||||
}
|
||||
s.sess.core.Submit(s.sess.ctx, "/agent "+input, s.makePublish())
|
||||
s.sess.mu.Lock()
|
||||
s.sess.mutableVers["agent"].update(s.sess.core.AgentName())
|
||||
s.sess.mu.Unlock()
|
||||
|
||||
case "model":
|
||||
if s.sess.core.IsRunning() {
|
||||
return fmt.Errorf("cannot switch model while agent is running")
|
||||
}
|
||||
s.sess.core.Submit(s.sess.ctx, "/model "+input, s.makePublish())
|
||||
s.sess.mu.Lock()
|
||||
s.sess.mutableVers["model"].update(s.sess.core.ModelName())
|
||||
s.sess.mu.Unlock()
|
||||
|
||||
case "cwd":
|
||||
if err := s.sess.core.SetCWD(input); err != nil {
|
||||
return err
|
||||
}
|
||||
s.sess.mu.Lock()
|
||||
s.sess.mutableVers["cwd"].update(s.sess.core.CWD())
|
||||
s.sess.mu.Unlock()
|
||||
|
||||
case "params":
|
||||
if s.sess.core.IsRunning() {
|
||||
|
|
@ -296,7 +282,6 @@ func (s *SessionFileStore) makePublish() func(agent.Event) {
|
|||
s.sess.chatOffset = len(s.sess.chatLog)
|
||||
s.sess.mu.Unlock()
|
||||
}
|
||||
s.sess.trackMutable()
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
Reference in New Issue