diff --git a/internal/process/manager.go b/internal/process/manager.go index cd91f85..50ec990 100644 --- a/internal/process/manager.go +++ b/internal/process/manager.go @@ -27,13 +27,20 @@ type Manager interface { } // LocalManager implements Manager using local OS processes. +// +// Concurrency model: +// - sessionsMu is an RWMutex. Read-only paths (Status, Subscribe lookup, +// Stdin lookup, ExitStatus lookup) take RLock; mutating paths (Start, +// Stop, Restart) take Lock. +// - exited is a sync.Map. waitExit is the only goroutine that calls +// sess.Cmd.Wait() and therefore the only one that writes here. Reads +// from anywhere (Start, Status, Subscribe, Restart) are lock-free. type LocalManager struct { - command string - args []string - mu sync.Mutex - sessions map[string]*Session - exitedMu sync.Mutex - exited map[string]struct{} + command string + args []string + sessionsMu sync.RWMutex + sessions map[string]*Session + exited sync.Map } // NewManager creates a LocalManager with the given command. @@ -43,23 +50,37 @@ func NewManager(command string, args []string) *LocalManager { command: normalizeCommand(command), args: args, sessions: make(map[string]*Session), - exited: make(map[string]struct{}), } } +// IsExited reports whether the process for the workspace has exited. +// Lock-free; safe to call from any goroutine. +func (m *LocalManager) IsExited(workspaceID string) bool { + _, exited := m.exited.Load(workspaceID) + return exited +} + +// MarkAsExited records that the process for the workspace has exited. +// Called only from waitExit after sess.Cmd.Wait() returns. +func (m *LocalManager) MarkAsExited(workspaceID string) { + m.exited.Store(workspaceID, struct{}{}) +} + +// ClearExited removes the exited marker for a workspace. Called when +// (re)starting a process so the next Status/Subscribe call won't see +// a stale exit. +func (m *LocalManager) ClearExited(workspaceID string) { + m.exited.Delete(workspaceID) +} + // Start launches a process for the given workspace. // Returns CodeConflict if a session is already running. func (m *LocalManager) Start(workspaceID string, workspaceRoot string) error { - m.mu.Lock() - defer m.mu.Unlock() + m.sessionsMu.Lock() + defer m.sessionsMu.Unlock() - if _, ok := m.sessions[workspaceID]; ok { - m.exitedMu.Lock() - _, exited := m.exited[workspaceID] - m.exitedMu.Unlock() - if !exited { - return util.New(util.CodeConflict, "process already running") - } + if _, ok := m.sessions[workspaceID]; ok && !m.IsExited(workspaceID) { + return util.New(util.CodeConflict, "process already running") } cmd := exec.Command(m.command, m.args...) @@ -101,10 +122,7 @@ func (m *LocalManager) Start(workspaceID string, workspaceRoot string) error { Subscribers: make(map[Subscription]struct{}), } m.sessions[workspaceID] = sess - - m.exitedMu.Lock() - delete(m.exited, workspaceID) - m.exitedMu.Unlock() + m.ClearExited(workspaceID) outputDone := make(chan struct{}) go m.captureOutput(sess, stdoutR, outputDone) @@ -144,13 +162,15 @@ func (m *LocalManager) captureOutput(sess *Session, stdoutR *os.File, done chan< // waitExit waits for the process to exit, then records the exit status, closes // the stdin writer, and closes every subscriber channel exactly once. +// +// This is the only goroutine that calls sess.Cmd.Wait() and the only one +// that writes to m.exited. All other code paths use m.IsExited instead of +// reading sess.Cmd.ProcessState directly (which would race with Wait). func (m *LocalManager) waitExit(sess *Session, outputDone <-chan struct{}) { _ = sess.Cmd.Wait() <-outputDone - m.exitedMu.Lock() - m.exited[sess.WorkspaceID] = struct{}{} - m.exitedMu.Unlock() + m.MarkAsExited(sess.WorkspaceID) sess.mu.Lock() defer sess.mu.Unlock() @@ -174,8 +194,8 @@ func (m *LocalManager) waitExit(sess *Session, outputDone <-chan struct{}) { // Stop kills the process for the given workspace. // Returns CodeNotFound if no session exists. func (m *LocalManager) Stop(workspaceID string) error { - m.mu.Lock() - defer m.mu.Unlock() + m.sessionsMu.Lock() + defer m.sessionsMu.Unlock() sess, ok := m.sessions[workspaceID] if !ok || sess.Cmd.Process == nil { @@ -187,46 +207,36 @@ func (m *LocalManager) Stop(workspaceID string) error { } delete(m.sessions, workspaceID) - m.exitedMu.Lock() - delete(m.exited, workspaceID) - m.exitedMu.Unlock() + m.ClearExited(workspaceID) return nil } // Restart stops (if running) then starts the process. func (m *LocalManager) Restart(workspaceID string, workspaceRoot string) error { - m.mu.Lock() + m.sessionsMu.Lock() if sess, ok := m.sessions[workspaceID]; ok { - m.exitedMu.Lock() - _, exited := m.exited[workspaceID] - m.exitedMu.Unlock() - if !exited && sess.Cmd.Process != nil { + if !m.IsExited(workspaceID) && sess.Cmd.Process != nil { _ = sess.Cmd.Process.Kill() } delete(m.sessions, workspaceID) - m.exitedMu.Lock() - delete(m.exited, workspaceID) - m.exitedMu.Unlock() + m.ClearExited(workspaceID) } - m.mu.Unlock() + m.sessionsMu.Unlock() return m.Start(workspaceID, workspaceRoot) } // Status returns the current process status for the workspace. func (m *LocalManager) Status(workspaceID string) Status { - m.mu.Lock() + m.sessionsMu.RLock() sess, ok := m.sessions[workspaceID] - m.mu.Unlock() + m.sessionsMu.RUnlock() if !ok || sess.Cmd.Process == nil { return Status{WorkspaceID: workspaceID, Running: false} } - m.exitedMu.Lock() - _, exited := m.exited[workspaceID] - m.exitedMu.Unlock() - if exited { + if m.IsExited(workspaceID) { return Status{WorkspaceID: workspaceID, Running: false} } @@ -240,9 +250,9 @@ func (m *LocalManager) Status(workspaceID string) Status { // Subscribe creates a new output subscription for the workspace. // Returns CodeNotFound if the workspace has no session. func (m *LocalManager) Subscribe(workspaceID string) (Subscription, error) { - m.mu.Lock() + m.sessionsMu.RLock() sess, ok := m.sessions[workspaceID] - m.mu.Unlock() + m.sessionsMu.RUnlock() if !ok { return nil, util.New(util.CodeNotFound, "no running process") @@ -253,10 +263,7 @@ func (m *LocalManager) Subscribe(workspaceID string) (Subscription, error) { sess.Subscribers[sub] = struct{}{} sess.mu.Unlock() - m.exitedMu.Lock() - _, exited := m.exited[workspaceID] - m.exitedMu.Unlock() - if exited { + if m.IsExited(workspaceID) { sub.closeChan() } @@ -275,8 +282,8 @@ func (m *LocalManager) ExitStatus(workspaceID string) (ExitInfo, error) { // stdinOf returns the session's stdin writer or CodeNotFound. func (m *LocalManager) stdinOf(workspaceID string) (io.WriteCloser, error) { - m.mu.Lock() - defer m.mu.Unlock() + m.sessionsMu.RLock() + defer m.sessionsMu.RUnlock() sess, ok := m.sessions[workspaceID] if !ok { return nil, util.New(util.CodeNotFound, "no running process") @@ -286,9 +293,9 @@ func (m *LocalManager) stdinOf(workspaceID string) (io.WriteCloser, error) { // exitOf returns the session's recorded exit info or CodeNotFound. func (m *LocalManager) exitOf(workspaceID string) (ExitInfo, error) { - m.mu.Lock() + m.sessionsMu.RLock() sess, ok := m.sessions[workspaceID] - m.mu.Unlock() + m.sessionsMu.RUnlock() if !ok { return ExitInfo{}, util.New(util.CodeNotFound, "no running process") }