Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
62 changes: 27 additions & 35 deletions meshsync/exec.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,28 +72,26 @@ func (h *Handler) processExecRequest(obj interface{}, cfg config.ListenerConfig)

for _, req := range reqs {
id := fmt.Sprintf("exec.%s.%s.%s.%s", req.Namespace, req.Name, req.Container, req.ID)
if _, ok := h.channelPool[id]; !ok {
// Subscribing the first time
if !bool(req.Stop) {
h.channelPool[id] = channels.NewStructChannel()
h.Log.Debug("Starting session")

err := h.Broker.Publish("active_sessions.exec", &broker.Message{
ObjectType: broker.ActiveExecObject,
Object: h.getActiveChannels(),
})
if err != nil {
h.Log.Error(ErrGetObject(err))
}
go h.streamSession(id, req, cfg)
}
} else {
// Already running subscription
if bool(req.Stop) {
// TODO: once we have a unsubscribe functionality, need to publish message to active sessions subject
execCleanup(h, id)
}
if bool(req.Stop) {
// Stop request: tear down the session if one is running (no-op otherwise).
// TODO: once we have unsubscribe functionality, publish to active sessions subject.
execCleanup(h, id)
continue
}
if _, created := h.addSession(id); !created {
// A session for this id is already running.
continue
}
h.Log.Debug("Starting session")

err := h.Broker.Publish("active_sessions.exec", &broker.Message{
ObjectType: broker.ActiveExecObject,
Object: h.getActiveChannels(),
})
if err != nil {
h.Log.Error(ErrGetObject(err))
}
go h.streamSession(id, req, cfg)
}

return nil
Expand All @@ -104,9 +102,10 @@ func (h *Handler) processActiveExecRequest() error {
return nil
}
func (h *Handler) getActiveChannels() []*string {
activeChannels := make([]*string, 0, len(h.channelPool))
for k := range h.channelPool {
activeChannels = append(activeChannels, &k)
ids := h.activeSessionIDs()
activeChannels := make([]*string, 0, len(ids))
for i := range ids {
activeChannels = append(activeChannels, &ids[i])
}

return activeChannels
Expand Down Expand Up @@ -164,7 +163,7 @@ func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.Lis
_ = tstdin.Close()
_ = stdout.Close()
_ = getStdout.Close()
delete(h.channelPool, id)
h.deleteSession(id)
})
}
defer terminate()
Expand Down Expand Up @@ -263,10 +262,8 @@ func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.Lis
}()

for {
// The session's StructChannel is asserted below; once terminate() has
// removed the pool entry that assertion would be on a nil interface and
// panic, so bail out first if the session is already gone.
sessionCh, ok := h.channelPool[id].(channels.StructChannel)
// If terminate() has already removed the session, bail out.
sessionCh, ok := h.getSession(id)
if !ok {
h.Log.Debugf("Session closed for: %s", id)
return
Expand All @@ -293,12 +290,7 @@ func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.Lis
}

func execCleanup(h *Handler, id string) {
ch, ok := h.channelPool[id]
if !ok {
return
}

structChan, ok := ch.(channels.StructChannel)
structChan, ok := h.getSession(id)
if !ok {
return
}
Expand Down
78 changes: 53 additions & 25 deletions meshsync/logstream.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,24 +27,31 @@ func (h *Handler) processLogRequest(obj interface{}, cfg config.ListenerConfig)

for _, req := range reqs {
id := fmt.Sprintf("logs.%s.%s.%s", req.Namespace, req.Name, req.Container)
if _, ok := h.channelPool[id]; !ok {
// Subscribing the first time
if !bool(req.Stop) {
h.channelPool[id] = channels.NewStructChannel()
go h.streamLogs(id, req, cfg)
}
} else {
// Already running subscription
if bool(req.Stop) {
h.channelPool[id].(channels.StructChannel) <- struct{}{}
if bool(req.Stop) {
// Stop request: signal the running stream, if any, to close.
// Non-blocking: the stream may already be closing and no longer
// receiving, so a plain send could freeze this loop.
if ch, ok := h.getSession(id); ok {
select {
case ch <- struct{}{}:
default:
}
}
Comment thread
leecalcote marked this conversation as resolved.
continue
}
if _, created := h.addSession(id); created {
go h.streamLogs(id, req, cfg)
}
}

return nil
}

func (h *Handler) streamLogs(id string, req model.LogRequest, cfg config.ListenerConfig) {
// Remove the session however this function exits (stream-open failure, EOF,
// read error, or a stop request).
defer h.deleteSession(id)

resp, err := h.kubeClient.KubeClient.CoreV1().Pods(req.Namespace).GetLogs(req.Name, &v1.PodLogOptions{
Container: req.Container,
Follow: req.Follow,
Expand All @@ -58,33 +65,35 @@ func (h *Handler) streamLogs(id string, req model.LogRequest, cfg config.Listene
}).Stream(context.TODO())
if err != nil {
h.Log.Error(ErrLogStream(err))
delete(h.channelPool, id)
return
}
defer resp.Close()

go func() {
<-h.channelPool[id].(channels.StructChannel)
h.Log.Debugf("Closing %s", id)
delete(h.channelPool, id)
resp.Close()
}()
// done unblocks the waiter goroutine when the stream ends on its own (EOF or
// read error), so it is not leaked once streamLogs returns.
done := make(chan struct{})
defer close(done)

go h.awaitLogStreamStop(id, resp, done)

for {
buf := make([]byte, 2000)
numBytes, err := resp.Read(buf)
if err == io.EOF {
break
}
if numBytes == 0 {
continue
}
if err != nil {
// A non-EOF read error ends the stream; breaking avoids an infinite
// read/log loop on a failed stream.
h.Log.Error(ErrCopyBuffer(err))
delete(h.channelPool, id)
break
}
if numBytes == 0 {
continue
}
Comment thread
leecalcote marked this conversation as resolved.

message := string(buf[:numBytes])
err = h.Broker.Publish(cfg.PublishTo, &broker.Message{
if pubErr := h.Broker.Publish(cfg.PublishTo, &broker.Message{
ObjectType: broker.LogStreamObject,
EventType: broker.Add,
Object: &model.LogObject{
Expand All @@ -93,10 +102,29 @@ func (h *Handler) streamLogs(id string, req model.LogRequest, cfg config.Listene
Primary: req.Name,
Secondary: req.Container,
},
})
if err != nil {
h.Log.Error(ErrCopyBuffer(err))
}); pubErr != nil {
h.Log.Error(ErrCopyBuffer(pubErr))
}
}
}

// awaitLogStreamStop closes resp when the session receives an explicit stop
// signal or the process is shutting down (channels.Stop), and returns without
// closing when the stream has already ended on its own (done is closed by
// streamLogs). channelPool holds only the fixed system channels and is
// read-only after init, so reading channels.Stop from it needs no lock.
func (h *Handler) awaitLogStreamStop(id string, resp io.ReadCloser, done <-chan struct{}) {
ch, ok := h.getSession(id)
if !ok {
return
}
select {
case <-ch:
h.Log.Debugf("Closing %s", id)
resp.Close()
case <-h.channelPool[channels.Stop].(channels.StopChannel):
h.Log.Debugf("Stopping session %s on global stop", id)
resp.Close()
case <-done:
}
}
14 changes: 11 additions & 3 deletions meshsync/meshsync.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
package meshsync

import (
"sync"

"github.com/meshery/meshkit/broker"
"github.com/meshery/meshkit/config"
"github.com/meshery/meshkit/logger"
Expand All @@ -22,10 +24,15 @@ type Handler struct {
Log logger.Handler
Broker broker.Handler

clusterID string
informer dynamicinformer.DynamicSharedInformerFactory
kubeClient *mesherykube.Client
clusterID string
informer dynamicinformer.DynamicSharedInformerFactory
kubeClient *mesherykube.Client
// channelPool holds the fixed system channels (Stop/OS/ReSync) and is
// read-only after construction. Dynamic exec/log-stream sessions live in
// sessions (guarded by sessionsMu), not here, so the two never race.
channelPool map[string]channels.GenericChannel
sessions map[string]channels.StructChannel
sessionsMu sync.Mutex
stores map[string]cache.Store
outputWriter output.Writer
outputFiltration internalconfig.OutputFiltrationContainer
Expand Down Expand Up @@ -65,6 +72,7 @@ func New(
kubeClient: kubeClient,
clusterID: clusterID,
channelPool: pool,
sessions: make(map[string]channels.StructChannel),
outputFiltration: outputFiltration,
}, nil
}
Expand Down
65 changes: 65 additions & 0 deletions meshsync/sessions.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
package meshsync

import (
"github.com/meshery/meshsync/internal/channels"
)

// Interactive exec and log-stream sessions are keyed by a per-request id.
// They were previously stored in the shared channelPool alongside the fixed
// system channels (Stop/OS/ReSync), which meant session goroutines mutated the
// same map that other goroutines read (system-channel selects, getActiveChannels),
// a data race that can panic the process. They now live in their own
// mutex-guarded map so channelPool stays read-only after initialization.
//
// Every helper returns the channel (if any) and releases the lock before the
// caller performs any channel send/receive, so the sessions mutex is never held
// across a blocking channel operation.

// addSession registers a new session channel for id and returns it with
// created=true. If a session already exists for id, the existing channel is
// returned with created=false so the caller does not start a duplicate.
func (h *Handler) addSession(id string) (ch channels.StructChannel, created bool) {
h.sessionsMu.Lock()
defer h.sessionsMu.Unlock()
// Defensive: a Handler constructed outside New has a nil map; writing to it
// would panic.
if h.sessions == nil {
h.sessions = make(map[string]channels.StructChannel)
}
if existing, ok := h.sessions[id]; ok {
return existing, false
}
// 1-buffered: exec/log-stream stop signals are sent non-blocking, so a stop
// that arrives before the session goroutine begins receiving must still be
// recorded rather than dropped (which would leave the session running).
ch = make(channels.StructChannel, 1)
h.sessions[id] = ch
return ch, true
}

// getSession returns the session channel for id, if present.
func (h *Handler) getSession(id string) (channels.StructChannel, bool) {
h.sessionsMu.Lock()
defer h.sessionsMu.Unlock()
ch, ok := h.sessions[id]
return ch, ok
}

// deleteSession removes the session for id. It is safe to call for an id that
// is not present.
func (h *Handler) deleteSession(id string) {
h.sessionsMu.Lock()
defer h.sessionsMu.Unlock()
delete(h.sessions, id)
}

// activeSessionIDs returns the ids of the currently active sessions.
func (h *Handler) activeSessionIDs() []string {
h.sessionsMu.Lock()
defer h.sessionsMu.Unlock()
ids := make([]string, 0, len(h.sessions))
for id := range h.sessions {
ids = append(ids, id)
}
return ids
}
Loading
Loading