Skip to content
Closed
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
7 changes: 4 additions & 3 deletions services/nvpair-cluster-manager/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ package main
import "nvpair-shared/jsonrpc"

type (
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
Codec = jsonrpc.Codec
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
DecodeError = jsonrpc.DecodeError
Codec = jsonrpc.Codec
)

var NewCodec = jsonrpc.NewCodec
12 changes: 10 additions & 2 deletions services/nvpair-cluster-manager/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ package main
import (
"context"
"encoding/json"
stderrors "errors"
"fmt"
"io"
"log"
Expand Down Expand Up @@ -372,8 +373,15 @@ func (m *Manager) readLoop(ctx context.Context) error {
if err == io.EOF || ctx.Err() != nil {
return nil
}
log.Printf("JSON-RPC read error: %v", err)
continue
var de *DecodeError
if stderrors.As(err, &de) {
log.Printf("JSON-RPC decode error (skipping frame): %v", err)
continue
}
// Terminal transport/scanner error (e.g. an over-long frame —
// bufio.Scanner cannot resync) — stop instead of spinning.
log.Printf("JSON-RPC read error (terminal): %v", err)
return err
}
m.handleMessage(msg)
if ctx.Err() != nil {
Expand Down
7 changes: 4 additions & 3 deletions services/nvpair-errors/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ package main
import "nvpair-shared/jsonrpc"

type (
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
Codec = jsonrpc.Codec
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
DecodeError = jsonrpc.DecodeError
Codec = jsonrpc.Codec
)

var NewCodec = jsonrpc.NewCodec
12 changes: 10 additions & 2 deletions services/nvpair-errors/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ package main
import (
"context"
"encoding/json"
stderrors "errors"
"fmt"
"io"
"log"
Expand Down Expand Up @@ -176,8 +177,15 @@ func (m *Manager) readLoop(ctx context.Context) error {
if err == io.EOF || ctx.Err() != nil {
return nil
}
log.Printf("JSON-RPC read error: %v", err)
continue
var de *DecodeError
if stderrors.As(err, &de) {
log.Printf("JSON-RPC decode error (skipping frame): %v", err)
continue
}
// Terminal transport/scanner error (e.g. an over-long frame —
// bufio.Scanner cannot resync) — stop instead of spinning.
log.Printf("JSON-RPC read error (terminal): %v", err)
return err
}
m.handleMessage(msg)
if ctx.Err() != nil {
Expand Down
7 changes: 4 additions & 3 deletions services/nvpair-job-scheduler/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ package main
import "nvpair-shared/jsonrpc"

type (
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
Codec = jsonrpc.Codec
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
DecodeError = jsonrpc.DecodeError
Codec = jsonrpc.Codec
)

var NewCodec = jsonrpc.NewCodec
12 changes: 10 additions & 2 deletions services/nvpair-job-scheduler/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
Expand Down Expand Up @@ -98,8 +99,15 @@ func (m *Manager) readLoop(ctx context.Context) error {
if err == io.EOF || ctx.Err() != nil {
return nil
}
log.Printf("JSON-RPC read error: %v", err)
continue
var de *DecodeError
if errors.As(err, &de) {
log.Printf("JSON-RPC decode error (skipping frame): %v", err)
continue
}
// Terminal transport/scanner error (e.g. an over-long frame —
// bufio.Scanner cannot resync) — stop instead of spinning.
log.Printf("JSON-RPC read error (terminal): %v", err)
return err
}
m.handleMessage(msg)
if ctx.Err() != nil {
Expand Down
7 changes: 4 additions & 3 deletions services/nvpair-manual-nodes/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ package main
import "nvpair-shared/jsonrpc"

type (
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
Codec = jsonrpc.Codec
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
DecodeError = jsonrpc.DecodeError
Codec = jsonrpc.Codec
)

var NewCodec = jsonrpc.NewCodec
12 changes: 10 additions & 2 deletions services/nvpair-manual-nodes/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ package main
import (
"context"
"encoding/json"
stderrors "errors"
"fmt"
"io"
"log"
Expand Down Expand Up @@ -606,8 +607,15 @@ func (m *Manager) readLoop(ctx context.Context) error {
if err == io.EOF || ctx.Err() != nil {
return nil
}
log.Printf("JSON-RPC read error: %v", err)
continue
var de *DecodeError
if stderrors.As(err, &de) {
log.Printf("JSON-RPC decode error (skipping frame): %v", err)
continue
}
// Terminal transport/scanner error (e.g. an over-long frame —
// bufio.Scanner cannot resync) — stop instead of spinning.
log.Printf("JSON-RPC read error (terminal): %v", err)
return err
}
m.handleMessage(msg)
if ctx.Err() != nil {
Expand Down
7 changes: 4 additions & 3 deletions services/nvpair-node-scanner/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ package main
import "nvpair-shared/jsonrpc"

type (
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
Codec = jsonrpc.Codec
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
DecodeError = jsonrpc.DecodeError
Codec = jsonrpc.Codec
)

var NewCodec = jsonrpc.NewCodec
12 changes: 10 additions & 2 deletions services/nvpair-node-scanner/scanner.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ package main

import (
"context"
"errors"
"fmt"
"io"
"log"
Expand Down Expand Up @@ -101,8 +102,15 @@ func (s *Scanner) readLoop(ctx context.Context) error {
if err == io.EOF || ctx.Err() != nil {
return nil
}
log.Printf("JSON-RPC read error: %v", err)
continue
var de *DecodeError
if errors.As(err, &de) {
log.Printf("JSON-RPC decode error (skipping frame): %v", err)
continue
}
// Terminal transport/scanner error (e.g. an over-long frame —
// bufio.Scanner cannot resync) — stop instead of spinning.
log.Printf("JSON-RPC read error (terminal): %v", err)
return err
}
s.handleMessage(msg)
}
Expand Down
7 changes: 4 additions & 3 deletions services/nvpair-node-settings/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ package main
import "nvpair-shared/jsonrpc"

type (
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
Codec = jsonrpc.Codec
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
DecodeError = jsonrpc.DecodeError
Codec = jsonrpc.Codec
)

var NewCodec = jsonrpc.NewCodec
11 changes: 9 additions & 2 deletions services/nvpair-node-settings/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -283,8 +283,15 @@ func (m *Manager) readLoop(ctx context.Context) error {
if err == io.EOF || ctx.Err() != nil {
return nil
}
log.Printf("JSON-RPC read error: %v", err)
continue
var de *DecodeError
if errors.As(err, &de) {
log.Printf("JSON-RPC decode error (skipping frame): %v", err)
continue
}
// Terminal transport/scanner error (e.g. an over-long frame —
// bufio.Scanner cannot resync) — stop instead of spinning.
log.Printf("JSON-RPC read error (terminal): %v", err)
return err
}
m.handleMessage(msg)
if ctx.Err() != nil {
Expand Down
7 changes: 4 additions & 3 deletions services/nvpair-proxy/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ package main
import "nvpair-shared/jsonrpc"

type (
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
Codec = jsonrpc.Codec
Message = jsonrpc.Message
RPCError = jsonrpc.RPCError
DecodeError = jsonrpc.DecodeError
Codec = jsonrpc.Codec
)

var NewCodec = jsonrpc.NewCodec
11 changes: 9 additions & 2 deletions services/nvpair-proxy/proxy.go
Original file line number Diff line number Diff line change
Expand Up @@ -2553,8 +2553,15 @@ func (p *Proxy) readLoop(ctx context.Context) error {
if err == io.EOF || ctx.Err() != nil {
return nil
}
log.Printf("JSON-RPC read error: %v", err)
continue
var de *DecodeError
if stderrors.As(err, &de) {
log.Printf("JSON-RPC decode error (skipping frame): %v", err)
continue
}
// Terminal transport/scanner error (e.g. an over-long frame —
// bufio.Scanner cannot resync) — stop instead of spinning.
log.Printf("JSON-RPC read error (terminal): %v", err)
return err
}
p.handleMessage(msg)
}
Expand Down
11 changes: 9 additions & 2 deletions services/nvpair-tui/rpc/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ package rpc
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"strconv"
Expand Down Expand Up @@ -68,9 +69,15 @@ func (c *Client) Run(ctx context.Context) error {
if err == io.EOF {
return nil
}
// A single malformed line should not kill the session; the
// A single malformed frame should not kill the session; the
// broker may emit a frame we don't model. Skip and continue.
continue
// Anything else (over-long frame, transport error) is terminal:
// bufio.Scanner cannot resync, so continuing would spin.
var de *DecodeError
if errors.As(err, &de) {
continue
}
return err
}
switch {
case msg.IsResponse():
Expand Down
16 changes: 13 additions & 3 deletions services/nvpair-tui/rpc/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,14 @@ type Codec struct {
wmu sync.Mutex
}

// DecodeError marks a recoverable per-frame failure (bad JSON or wrong
// version): the stream is still positioned at the next line, so callers
// may skip the frame and continue reading. Mirrors nvpair-shared/jsonrpc.
type DecodeError struct{ Err error }

func (e *DecodeError) Error() string { return e.Err.Error() }
func (e *DecodeError) Unwrap() error { return e.Err }

// NewCodec wraps a reader/writer pair (the broker's stdout/stdin) in a
// framing codec.
func NewCodec(r io.Reader, w io.Writer) *Codec {
Expand All @@ -79,7 +87,9 @@ func NewCodec(r io.Reader, w io.Writer) *Codec {
return &Codec{scanner: scanner, writer: w}
}

// Read returns the next frame, or io.EOF when the stream closes.
// Read returns the next frame, or io.EOF when the stream closes. A
// malformed frame (bad JSON or wrong version) is returned as a recoverable
// *DecodeError; a terminal scanner/transport error is a plain error.
func (c *Codec) Read() (*Message, error) {
if !c.scanner.Scan() {
if err := c.scanner.Err(); err != nil {
Expand All @@ -89,10 +99,10 @@ func (c *Codec) Read() (*Message, error) {
}
var msg Message
if err := json.Unmarshal(c.scanner.Bytes(), &msg); err != nil {
return nil, fmt.Errorf("invalid JSON-RPC message: %w", err)
return nil, &DecodeError{fmt.Errorf("invalid JSON-RPC message: %w", err)}
}
if msg.JSONRPC != "2.0" {
return nil, fmt.Errorf("unsupported JSON-RPC version: %q", msg.JSONRPC)
return nil, &DecodeError{fmt.Errorf("unsupported JSON-RPC version: %q", msg.JSONRPC)}
}
return &msg, nil
}
Expand Down
37 changes: 31 additions & 6 deletions services/nvpair-ui-broker/broker.go
Original file line number Diff line number Diff line change
Expand Up @@ -2585,6 +2585,20 @@ func setNodeIDIfEmpty(m map[string]json.RawMessage, key, nodeID string) bool {
return true
}

// errTerminalRead reports that the client stdin read loop ended on a
// non-recoverable scanner/transport error (e.g. an over-long frame that
// bufio.Scanner cannot resync past), distinct from a clean EOF.
var errTerminalRead = stderrors.New("terminal read error")

// recoverableDecode reports whether a codec Read error is a recoverable
// per-frame decode failure (bad JSON / wrong version): the scanner advances
// past the bad frame, so both the producer and the consumer keep pumping
// instead of tearing the connection down.
func recoverableDecode(err error) bool {
var de *DecodeError
return stderrors.As(err, &de)
}

func (b *Broker) readLoop(ctx context.Context) error {
// codec.Read() blocks on stdin, so we run it on its own goroutine and
// select against ctx.Done(). Otherwise a SIGINT/SIGTERM (which cancels
Expand All @@ -2604,11 +2618,18 @@ func (b *Broker) readLoop(ctx context.Context) error {
case <-ctx.Done():
return
}
// EOF is terminal (stream closed); other errors are per-line
// (e.g. a bad JSON frame) and the next Read advances past them.
if err == io.EOF {
return
// A decoded message (err nil) and a recoverable decode error
// (bad frame; the next Read advances past it) both keep the pump
// running. EOF is terminal (stream closed), and any other error
// is a terminal scanner/transport error: stop feeding the
// channel so the consumer exits instead of spinning.
if err == nil {
continue
}
if recoverableDecode(err) {
continue
}
return
}
}()

Expand All @@ -2621,8 +2642,12 @@ func (b *Broker) readLoop(ctx context.Context) error {
if r.err == io.EOF || ctx.Err() != nil {
return nil
}
slog.Warn("JSON-RPC read error", "err", r.err)
continue
if recoverableDecode(r.err) {
slog.Warn("JSON-RPC decode error (skipping frame)", "err", r.err)
continue
}
slog.Warn("JSON-RPC read error (terminal)", "err", r.err)
return errTerminalRead
}
b.handleMessage(r.msg)
if ctx.Err() != nil {
Expand Down
Loading