9p: snapshot *wait baseline at open time to eliminate inter-read race
The previous design read the current value at read time. If a state change fired between one cat returning and the next cat opening, the new cat would see the changed value as current and wait forever for the next change. Fix: snapshot the baseline in the fid at Topen time. Each read blocks until the value changes from the open-time snapshot, then updates the fid's baseline for the next read. The race window is closed. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
a6b6e911dd
commit
2bc30d275c
|
|
@ -91,6 +91,7 @@ type fid struct {
|
|||
qid plan9.Qid
|
||||
mode uint8
|
||||
writeBuf []byte
|
||||
waitBase string // for *wait files: value snapshotted at open time
|
||||
}
|
||||
|
||||
// connState tracks all open fids for a single 9P connection.
|
||||
|
|
@ -515,6 +516,17 @@ func (s *Server) open(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
}
|
||||
|
||||
f.mode = fc.Mode
|
||||
// Snapshot the current value for *wait files at open time so that the
|
||||
// baseline is fixed before any read — eliminating the race between a
|
||||
// prior read returning and the next read opening.
|
||||
if strings.HasSuffix(pathBase(f.path), "wait") {
|
||||
parts := strings.SplitN(strings.TrimPrefix(f.path, "/"), "/", 3)
|
||||
if len(parts) == 3 {
|
||||
if store, ok := s.sessionFileStore(parts[1]); ok {
|
||||
f.waitBase = store.CurrentWaitValue(parts[2])
|
||||
}
|
||||
}
|
||||
}
|
||||
plog.Debug("Topen fid=%d path=%q mode=%d", fc.Fid, f.path, fc.Mode)
|
||||
return &plan9.Fcall{Type: plan9.Ropen, Tag: fc.Tag, Qid: f.qid}
|
||||
}
|
||||
|
|
@ -751,10 +763,25 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
}
|
||||
// *wait files block until a value changes; use connection context.
|
||||
if strings.HasSuffix(parts[2], "wait") {
|
||||
content, err := store.Wait(cs.ctx, parts[2])
|
||||
cs.mu.RLock()
|
||||
f, fidOK := cs.fids[fc.Fid]
|
||||
var base string
|
||||
if fidOK {
|
||||
base = f.waitBase
|
||||
}
|
||||
cs.mu.RUnlock()
|
||||
content, err := store.Wait(cs.ctx, parts[2], base)
|
||||
if err != nil {
|
||||
return errFcall(fc, err.Error())
|
||||
}
|
||||
// Update the fid's baseline to the returned value for subsequent reads.
|
||||
if content != nil {
|
||||
cs.mu.Lock()
|
||||
if f, ok := cs.fids[fc.Fid]; ok {
|
||||
f.waitBase = strings.TrimSuffix(string(content), "\n")
|
||||
}
|
||||
cs.mu.Unlock()
|
||||
}
|
||||
return s.readSlice(fc, content)
|
||||
}
|
||||
content, err := store.Get(parts[2])
|
||||
|
|
|
|||
|
|
@ -272,30 +272,49 @@ 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) {
|
||||
// CurrentWaitValue returns the current value for the named *wait file.
|
||||
func (s *SessionFileStore) CurrentWaitValue(name string) string {
|
||||
switch name {
|
||||
case "statewait":
|
||||
return s.sess.core.State()
|
||||
case "usagewait":
|
||||
return s.sess.core.Usage()
|
||||
case "ctxszwait":
|
||||
return s.sess.core.CtxSz()
|
||||
case "cwdwait":
|
||||
return s.sess.core.CWD()
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// Wait blocks until the named *wait file's underlying value changes from base,
|
||||
// then returns the new value. If base is empty, the current value is used.
|
||||
// Unblocks when connCtx or the session context is cancelled.
|
||||
func (s *SessionFileStore) Wait(connCtx context.Context, name, base 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
|
||||
var field string
|
||||
switch name {
|
||||
case "statewait":
|
||||
field, current = agent.WatchState, s.sess.core.State()
|
||||
field = agent.WatchState
|
||||
case "usagewait":
|
||||
field, current = agent.WatchUsage, s.sess.core.Usage()
|
||||
field = agent.WatchUsage
|
||||
case "ctxszwait":
|
||||
field, current = agent.WatchCtxSz, s.sess.core.CtxSz()
|
||||
field = agent.WatchCtxSz
|
||||
case "cwdwait":
|
||||
field, current = agent.WatchCWD, s.sess.core.CWD()
|
||||
field = agent.WatchCWD
|
||||
default:
|
||||
return nil, fmt.Errorf("%s: not a wait file", name)
|
||||
}
|
||||
|
||||
v, ok := s.sess.core.WaitChange(ctx, field, current)
|
||||
if base == "" {
|
||||
base = s.CurrentWaitValue(name)
|
||||
}
|
||||
|
||||
v, ok := s.sess.core.WaitChange(ctx, field, base)
|
||||
if !ok {
|
||||
return nil, nil
|
||||
}
|
||||
|
|
|
|||
Reference in New Issue