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
17 changes: 17 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,14 @@ Detailed per-release notes are on the
## [Unreleased]

### Added
- **`pilotctl send-message` can take its payload from stdin or a file**:
`--data -` and `--data-file <path>`. A payload passed as an argument is
capped by the OS — a 1 MB `--data` failed with "Argument list too long",
and on Linux a single argument stops at 128 KiB — which is why senders of
large bodies needed a separate stdin helper.
`--data -` now means stdin, so a message that is just `-` has to come from
`--data-file`. One message can carry up to the 64 MiB data-exchange frame
limit; a larger payload is refused before anything is sent.
- **A client can ask the daemon whether a datagram was actually sent.** The
IPC `SendTo` command is fire-and-forget: when the daemon could not send a
datagram (no route to the node, port policy, ephemeral ports exhausted) it
Expand Down Expand Up @@ -356,6 +364,15 @@ Detailed per-release notes are on the
with 2 of 5 failing before, 10.3–11.8s with none failing after; at 0.5%
loss 9.4–10.5s before, 7.3–7.6s after; without loss 8.2–8.4s before,
7.6–7.7s after.
- **A message larger than 256 KB is delivered instead of silently dropped.**
The daemon refused any single stream write bigger than its send buffer, and
an IPC send has no reply, so the client never knew: `send-message` printed
`"status":"ok"` for a 1 MB message that never left the node. A large write is
now fed through the buffer in pieces, blocking on the window like any other
write. The buffer's size cap is unchanged.
- **`pilotctl send-message` fails when the receiver does not acknowledge the
message.** Every receiver answers a stored message with an ACK; with none,
the command used to exit 0.
- **`-advertise-endpoint` survives a re-registration.** When the daemon
re-registered (registry reconnect, transport watchdog recovery) it sent
the tunnel socket's local address instead of the advertised endpoint, so
Expand Down
172 changes: 158 additions & 14 deletions cmd/pilotctl/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,10 @@ func classifyDaemonError(err error) string {
return ""
}

// fatalResults, when set, goes into fatalHint's JSON envelope as "results",
// so a multi-message send that fails still says what became of each message.
var fatalResults []map[string]interface{}

// fatalHint is like fatalCode but adds an actionable hint telling the user what to do next.
func fatalHint(code, hint, format string, args ...interface{}) {
msg := fmt.Sprintf(format, args...)
Expand All @@ -204,6 +208,9 @@ func fatalHint(code, hint, format string, args ...interface{}) {
if ns := nextStepsEnvelope(exitNextSteps); ns != nil {
env["next_steps"] = ns
}
if fatalResults != nil {
env["results"] = fatalResults
}
b, _ := json.Marshal(env)
fmt.Fprintln(os.Stderr, string(b))
} else {
Expand Down Expand Up @@ -895,7 +902,9 @@ var commandHelp = map[string]string{
Send a message to a remote agent and optionally wait for the reply.

Flags:
--data <text> message payload (required)
--data <text> message payload; "-" reads it from stdin
--data-file <path> read the payload from a file (for payloads too large
for a command-line argument)
--type text|json|binary payload encoding (default: text)
--count <n> send N times (default: 1)
--reuse-conn reuse the connection across --count sends (saves ~1 RTT)
Expand All @@ -912,6 +921,8 @@ Examples:
pilotctl send-message list-agents --data '/data {"search":"weather","limit":5}'
pilotctl send-message my-peer --data "hello" --wait
pilotctl send-message 0:0000.0000.400E --data "ping" --trace
pilotctl send-message my-peer --data-file report.json --type json
generate-report | pilotctl send-message my-peer --data -
`,
"ping": `Usage: pilotctl ping <address|hostname> [flags]

Expand Down Expand Up @@ -1628,7 +1639,7 @@ Communication commands:
pilotctl send <address|hostname> <port> --data <msg> [--timeout <dur>]
pilotctl recv <port> [--count <n>] [--timeout <dur>]
pilotctl send-file <address|hostname> <filepath>
pilotctl send-message <address|hostname> --data <text> [--type text|json|binary] [--count <n>] [--reuse-conn] [--wait <dur>]
pilotctl send-message <address|hostname> --data <text> | --data - | --data-file <path> [--type text|json|binary] [--count <n>] [--reuse-conn] [--wait <dur>]
pilotctl dgram <address|hostname> <port> --data <msg>
pilotctl subscribe <address|hostname> <topic> [--count <n>] [--timeout <dur>]
pilotctl publish <address|hostname> <topic> --data <message>
Expand Down Expand Up @@ -2487,8 +2498,8 @@ func contextCatalog() map[string]interface{} {

// Messaging
"send-message": map[string]interface{}{
"args": []string{"<address|hostname>", "--data <text>", "[--type text|json|binary]", "[--count <n>]", "[--reuse-conn]", "[--wait <dur>]", "[--timeout <dur>]"},
"description": "Send a typed message to a node via data exchange (port 1001). --count N sends N messages; --reuse-conn shares one connection across all N (env: PILOT_SENDMSG_REUSE_CONN=1). Default type: text",
"args": []string{"<address|hostname>", "--data <text> | --data - | --data-file <path>", "[--type text|json|binary]", "[--count <n>]", "[--reuse-conn]", "[--wait <dur>]", "[--timeout <dur>]"},
"description": "Send a typed message to a node via data exchange (port 1001). --data - reads the payload from stdin, --data-file from a file (up to 64 MiB; a command-line argument is capped by the OS). --count N sends N messages; --reuse-conn shares one connection across all N (env: PILOT_SENDMSG_REUSE_CONN=1). Fails unless every message is acknowledged. Default type: text",
"returns": "target, to, type, bytes, ack, reuse_conn",
},
"send-file": map[string]interface{}{
Expand Down Expand Up @@ -4744,11 +4755,117 @@ func streamSendFile(d *driver.Driver, target protocol.Addr, filePath, filename s
return result, nil
}

// messagePayload returns the body for send-message: --data, or the whole of
// stdin for "--data -", or the contents of --data-file. A payload passed as an
// argument is capped by the OS (ARG_MAX; 128 KiB per argument on Linux), which
// is why callers with large bodies needed a separate stdin helper.
func messagePayload(flags map[string]string, stdin io.Reader) (string, error) {
data, hasData := flags["data"]
file, hasFile := flags["data-file"]
switch {
case hasData && hasFile:
return "", fmt.Errorf("give --data or --data-file, not both")
case hasFile:
if file == "" || file == "true" {
return "", fmt.Errorf("--data-file needs a path")
}
b, err := os.ReadFile(file) // #nosec G304 -- the operator names the file to send
if err != nil {
return "", fmt.Errorf("read --data-file: %v", err)
}
data = string(b)
case data == "-":
b, err := io.ReadAll(stdin)
if err != nil {
return "", fmt.Errorf("read payload from stdin: %v", err)
}
data = string(b)
}
if data == "" {
if hasFile {
return "", fmt.Errorf("--data-file %s is empty", file)
}
return "", fmt.Errorf("--data is required")
}
// One message is one data-exchange frame. A receiver drops a frame over
// its limit at the header, which the sender would only see as a missing
// acknowledgement, and past 4 GiB the length prefix would wrap.
if limit := maxMessageBytes(); len(data) > limit {
return "", fmt.Errorf("message is %d bytes; one message can carry at most %d (use send-file for larger payloads)", len(data), limit)
}
return data, nil
}

// maxMessageBytes is the most one send-message payload can be: the frame
// limit, less room for the frame and message-tag headers.
func maxMessageBytes() int {
return int(dataexchange.MaxFrameSize) - 4096
}

// failIfUndelivered ends a multi-message send with an error when any message
// failed to send, was not acknowledged, or was refused by the receiver:
// reporting "ok" for those hid messages that never arrived. The error carries
// every message's result, so a caller can tell which ones to send again.
func failIfUndelivered(target string, results []map[string]interface{}) {
failed, allRefused, first := undelivered(results)
if failed == 0 {
return
}
code := "connection_failed"
if allRefused {
code = "internal" // as for a single message the receiver refuses
}
fatalResults = results
fatalHint(code, "the results list which messages were delivered; check `pilotctl peers` and the daemon log before sending the others again",
"%d of %d messages to %s were not delivered (first, %s)", failed, len(results), target, first)
}

// undelivered counts the send results that did not end in a stored message,
// says whether all of those were refused by the receiver, and describes the
// first of them.
func undelivered(results []map[string]interface{}) (failed int, allRefused bool, first string) {
allRefused = true
for _, r := range results {
why, refused := "", false
_, acked := r["ack"].(string)
switch {
case r["error"] != nil:
why = fmt.Sprint(r["error"])
case !acked:
why = "not acknowledged"
if e, ok := r["ack_error"].(string); ok {
why += " (" + e + ")"
}
case refusal(r) != "":
why, refused = "receiver refused it: "+refusal(r), true
}
if why != "" {
failed++
allRefused = allRefused && refused
if first == "" {
first = fmt.Sprintf("message %v: %s", r["seq"], why)
}
}
}
return failed, failed > 0 && allRefused, first
}

// refusal returns the receiver's "ERR ..." answer in a send result: its ACK,
// or with --trace the inner ACK that the timing reply carries.
func refusal(r map[string]interface{}) string {
for _, k := range []string{"ack", "inner_ack"} {
if s, _ := r[k].(string); strings.HasPrefix(s, "ERR ") {
return s
}
}
return ""
}

func cmdSendMessage(args []string) {
flags, pos := parseFlags(args)
for name := range flags {
switch name {
case "data", "type", "count", "reuse-conn", "wait", "timeout", "no-resend", "trace", "no-auto-handshake":
case "data", "data-file", "type", "count", "reuse-conn", "wait", "timeout", "no-resend", "trace", "no-auto-handshake":
default:
fatalCode("invalid_argument", "send-message: unknown flag --%s", name)
}
Expand All @@ -4768,7 +4885,7 @@ func cmdSendMessage(args []string) {
defer timer.Stop()
}
if len(pos) != 1 {
fatalCode("invalid_argument", "usage: pilotctl send-message <address|hostname> --data <text> [--type text|json|binary] [--trace] [--count <n>] [--reuse-conn] [--wait <dur>] [--timeout <dur>] [--no-resend]")
fatalCode("invalid_argument", "usage: pilotctl send-message <address|hostname> --data <text> | --data - | --data-file <path> [--type text|json|binary] [--trace] [--count <n>] [--reuse-conn] [--wait <dur>] [--timeout <dur>] [--no-resend]")
}

sendCount := flagInt(flags, "count", 1)
Expand Down Expand Up @@ -4812,6 +4929,13 @@ func cmdSendMessage(args []string) {
}
}

// Read and check the payload first, so a missing file or an oversized
// message is reported before the daemon or the network is touched.
data, err := messagePayload(flags, os.Stdin)
if err != nil {
fatalCode("invalid_argument", "%v", err)
}

d := connectDriver()
tracef("connectDriver")
// d may be replaced by a fresh connection when a dial is retried.
Expand All @@ -4823,10 +4947,6 @@ func cmdSendMessage(args []string) {
fatalCode("not_found", "%v", err)
}

data := flagString(flags, "data", "")
if data == "" {
fatalCode("invalid_argument", "--data is required")
}
msgType := flagString(flags, "type", "text")

// First contact: the daemon holds no session with the peer yet, so this
Expand Down Expand Up @@ -4903,6 +5023,9 @@ func cmdSendMessage(args []string) {
if ack != nil {
r["ack"] = string(ack.Payload)
}
if ackErr != nil {
r["ack_error"] = ackErr.Error()
}
if traceTime {
r["total_ms"] = float64(time.Duration(ackRecvAtNs-sentAtNs).Microseconds()) / 1000.0
if ack != nil && ack.Type == dataexchange.TypeJSON {
Expand Down Expand Up @@ -4974,11 +5097,31 @@ func cmdSendMessage(args []string) {
defer cl.Close()
r := sendOne(cl, 0, false)
ackAt := time.Now()
// Every receiver answers a stored message with an ACK frame. No ACK
// means the message was not stored, or was never sent — the daemon
// drops a write it cannot buffer without telling the client — so
// "ok" here was reporting messages that never arrived. A send that
// failed outright was reported as "ok" with an error field; it is a
// failure too.
if e, failed := r["error"].(string); failed {
fatalHint("connection_failed",
"the message was not sent; check `pilotctl peers` and the daemon log, then send again",
"sending to %s failed: %s", target, e)
}
if _, acked := r["ack"]; !acked {
why := ""
if e, ok := r["ack_error"].(string); ok {
why = ": " + e
}
fatalHint("connection_failed",
"the receiver did not confirm it stored the message; check `pilotctl peers` and the daemon log, then send again",
"%s did not acknowledge the message (%d bytes)%s", target, len(data), why)
}
// The receiver answers "ERR ..." when it could not store the
// message (disk full, inbox unwritable). That is a failed send, not
// a delivered one — same rule send-file applies.
if ackText, _ := r["ack"].(string); strings.HasPrefix(ackText, "ERR ") {
fatalCode("internal", "receiver rejected message: %s", ackText)
if rej := refusal(r); rej != "" {
fatalCode("internal", "receiver rejected message: %s", rej)
}
result := map[string]interface{}{
"target": target.String(),
Expand Down Expand Up @@ -5012,8 +5155,7 @@ func cmdSendMessage(args []string) {
// more on a new stream if nothing arrived by mid-window.
// --no-resend opts out for requests that are not safe to
// repeat.
_, sendFailed := r["error"]
if firstContact && !sendFailed && !flagBool(flags, "no-resend") {
if firstContact && !flagBool(flags, "no-resend") {
if after := resendDelay(waitDur); after > 0 {
cfg.resendAfter = after
cfg.resend = func() (time.Time, error) {
Expand Down Expand Up @@ -5084,6 +5226,7 @@ func cmdSendMessage(args []string) {
time.Sleep(50 * time.Millisecond)
}
}
failIfUndelivered(target.String(), results)
outputOK(map[string]interface{}{
"target": target.String(),
"to": target.String(),
Expand All @@ -5104,6 +5247,7 @@ func cmdSendMessage(args []string) {
time.Sleep(50 * time.Millisecond)
}
}
failIfUndelivered(target.String(), results)
outputOK(map[string]interface{}{
"target": target.String(),
"to": target.String(),
Expand Down
Loading
Loading