Skip to content
Open
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
5 changes: 4 additions & 1 deletion .vscode/settings.json
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
{
"editor.aiStats.enabled": true,
"go.lintTool": "golangci-lint",
"go.lintOnSave": "package"
"go.lintOnSave": "package",
"git.scanRepositories": [
"common"
]
}
6 changes: 6 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,12 @@ image: vendor
coverage: ##@ Calculate test coverage percentage from coverage.out
@go tool cover -func=$(REPORTS_DIR)/coverage.out | grep total | awk '{print $$3}'

debug-setup: ##@ Set up local debug environment
##@ Generates go.work (Go version taken from go.mod) and symlinks the common module for local debugging
@GO_VERSION=$$(grep -m1 '^go ' go.mod | awk '{print $$2}') && \
printf 'go %s\n\nuse (\n\t.\n\t/opt/shared/common\n)\n' "$$GO_VERSION" > go.work
ln -sfn /opt/shared/common common

##@
##@ Misc commands
##@
Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ require (
github.com/fatih/color v1.18.0
github.com/google/shlex v0.0.0-20191202100458-e7afc7fbc510
github.com/gorilla/mux v1.8.1
github.com/k8shell-io/common v0.40.0
github.com/k8shell-io/common v0.45.0
github.com/k8shell-io/k8shell-go v0.2.3
github.com/pkg/sftp v1.13.10
github.com/rs/zerolog v1.34.0
Expand Down
16 changes: 4 additions & 12 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -42,18 +42,10 @@ github.com/gorilla/mux v1.8.1 h1:TuBL49tXwgrFYWhqrNgrUNEY92u81SPhu7sTdzQEiWY=
github.com/gorilla/mux v1.8.1/go.mod h1:AKf9I4AEqPTmMytcMc0KkNouC66V3BtZ4qD5fmWSiMQ=
github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw=
github.com/k8shell-io/common v0.36.0 h1:fkMH1XfYRLzDxqhIq5/luHusWWPGnCGJUXSTEIhEDzI=
github.com/k8shell-io/common v0.36.0/go.mod h1:E8dsb9ta4v3ne61AJgtRyTTbTkMMmKeCMAcXD+/9+cY=
github.com/k8shell-io/common v0.37.0 h1:whq66WosIJECKErUKZF1RQep7tdpOfI6GtP4hXREpsQ=
github.com/k8shell-io/common v0.37.0/go.mod h1:E8dsb9ta4v3ne61AJgtRyTTbTkMMmKeCMAcXD+/9+cY=
github.com/k8shell-io/common v0.39.0 h1:hfrKZYX2lBonornGrfK35rSY/+5o2sK+SakgUPKDL24=
github.com/k8shell-io/common v0.39.0/go.mod h1:E8dsb9ta4v3ne61AJgtRyTTbTkMMmKeCMAcXD+/9+cY=
github.com/k8shell-io/common v0.40.0 h1:MhQPVI5oe+JSdRJhWrtvCTJf0MaJDjM58KR7Nq3+lsM=
github.com/k8shell-io/common v0.40.0/go.mod h1:E8dsb9ta4v3ne61AJgtRyTTbTkMMmKeCMAcXD+/9+cY=
github.com/k8shell-io/k8shell-go v0.2.1 h1:6n88ijXkzP39//lIy4ai3XqtpSUXzoa/dVaWogHQYf4=
github.com/k8shell-io/k8shell-go v0.2.1/go.mod h1:j1JHgUIKIbaiRaitx6Pzw37ahqS4Hu9OcM4uvJ7BP4g=
github.com/k8shell-io/k8shell-go v0.2.2 h1:rwLOeIfyq1+l2Jyv0ak/lXZS7x6xbA5ye0yMsRkTonw=
github.com/k8shell-io/k8shell-go v0.2.2/go.mod h1:ZShnaWs7zxUlNwAkIn4lJodFqaB+PB8O+gn2EIscxq8=
github.com/k8shell-io/common v0.44.0 h1:cry62gDIXW8oyRWCzc2M+Hx62mpnarGdKZRTZ9TV8J8=
github.com/k8shell-io/common v0.44.0/go.mod h1:E8dsb9ta4v3ne61AJgtRyTTbTkMMmKeCMAcXD+/9+cY=
github.com/k8shell-io/common v0.45.0 h1:1f3QUIU1DhQ/HpBZVO/5BDcKYjf/EO4RKdOXZk33FzI=
github.com/k8shell-io/common v0.45.0/go.mod h1:E8dsb9ta4v3ne61AJgtRyTTbTkMMmKeCMAcXD+/9+cY=
github.com/k8shell-io/k8shell-go v0.2.3 h1:gL7dXDYN4EhWdQvvnY4B2On9Tpb1sZ7G5lO8RtI/nr4=
github.com/k8shell-io/k8shell-go v0.2.3/go.mod h1:wWb5gq693qqb48/p5iYrosLG4uNeGOr5dQJiOClIbE8=
github.com/kr/fs v0.1.0 h1:Jskdu9ieNAYnjxsi0LbQp1ulIKZV1LAFgK1tWhpZgl8=
Expand Down
59 changes: 59 additions & 0 deletions internal/grpc/acquire.go
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,65 @@ func (s *ShellHandler) AcquireSession(ctx context.Context, req *k8shelldv1.Acqui
}, nil
}

// ListSessions implements SshServiceServer.ListSessions.
// It returns the set of live PTY sessions that AcquireSession would currently
// accept: no client attached and no unexpired lock held, along with the OS
// user each session runs as.
func (s *ShellHandler) ListSessions(_ context.Context, _ *k8shelldv1.ListSessionsRequest) (*k8shelldv1.ListSessionsResponse, error) {
if !s.grpcApi.allowSessionDetach {
return nil, status.Errorf(codes.PermissionDenied, "session attachment is not enabled on this server")
}

now := time.Now()
locked := make(map[string]bool)
s.grpcApi.SessionLockStore.Range(func(_, v any) bool {
lk := v.(*sessionLock)
if now.Before(lk.expiresAt) {
locked[lk.sessionId] = true
}
return true
})

resp := &k8shelldv1.ListSessionsResponse{}
s.grpcApi.SessionStore.Range(func(_, value any) bool {
session, ok := value.(*SessionData)
if !ok || session.ptyDone == nil || !session.Deleted.IsZero() {
return true
}
select {
case <-session.ptyDone:
return true
default:
}

session.mu.Lock()
attached := session.attachedSender != nil
detachedAt := session.DetachedAt
session.mu.Unlock()

if attached || locked[session.Id] {
return true
}

var detachedAtStr string
if !detachedAt.IsZero() {
detachedAtStr = detachedAt.Format(timeFormat)
}

resp.Sessions = append(resp.Sessions, &k8shelldv1.AcquirableSession{
SessionId: session.Id,
Owner: session.user.Username,
CmdShell: session.CmdShell,
Pid: int32(session.Pid),
Created: session.Created.Format(timeFormat),
DetachedAt: detachedAtStr,
})
return true
})

return resp, nil
}

// cleanupExpiredLocks removes session locks that have passed their TTL.
func (a *GRPCService) cleanupExpiredLocks() {
now := time.Now()
Expand Down
4 changes: 4 additions & 0 deletions internal/grpc/ssh.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,10 @@ func (s *SshServiceServer) AcquireSession(ctx context.Context, req *k8shelldv1.A
return s.shell.AcquireSession(ctx, req)
}

func (s *SshServiceServer) ListSessions(ctx context.Context, req *k8shelldv1.ListSessionsRequest) (*k8shelldv1.ListSessionsResponse, error) {
return s.shell.ListSessions(ctx, req)
}

func (s *SshServiceServer) Exec(stream grpc.BidiStreamingServer[k8shelldv1.ExecRequest, k8shelldv1.ExecResponse]) error {
return s.exec.Exec(stream)
}
Expand Down
106 changes: 106 additions & 0 deletions internal/grpc/system.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,14 @@ import (
"google.golang.org/grpc/status"
)

// logStreamPollInterval is how often GetLogsStream polls the in-memory log
// store for new entries while following (mirrors the REST /logs handler).
const logStreamPollInterval = 100 * time.Millisecond

// defaultLogPageLimit is the page size GetLogsPage falls back to when the
// caller doesn't specify one (mirrors the REST /logs handler's default).
const defaultLogPageLimit = 100

// SystemServiceServer is the gRPC server for the system service
type SystemServiceServer struct {
grpcApi *GRPCService
Expand Down Expand Up @@ -124,3 +132,101 @@ func (s *SystemServiceServer) SystemInfo(ctx context.Context,

return k8shelld.SystemInfoToProto(&systemInfo), nil
}

// GetLogsStream streams k8shelld daemon logs (the same logs shown by
// `kbox logs`). With Follow=false it sends the currently buffered entries
// and closes the stream; with Follow=true it keeps streaming new entries as
// they are produced until the client cancels.
func (s *SystemServiceServer) GetLogsStream(req *k8shelldv1.SystemLogsStreamRequest,
stream k8shelldv1.SystemService_GetLogsStreamServer) error {

component := req.GetComponent()
level := k8shelld.LogLevelFromProto(req.GetLevel())
follow := req.GetFollow()

send := func(entries []logger.LogEntry) error {
for _, entry := range entries {
if sendErr := stream.Send(logEntryToProto(entry)); sendErr != nil {
return status.Errorf(codes.Canceled, "client canceled")
}
}
return nil
}

var backlog []logger.LogEntry
var sinceID int64
if n := req.GetLastN(); n > 0 {
backlog, _ = logger.GetLogsBefore(0, int(n), component, level)
if len(backlog) > 0 {
sinceID = backlog[len(backlog)-1].ID
}
} else {
backlog, sinceID = logger.GetLogsSince(0, component, level)
}
if err := send(backlog); err != nil {
return err
}

if !follow {
return nil
}

ctx := stream.Context()
for {
select {
case <-ctx.Done():
return status.Errorf(codes.Canceled, "client canceled")
default:
entries, newSinceID := logger.GetLogsSince(sinceID, component, level)
if err := send(entries); err != nil {
return err
}
sinceID = newSinceID

time.Sleep(logStreamPollInterval)
}
}
}

// GetLogsPage returns one page of k8shelld daemon logs strictly older than
// the requested BeforeId, for "load more" / infinite-scroll style backward
// pagination independent of GetLogsStream's live tail.
func (s *SystemServiceServer) GetLogsPage(ctx context.Context,
req *k8shelldv1.GetLogsPageRequest) (*k8shelldv1.GetLogsPageResponse, error) {

if req.GetBeforeId() < 0 {
return nil, status.Errorf(codes.InvalidArgument, "before_id must not be negative")
}
if req.GetLimit() < 0 {
return nil, status.Errorf(codes.InvalidArgument, "limit must not be negative")
}

limit := int(req.GetLimit())
if limit == 0 {
limit = defaultLogPageLimit
}

component := req.GetComponent()
level := k8shelld.LogLevelFromProto(req.GetLevel())

entries, hasMore := logger.GetLogsBefore(req.GetBeforeId(), limit, component, level)

resp := &k8shelldv1.GetLogsPageResponse{
Entries: make([]*k8shelldv1.SystemLogsStreamResponse, 0, len(entries)),
HasMore: hasMore,
}
for _, entry := range entries {
resp.Entries = append(resp.Entries, logEntryToProto(entry))
}
return resp, nil
}

func logEntryToProto(entry logger.LogEntry) *k8shelldv1.SystemLogsStreamResponse {
return &k8shelldv1.SystemLogsStreamResponse{
Id: entry.ID,
Time: entry.Timestamp,
Component: entry.Component,
Level: entry.Level,
Message: entry.Message,
}
}
82 changes: 66 additions & 16 deletions internal/logger/logger.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"encoding/json"
"fmt"
"io"
"sort"
"sync"

clogger "github.com/k8shell-io/common/pkg/logger"
Expand All @@ -17,19 +18,26 @@ const LOGSTORE_CAPACITY = 10000

// Shared memory log store
var logStore = &MemoryLogStore{
entries: make([]logEntry, 0, LOGSTORE_CAPACITY),
entries: make([]LogEntry, 0, LOGSTORE_CAPACITY),
cap: LOGSTORE_CAPACITY,
}

// MemoryLogStore is an in-memory log store that implements io.Writer
type MemoryLogStore struct {
mu sync.Mutex
entries []logEntry
entries []LogEntry
cap int
nextID int64
}

// logEntry represents a single log entry
type logEntry struct {
// LogEntry represents a single log entry.
//
// ID is a per-store, monotonically increasing sequence number assigned on
// write (starting at 1). Unlike a slice index it never shifts as older
// entries are evicted from the buffer, so it's safe to hold onto as a
// pagination cursor across calls.
type LogEntry struct {
ID int64 `json:"id"`
Timestamp string `json:"time"`
Component string `json:"component"`
Level string `json:"level"`
Expand All @@ -45,7 +53,7 @@ func NewLogger(component string) *zerolog.Logger {

// Write implements the io.Writer interface for MemoryLogStore
func (s *MemoryLogStore) Write(p []byte) (int, error) {
var entry logEntry
var entry LogEntry

if err := json.Unmarshal(p, &entry); err != nil {
return 0, fmt.Errorf("failed to unmarshal log entry: %w", err)
Expand All @@ -54,6 +62,9 @@ func (s *MemoryLogStore) Write(p []byte) (int, error) {
s.mu.Lock()
defer s.mu.Unlock()

s.nextID++
entry.ID = s.nextID

if len(s.entries) >= s.cap {
s.entries = s.entries[len(s.entries)-s.cap:]
}
Expand Down Expand Up @@ -82,28 +93,67 @@ func InitLogLevel(level string) error {
return nil
}

// GetLogsSince returns new log entries from the given offset.
func GetLogsSince(offset int, component, level string) ([]logEntry, int) {
// GetLogsSince returns log entries with ID > sinceID (in ID order), along
// with the highest entry ID currently in the store — pass that back in as
// sinceID on the next call to continue tailing without gaps or repeats,
// regardless of how many entries have since been evicted from the buffer.
// sinceID <= 0 returns every entry currently buffered.
func GetLogsSince(sinceID int64, component, level string) ([]LogEntry, int64) {
logStore.mu.Lock()
defer logStore.mu.Unlock()

if offset >= len(logStore.entries) {
return nil, len(logStore.entries)
lastID := sinceID
if n := len(logStore.entries); n > 0 && logStore.entries[n-1].ID > lastID {
lastID = logStore.entries[n-1].ID
}

if offset < 0 {
offset = len(logStore.entries) + offset
if offset < 0 {
offset = 0
var logs []LogEntry
for _, entry := range logStore.entries {
if entry.ID <= sinceID {
continue
}
if (component == "" || entry.Component == component) && (level == "" || entry.Level == level) {
logs = append(logs, entry)
}
}
return logs, lastID
}

var logs []logEntry
for i := offset; i < len(logStore.entries); i++ {
// GetLogsBefore returns up to limit log entries with ID < beforeID, oldest
// matching entry first — the "before" half of cursor-based pagination:
// pass the ID of the oldest entry from the previous page back in as
// beforeID to load the page before it. beforeID <= 0 starts from the most
// recent entry (there is no valid entry ID 0, since IDs are assigned
// starting at 1, so it doubles as the "no cursor yet" sentinel for the
// first page). The returned bool reports whether older entries remain
// unscanned in the buffer, i.e. whether a further "load more" call could
// return anything.
func GetLogsBefore(beforeID int64, limit int, component, level string) ([]LogEntry, bool) {
logStore.mu.Lock()
defer logStore.mu.Unlock()

end := len(logStore.entries)
if beforeID > 0 {
end = sort.Search(end, func(i int) bool {
return logStore.entries[i].ID >= beforeID
})
}

var logs []LogEntry
hasMore := false
for i := end - 1; i >= 0; i-- {
if len(logs) == limit {
hasMore = true
break
}
entry := logStore.entries[i]
if (component == "" || entry.Component == component) && (level == "" || entry.Level == level) {
logs = append(logs, entry)
}
}
return logs, len(logStore.entries)

for i, j := 0, len(logs)-1; i < j; i, j = i+1, j-1 {
logs[i], logs[j] = logs[j], logs[i]
}
return logs, hasMore
}
Loading
Loading