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
12 changes: 12 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,18 @@ Published to 34.126.161.115:33211 (Singapore) in 259ms
mump2p publish --topic test/data --file ./payload.json
```

### From stdin

When stdin is a pipe or redirect, its contents are published. Non-empty stdin takes precedence over `--message`; an empty pipe falls back to `--message`. This makes `publish` work like any other Unix tool in a pipe:

```bash
echo "Hello World" | mump2p publish --topic test
cat ./payload.json | mump2p publish --topic test/data
curl -s https://api.example.com/status | mump2p publish --topic status
```

`--file -` also reads from stdin explicitly. Stdin is capped at your account's maximum message size.

## Debug Mode

Use `--debug` to see session details, node scores, timing breakdowns, message IDs, and peer paths.
Expand Down
137 changes: 121 additions & 16 deletions cmd/publish.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"encoding/json"
"errors"
"fmt"
"io"
"os"
"time"

Expand Down Expand Up @@ -64,20 +65,127 @@ func shortMsgID(resp *pb.Response) string {
return ""
}

// payloadSource describes where the publish payload should come from, as
// derived from the command flags.
type payloadSource struct {
message string
messageSet bool // --message was given explicitly (possibly empty)
filePath string
fileSet bool // --file was given explicitly (possibly empty)
}

// validatePayloadFlags rejects flag combinations that cannot be resolved.
func validatePayloadFlags(src payloadSource) error {
if src.messageSet && src.fileSet {
return errors.New("only one of --message or --file should be used at a time")
}
if src.fileSet && src.filePath == "" {
return errors.New("--file requires a path (use - for stdin)")
}
return nil
}

// resolvePublishPayload returns the payload from --file, stdin, or --message.
// Stdin (a pipe or redirect, never an interactive terminal) takes precedence
// over --message when it carries data, as requested in #92; --file=- reads
// stdin explicitly. Reads are capped at maxBytes to bound memory use.
func resolvePublishPayload(src payloadSource, stdin *os.File, maxBytes int64) ([]byte, error) {
if err := validatePayloadFlags(src); err != nil {
return nil, err
}

if src.fileSet && src.filePath != "-" {
content, err := os.ReadFile(src.filePath)
if err != nil {
return nil, fmt.Errorf("failed to read file: %v", err)
}
return content, nil
}

stdinAvailable := src.filePath == "-" || !isTerminal(stdin)
if stdinAvailable {
content, err := readStdinBounded(stdin, maxBytes)
if err != nil {
return nil, err
}
if len(content) > 0 {
return content, nil
}
if src.filePath == "-" || !src.messageSet {
return nil, errors.New("stdin is empty: nothing to publish")
}
}

if src.messageSet {
if src.message == "" {
return nil, errors.New("--message is empty: nothing to publish")
}
return []byte(src.message), nil
}

return nil, errors.New("no message provided: use --message, --file, or pipe data via stdin")
}

// readStdinBounded reads stdin to EOF, failing if it exceeds maxBytes.
func readStdinBounded(stdin *os.File, maxBytes int64) ([]byte, error) {
var r io.Reader = stdin
if maxBytes > 0 {
r = io.LimitReader(stdin, maxBytes+1)
}
content, err := io.ReadAll(r)
if err != nil {
return nil, fmt.Errorf("failed to read stdin: %v", err)
}
if maxBytes > 0 && int64(len(content)) > maxBytes {
return nil, fmt.Errorf("stdin exceeds the maximum message size of %d bytes", maxBytes)
}
return content, nil
}

// isTerminal reports whether f is an interactive terminal. A nil or
// unreadable file counts as one so we never block on input that cannot arrive.
func isTerminal(f *os.File) bool {
if f == nil {
return true
}
info, err := f.Stat()
if err != nil {
return true
}
return info.Mode()&os.ModeCharDevice != 0
}

var publishCmd = &cobra.Command{
Use: "publish",
Short: "Publish a message to the Optimum Network",
Long: `Publish a message to the Optimum Network.

The payload comes from --file, stdin, or --message. Piped stdin is used
whenever it carries data (it takes precedence over --message), so the
command composes in pipes:

echo "hello" | mump2p publish --topic=test
curl -s https://api.example.com/status | mump2p publish --topic=status

Use --file=- to read from stdin explicitly.`,
Example: ` mump2p publish --topic=test --message="Hello World"
mump2p publish --topic=test/data --file=./payload.json
cat payload.json | mump2p publish --topic=test/data`,
RunE: func(cmd *cobra.Command, args []string) error {
if pubMessage == "" && file == "" {
return errors.New("either --message or --file must be provided")
src := payloadSource{
message: pubMessage,
messageSet: cmd.Flags().Changed("message"),
filePath: file,
fileSet: cmd.Flags().Changed("file"),
}
if pubMessage != "" && file != "" {
return errors.New("only one of --message or --file should be used at a time")
if err := validatePayloadFlags(src); err != nil {
return err
}

var claims *auth.TokenClaims
var clientIDToUse string
var accessToken string
maxMessageSize := int64(config.DefaultMaxMessageSize)

if !IsAuthDisabled() {
authClient := auth.NewClient()
Expand All @@ -96,23 +204,20 @@ var publishCmd = &cobra.Command{
return fmt.Errorf("your account is inactive, please contact support")
}
clientIDToUse = claims.ClientID
if claims.MaxMessageSize > 0 {
maxMessageSize = claims.MaxMessageSize
}
} else {
clientIDToUse = GetClientID()
if clientIDToUse == "" {
return fmt.Errorf("--client-id is required when using --disable-auth")
}
}

var data []byte

if file != "" {
content, err := os.ReadFile(file)
if err != nil {
return fmt.Errorf("failed to read file: %v", err)
}
data = content
} else {
data = []byte(pubMessage)
// Read the payload only after the size limit is known so stdin is bounded.
data, err := resolvePublishPayload(src, os.Stdin, maxMessageSize)
if err != nil {
return err
}

messageSize := int64(len(data))
Expand Down Expand Up @@ -244,8 +349,8 @@ var publishCmd = &cobra.Command{

func init() {
publishCmd.Flags().StringVar(&pubTopic, "topic", "", "Topic to publish to")
publishCmd.Flags().StringVar(&pubMessage, "message", "", "Message string to publish")
publishCmd.Flags().StringVar(&file, "file", "", "Path of the file to publish")
publishCmd.Flags().StringVar(&pubMessage, "message", "", "Message string to publish (reads stdin if neither --message nor --file is set)")
publishCmd.Flags().StringVar(&file, "file", "", "Path of the file to publish (use - for stdin)")
publishCmd.Flags().StringVar(&serviceURL, "service-url", "", "Override the default proxy URL")
publishCmd.Flags().Uint32Var(&pubExposeAmount, "expose-amount", 1, "Number of nodes to request from proxy")
publishCmd.MarkFlagRequired("topic") //nolint:errcheck
Expand Down
Loading