diff --git a/cmd/kcd/cli_connectivity.go b/cmd/kcd/cli_connectivity.go index 1c068e5..4eaf136 100644 --- a/cmd/kcd/cli_connectivity.go +++ b/cmd/kcd/cli_connectivity.go @@ -51,19 +51,42 @@ var connectivityCmd = &cli.Command{ }, } -// formatConnectivity renders one line per SIM, e.g. "LTE [███░] (3/4)". -// The primary SIM ("0", else lowest key) comes first; keys are sorted for -// stable output. Level is clamped to 0-4 so the bar always parses. +// signalDots renders signal level 0-4 as a four-position dot bar. The range is +// clamped so a bogus level from the phone still produces a parseable bar, and +// four positions means the maximum level reads as genuinely full. +// These are text-presentation geometric shapes (single cell width), not +// pictographs, so they align like the block-element bar they replace. +func signalDots(level int) string { + if level < 0 { + level = 0 + } + if level > 4 { + level = 4 + } + dots := []string{"○○○○", "●○○○", "●●○○", "●●●○", "●●●●"} + return dots[level] +} + +func formatSIMStatus(networkType string, level int) string { + if level <= 0 { + return "No service" + } + label := networkType + if label == "" || strings.EqualFold(label, "unknown") { + label = "Cellular" + } + return fmt.Sprintf("%-8s %s", label, signalDots(level)) +} + func formatConnectivity(body connectivity.ConnectivityBody) []string { if len(body.SignalStrengths) == 0 { - return []string{"No signal data reported"} + return []string{"No SIM or cellular data available"} } keys := make([]string, 0, len(body.SignalStrengths)) for k := range body.SignalStrengths { keys = append(keys, k) } sort.Strings(keys) - // Primary SIM first. if _, ok := body.SignalStrengths["0"]; ok { ordered := []string{"0"} for _, k := range keys { @@ -77,26 +100,11 @@ func formatConnectivity(body connectivity.ConnectivityBody) []string { lines := make([]string, 0, len(keys)) for _, k := range keys { sig := body.SignalStrengths[k] - level := sig.SignalStrength - if level < 0 { - level = 0 - } - if level > 4 { - level = 4 - } netType := sig.NetworkDetailedType if netType == "" { netType = sig.NetworkType } - if netType == "" { - netType = "CELL" - } - bar := strings.Repeat("█", level) + strings.Repeat("░", 4-level) - if len(keys) == 1 { - lines = append(lines, fmt.Sprintf("%s [%s] (%d/4)", strings.ToUpper(netType), bar, level)) - } else { - lines = append(lines, fmt.Sprintf("SIM %s: %s [%s] (%d/4)", k, strings.ToUpper(netType), bar, level)) - } + lines = append(lines, fmt.Sprintf("SIM %s: %s", k, formatSIMStatus(netType, sig.SignalStrength))) } return lines } diff --git a/cmd/kcd/cli_connectivity_test.go b/cmd/kcd/cli_connectivity_test.go index 7d62c0b..ebd0c12 100644 --- a/cmd/kcd/cli_connectivity_test.go +++ b/cmd/kcd/cli_connectivity_test.go @@ -15,36 +15,72 @@ func TestFormatConnectivity(t *testing.T) { { name: "empty", body: connectivity.ConnectivityBody{}, - want: []string{"No signal data reported"}, + want: []string{"No SIM or cellular data available"}, }, { name: "single sim detailed type", body: connectivity.ConnectivityBody{SignalStrengths: map[string]connectivity.SignalStrength{ "0": {NetworkType: "LTE", NetworkDetailedType: "LTE", SignalStrength: 3}, }}, - want: []string{"LTE [███░] (3/4)"}, + want: []string{"SIM 0: LTE ●●●○"}, }, { name: "falls back to network type", body: connectivity.ConnectivityBody{SignalStrengths: map[string]connectivity.SignalStrength{ "0": {NetworkType: "5G", SignalStrength: 4}, }}, - want: []string{"5G [████] (4/4)"}, + want: []string{"SIM 0: 5G ●●●●"}, }, { - name: "clamps out of range", + name: "clamps above range to a full bar", body: connectivity.ConnectivityBody{SignalStrengths: map[string]connectivity.SignalStrength{ "0": {NetworkType: "GSM", SignalStrength: 9}, }}, - want: []string{"GSM [████] (4/4)"}, + want: []string{"SIM 0: GSM ●●●●"}, }, { - name: "dual sim primary first", + name: "clamps below range reads as no service", + body: connectivity.ConnectivityBody{SignalStrengths: map[string]connectivity.SignalStrength{ + "0": {NetworkType: "GSM", SignalStrength: -2}, + }}, + want: []string{"SIM 0: No service"}, + }, + { + name: "zero level is no service, not an empty bar", + body: connectivity.ConnectivityBody{SignalStrengths: map[string]connectivity.SignalStrength{ + "0": {NetworkType: "LTE", SignalStrength: 0}, + }}, + want: []string{"SIM 0: No service"}, + }, + { + name: "unknown network type reads as cellular", + body: connectivity.ConnectivityBody{SignalStrengths: map[string]connectivity.SignalStrength{ + "0": {NetworkType: "unknown", SignalStrength: 2}, + }}, + want: []string{"SIM 0: Cellular ●●○○"}, + }, + { + name: "empty network type reads as cellular", + body: connectivity.ConnectivityBody{SignalStrengths: map[string]connectivity.SignalStrength{ + "0": {SignalStrength: 1}, + }}, + want: []string{"SIM 0: Cellular ●○○○"}, + }, + { + name: "dual sim primary first and labels stay aligned", body: connectivity.ConnectivityBody{SignalStrengths: map[string]connectivity.SignalStrength{ "1": {NetworkType: "EDGE", SignalStrength: 2}, "0": {NetworkType: "LTE", SignalStrength: 4}, }}, - want: []string{"SIM 0: LTE [████] (4/4)", "SIM 1: EDGE [██░░] (2/4)"}, + want: []string{"SIM 0: LTE ●●●●", "SIM 1: EDGE ●●○○"}, + }, + { + name: "sparse subscription ids are not renumbered", + body: connectivity.ConnectivityBody{SignalStrengths: map[string]connectivity.SignalStrength{ + "3": {NetworkType: "NR", SignalStrength: 3}, + "7": {NetworkType: "LTE", SignalStrength: 1}, + }}, + want: []string{"SIM 3: NR ●●●○", "SIM 7: LTE ●○○○"}, }, } for _, tc := range cases { diff --git a/cmd/kcd/cli_watch.go b/cmd/kcd/cli_watch.go index cadbd1c..0802aa7 100644 --- a/cmd/kcd/cli_watch.go +++ b/cmd/kcd/cli_watch.go @@ -56,62 +56,7 @@ var watchCmd = &cli.Command{ b, _ := json.Marshal(ev) fmt.Println(string(b)) } else { - switch ev.Type { - case events.TypeBatteryUpdate: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("[%s] battery: %v%% (charging: %v)\n", ev.DeviceID, payload["charge"], payload["charging"]) - case events.TypeNotification: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("[%s] notification: %s - %s\n", ev.DeviceID, payload["appName"], payload["title"]) - case events.TypeShareProgress: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("\r[%s] transfer: %s... %v/%v bytes", ev.DeviceID, payload["file"], payload["current"], payload["total"]) - case events.TypeShareComplete: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("\n[%s] transfer complete: %s\n", ev.DeviceID, payload["file"]) - case events.TypeShareText: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("[%s] share text: %s\n", ev.DeviceID, payload["text"]) - case events.TypeShareURL: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("[%s] share url: %s\n", ev.DeviceID, payload["url"]) - case events.TypeSftpMount: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("[%s] SFTP credentials received: %s\n", ev.DeviceID, payload["uri"]) - case events.TypePairRequested: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("[%s] pair request from %s (%s). code: %v\n", ev.DeviceID, payload["name"], payload["type"], payload["verificationKey"]) - case events.TypePairAccepted: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("[%s] paired with %s\n", ev.DeviceID, payload["name"]) - case events.TypePairRejected: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("[%s] pairing rejected or cancelled by %s\n", ev.DeviceID, payload["name"]) - case events.TypeNotificationCanceled: - payload, _ := ev.Payload.(map[string]interface{}) - fmt.Printf("[%s] notification cancelled: %s\n", ev.DeviceID, payload["id"]) - case events.TypeMprisUpdate: - payload, _ := ev.Payload.(map[string]interface{}) - state := "⏹" - if isPlaying, _ := payload["isPlaying"].(bool); isPlaying { - state = "▶" - } else if ps, _ := payload["playbackStatus"].(string); ps == "Paused" { - state = "⏸" - } - player, _ := payload["player"].(string) - title, _ := payload["title"].(string) - artist, _ := payload["artist"].(string) - fmt.Printf("[%s] %s %s", ev.DeviceID, state, player) - if title != "" { - fmt.Printf(" - %s", title) - } - if artist != "" { - fmt.Printf(" (%s)", artist) - } - fmt.Println() - default: - fmt.Printf("[%s] %s\n", ev.DeviceID, ev.Type) - } + fmt.Print(formatEvent(ev)) } } }() diff --git a/cmd/kcd/format_event.go b/cmd/kcd/format_event.go new file mode 100644 index 0000000..6d023f9 --- /dev/null +++ b/cmd/kcd/format_event.go @@ -0,0 +1,321 @@ +package main + +import ( + "encoding/json" + "fmt" + "strings" + + "github.com/bethropolis/kcd/internal/events" + "github.com/bethropolis/kcd/internal/plugins/connectivity" + "github.com/bethropolis/kcd/internal/plugins/remotesystemvolume" + "github.com/bethropolis/kcd/internal/plugins/telephony" +) + +// truncate shortens free-form text to keep one event on one line. +func truncate(s string, max int) string { + if len(s) > max { + return s[:max-1] + "…" + } + return s +} + +// oneLine collapses newlines so a multi-line message cannot break the stream's +// line-per-event shape. +func oneLine(s string) string { + return strings.Join(strings.Fields(s), " ") +} + +// decodePayload re-decodes an event payload into a concrete type. The daemon +// hands the CLI generic JSON, so a payload that was published as a struct +// arrives as an untyped map and has to be round-tripped to reuse the same +// formatter the dedicated commands use. +func decodePayload(payload any, out any) bool { + b, err := json.Marshal(payload) + if err != nil { + return false + } + return json.Unmarshal(b, out) == nil +} + +func str(payload map[string]any, key string) string { + s, _ := payload[key].(string) + return s +} + +// num reads a numeric payload field. Values normally arrive as float64 from the +// JSON stream, but the same event can be produced in-process, so accept the +// other integer widths rather than silently reading zero. +func num(payload map[string]any, key string) (int, bool) { + switch v := payload[key].(type) { + case float64: + return int(v), true + case int: + return v, true + case int64: + return int(v), true + case json.Number: + n, err := v.Int64() + return int(n), err == nil + default: + return 0, false + } +} + +// formatEvent renders one event as a human-readable line, including its +// trailing newline (or leading carriage return, for the in-place transfer +// progress line). The trailing newline is part of the return value so callers +// can print without adding one. +func formatEvent(ev events.Event) string { + // A nil or non-map payload still renders for the types below that read + // nothing; the map assertion simply yields nil and lookups return "". + payload, _ := ev.Payload.(map[string]interface{}) + + switch ev.Type { + case events.TypeBatteryUpdate: + return fmt.Sprintf("[%s] battery: %v%% (charging: %v)\n", ev.DeviceID, payload["charge"], payload["charging"]) + + case events.TypeBatteryThreshold: + level := "full" + if e, ok := num(payload, "event"); ok && e == 1 { + level = "low" + } + return fmt.Sprintf("[%s] battery %s: %v%%%s\n", ev.DeviceID, level, payload["charge"], chargingSuffix(payload)) + + case events.TypeNotification: + return fmt.Sprintf("[%s] notification: %s - %s\n", ev.DeviceID, payload["appName"], payload["title"]) + + case events.TypeNotificationCanceled: + return fmt.Sprintf("[%s] notification cancelled: %s\n", ev.DeviceID, payload["id"]) + + case events.TypeShareProgress: + return fmt.Sprintf("\r[%s] transfer: %s... %v/%v bytes", ev.DeviceID, payload["file"], payload["current"], payload["total"]) + + case events.TypeShareComplete: + return fmt.Sprintf("\n[%s] transfer complete: %s\n", ev.DeviceID, payload["file"]) + + case events.TypeShareText: + return fmt.Sprintf("[%s] share text: %s\n", ev.DeviceID, payload["text"]) + + case events.TypeShareURL: + return fmt.Sprintf("[%s] share url: %s\n", ev.DeviceID, payload["url"]) + + case events.TypeSftpMount: + return fmt.Sprintf("[%s] SFTP credentials received: %s\n", ev.DeviceID, payload["uri"]) + + case events.TypePairRequested: + return fmt.Sprintf("[%s] pair request from %s (%s). code: %v\n", ev.DeviceID, payload["name"], payload["type"], payload["verificationKey"]) + + case events.TypePairAccepted: + return fmt.Sprintf("[%s] paired with %s\n", ev.DeviceID, payload["name"]) + + case events.TypePairRejected: + return fmt.Sprintf("[%s] pair rejected or cancelled by %s\n", ev.DeviceID, payload["name"]) + + case events.TypeMprisUpdate: + state := "⏹" + if isPlaying, _ := payload["isPlaying"].(bool); isPlaying { + state = "▶" + } else if ps, _ := payload["playbackStatus"].(string); ps == "Paused" { + state = "⏸" + } + player, _ := payload["player"].(string) + title, _ := payload["title"].(string) + artist, _ := payload["artist"].(string) + var b strings.Builder + fmt.Fprintf(&b, "[%s] %s %s", ev.DeviceID, state, player) + if title != "" { + fmt.Fprintf(&b, " - %s", title) + } + if artist != "" { + fmt.Fprintf(&b, " (%s)", artist) + } + b.WriteString("\n") + return b.String() + + case events.TypeSMSIncoming: + return fmt.Sprintf("[%s] sms from %s: %s\n", ev.DeviceID, payload["sender"], oneLine(truncate(str(payload, "body"), 72))) + + case events.TypeSMSAttachment: + return fmt.Sprintf("[%s] sms attachment saved: %s\n", ev.DeviceID, oneLine(truncate(str(payload, "filename"), 48))) + + case events.TypePingReceived: + return fmt.Sprintf("[%s] ping: %s\n", ev.DeviceID, oneLine(truncate(str(payload, "message"), 60))) + + case events.TypeConnectivityUpdate: + var body connectivity.ConnectivityBody + if !decodePayload(ev.Payload, &body) { + break + } + lines := formatConnectivity(body) + if len(lines) == 0 { + break + } + return fmt.Sprintf("[%s] connectivity: %s\n", ev.DeviceID, strings.Join(lines, ", ")) + + case events.TypeTelephonyRinging, events.TypeTelephonyTalking, events.TypeTelephonyMissed, events.TypeTelephonyCanceled: + var body telephony.TelephonyBody + if !decodePayload(ev.Payload, &body) { + break + } + // Wording matches the desktop notification this same event raises, so + // the terminal line and the popup read the same way. + caller := callerLabel(&body) + if caller == "" { + return fmt.Sprintf("[%s] %s\n", ev.DeviceID, telephonyLabel(ev.Type)) + } + return fmt.Sprintf("[%s] %s: %s\n", ev.DeviceID, telephonyLabel(ev.Type), caller) + + case events.TypeVolumeUpdate: + // Two publishers share this type: the systemvolume plugin reports a + // single stream, remotesystemvolume can report a whole sink list. + if _, ok := payload["sinks"]; ok { + return fmt.Sprintf("[%s] volume: %s\n", ev.DeviceID, sinkSummary(payload["sinks"])) + } + name := str(payload, "name") + if name == "" { + name = "output" + } + name = oneLine(truncate(name, 40)) + if m, _ := payload["muted"].(bool); m { + return fmt.Sprintf("[%s] volume: %s (muted)\n", ev.DeviceID, name) + } + return fmt.Sprintf("[%s] volume: %s %v%%\n", ev.DeviceID, name, payload["volume"]) + + case events.TypeContactsUpdated: + if str(payload, "phase") == "vcards" { + stored, _ := num(payload, "stored") + skipped, _ := num(payload, "skipped") + if skipped == 0 { + return fmt.Sprintf("[%s] contacts: %d saved\n", ev.DeviceID, stored) + } + return fmt.Sprintf("[%s] contacts: %d saved, %d skipped\n", ev.DeviceID, stored, skipped) + } + return fmt.Sprintf("[%s] contacts: %s\n", ev.DeviceID, contactDelta(payload)) + + // device.connected / device.added keep the bare type token as the first + // field so anything grepping for it still matches; the detail follows. + case events.TypeDeviceAdded: + name, _ := ev.Payload.(string) + if name == "" { + break + } + return fmt.Sprintf("[%s] device.added: %s\n", ev.DeviceID, oneLine(truncate(name, 40))) + + case events.TypeDeviceConnected: + name, typeName := str(payload, "name"), str(payload, "type") + if name == "" && typeName == "" { + break + } + return fmt.Sprintf("[%s] device.connected: %s (%s)\n", ev.DeviceID, oneLine(truncate(name, 32)), typeName) + + default: + // state.snapshot deliberately lands here: the IPC layer writes it + // straight to the socket with a full device+plugin dump as payload, + // so rendering it would flood the terminal. ring.received, + // device.removed and device.disconnected carry no payload. + return fmt.Sprintf("[%s] %s\n", ev.DeviceID, ev.Type) + } + + return fmt.Sprintf("[%s] %s\n", ev.DeviceID, ev.Type) +} + +// callerLabel prefers a contact name and falls back to the raw number, the way +// the desktop notification does. +func callerLabel(body *telephony.TelephonyBody) string { + if body.ContactName != "" { + if body.PhoneNumber != "" { + return fmt.Sprintf("%s (%s)", oneLine(truncate(body.ContactName, 32)), body.PhoneNumber) + } + return oneLine(truncate(body.ContactName, 40)) + } + if body.PhoneNumber != "" { + return body.PhoneNumber + } + return "" +} + +// telephonyLabel names the call event in plain language. It deliberately +// mirrors the wording telephony.go puts in the desktop notification so the two +// do not describe the same event differently. +func telephonyLabel(typ events.EventType) string { + switch typ { + case events.TypeTelephonyRinging: + return "incoming call" + case events.TypeTelephonyTalking: + return "call answered" + case events.TypeTelephonyMissed: + return "missed call" + default: + return "call ended" + } +} + +// chargingSuffix only mentions charging when it is actually happening. +func chargingSuffix(payload map[string]any) string { + if c, _ := payload["charging"].(bool); c { + return " (charging)" + } + return "" +} + +// contactDelta describes a contacts sync pass, dropping zero counts so a +// no-op phase does not read as "+0 ~0 -0". +func contactDelta(payload map[string]any) string { + added, _ := num(payload, "added") + updated, _ := num(payload, "updated") + deleted, _ := num(payload, "deleted") + pending, _ := num(payload, "pending") + + var parts []string + if added > 0 { + parts = append(parts, fmt.Sprintf("%d added", added)) + } + if updated > 0 { + parts = append(parts, fmt.Sprintf("%d updated", updated)) + } + if deleted > 0 { + parts = append(parts, fmt.Sprintf("%d deleted", deleted)) + } + if pending > 0 { + parts = append(parts, fmt.Sprintf("%d pending", pending)) + } + if len(parts) == 0 { + return "no changes" + } + return strings.Join(parts, ", ") +} + +// sinkSummary renders a sink list as "Speaker 42%, Headset 80%". Raw %v on the +// slice would print Go slice syntax, which is not something a person should +// have to read. +func sinkSummary(raw any) string { + var sinks []remotesystemvolume.SinkInfo + if !decodePayload(raw, &sinks) { + return "unknown" + } + if len(sinks) == 0 { + return "no streams" + } + + const maxSinks = 3 + parts := make([]string, 0, maxSinks+1) + for i, s := range sinks { + if i == maxSinks { + parts = append(parts, fmt.Sprintf("+%d more", len(sinks)-maxSinks)) + break + } + label := oneLine(truncate(s.Name, 28)) + if label == "" { + label = s.Description + } + if label == "" { + label = "stream" + } + if s.Muted { + parts = append(parts, fmt.Sprintf("%s (muted)", label)) + continue + } + parts = append(parts, fmt.Sprintf("%s %d%%", label, s.Volume)) + } + return strings.Join(parts, ", ") +} diff --git a/cmd/kcd/format_event_test.go b/cmd/kcd/format_event_test.go new file mode 100644 index 0000000..c83daf6 --- /dev/null +++ b/cmd/kcd/format_event_test.go @@ -0,0 +1,286 @@ +package main + +import ( + "strings" + "testing" + + "github.com/bethropolis/kcd/internal/events" +) + +// These tests pin the exact output of every format that existed before +// formatEvent was extracted. The extraction is a pure refactor, so any change +// here means a line a script may already parse has moved. New formats are +// covered separately in TestFormatEventNewTypes. +func TestFormatEventExistingFormats(t *testing.T) { + tests := []struct { + name string + ev events.Event + want string + }{ + { + name: "battery", + ev: events.Event{Type: events.TypeBatteryUpdate, DeviceID: "d1", Payload: map[string]any{"charge": 62, "charging": false}}, + want: "[d1] battery: 62% (charging: false)\n", + }, + { + name: "notification", + ev: events.Event{Type: events.TypeNotification, DeviceID: "d1", Payload: map[string]any{"appName": "WhatsApp", "title": "Alice"}}, + want: "[d1] notification: WhatsApp - Alice\n", + }, + { + name: "notification canceled", + ev: events.Event{Type: events.TypeNotificationCanceled, DeviceID: "d1", Payload: map[string]any{"id": "notif-abc"}}, + want: "[d1] notification cancelled: notif-abc\n", + }, + { + name: "share progress keeps leading carriage return and no newline", + ev: events.Event{Type: events.TypeShareProgress, DeviceID: "d1", Payload: map[string]any{"file": "a.png", "current": 10, "total": 20}}, + want: "\r[d1] transfer: a.png... 10/20 bytes", + }, + { + name: "share complete keeps leading newline", + ev: events.Event{Type: events.TypeShareComplete, DeviceID: "d1", Payload: map[string]any{"file": "a.png"}}, + want: "\n[d1] transfer complete: a.png\n", + }, + { + name: "share text", + ev: events.Event{Type: events.TypeShareText, DeviceID: "d1", Payload: map[string]any{"text": "hi"}}, + want: "[d1] share text: hi\n", + }, + { + name: "share url", + ev: events.Event{Type: events.TypeShareURL, DeviceID: "d1", Payload: map[string]any{"url": "https://kdeconnect.org"}}, + want: "[d1] share url: https://kdeconnect.org\n", + }, + { + name: "sftp mount", + ev: events.Event{Type: events.TypeSftpMount, DeviceID: "d1", Payload: map[string]any{"uri": "sftp://x"}}, + want: "[d1] SFTP credentials received: sftp://x\n", + }, + { + name: "pair requested", + ev: events.Event{Type: events.TypePairRequested, DeviceID: "d1", Payload: map[string]any{"name": "Pixel", "type": "phone", "verificationKey": "12345"}}, + want: "[d1] pair request from Pixel (phone). code: 12345\n", + }, + { + name: "pair accepted", + ev: events.Event{Type: events.TypePairAccepted, DeviceID: "d1", Payload: map[string]any{"name": "Pixel"}}, + want: "[d1] paired with Pixel\n", + }, + { + name: "pair rejected", + ev: events.Event{Type: events.TypePairRejected, DeviceID: "d1", Payload: map[string]any{"name": "Pixel"}}, + want: "[d1] pair rejected or cancelled by Pixel\n", + }, + { + name: "mpris playing with title and artist", + ev: events.Event{Type: events.TypeMprisUpdate, DeviceID: "d1", Payload: map[string]any{"isPlaying": true, "player": "Firefox", "title": "Song", "artist": "Band"}}, + want: "[d1] ▶ Firefox - Song (Band)\n", + }, + { + name: "mpris paused", + ev: events.Event{Type: events.TypeMprisUpdate, DeviceID: "d1", Payload: map[string]any{"playbackStatus": "Paused", "player": "Firefox"}}, + want: "[d1] ⏸ Firefox\n", + }, + { + name: "mpris stopped omits empty title and artist", + ev: events.Event{Type: events.TypeMprisUpdate, DeviceID: "d1", Payload: map[string]any{"isPlaying": false, "player": "Firefox", "title": "", "artist": ""}}, + want: "[d1] ⏹ Firefox\n", + }, + { + name: "sms incoming", + ev: events.Event{Type: events.TypeSMSIncoming, DeviceID: "d1", Payload: map[string]any{"sender": "+1555", "body": "hello"}}, + want: "[d1] sms from +1555: hello\n", + }, + { + name: "unrendered type falls back to the bare type token", + ev: events.Event{Type: events.EventType("something.new"), DeviceID: "d1", Payload: nil}, + want: "[d1] something.new\n", + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + if got := formatEvent(tc.ev); got != tc.want { + t.Errorf("formatEvent() =\n%q\nwant\n%q", got, tc.want) + } + }) + } +} + +func TestFormatEventNewTypes(t *testing.T) { + tests := []struct { + name string + ev events.Event + want string + }{ + { + name: "battery low threshold omits the charging noise", + ev: events.Event{Type: events.TypeBatteryThreshold, DeviceID: "d1", Payload: map[string]any{"charge": 15, "charging": false, "event": 1}}, + want: "[d1] battery low: 15%\n", + }, + { + name: "battery full threshold mentions charging when true", + ev: events.Event{Type: events.TypeBatteryThreshold, DeviceID: "d1", Payload: map[string]any{"charge": 100, "charging": true, "event": 0}}, + want: "[d1] battery full: 100% (charging)\n", + }, + { + name: "sms attachment", + ev: events.Event{Type: events.TypeSMSAttachment, DeviceID: "d1", Payload: map[string]any{"filename": "photo.jpg", "path": "/tmp/x", "thread_id": 3}}, + want: "[d1] sms attachment saved: photo.jpg\n", + }, + { + name: "ping", + ev: events.Event{Type: events.TypePingReceived, DeviceID: "d1", Payload: map[string]any{"message": "pong"}}, + want: "[d1] ping: pong\n", + }, + { + name: "telephony ringing names the event in plain language", + ev: events.Event{Type: events.TypeTelephonyRinging, DeviceID: "d1", Payload: map[string]any{"contactName": "Bob", "phoneNumber": "+15550001234"}}, + want: "[d1] incoming call: Bob (+15550001234)\n", + }, + { + name: "telephony missed falls back to the number", + ev: events.Event{Type: events.TypeTelephonyMissed, DeviceID: "d1", Payload: map[string]any{"phoneNumber": "+15550001234"}}, + want: "[d1] missed call: +15550001234\n", + }, + { + name: "telephony talking with no caller still names the event", + ev: events.Event{Type: events.TypeTelephonyTalking, DeviceID: "d1", Payload: map[string]any{}}, + want: "[d1] call answered\n", + }, + { + name: "telephony canceled", + ev: events.Event{Type: events.TypeTelephonyCanceled, DeviceID: "d1", Payload: map[string]any{"contactName": "Bob"}}, + want: "[d1] call ended: Bob\n", + }, + { + name: "volume keeps the stream name", + ev: events.Event{Type: events.TypeVolumeUpdate, DeviceID: "d1", Payload: map[string]any{"name": "Speaker", "volume": 42, "muted": false}}, + want: "[d1] volume: Speaker 42%\n", + }, + { + name: "volume muted", + ev: events.Event{Type: events.TypeVolumeUpdate, DeviceID: "d1", Payload: map[string]any{"name": "Speaker", "volume": 42, "muted": true}}, + want: "[d1] volume: Speaker (muted)\n", + }, + { + name: "volume sink list renders names not Go slice syntax", + ev: events.Event{Type: events.TypeVolumeUpdate, DeviceID: "d1", Payload: map[string]any{"sinks": []map[string]any{ + {"name": "Speaker", "volume": 42, "muted": false}, + {"name": "Headset", "volume": 80, "muted": true}, + }}}, + want: "[d1] volume: Speaker 42%, Headset (muted)\n", + }, + { + name: "volume sink list caps the number shown", + ev: events.Event{Type: events.TypeVolumeUpdate, DeviceID: "d1", Payload: map[string]any{"sinks": []map[string]any{ + {"name": "a", "volume": 1}, {"name": "b", "volume": 2}, + {"name": "c", "volume": 3}, {"name": "d", "volume": 4}, + {"name": "e", "volume": 5}, + }}}, + want: "[d1] volume: a 1%, b 2%, c 3%, +2 more\n", + }, + { + name: "contacts uids phase drops zero counts", + ev: events.Event{Type: events.TypeContactsUpdated, DeviceID: "d1", Payload: map[string]any{"phase": "uids", "added": 2, "updated": 0, "deleted": 0, "pending": 5}}, + want: "[d1] contacts: 2 added, 5 pending\n", + }, + { + name: "contacts vcards phase drops a zero skip count", + ev: events.Event{Type: events.TypeContactsUpdated, DeviceID: "d1", Payload: map[string]any{"phase": "vcards", "stored": 4, "skipped": 0}}, + want: "[d1] contacts: 4 saved\n", + }, + { + name: "contacts vcards phase keeps a real skip count", + ev: events.Event{Type: events.TypeContactsUpdated, DeviceID: "d1", Payload: map[string]any{"phase": "vcards", "stored": 4, "skipped": 1}}, + want: "[d1] contacts: 4 saved, 1 skipped\n", + }, + { + name: "contacts with nothing to report", + ev: events.Event{Type: events.TypeContactsUpdated, DeviceID: "d1", Payload: map[string]any{"phase": "uids", "added": 0, "updated": 0, "deleted": 0, "pending": 0}}, + want: "[d1] contacts: no changes\n", + }, + { + name: "connectivity reuses the dedicated command's formatting", + ev: events.Event{Type: events.TypeConnectivityUpdate, DeviceID: "d1", Payload: map[string]any{ + "signalStrengths": map[string]any{ + "0": map[string]any{"networkType": "LTE", "signalStrength": 3}, + }, + }}, + want: "[d1] connectivity: SIM 0: LTE ●●●○\n", + }, + { + name: "device added keeps type token first", + ev: events.Event{Type: events.TypeDeviceAdded, DeviceID: "d1", Payload: "Pixel 8"}, + want: "[d1] device.added: Pixel 8\n", + }, + { + name: "device connected keeps type token first", + ev: events.Event{Type: events.TypeDeviceConnected, DeviceID: "d1", Payload: map[string]any{"name": "Pixel 8", "type": "phone"}}, + want: "[d1] device.connected: Pixel 8 (phone)\n", + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + if got := formatEvent(tc.ev); got != tc.want { + t.Errorf("formatEvent() =\n%q\nwant\n%q", got, tc.want) + } + }) + } +} + +// state.snapshot is written straight to the socket by the IPC layer with a full +// device+plugin dump. Rendering it would flood the terminal, so it must stay on +// the bare type-token line. +func TestFormatEventStateSnapshotStaysBare(t *testing.T) { + ev := events.Event{ + Type: events.TypeStateSnapshot, + DeviceID: "d1", + Payload: map[string]any{"devices": []any{"a", "b", "c"}, "plugins": []any{"Battery", "MPRIS"}}, + } + got := formatEvent(ev) + if got != "[d1] state.snapshot\n" { + t.Errorf("state.snapshot rendered with payload: %q", got) + } +} + +// A malformed or missing payload must degrade to the bare line, never panic and +// never print Go syntax like map[...] or for a whole struct. +func TestFormatEventSurvivesBadPayloads(t *testing.T) { + types := []events.EventType{ + events.TypeBatteryUpdate, events.TypeBatteryThreshold, events.TypeNotification, + events.TypeSftpMount, events.TypePairAccepted, events.TypeMprisUpdate, + events.TypeSMSIncoming, events.TypeSMSAttachment, events.TypePingReceived, + events.TypeConnectivityUpdate, events.TypeTelephonyRinging, events.TypeTelephonyMissed, + events.TypeTelephonyTalking, events.TypeTelephonyCanceled, events.TypeVolumeUpdate, + events.TypeContactsUpdated, events.TypeDeviceAdded, events.TypeDeviceConnected, + } + + for _, typ := range types { + t.Run(string(typ), func(t *testing.T) { + for _, payload := range []any{nil, "a bare string", 42, []any{1, 2}} { + got := formatEvent(events.Event{Type: typ, DeviceID: "d1", Payload: payload}) + if !strings.HasPrefix(got, "[d1] ") { + t.Fatalf("payload %#v: line does not start with the device prefix: %q", payload, got) + } + if !strings.HasSuffix(got, "\n") { + t.Fatalf("payload %#v: line is not newline terminated: %q", payload, got) + } + } + }) + } +} + +func TestTruncateAndOneLine(t *testing.T) { + if got := truncate("hello", 10); got != "hello" { + t.Errorf("short string was altered: %q", got) + } + if got := truncate("hello world", 8); got != "hello w…" { + t.Errorf("truncate = %q, want %q", got, "hello w…") + } + if got := oneLine("a\n b\tc "); got != "a b c" { + t.Errorf("oneLine = %q, want %q", got, "a b c") + } +} diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index e831e6d..58415ad 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -91,23 +91,56 @@ Connected steady state (all pairs connected, nothing playing, no transfers, no p | mDNS browse | Same owned lifetime as broadcast (probes are periodic by library design) | | mDNS advertise | Lifetime-on, responder-only (no timers) | | UDP/TCP/IPC listeners, D-Bus signals, bus subscriptions | Blocking waits, zero CPU until an event arrives | -| Local position poller | Exists only while ≥1 local player `IsPlaying` (`[mpris] poll_while_playing`, `position_interval`); ticks re-broadcast only on metadata change or position drift >3s off the anchor extrapolation | -| MPRIS watchdog | 10s re-check, alive only while ≥1 local player is tracked; restarts the position poller if it finds unpolled playback (see below) | +| Local position poller | Exists only while ≥1 local player `IsPlaying` (`[mpris] poll_while_playing`, `position_interval`); armed by D-Bus signals, ticks re-broadcast only on metadata change or position drift >3s off the anchor extrapolation | | Remote state poller | Ticker itself exists only while a client subscribes to `mpris.update` (bus subscriber-change hook starts/stops it) | | Reconnect redial | Parked on discovery sightings; fallback escalates to `fallback_max`, then gives up past `stale_after` until the next sighting | | TCP keepalive | Kernel probes, first delay `[network] keepalive_idle` (default 30s, minimum 10s) | -### The one deliberate exception: the MPRIS watchdog +There are no standing timers. The position poller is armed by D-Bus signals and self-stops on a confirmed pause, so a tracked-but-paused player costs zero wakeups — exactly like an untracked one. -The position poller arms on an observed state change and stops as soon as a live read confirms nothing is playing. Restarting it therefore depends on a D-Bus `PlaybackStatus` signal arriving — and a missed signal strands the poller for the rest of the session, leaving the phone's now-playing frozen while audio plays. Firefox's MPRIS endpoint answers intermittently, so a dropped edge is routine rather than a corner case. +Measured 2026-09-25 (phone connected, Firefox playing): 7.5 CPU ticks/min, 0 voluntary context switches. Idle with no MPRIS player tracked: 0.00 CPU ticks/min, 0 `GetAll`/min. Methodology note: Go timers are runtime-managed (no timerfds to count) and `ptrace` is restricted by Yama, so `/proc` CPU deltas + `dbus-monitor` call rates are the working proxies. -So while any local player is tracked, a 10s watchdog samples live state and re-arms the poller if it finds unpolled playback. This bounds the stale window regardless of signal reliability. +### The signal path is load-bearing -The cost: with an MPRIS application open but paused, the daemon performs one D-Bus read per 10s. **"Zero timers at idle" means zero when no MPRIS player is tracked**, not zero on a desktop with a media player merely running. This is a deliberate trade — a bounded ~6 reads/min beats an unbounded frozen now-playing display. +The poller has no self-healing timer, so every arm and disarm rides on D-Bus signal delivery. Two details are easy to get wrong and both have already caused a total outage: -Measured 2026-09-25 (phone connected, Firefox playing): 7.5 CPU ticks/min, 0 voluntary context switches. Idle with no MPRIS player tracked: 0.00 CPU ticks/min, 0 `GetAll`/min. Methodology note: Go timers are runtime-managed (no timerfds to count) and `ptrace` is restricted by Yama, so `/proc` CPU deltas + `dbus-monitor` call rates are the working proxies. +- **godbus reports `Signal.Name` as `.`, not the bare member.** Dispatch compares the fully qualified name (`org.mpris.MediaPlayer2.Player.Seeked`); match rules passed to `AddMatchSignal` use the bare member (`Seeked`). The two spellings are not interchangeable, and matching the wrong one silently drops every signal with no error anywhere. +- **Match rules use `WithMatchInterface` + `WithMatchMember`, never the unfiltered rule.** The fully qualified name is not a legal value for `WithMatchMember`, and the daemon closes with an invalid-argument error, which surfaces as a crash-looping watcher rather than a dropped signal. + +`classifySignal` in `internal/plugins/mpris/signals.go` is the single place that maps a signal name to a handler, and it is table-tested against both the qualified and the bare spellings. Keep it that way; a `switch` inline in the watcher loop cannot be unit-tested and that is precisely how the bare-name bug survived so long. + +**Contributor invariant: new periodic work must be owner-gated or activity-gated, never standing.** A ticker that fires while nothing is happening is a bug — gate it on owners (discovery), playback state (MPRIS), subscribers (remote refresh), or sightings (reconnect). Do not add a slow self-healing timer to paper over unreliable event delivery: fix the event path instead. An earlier build carried a 10s watchdog that re-armed the poller on missed signals, which cost ~6 D-Bus reads/min on every desktop with a paused media player. It was removed once signal routing was fixed and measured to never fire. + +### Known optimization: the position poller's 2s tick + +Accepted for now, not because it is required, but because it is cheap and correct. While a player plays, the poller issues one `GetAll` every 2s — measured 2026-09-27 at **28.5 `GetAll`/min, 6 CPU ticks/min (~0.1% of one core)**. Paused and untracked both cost zero. + +Seeking is **not** polled. `Seeked` is handled as a free D-Bus signal that updates `pos` directly, so scrubbing costs nothing regardless of this interval. The 2s ticker exists for one reason only: the phone extrapolates from `posAnchorMs`, and that extrapolation needs periodic correction against real metadata and position. + +That makes the poll pure drift correction, which is why it is the right next target. Three options, cheapest first: + +1. **Owner-gate it on subscribers.** The remote poller already does this — a bus subscriber-change hook starts and stops its ticker. The local poller runs unconditionally even when no client is watching local now-playing. Gating it the same way costs nothing when nobody is looking and changes nothing when someone is. Preferred, because it reuses a pattern already in the codebase rather than introducing a new policy. +2. **Stretch `position_interval`.** Drift tolerance is already 3s, so extrapolation can cover a wider gap; 2s → 10s drops the rate to ~6/min at the cost of scrubber latency on the phone. +3. **Stop trusting `GetAll` for `Position`.** Worth verifying what Firefox actually returns there. If it is stale or zero, the code already issues a second `Get` in `queryPositionAndCanSeek` (`state.go`), meaning each tick performs two round-trips to obtain one number. Fixing that halves the rate with no behavioural change. -**Contributor invariant: new periodic work must be owner-gated or activity-gated, never standing.** A ticker that fires while nothing is happening is a bug — gate it on owners (discovery), playback state (MPRIS), subscribers (remote refresh), or sightings (reconnect). Where a signal-driven design cannot be made reliable, add a slow self-healing check scoped to the thing it watches, and document it here as a known cost rather than quietly reintroducing a standing timer. +Any of these should be measured with the same method before and after: `dbus-monitor` call rate for round-trips, `/proc//stat` deltas for CPU, and `voluntary_ctxt_switches` to confirm the ticker actually stopped rather than merely slowing down. + +### SMS push: arming is a one-way door + +The phone suppresses every SMS push until the desktop sends `request_conversations` or `request_conversation` once, which sets the plugin's `haveMessagesBeenRequested` flag. **There is no packet that clears it.** A phone that has been armed keeps pushing for the rest of its app's lifetime, whether or not anyone still wants to hear it. + +That asymmetry drives the whole design, and it is the part that is easy to get wrong: + +- **Arming is gated, not the notifications alone.** `armed()` is true when `[sms] always_arm` is set or a client subscribes to `sms.incoming`. `always_arm` is **off by default**, because an armed phone keeps streaming for the rest of its app's lifetime and opting in should be deliberate. So the bus subscriber-change hook is the default opt-in, and the phone is only asked while somebody is watching. +- **The arm time is a second, independent gate.** The reply to `request_conversations` is one `kdeconnect.sms.messages` packet *per thread* carrying that thread's head message, so an ungated notify would fire one desktop popup per existing conversation on every connect. `shouldNotify` drops anything older than the arm time. +- **Notifications stop when clients do, even though packets do not.** Because the ratchet cannot be undone, a single transient `kcd watch` would otherwise silently become a permanent notifier. The armed check in `shouldNotify` is what makes the residual stream a no-op. +- **Arming is idempotent per connection.** `watch` reconnects with backoff and re-subscribes each time; without the `armedAt` guard every reconnect would re-trigger a full conversation-head burst. +- **The arming goroutine must stay off the hook.** `bus.Subscribe` invokes hooks inline, and `dev.Send` can block for up to `writeTimeout` (10s) when a peer's send channel is full, so `syncArming` collects targets under the lock and sends in a goroutine. +- **A reconnect re-arms**, because that is the only recovery available after the phone's own app restarts and resets the flag. A phone-side restart mid-session is therefore the one case that silently stops push until the connection drops. +- **The phone emits empty batches.** Its content observer fires on any SMS database change with no empty guard of its own, so `handleMessages` returns early on a zero-length batch. +- **Message bodies never reach the logger**, only the event bus and the notification text. Journals are routinely collected and shipped off-box, and a message is the most sensitive thing this daemon handles. `TestMessageBodyNeverLogged` guards it. + +Known gap, accepted: the phone applies its blocked-numbers list only on the deprecated `kdeconnect.telephony` push, not on the content-observer path, so blocked senders can still arrive over `kdeconnect.sms.messages`. --- @@ -173,17 +206,17 @@ A `Device` wraps an active TCP connection with: ```mermaid stateDiagram-v2 direction TB - + [*] --> Unpaired - + Unpaired --> PairRequested : Local initiates PairRequested --> Paired : Peer accepts PairRequested --> Unpaired : Peer rejects / timeout - + Unpaired --> PairRequestedByPeer : Peer initiates PairRequestedByPeer --> Paired : Local accepts PairRequestedByPeer --> Unpaired : Local rejects - + Paired --> Unpaired : Unpair ``` @@ -257,7 +290,7 @@ goroutine leak when the child wedges. | `notification` | `kdeconnect.notification` | Downloads icon payload over TLS side-channel; per-app filter via `SetFilters()`; `notify-send --help` probe for `--print-id` support; `tlsConfig` + `logger` required in constructor | | `pair` | `kdeconnect.pair` | Manages the pairing handshake and certificate fingerprint verification | | `ping` | `kdeconnect.ping` | Fires `ping.received`; can be sent outbound | -| `runcommand` | `kdeconnect.runcommand` | Executes commands from the `[commands]` config table | +| `runcommand` | `kdeconnect.runcommand`, `kdeconnect.runcommand.output` | Executes commands from the `[commands]` config table; results stream to the phone's output card via `runcommand.output` (`commandStarted` → batched `commandOutput` → `commandFinished`, all sharing one 32-bit id). A capped notification is still sent as a fallback. A 15s bound per execution; `{"stop":true}` from the phone cancels it, as does disconnect. | | `sms` | `kdeconnect.sms.messages`, `kdeconnect.sms.attachment_file` | Sends `kdeconnect.sms.request`, `kdeconnect.sms.request_conversations`, `kdeconnect.sms.request_conversation`, `kdeconnect.sms.request_attachment` | | `sftp` | `kdeconnect.sftp` | Parses `multiPaths`, `pathNames`, and `errorMessage` from the phone's response. `Info()` returns cached credentials + `StorageVolume` slices; `Volumes()` lists storage roots with human-readable names. `Handle()` logs errors when the phone returns `errorMessage` (e.g. missing storage permission). Mounts at server root to avoid chroot double-path bug; tracks mounts in `mountPoints` map; `Unmount()` calls `fusermount3`/`fusermount` | | `share` | `kdeconnect.share.request` | Streaming file receive + URL/text handling; fires progress events | diff --git a/docs/CLI.md b/docs/CLI.md index 0f30242..77d392e 100644 --- a/docs/CLI.md +++ b/docs/CLI.md @@ -357,14 +357,19 @@ If `device-id` is omitted, `kcd` automatically targets the first paired and conn **Example output** ``` -LTE [███░] (3/4) +SIM 0: LTE ●●●○ ``` -Dual-SIM phones print one line per SIM (`SIM 0: …`), primary first. +Dual-SIM phones print one line per SIM (`SIM 0: …`, `SIM 1: …`), primary first. +The network label is padded to eight columns so the dot bars line up. A SIM +that is registered but has no usable signal reads `No service` rather than +showing an empty bar, and an `unknown` or absent network type reads `Cellular`. + `--json` prints the raw report (same shape as `connectivity.update` event payloads) for scripting. Exits non-zero with `no connectivity data` when -the device is offline or never reported — reports are requested fresh on -every connect. +the device is offline or has not reported yet. The phone pushes a report +whenever its signal state changes; the daemon does not ask, because Android's +plugin cannot receive a request. > For continuous monitoring, use `kcd watch --events=connectivity.update` instead. @@ -1027,6 +1032,39 @@ kcd watch [--events ] [--json] | `sms.attachment` | MMS attachment downloaded: `{filename, path, thread_id}` | | `ring.received` | Phone wants this PC to ring | +Text mode is written for people, not for parsing. Use `--json` if you need the +raw payload; `kcd watch` lines are not a stable interface. + +| Event | Rendered as | +|---|---| +| `battery.update` | `battery: 62% (charging: false)` | +| `battery.threshold` | `battery low: 15%` (` (charging)` only while charging) | +| `notification` | `notification: WhatsApp - Alice` | +| `share.progress` | `transfer: a.png... 10/20 bytes` (rewritten in place) | +| `share.complete` | `transfer complete: a.png` | +| `mpris.update` | `▶ Firefox - Song (Band)` | +| `sms.incoming` | `sms from +1555: hello` | +| `sms.attachment` | `sms attachment saved: photo.jpg` | +| `ping.received` | `ping: pong` | +| `connectivity.update` | `connectivity: SIM 0: LTE ●●●○` | +| `telephony.ringing` | `incoming call: Bob (+15550001234)` | +| `telephony.talking` | `call answered: Bob (+15550001234)` | +| `telephony.missed` | `missed call: Bob (+15550001234)` | +| `telephony.canceled` | `call ended: Bob (+15550001234)` | +| `volume.update` | `volume: Speaker 42%` or `volume: Speaker (muted)` | +| `contacts.updated` | `contacts: 2 added, 1 updated, 5 pending` | +| `device.added` | `device.added: Pixel 8` | +| `device.connected` | `device.connected: Pixel 8 (phone)` | + +Four types intentionally stay on the bare `[] ` line. +`ring.received`, `device.removed` and `device.disconnected` carry no payload to +show, and `state.snapshot` carries a full device and plugin dump that would +flood the terminal. + +Long values (message bodies, contact names, file names) are truncated with an +ellipsis and collapsed onto one line, so a multi-line SMS cannot break the +one-event-per-line shape. + ### Examples **Watch everything, human-readable** @@ -1038,7 +1076,11 @@ kcd watch ``` [a1b2...] battery: 62% (charging: false) [a1b2...] notification: WhatsApp - Alice: "Hey, are you free?" -[a1b2...] telephony.ringing — Bob (+15550001234) +[a1b2...] incoming call: Bob (+15550001234) +[a1b2...] connectivity: SIM 0: LTE ●●●○ +[a1b2...] volume: Speaker 42% +[a1b2...] sms from +15550001234: Running about 10 minutes late +[a1b2...] contacts: 2 added, 1 updated, 5 pending ``` **Filter to battery and calls only** @@ -1047,6 +1089,22 @@ kcd watch kcd watch --events=battery.update,telephony.ringing ``` +**Receive SMS as they arrive** + +```bash +kcd watch --events=sms.incoming +``` + +``` +[a1b2...] sms from +15550001234: Running about 10 minutes late, order without me +[a1b2...] sms from +15550009999: Your code is 481920. Do not share it. +``` + +The phone pushes messages as they arrive, so this needs no polling and no +request command. Since `[sms] always_arm` is off by default, this +subscription is also what arms the push — nothing is asked of the phone +until it is running. + **Raw NDJSON for scripting** ```bash diff --git a/docs/CLIENT_GUIDE.md b/docs/CLIENT_GUIDE.md index ab96150..d652420 100644 --- a/docs/CLIENT_GUIDE.md +++ b/docs/CLIENT_GUIDE.md @@ -290,6 +290,16 @@ except KeyboardInterrupt: | `share.complete` | File transfer finished | | `mpris.update` | Now-playing state changed (deduplicated — only on real changes) | | `sms.incoming` | SMS/MMS received | + +> **SMS freshness:** the phone only pushes new messages after the daemon +> has asked once, which by default happens only while a client subscribes to +> `sms.incoming` — so `kcd watch --events sms.incoming` is all it takes to +> receive messages live. Set `[sms] always_arm = true` to ask on every +> connect instead and notify with no client attached. The phone cannot be un-asked, +> so a client that subscribes leaves the phone pushing afterwards; the +> daemon keeps publishing events in that case but stays quiet on desktop +> notifications. Messages that predate the ask arrive as a one-off burst of +> per-thread history and are not notified. | `contacts.updated` | Contacts sync progress (counts only; call `contacts_list` for data) | | `pair.requested` | Remote device wants to pair | | `ping.received` | Ping from device | @@ -324,11 +334,9 @@ except KeyboardInterrupt: > stays silent — the phone extrapolates from `posAnchorMs` — and a tick > re-broadcasts only on a metadata change or when the true position drifts > more than 3s off the extrapolation (seek, missed signal, clock drift). -> With no player running at all, nothing is polled — a silent desktop costs -> zero wakeups. While a player is merely paused, a 10s watchdog re-checks -> live state and restarts the poller if playback resumed without the daemon -> seeing the signal, so a dropped D-Bus edge cannot leave the phone's -> display frozen. Set +> A paused or absent player costs zero wakeups: D-Bus signals arm the poller +> on playback and it stops itself on a confirmed pause, so there is no +> background timer once nothing is playing. Set > `poll_while_playing = false` for pure event-driven mode (position then > extrapolates from `posAnchorMs` between D-Bus signals). @@ -769,6 +777,12 @@ Usage: `python3 monitor.py '["battery.update","mpris.update"]'` ### 9.3 GTK4/Shell Proxy +The plain-text output of `kcd watch` is formatted for people, not for parsing: +long values are truncated, zero counts are dropped, and the wording is prose +("incoming call", "battery low"). It is not a stable interface. Anything that +needs the full payload should use `--json`, whose event stream is the +documented contract in [`IPC_PROTOCOL.md §5`](IPC_PROTOCOL.md#5-event-types). + For desktop shell widgets (eww, ags, quickshell), run `kcd watch` in the background and pipe the JSON output to a named pipe or parse it directly: @@ -780,6 +794,14 @@ kcd watch --json '["battery.update","mpris.update"]' | while read -r line; do done ``` +To read events by eye instead: + +```bash +kcd watch # everything, rendered for a terminal +kcd watch --events=sms.incoming +kcd watch --events=connectivity.update +``` + --- ## 10. Going Further diff --git a/docs/IPC_PROTOCOL.md b/docs/IPC_PROTOCOL.md index 6700f9f..c3f1c90 100644 --- a/docs/IPC_PROTOCOL.md +++ b/docs/IPC_PROTOCOL.md @@ -283,9 +283,11 @@ shape as `connectivity.update` event payloads). | `signalStrength` | number | Level 0 (no signal) – 4 (full) | Errors: `device not found`, `connectivity plugin not enabled`, -`no connectivity data (device offline or never reported)`. Reports are -requested fresh on every connect; `kcd watch` also emits a cached -`connectivity.update` on subscribe so clients never boot blind. +`no connectivity data (device offline or never reported)`. The daemon +never requests a report: Android's connectivity plugin declares no incoming +packet types, so a request would be discarded. The phone pushes one whenever +its signal state changes and the daemon caches the last; `kcd watch` also +emits a cached `connectivity.update` on subscribe so clients never boot blind. #### `clipboard_push` @@ -1207,6 +1209,17 @@ An SMS or MMS message was received. } ``` +**Delivery:** the phone pushes these once the daemon has asked for messages, +which it does per connection when `[sms] always_arm` is set, or while a +client is subscribed to this event type. `always_arm` is off by default, so +subscribing is the opt-in — no request command is needed to receive +messages. The reply to that ask is a one-off burst of per-thread history +which is published but not notified; `type` 2 marks an outbound message the +phone echoes back. + +Note: `attachments` carries only descriptors. Fetch the bytes with +`sms_request_attachment`, then read `sms.attachment`. + #### `sms.attachment` An MMS attachment has been downloaded. @@ -1341,6 +1354,7 @@ who may want to implement a full network-level implementation. | `kdeconnect.notification.request` | Notification | Clear a notification on the phone (`{"cancel": ""}`) | | `kdeconnect.notification` | RunCommand | Command output notification pushed to phone | | `kdeconnect.runcommand` | RunCommand | Send command list to phone | +| `kdeconnect.runcommand.output` | RunCommand | Stream execution results to the phone's output card | | `kdeconnect.runcommand.request` | RunCommand | Request phone's command list / execute command | | `kdeconnect.share.request` | Share | File transfer invitation (side-channel) | | `kdeconnect.sftp.request` | SFTP | Request the phone to start its SFTP server | @@ -1390,7 +1404,8 @@ plugin processes it and a link to the body struct definition. | `kdeconnect.mpris` | MPRIS | `MPRISRequest{RequestPlayerList, RequestNowPlaying, RequestVolume, Player, Action, AlbumArtUrl, TransferringAlbumArt, ...}` — inbound packets with `transferringAlbumArt: true` + `payloadTransferInfo` carry album art bytes (side channel) that the daemon caches to `$XDG_CACHE_HOME/kcd/art/` | | `kdeconnect.mpris.request` | MPRIS | `MPRISRequest{}` (same struct, different semantics) — an outbound `kdeconnect.mpris.request` with `player` + `albumArtUrl` asks the phone to stream art back | | `kdeconnect.runcommand` | RunCommand | `{CommandList string}` — the phone's reply to a command-list request, holding a JSON object of label → `{name, command}` | -| `kdeconnect.runcommand.request` | RunCommand | `RequestBody{RequestCommandList bool, Key string}` | +| `kdeconnect.runcommand.request` | RunCommand | `RequestBody{RequestCommandList bool, Key string, Stop bool, ID int32}` — `Stop`+`ID` cancels a running execution | +| `kdeconnect.runcommand.output` | RunCommand | Execution results. One packet type, three shapes distinguished by which key is present: `{"commandStarted":true,"id":N,"command":"label"}`, `{"commandOutput":true,"id":N,"stdout":[...],"stderr":[...]}`, `{"commandFinished":true,"id":N,"success":bool}`. Order is mandatory — the phone registers a display row on `commandStarted` and keys every later packet for that execution to the same `id`. `id` is read with `getInt`, so it must fit a 32-bit int. Both `stdout` and `stderr` must be present on every `commandOutput` batch; the phone iterates both lists without a null check. | | `kdeconnect.presenter` | Presenter | `PresenterBody{Dx, Dy *float64, Stop *bool}` | | `kdeconnect.systemvolume` | RemoteSystemVolume | `VolumeBody{SinkList, Name, Volume, Muted}` | diff --git a/hk.pkl b/hk.pkl new file mode 100644 index 0000000..5b1da5f --- /dev/null +++ b/hk.pkl @@ -0,0 +1,56 @@ +amends "package://github.com/jdx/hk/releases/download/v2.4.0/hk@2.4.0#/Config.pkl" +import "package://github.com/jdx/hk/releases/download/v2.4.0/hk@2.4.0#/Builtins.pkl" +// Using a coding agent? See https://hk.jdx.dev/agents + +// Declaring top-level steps creates the implicit check, fix, and pre-commit +// hooks. `hk install` wires the git pre-commit hook to `hk run pre-commit`, +// which applies fixes and stages the result. +// +// These mirror what CI actually enforces, and every tool here is already +// available, so a clean checkout passes. +steps { + ["go_fmt"] = Builtins.go_fmt + ["go_vet"] = Builtins.go_vet + ["golangci_lint"] = Builtins.golangci_lint + ["golangci_lint_fmt"] = Builtins.golangci_lint_fmt + ["gomod_tidy"] = Builtins.gomod_tidy + ["trailing_whitespace"] = Builtins.trailing_whitespace +} + +// The repo ships shell scripts and GitHub workflows that nothing currently +// checks. These are worth having, but the tools are not installed and CI does +// not run them, so they stay opt-in: `hk check --profile shell`. +// +// Promote one to the default set once its tool is added to mise.toml and CI. +local extra = new Mapping { + ["shellcheck"] = (Builtins.shellcheck) { profiles = List("shell") } + ["shellharden"] = (Builtins.shellharden) { profiles = List("shell") } + ["shfmt"] = (Builtins.shfmt) { profiles = List("shell") } + ["actionlint"] = (Builtins.actionlint) { profiles = List("shell") } + ["zizmor"] = (Builtins.zizmor) { profiles = List("shell") } + ["go_vuln_check"] = (Builtins.go_vuln_check) { profiles = List("slow") } +} + +hooks { + ["check"] { + steps = extra + } + ["fix"] { + steps = extra + } + ["pre-commit"] { + steps = extra + } +} + +// The upstream KDE Connect trees are checked out for reference only. They are +// gitignored, so hk skips them already; excluding them keeps that true even if +// the ignore entry goes away. +exclude = List("kdeconnect-kde", "flake.lock") + +// Let hk size the pool from the machine instead of a fixed number. +jobs = 0 + +// Report every failure in one pass instead of stopping at the first, so a +// single missing tool does not hide the rest of the run. +fail_fast = false diff --git a/internal/config/plugins.go b/internal/config/plugins.go index 4f87180..69d302d 100644 --- a/internal/config/plugins.go +++ b/internal/config/plugins.go @@ -99,6 +99,10 @@ type MousepadConfig struct { type SMSConfig struct { // NotifyIncoming shows a desktop notification when an SMS is received. NotifyIncoming bool `toml:"notify_incoming"` + + // AlwaysArm asks the phone to push new SMS on every connect. The phone + // cannot be un-asked, so this is opt-in; see internal/plugins/sms. + AlwaysArm bool `toml:"always_arm"` } func (p *PluginConfig) Defaults() { @@ -186,4 +190,5 @@ func (c *MousepadConfig) Defaults() { func (c *SMSConfig) Defaults() { c.NotifyIncoming = true + c.AlwaysArm = false } diff --git a/internal/log/log_test.go b/internal/log/log_test.go index c2548db..310faa6 100644 --- a/internal/log/log_test.go +++ b/internal/log/log_test.go @@ -77,3 +77,20 @@ func TestSetLevelFilters(t *testing.T) { t.Errorf("surviving entry = %q, want %q", got, "kept") } } + +func TestObserveCapturesFieldsAndLevels(t *testing.T) { + l, snapshot := Observe() + l.Debug("d", String("body", "secret")) + l.Info("i", Int("n", 7)) + + entries := snapshot() + if len(entries) != 2 { + t.Fatalf("captured %d entries, want 2: %v", len(entries), entries) + } + if !strings.Contains(entries[0], "body=secret") { + t.Errorf("entry missing field: %q", entries[0]) + } + if !strings.Contains(entries[1], "n=7") { + t.Errorf("entry missing field: %q", entries[1]) + } +} diff --git a/internal/log/testing.go b/internal/log/testing.go index 9a99c8c..e60385c 100644 --- a/internal/log/testing.go +++ b/internal/log/testing.go @@ -1,10 +1,15 @@ package log import ( + "fmt" + "sort" + "strings" "testing" "go.uber.org/zap" + "go.uber.org/zap/zapcore" "go.uber.org/zap/zaptest" + "go.uber.org/zap/zaptest/observer" ) // NewDevelopment returns a human-readable logger for tests. @@ -17,3 +22,35 @@ func NewDevelopment() Logger { func NewTest(t *testing.T) Logger { return Logger{zap: zaptest.NewLogger(t, zaptest.WrapOptions(zap.AddCallerSkip(1))), level: zap.NewAtomicLevel()} } + +// Observe returns a debug-level Logger and a snapshot func rendering every +// captured entry as "message key=value ...". Tests outside this package need +// it because depguard keeps zap imports here, so they cannot build an +// observing core themselves. Use it to assert what does and does not reach +// the log, which is the only way to guard a "never log this" rule. +func Observe() (Logger, func() []string) { + level := zap.NewAtomicLevel() + level.SetLevel(zapcore.DebugLevel) + core, logs := observer.New(level) + l := Logger{zap: zap.New(core, zap.AddCallerSkip(1)), level: level} + + return l, func() []string { + entries := logs.All() + out := make([]string, 0, len(entries)) + for _, e := range entries { + var b strings.Builder + b.WriteString(e.Message) + fields := e.ContextMap() + keys := make([]string, 0, len(fields)) + for k := range fields { + keys = append(keys, k) + } + sort.Strings(keys) + for _, k := range keys { + fmt.Fprint(&b, " ", k, "=", fields[k]) + } + out = append(out, b.String()) + } + return out + } +} diff --git a/internal/plugins/connectivity/connectivity.go b/internal/plugins/connectivity/connectivity.go index 8015a54..a66581e 100644 --- a/internal/plugins/connectivity/connectivity.go +++ b/internal/plugins/connectivity/connectivity.go @@ -83,11 +83,13 @@ func connectivityEqual(a, b ConnectivityBody) bool { return true } -func (p *ConnectivityPlugin) OnConnect(dev device.Sender) { - // Request an immediate connectivity report upon connection - pkt, _ := protocol.NewPacket("kdeconnect.connectivity_report.request", map[string]interface{}{}) - dev.Send(pkt) -} +// OnConnect intentionally sends nothing. The phone's ConnectivityReportPlugin +// declares supportedPacketTypes = emptyArray() and its onPacketReceived returns +// false unconditionally, so kdeconnect.connectivity_report.request can never be +// processed on Android. It pushes reports on its own whenever its telephony +// state listener fires, so the only thing a connect-time request accomplished +// was a discarded packet. +func (p *ConnectivityPlugin) OnConnect(_ device.Sender) {} // Report returns the last connectivity report received from a device. // The second return value is false when the device never reported (or the diff --git a/internal/plugins/connectivity/connectivity_test.go b/internal/plugins/connectivity/connectivity_test.go index 0c086a8..e1a1231 100644 --- a/internal/plugins/connectivity/connectivity_test.go +++ b/internal/plugins/connectivity/connectivity_test.go @@ -4,6 +4,7 @@ import ( "context" "crypto/x509" "net" + "sync" "testing" "time" @@ -30,6 +31,27 @@ func (s testSender) HasCapability(string) bool { return false } func (s testSender) UpdateBattery(charge int, charging bool) {} func (s testSender) GetBattery() (int, bool) { return 0, false } +// countingSender records outbound packets so a no-op can be asserted rather +// than inferred from the absence of output. +type countingSender struct { + testSender + mu sync.Mutex + sent []*protocol.Packet +} + +func (s *countingSender) Send(p *protocol.Packet) error { + s.mu.Lock() + defer s.mu.Unlock() + s.sent = append(s.sent, p) + return nil +} + +func (s *countingSender) count() int { + s.mu.Lock() + defer s.mu.Unlock() + return len(s.sent) +} + func TestHandleDeduplicatesConnectivityReports(t *testing.T) { bus := events.NewBus(log.Nop()) sub := bus.Subscribe(4, events.TypeConnectivityUpdate) @@ -135,3 +157,20 @@ func TestReportCachesLastHandle(t *testing.T) { t.Error("expected miss after disconnect clears the cache") } } + +// The phone's ConnectivityReportPlugin declares supportedPacketTypes as empty +// and rejects every packet in onPacketReceived, so asking for a report is +// pointless. The phone pushes on its own when its telephony listener fires. +// This pins that OnConnect stays silent: without it, a reader has no way to tell +// a deliberate no-op from a forgotten request, and someone will re-add it. +func TestOnConnectSendsNothing(t *testing.T) { + bus := events.NewBus(log.Nop()) + plugin := NewConnectivityPlugin(bus) + dev := &countingSender{testSender: testSender{id: "device-1"}} + + plugin.OnConnect(dev) + + if n := dev.count(); n != 0 { + t.Errorf("OnConnect sent %d packets, want 0", n) + } +} diff --git a/internal/plugins/mpris/local.go b/internal/plugins/mpris/local.go index a74bf04..7ac00de 100644 --- a/internal/plugins/mpris/local.go +++ b/internal/plugins/mpris/local.go @@ -123,11 +123,6 @@ func (p *MPRISPlugin) removePlayer(displayName string) { delete(p.players, displayName) delete(p.lastTracks, displayName) delete(p.lastStates, displayName) - // With no players left there is nothing to poll, so stop the ticker - // and the watchdog rather than let them discover it on their own. - if len(p.players) == 0 { - p.stopWatchdogLocked() - } p.syncPlayingPollerLocked() p.mu.Unlock() diff --git a/internal/plugins/mpris/mpris.go b/internal/plugins/mpris/mpris.go index 0252613..14efefb 100644 --- a/internal/plugins/mpris/mpris.go +++ b/internal/plugins/mpris/mpris.go @@ -37,12 +37,11 @@ type MPRISPlugin struct { // ticker exists only while at least one local player IsPlaying. // pollCancel stops it; nil means no poller is running. pollGen // identifies the current poller so a self-stopping one does not - // clear its successor's handle. watchdogCancel stops the slow - // re-check that restarts a poller stranded by a missed signal. - mprisCfg config.MPRISConfig - pollCancel context.CancelFunc - pollGen uint64 - watchdogCancel context.CancelFunc + // clear its successor's handle. The poller is armed purely by + // observed D-Bus signals, so a paused desktop holds no timers. + mprisCfg config.MPRISConfig + pollCancel context.CancelFunc + pollGen uint64 // remotePollCancel stops the remote-state poller; nil means it is // not running. The poller is demand-driven (see syncRemotePoller): diff --git a/internal/plugins/mpris/playing_poller.go b/internal/plugins/mpris/playing_poller.go index a59005d..f5496c0 100644 --- a/internal/plugins/mpris/playing_poller.go +++ b/internal/plugins/mpris/playing_poller.go @@ -44,77 +44,6 @@ func (p *MPRISPlugin) syncPlayingPollerLocked() { return } p.armPlayingPollerLocked() - p.startWatchdogLocked() -} - -// startWatchdogLocked starts the slow re-check if it is not already -// running. Callers must hold p.mu. -func (p *MPRISPlugin) startWatchdogLocked() { - if p.watchdogCancel != nil { - return - } - ctx, cancel := context.WithCancel(p.watchCtx) - p.watchdogCancel = cancel - go p.runWatchdog(ctx) -} - -// stopWatchdogLocked stops the slow re-check. Callers must hold p.mu. -func (p *MPRISPlugin) stopWatchdogLocked() { - if p.watchdogCancel != nil { - p.watchdogCancel() - p.watchdogCancel = nil - } -} - -// runWatchdog restarts the position poller when it finds playback the -// signal path failed to announce. It stays parked — one read per -// interval — whenever the poller is already doing that job itself. -func (p *MPRISPlugin) runWatchdog(ctx context.Context) { - ticker := time.NewTicker(watchdogInterval) - defer ticker.Stop() - - for { - select { - case <-ctx.Done(): - return - case <-ticker.C: - p.mu.RLock() - armed := p.pollCancel != nil - tracked := len(p.players) - p.mu.RUnlock() - - if armed || tracked == 0 { - continue - } - if !p.anyPlayerPlaying() { - continue - } - p.logger.Debug("mpris: watchdog found unpolled playback, arming position poller") - p.mu.Lock() - p.armPlayingPollerLocked() - p.mu.Unlock() - } - } -} - -// anyPlayerPlaying reports whether any tracked player is playing right -// now. It reads live state without touching the cache or broadcasting — -// this only decides whether to restart the poller. -func (p *MPRISPlugin) anyPlayerPlaying() bool { - p.mu.RLock() - players := make([]*trackedPlayer, 0, len(p.players)) - for _, pl := range p.players { - players = append(players, pl) - } - p.mu.RUnlock() - - for _, pl := range players { - state, err := p.playerState(pl.displayName) - if err == nil && state.IsPlaying { - return true - } - } - return false } // armPlayingPollerLocked starts the position ticker if it is not already @@ -151,23 +80,6 @@ func (p *MPRISPlugin) stopPlayingPollerLocked() { // backstop for a name that lingers with a dead object behind it. const maxConsecutiveReadFailures = 5 -// watchdogInterval is how often the slow re-check looks for playback that -// the signal path missed. -// -// The poller stops the moment a live read reports nothing playing, and it -// only restarts when an observed change re-arms it. That makes it hostage -// to signal delivery: lose the PlaybackStatus=Playing edge and the poller -// stays down for the rest of the session, so the phone's now-playing -// freezes even though audio is playing. Firefox's MPRIS endpoint answers -// intermittently, which makes that a routine event, not a corner case. -// -// So the watchdog samples once per interval and restarts the poller if it -// finds something playing. 10s bounds how stale the phone's now-playing -// can get after a resume, at six reads per minute while an MPRIS app sits -// paused — still far below the 30/min the poller itself costs while -// playing, and it disappears entirely once no player is tracked. -const watchdogInterval = 10 * time.Second - // runPlayingPoller re-reads D-Bus state for every tracked player at // PositionInterval. It returns once live reads say nothing is playing, so // its lifetime follows the player rather than the cache that armed it. diff --git a/internal/plugins/mpris/playing_poller_test.go b/internal/plugins/mpris/playing_poller_test.go index d709c2b..4e37b94 100644 --- a/internal/plugins/mpris/playing_poller_test.go +++ b/internal/plugins/mpris/playing_poller_test.go @@ -148,48 +148,6 @@ func TestPollDisarmsOnPlayerRemoval(t *testing.T) { } } -// The watchdog exists because signal delivery is unreliable: the poller -// stops on a confirmed pause and must be able to come back on its own -// when playback resumes. It starts with the first tracked player and is -// torn down with the last, so a desktop with no MPRIS app keeps no -// timers at all. -func TestWatchdogLifecycleFollowsTrackedPlayers(t *testing.T) { - p := NewMPRISPlugin(nil, events.NewBus(log.Nop()), false, testMPRISConfig(), log.Nop()) - defer p.watchCancel() - - watchdogRunning := func() bool { - p.mu.RLock() - defer p.mu.RUnlock() - return p.watchdogCancel != nil - } - - if watchdogRunning() { - t.Fatal("watchdog running with zero tracked players") - } - - trackPlayer(p, "Nightdrive") - p.storeLocalState("Nightdrive", &NowPlaying{Player: "Nightdrive", IsPlaying: false}) - if !watchdogRunning() { - t.Fatal("watchdog not started alongside the first tracked player") - } - - trackPlayer(p, "Second FM") - p.storeLocalState("Second FM", &NowPlaying{Player: "Second FM", IsPlaying: false}) - if !watchdogRunning() { - t.Fatal("watchdog stopped while players remain") - } - - p.removePlayer("Nightdrive") - if !watchdogRunning() { - t.Fatal("watchdog stopped while one player remains") - } - - p.removePlayer("Second FM") - if watchdogRunning() { - t.Fatal("watchdog still running after the last player was removed") - } -} - // With PollWhilePlaying=false the poller must never arm (pure // event-driven mode). func TestPollNeverArmsWhenDisabled(t *testing.T) { diff --git a/internal/plugins/mpris/signals.go b/internal/plugins/mpris/signals.go index bbb5f1b..a8455f4 100644 --- a/internal/plugins/mpris/signals.go +++ b/internal/plugins/mpris/signals.go @@ -9,6 +9,47 @@ import ( "github.com/godbus/dbus/v5" ) +// signalKind identifies which handler owns an incoming D-Bus signal. +type signalKind int + +const ( + signalKindNone signalKind = iota + signalKindNameOwnerChanged + signalKindSeeked + signalKindPropertiesChanged +) + +func (k signalKind) String() string { + switch k { + case signalKindNameOwnerChanged: + return "NameOwnerChanged" + case signalKindSeeked: + return "Seeked" + case signalKindPropertiesChanged: + return "PropertiesChanged" + default: + return "None" + } +} + +// classifySignal routes a signal to its handler. +// +// godbus fills Signal.Name as ".", so these must be fully +// qualified. Match rules passed to AddMatchSignal use the bare member instead; +// the two spellings are not interchangeable. +func classifySignal(name string) signalKind { + switch name { + case "org.freedesktop.DBus.NameOwnerChanged": + return signalKindNameOwnerChanged + case "org.mpris.MediaPlayer2.Player.Seeked": + return signalKindSeeked + case "org.freedesktop.DBus.Properties.PropertiesChanged": + return signalKindPropertiesChanged + default: + return signalKindNone + } +} + func (p *MPRISPlugin) handleNameOwnerChanged(sig *dbus.Signal, conn *dbus.Conn, uniqueToDisplay map[string]string) { if len(sig.Body) < 3 { return diff --git a/internal/plugins/mpris/signals_test.go b/internal/plugins/mpris/signals_test.go new file mode 100644 index 0000000..e8c64d6 --- /dev/null +++ b/internal/plugins/mpris/signals_test.go @@ -0,0 +1,33 @@ +package mpris + +import "testing" + +func TestClassifySignalMatchesQualifiedNames(t *testing.T) { + tests := []struct { + name string + signal string + want signalKind + }{ + {"NameOwnerChanged", "org.freedesktop.DBus.NameOwnerChanged", signalKindNameOwnerChanged}, + {"Seeked", "org.mpris.MediaPlayer2.Player.Seeked", signalKindSeeked}, + {"PropertiesChanged", "org.freedesktop.DBus.Properties.PropertiesChanged", signalKindPropertiesChanged}, + + // godbus reports ".". Matching the bare member + // silently dropped every signal, so these must stay unrouted. + {"bare NameOwnerChanged", "NameOwnerChanged", signalKindNone}, + {"bare Seeked", "Seeked", signalKindNone}, + {"bare PropertiesChanged", "PropertiesChanged", signalKindNone}, + + {"other interface", "org.freedesktop.DBus.ObjectManager.InterfacesAdded", signalKindNone}, + {"player interface other member", "org.mpris.MediaPlayer2.Player.VolumeChanged", signalKindNone}, + {"empty", "", signalKindNone}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + if got := classifySignal(tc.signal); got != tc.want { + t.Errorf("classifySignal(%q) = %v, want %v", tc.signal, got, tc.want) + } + }) + } +} diff --git a/internal/plugins/mpris/watcher.go b/internal/plugins/mpris/watcher.go index 273b946..b4465d2 100644 --- a/internal/plugins/mpris/watcher.go +++ b/internal/plugins/mpris/watcher.go @@ -97,12 +97,12 @@ func (p *MPRISPlugin) runDBusWatcher(ctx context.Context) error { continue } - switch sig.Name { - case "org.freedesktop.DBus.NameOwnerChanged": + switch classifySignal(sig.Name) { + case signalKindNameOwnerChanged: p.handleNameOwnerChanged(sig, conn, uniqueToDisplay) - case "org.mpris.MediaPlayer2.Player.Seeked": + case signalKindSeeked: p.handleSeeked(sig, uniqueToDisplay) - case "org.freedesktop.DBus.Properties.PropertiesChanged": + case signalKindPropertiesChanged: p.handlePropertiesChanged(sig, uniqueToDisplay) } } diff --git a/internal/plugins/runcommand/output.go b/internal/plugins/runcommand/output.go new file mode 100644 index 0000000..c54a394 --- /dev/null +++ b/internal/plugins/runcommand/output.go @@ -0,0 +1,380 @@ +package runcommand + +import ( + "bufio" + "context" + "io" + "os/exec" + "strings" + "sync" + "time" + + "github.com/bethropolis/kcd/internal/device" + "github.com/bethropolis/kcd/internal/log" + "github.com/bethropolis/kcd/internal/protocol" +) + +// PacketTypeOutput carries command execution results to the phone's in-app +// output card. Its presence is what makes the daemon advertise the streaming +// path rather than relying on a desktop notification, which Android's +// ReceiveNotificationsPlugin ignores unless the user explicitly enables it. +const PacketTypeOutput = "kdeconnect.runcommand.output" + +const ( + // streamFlushInterval bounds how long a line waits before the phone sees + // it. It is the only knob trading latency against packet count: a chatty + // command emitting thousands of lines would otherwise produce one packet + // per line. + streamFlushInterval = 250 * time.Millisecond + + // streamMaxLines caps a single execution's output. The phone keeps every + // line in an unbounded Compose list, so the bound has to live here. + streamMaxLines = 2000 + + // streamMaxLineLen truncates a single pathological line (a minified + // blob, a base64 dump) before it reaches the wire. + streamMaxLineLen = 1024 + + // streamChanBuffer absorbs a burst between flushes. + streamChanBuffer = 256 + + // streamMaxScanBuffer is the largest single line the scanner will read. + // Anything longer is truncated by add rather than failing the scan. + streamMaxScanBuffer = 1024 * 1024 + + // streamNotifMaxBytes bounds the notification transcript. It is separate + // from the streamed cap because a notification is a summary, not a log. + streamNotifMaxBytes = 4000 +) + +// outputStream accumulates one execution's output and emits it as batched +// packets. Both stdout and stderr are sent on every batch, even when one is +// empty: the phone's handler iterates both lists without a null check. +type outputStream struct { + mu sync.Mutex + dev device.Sender + logger log.Logger + id int32 + command string + + stdout []string + stderr []string + total int + // notif accumulates a combined transcript for the notification fallback, + // which has a much smaller budget than the streamed output. + notif []string + notifBytes int + // dropped records that the line cap was hit, so the final batch can say so + // rather than silently ending early. + dropped bool + // stopped records that the execution was cancelled, so the finish packet + // reports failure even if the process happened to exit zero. + stopped bool +} + +func (s *outputStream) add(isStderr bool, line string) { + if len(line) > streamMaxLineLen { + line = line[:streamMaxLineLen] + "..." + } + + s.mu.Lock() + defer s.mu.Unlock() + + if s.total >= streamMaxLines { + s.dropped = true + return + } + s.total++ + if isStderr { + s.stderr = append(s.stderr, line) + } else { + s.stdout = append(s.stdout, line) + } + + // The notification has its own much smaller budget, so it tracks bytes + // rather than lines and stops growing once full. + if s.notifBytes < streamNotifMaxBytes { + s.notif = append(s.notif, line) + s.notifBytes += len(line) + 1 + } +} + +// summary renders the notification transcript, stdout and stderr interleaved +// in arrival order. +func (s *outputStream) summary() string { + s.mu.Lock() + defer s.mu.Unlock() + text := strings.Join(s.notif, "\n") + if s.dropped { + text += "\n...[output truncated]" + } + return text +} + +// pending reports whether anything is buffered, so the caller can skip +// sending an empty batch. +func (s *outputStream) pending() bool { + s.mu.Lock() + defer s.mu.Unlock() + return len(s.stdout) > 0 || len(s.stderr) > 0 +} + +// take swaps out the buffered lines, leaving empty slices behind. Both keys +// are always present in the packet even when one is empty. +func (s *outputStream) take() (stdout, stderr []string, dropped bool) { + s.mu.Lock() + defer s.mu.Unlock() + stdout, stderr, dropped = s.stdout, s.stderr, s.dropped + s.stdout, s.stderr = nil, nil + return stdout, stderr, dropped +} + +// send emits one commandOutput packet. +func (s *outputStream) send(stdout, stderr []string, dropped bool) { + if len(stdout) == 0 { + stdout = []string{} + } + if len(stderr) == 0 { + stderr = []string{} + } + if dropped { + stderr = append(stderr, "[output truncated: over 2000 lines]") + } + + body := map[string]any{ + "commandOutput": true, + "id": s.id, + "stdout": stdout, + "stderr": stderr, + } + pkt, err := protocol.NewPacket(PacketTypeOutput, body) + if err != nil { + s.logger.Warn("runcommand: build output packet", log.Error(err)) + return + } + if err := s.dev.Send(pkt); err != nil { + s.logger.Warn("runcommand: send output", log.Error(err), log.Int("id", int(s.id))) + } +} + +// flush sends a batch only when lines are buffered, so an idle command does not +// generate traffic. +func (s *outputStream) flush() { + if !s.pending() { + return + } + stdout, stderr, dropped := s.take() + s.send(stdout, stderr, dropped) +} + +// line is one scanned line tagged with the stream it came from. +type line struct { + text string + isStderr bool +} + +// streamOutput runs cmd and reports its output to the phone as it arrives. +// The three packet types must be sent in order -- started, then output, then +// finished -- because the phone keys its output rows off the id registered by +// commandStarted; a finished packet for an unknown id only flips its spinner. +// +// sendFinished is deferred so every exit path, including a panic, still +// releases the phone's running indicator. +func (p *RunCommandPlugin) streamOutput( + ctx context.Context, + dev device.Sender, + id int32, + key string, + cmd *exec.Cmd, +) string { + var stream *outputStream + // The finish packet reports the outcome, so it is always sent exactly + // once, including on the early-return paths below. + var success bool + defer func() { + if stream != nil && stream.stopped { + success = false + } + p.finishOutput(dev, id, success) + }() + + stdoutPipe, err := cmd.StdoutPipe() + if err != nil { + p.logger.Warn("runcommand: stdout pipe", log.Error(err)) + return "" + } + stderrPipe, err := cmd.StderrPipe() + if err != nil { + p.logger.Warn("runcommand: stderr pipe", log.Error(err)) + return "" + } + if err := cmd.Start(); err != nil { + p.logger.Warn("runcommand: start", log.Error(err)) + return "" + } + + stream = &outputStream{dev: dev, logger: p.logger, id: id, command: key} + + // exec.CommandContext kills only the direct child. If the shell forks + // rather than execs, the orphan keeps the write end of the pipes open, so + // the scanners would never see EOF and Wait would never be reached -- + // cancelling would hang instead of stopping. Closing our read ends when + // the context ends makes that unblock deterministically, whatever the + // shell does. + scanDone := make(chan struct{}) + defer close(scanDone) + go func() { + select { + case <-ctx.Done(): + _ = stdoutPipe.Close() + _ = stderrPipe.Close() + case <-scanDone: + } + }() + + lines := make(chan line, streamChanBuffer) + var readers sync.WaitGroup + readers.Add(2) + go func() { defer readers.Done(); scanInto(lines, stdoutPipe, false, p.logger) }() + go func() { defer readers.Done(); scanInto(lines, stderrPipe, true, p.logger) }() + go func() { + readers.Wait() + close(lines) + }() + + ticker := time.NewTicker(streamFlushInterval) + defer ticker.Stop() + + // The loop runs until both pipes hit EOF rather than returning on ctx. + // Cancelling the context kills the process, which closes the pipes, which + // ends the scanners. Tailing the channel to completion keeps the scanner + // goroutines from blocking forever on a send nobody is reading, and lets + // the final flush carry whatever output arrived before the kill. + for { + closed := drainLines(lines, stream) + stream.flush() + if closed { + break + } + + select { + case l, ok := <-lines: + if ok { + stream.add(l.isStderr, l.text) + } + case <-ticker.C: + stream.flush() + case <-ctx.Done(): + // The process is being killed; keep draining until EOF so the + // last lines are not lost, then report failure. + stream.stopped = true + } + } + + readers.Wait() + waitErr := cmd.Wait() + success = waitErr == nil + + return stream.summary() +} + +// drainLines consumes everything currently buffered and reports whether the +// channel has closed and gone empty. Draining first is what collapses a burst +// into one packet instead of one packet per line. +func drainLines(lines <-chan line, s *outputStream) (closed bool) { + for { + select { + case l, ok := <-lines: + if !ok { + return true + } + s.add(l.isStderr, l.text) + default: + return false + } + } +} + +// scanInto splits a pipe into lines. Sends block so a slow consumer applies +// backpressure to the command rather than silently discarding output; the +// line cap in add is the authoritative bound. +func scanInto(out chan<- line, r io.Reader, isStderr bool, logger log.Logger) { + scanner := bufio.NewScanner(r) + scanner.Buffer(make([]byte, 0, 64*1024), streamMaxScanBuffer) + for scanner.Scan() { + out <- line{text: scanner.Text(), isStderr: isStderr} + } + if err := scanner.Err(); err != nil { + logger.Debug("runcommand: output scan ended", log.Error(err)) + } +} + +// finishOutput emits commandFinished and clears the cancel registration. +// The phone uses this to turn the command line green or red, so it must be +// sent exactly once per started execution. +func (p *RunCommandPlugin) finishOutput(dev device.Sender, id int32, success bool) { + p.Mu.Lock() + if byID, ok := p.running[dev.ID()]; ok { + delete(byID, id) + if len(byID) == 0 { + delete(p.running, dev.ID()) + } + } + p.Mu.Unlock() + + body := map[string]any{ + "commandFinished": true, + "id": id, + "success": success, + } + pkt, err := protocol.NewPacket(PacketTypeOutput, body) + if err != nil { + p.logger.Warn("runcommand: build finished packet", log.Error(err)) + return + } + if err := dev.Send(pkt); err != nil { + p.logger.Warn("runcommand: send finished", log.Error(err), log.Int("id", int(id))) + } +} + +// nextExecID hands out execution ids. They must fit a Java int because the +// phone reads them with getInt; a nanosecond timestamp does not, so this is a +// plain counter. +func (p *RunCommandPlugin) nextExecID() int32 { + p.Mu.Lock() + defer p.Mu.Unlock() + p.execSeq++ + return p.execSeq +} + +// stopRunning cancels a running execution, if the phone asked for it. The +// phone's stop button sends kdeconnect.runcommand.request {"stop": true}. +func (p *RunCommandPlugin) stopRunning(deviceID string, id int32) { + p.Mu.Lock() + cancel := p.running[deviceID][id] + p.Mu.Unlock() + + if cancel == nil { + p.logger.Debug("runcommand: stop for unknown execution", + log.String("device_id", deviceID), log.Int("id", int(id))) + return + } + p.logger.Info("runcommand: stopping execution on request", + log.String("device_id", deviceID), log.Int("id", int(id))) + cancel() +} + +// stopAll cancels everything running for a device, used when it disconnects. +func (p *RunCommandPlugin) stopAll(deviceID string) { + p.Mu.Lock() + cancels := make([]context.CancelFunc, 0, len(p.running[deviceID])) + for _, cancel := range p.running[deviceID] { + cancels = append(cancels, cancel) + } + delete(p.running, deviceID) + p.Mu.Unlock() + + for _, cancel := range cancels { + cancel() + } +} diff --git a/internal/plugins/runcommand/output_test.go b/internal/plugins/runcommand/output_test.go new file mode 100644 index 0000000..e8a370f --- /dev/null +++ b/internal/plugins/runcommand/output_test.go @@ -0,0 +1,477 @@ +package runcommand + +import ( + "context" + "crypto/x509" + "encoding/json" + "net" + "sync" + "testing" + "time" + + "github.com/bethropolis/kcd/internal/device" + "github.com/bethropolis/kcd/internal/log" + "github.com/bethropolis/kcd/internal/protocol" +) + +// outputSender records every packet a run emits so ordering and payload shape +// can be asserted. The execution goroutine calls Send concurrently, so the +// recording is mutex-guarded. +type outputSender struct { + id string + mu sync.Mutex + sent []*protocol.Packet +} + +func (s *outputSender) ID() string { return s.id } +func (s *outputSender) Name() string { return "Test" } +func (s *outputSender) SetName(string) {} +func (s *outputSender) State() device.PairingState { return device.StatePaired } +func (s *outputSender) SetState(device.PairingState) {} +func (s *outputSender) IsConnected() bool { return true } +func (s *outputSender) RemoteIP() net.IP { return nil } +func (s *outputSender) PeerCert() *x509.Certificate { return nil } +func (s *outputSender) HasCapability(string) bool { return true } +func (s *outputSender) UpdateBattery(int, bool) {} +func (s *outputSender) GetBattery() (int, bool) { return 0, false } + +func (s *outputSender) Send(p *protocol.Packet) error { + s.mu.Lock() + defer s.mu.Unlock() + s.sent = append(s.sent, p) + return nil +} + +func (s *outputSender) packets() []*protocol.Packet { + s.mu.Lock() + defer s.mu.Unlock() + return append([]*protocol.Packet(nil), s.sent...) +} + +func (s *outputSender) outputs() []*protocol.Packet { + var out []*protocol.Packet + for _, p := range s.packets() { + if p.Type == PacketTypeOutput { + out = append(out, p) + } + } + return out +} + +// body decodes a captured packet's body into a generic map. +func body(t *testing.T, p *protocol.Packet) map[string]any { + t.Helper() + var m map[string]any + if err := json.Unmarshal(p.Body, &m); err != nil { + t.Fatalf("decode %s body: %v", p.Type, err) + } + return m +} + +// runKey sends a command request for key and waits for the plugin's background +// goroutines to finish. +func runKey(t *testing.T, p *RunCommandPlugin, dev *outputSender, key string) { + t.Helper() + pkt := &protocol.Packet{ + Type: "kdeconnect.runcommand.request", + Body: json.RawMessage(`{"key":"` + key + `"}`), + } + if err := p.Handle(context.Background(), dev, pkt); err != nil { + t.Fatalf("Handle: %v", err) + } + p.wg.Wait() +} + +func newTestPlugin(commands map[string]string) *RunCommandPlugin { + return NewRunCommandPlugin(commands, nil, log.Nop()) +} + +// The phone keys its output rows off the id registered by commandStarted, and +// a finished packet for an unknown id only flips its spinner. The three packet +// types must therefore arrive started -> output -> finished, all sharing one +// id that fits a Java int. +func TestOutputPacketOrdering(t *testing.T) { + p := newTestPlugin(map[string]string{"hi": "printf 'hello\\n'"}) + dev := &outputSender{id: "dev1"} + + runKey(t, p, dev, "hi") + + outs := dev.outputs() + if len(outs) < 2 { + t.Fatalf("got %d output packets, want started+output+finished", len(outs)) + } + + first := body(t, outs[0]) + if first["commandStarted"] != true { + t.Errorf("first packet is not commandStarted: %v", first) + } + if first["command"] != "hi" { + t.Errorf("commandStarted command = %v, want %q", first["command"], "hi") + } + + last := body(t, outs[len(outs)-1]) + if last["commandFinished"] != true { + t.Errorf("last packet is not commandFinished: %v", last) + } + + id, ok := first["id"].(float64) + if !ok { + t.Fatalf("id is not a number: %T", first["id"]) + } + if int64(id) > 2147483647 || int64(id) < -2147483648 { + t.Errorf("id %v does not fit a Java int", id) + } + if last["id"] != first["id"] { + t.Errorf("finished id %v != started id %v", last["id"], first["id"]) + } + if last["success"] != true { + t.Errorf("success = %v, want true for a command that worked", last["success"]) + } +} + +// getStringList returns null for a missing key and the phone iterates both +// lists unguarded, so every output batch must carry both keys even when one +// stream produced nothing. +func TestOutputBatchesAlwaysCarryBothStreams(t *testing.T) { + p := newTestPlugin(map[string]string{"out": "printf 'only stdout\\n'"}) + dev := &outputSender{id: "dev1"} + + runKey(t, p, dev, "out") + + var batches int + for _, pkt := range dev.outputs() { + m := body(t, pkt) + if m["commandOutput"] != true { + continue + } + batches++ + if _, ok := m["stdout"]; !ok { + t.Errorf("batch missing stdout key: %v", m) + } + if _, ok := m["stderr"]; !ok { + t.Errorf("batch missing stderr key: %v", m) + } + } + if batches == 0 { + t.Fatal("no commandOutput packet was sent") + } +} + +func TestOutputStreamsStdoutAndStderrSeparately(t *testing.T) { + p := newTestPlugin(map[string]string{"mix": "printf 'to out\\n'; printf 'to err\\n' >&2"}) + dev := &outputSender{id: "dev1"} + + runKey(t, p, dev, "mix") + + var stdout, stderr []string + for _, pkt := range dev.outputs() { + m := body(t, pkt) + if m["commandOutput"] != true { + continue + } + for _, v := range m["stdout"].([]any) { + stdout = append(stdout, v.(string)) + } + for _, v := range m["stderr"].([]any) { + stderr = append(stderr, v.(string)) + } + } + + if !contains(stdout, "to out") { + t.Errorf("stdout = %v, want it to contain %q", stdout, "to out") + } + if !contains(stderr, "to err") { + t.Errorf("stderr = %v, want it to contain %q", stderr, "to err") + } + if contains(stdout, "to err") { + t.Errorf("stderr line leaked into stdout: %v", stdout) + } +} + +// A command that prints nothing still has to be bracketed, otherwise the +// phone's spinner never clears. +func TestSilentCommandStillBracketsExecution(t *testing.T) { + p := newTestPlugin(map[string]string{"quiet": "true"}) + dev := &outputSender{id: "dev1"} + + runKey(t, p, dev, "quiet") + + outs := dev.outputs() + if len(outs) < 2 { + t.Fatalf("got %d output packets, want started+finished", len(outs)) + } + if body(t, outs[0])["commandStarted"] != true { + t.Error("missing commandStarted") + } + if body(t, outs[len(outs)-1])["commandFinished"] != true { + t.Error("missing commandFinished") + } +} + +func TestFailedCommandReportsFailure(t *testing.T) { + p := newTestPlugin(map[string]string{"bad": "exit 3"}) + dev := &outputSender{id: "dev1"} + + runKey(t, p, dev, "bad") + + outs := dev.outputs() + if len(outs) == 0 { + t.Fatal("no output packets") + } + last := body(t, outs[len(outs)-1]) + if last["success"] != false { + t.Errorf("success = %v, want false for a failing command", last["success"]) + } +} + +// The line cap is the only bound protecting the phone's unbounded output list. +func TestOutputLineCapTruncates(t *testing.T) { + p := newTestPlugin(map[string]string{"spam": "seq 1 5000"}) + dev := &outputSender{id: "dev1"} + + runKey(t, p, dev, "spam") + + var lines int + var sawTruncation bool + for _, pkt := range dev.outputs() { + m := body(t, pkt) + if m["commandOutput"] != true { + continue + } + for _, v := range m["stderr"].([]any) { + lines++ + if s, ok := v.(string); ok && s != "" && s[len(s)-1:] == "]" { + sawTruncation = true + } + } + } + if lines == 0 { + t.Fatal("no output reached the phone") + } + if lines > streamMaxLines+10 { + t.Errorf("sent %d lines, want the cap of %d to hold", lines, streamMaxLines) + } + if !sawTruncation { + t.Error("truncation was not reported to the phone") + } +} + +// Handle must return before the command finishes, or it stalls the whole +// per-device read loop. +func TestHandleDoesNotBlockOnCommand(t *testing.T) { + p := newTestPlugin(map[string]string{"slow": "sleep 2; printf done"}) + dev := &outputSender{id: "dev1"} + + pkt := &protocol.Packet{ + Type: "kdeconnect.runcommand.request", + Body: json.RawMessage(`{"key":"slow"}`), + } + + done := make(chan struct{}) + go func() { + defer close(done) + if err := p.Handle(context.Background(), dev, pkt); err != nil { + t.Errorf("Handle: %v", err) + } + }() + + select { + case <-done: + case <-time.After(500 * time.Millisecond): + t.Fatal("Handle blocked on command execution") + } + p.wg.Wait() +} + +// The phone's stop button sends {"stop": true} with an execution id. +func TestStopRequestCancelsExecution(t *testing.T) { + p := newTestPlugin(map[string]string{"long": "sleep 30"}) + dev := &outputSender{id: "dev1"} + + pkt := &protocol.Packet{ + Type: "kdeconnect.runcommand.request", + Body: json.RawMessage(`{"key":"long"}`), + } + if err := p.Handle(context.Background(), dev, pkt); err != nil { + t.Fatalf("Handle: %v", err) + } + + // Wait for the execution to register before stopping it. + var id int32 + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + p.Mu.RLock() + for candidate := range p.running["dev1"] { + id = candidate + } + p.Mu.RUnlock() + if id != 0 { + break + } + time.Sleep(10 * time.Millisecond) + } + if id == 0 { + p.wg.Wait() + t.Fatal("execution never registered") + } + + stopBody, err := json.Marshal(map[string]any{"stop": true, "id": id}) + if err != nil { + t.Fatalf("marshal stop: %v", err) + } + stop := &protocol.Packet{Type: "kdeconnect.runcommand.request", Body: stopBody} + if err := p.Handle(context.Background(), dev, stop); err != nil { + t.Fatalf("stop Handle: %v", err) + } + + done := make(chan struct{}) + go func() { p.wg.Wait(); close(done) }() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("stop did not cancel the execution") + } + + p.Mu.RLock() + remaining := len(p.running["dev1"]) + p.Mu.RUnlock() + if remaining != 0 { + t.Errorf("%d executions still registered after stop", remaining) + } +} + +func TestOnDisconnectCancelsRunning(t *testing.T) { + p := newTestPlugin(map[string]string{"long": "sleep 30"}) + dev := &outputSender{id: "dev1"} + + pkt := &protocol.Packet{ + Type: "kdeconnect.runcommand.request", + Body: json.RawMessage(`{"key":"long"}`), + } + if err := p.Handle(context.Background(), dev, pkt); err != nil { + t.Fatalf("Handle: %v", err) + } + + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + p.Mu.RLock() + n := len(p.running["dev1"]) + p.Mu.RUnlock() + if n > 0 { + break + } + time.Sleep(10 * time.Millisecond) + } + + p.OnDisconnect(dev) + + done := make(chan struct{}) + go func() { p.wg.Wait(); close(done) }() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("disconnect did not cancel the running execution") + } +} + +// The output packet type must be advertised, or the phone will not expect it. +func TestOutputTypeIsAdvertised(t *testing.T) { + p := newTestPlugin(nil) + var found bool + for _, typ := range p.OutgoingTypes() { + if typ == PacketTypeOutput { + found = true + } + } + if !found { + t.Errorf("OutgoingTypes = %v, want it to include %q", p.OutgoingTypes(), PacketTypeOutput) + } +} + +// The notification fallback stays: it is the only channel for users who have +// not enabled the phone's in-app output card. +func TestNotificationFallbackStillSent(t *testing.T) { + p := newTestPlugin(map[string]string{"hi": "printf 'fallback text\\n'"}) + dev := &outputSender{id: "dev1"} + + runKey(t, p, dev, "hi") + + var notif *protocol.Packet + for _, pkt := range dev.packets() { + if pkt.Type == "kdeconnect.notification" { + notif = pkt + } + } + if notif == nil { + t.Fatal("no notification packet was sent") + } + m := body(t, notif) + if ticker, _ := m["ticker"].(string); !contains([]string{ticker}, "fallback text") { + t.Errorf("notification ticker = %q, want the command output", ticker) + } +} + +func contains(haystack []string, needle string) bool { + for _, s := range haystack { + if s == needle { + return true + } + } + return false +} + +// A command that forks a child and waits is the shape of any real workload +// that spawns a subprocess. exec.CommandContext kills only the shell, so the +// orphan keeps the pipe write end open: without closing our read ends on +// cancellation the scanners never see EOF and the plugin's WaitGroup never +// drains. This is the case that hung in CI, where /bin/sh forks where the +// local one happens to exec. +func TestStopCancelsCommandThatForksChildren(t *testing.T) { + p := newTestPlugin(map[string]string{"fork": "{ sleep 30 & wait; }"}) + dev := &outputSender{id: "dev1"} + + pkt := &protocol.Packet{ + Type: "kdeconnect.runcommand.request", + Body: json.RawMessage(`{"key":"fork"}`), + } + if err := p.Handle(context.Background(), dev, pkt); err != nil { + t.Fatalf("Handle: %v", err) + } + + var id int32 + deadline := time.Now().Add(3 * time.Second) + for time.Now().Before(deadline) { + p.Mu.RLock() + for candidate := range p.running["dev1"] { + id = candidate + } + p.Mu.RUnlock() + if id != 0 { + break + } + time.Sleep(10 * time.Millisecond) + } + if id == 0 { + p.wg.Wait() + t.Fatal("execution never registered") + } + + stopBody, err := json.Marshal(map[string]any{"stop": true, "id": id}) + if err != nil { + t.Fatalf("marshal stop: %v", err) + } + if err := p.Handle(context.Background(), dev, &protocol.Packet{ + Type: "kdeconnect.runcommand.request", + Body: stopBody, + }); err != nil { + t.Fatalf("stop Handle: %v", err) + } + + done := make(chan struct{}) + go func() { p.wg.Wait(); close(done) }() + select { + case <-done: + case <-time.After(10 * time.Second): + t.Fatal("stop did not cancel a forking command") + } +} diff --git a/internal/plugins/runcommand/runcommand.go b/internal/plugins/runcommand/runcommand.go index 628e2fe..eb47cbc 100644 --- a/internal/plugins/runcommand/runcommand.go +++ b/internal/plugins/runcommand/runcommand.go @@ -4,13 +4,13 @@ import ( "context" "encoding/json" "fmt" + "os/exec" "strings" "sync" "time" "github.com/bethropolis/kcd/internal/device" "github.com/bethropolis/kcd/internal/log" - "github.com/bethropolis/kcd/internal/plugin" "github.com/bethropolis/kcd/internal/protocol" ) @@ -22,8 +22,16 @@ type RunCommandPlugin struct { // pendingLists holds one buffered waiter per device awaiting that // device's command-list reply, keyed by device ID. pendingLists map[string]chan []Command - logger log.Logger - wg sync.WaitGroup // exported for tests to synchronize with background goroutines + // running holds the cancel func of every live execution, keyed by device + // ID then execution id, so the phone's stop button can reach it. Entries + // are removed when the execution finishes. + running map[string]map[int32]context.CancelFunc + // execSeq hands out execution ids. It is a counter rather than a + // timestamp because the phone reads the id with getInt, which would + // overflow on a nanosecond value. + execSeq int32 + logger log.Logger + wg sync.WaitGroup // exported for tests to synchronize with background goroutines } func NewRunCommandPlugin(commands map[string]string, commandsPerDevice map[string]map[string]string, logger log.Logger) *RunCommandPlugin { @@ -34,6 +42,7 @@ func NewRunCommandPlugin(commands map[string]string, commandsPerDevice map[strin Commands: commands, CommandsPerDevice: commandsPerDevice, pendingLists: make(map[string]chan []Command), + running: make(map[string]map[int32]context.CancelFunc), logger: logger.With(log.String("plugin", "runcommand")), } } @@ -42,6 +51,10 @@ func NewRunCommandPlugin(commands map[string]string, commandsPerDevice map[strin type RequestBody struct { RequestCommandList bool `json:"requestCommandList,omitempty"` Key string `json:"key,omitempty"` + // Stop cancels a running execution. The phone's stop button sends this + // with the execution id. + Stop bool `json:"stop,omitempty"` + ID int32 `json:"id,omitempty"` } // Name returns the plugin name. @@ -57,9 +70,11 @@ func (p *RunCommandPlugin) IncomingTypes() []string { return []string{"kdeconnect.runcommand.request", "kdeconnect.runcommand"} } -// OutgoingTypes returns the packet types this plugin may send. +// OutgoingTypes returns the packet types this plugin may send. The output type +// is what tells the phone it can render results in its in-app output card +// rather than relying on a notification it ignores by default. func (p *RunCommandPlugin) OutgoingTypes() []string { - return []string{"kdeconnect.runcommand", "kdeconnect.notification"} + return []string{"kdeconnect.runcommand", PacketTypeOutput, "kdeconnect.notification"} } // Handle processes incoming command requests. @@ -74,6 +89,12 @@ func (p *RunCommandPlugin) Handle(ctx context.Context, dev device.Sender, pkt *p return err } + // The phone's stop button sends {"stop": true} with the execution id. + if body.Stop { + p.stopRunning(dev.ID(), body.ID) + return nil + } + if body.RequestCommandList { p.Mu.RLock() cmds := p.Commands @@ -120,26 +141,37 @@ func (p *RunCommandPlugin) Handle(ctx context.Context, dev device.Sender, pkt *p return nil } - // Handlers must not block. Spawning goroutine to run the command - // and optionally send a notification with the output. + // Handlers must not block. Spawning goroutine to run the command, + // stream its output to the phone's output card, and fall back to a + // notification for users who have not enabled the output card. + execID := p.nextExecID() p.wg.Add(1) go func() { defer p.wg.Done() execCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() - out, err := plugin.RunCommandSync(execCtx, "sh", "-c", cmdStr) + p.Mu.Lock() + if p.running[dev.ID()] == nil { + p.running[dev.ID()] = make(map[int32]context.CancelFunc) + } + p.running[dev.ID()][execID] = cancel + p.Mu.Unlock() + + // Registered before the command starts so the phone's stop + // button can reach an execution that has just been launched. + if err := p.sendStarted(dev, execID, body.Key); err != nil { + p.logger.Warn("failed to send command started", log.Error(err)) + } + + cmd := exec.CommandContext(execCtx, "sh", "-c", cmdStr) + text := p.streamOutput(execCtx, dev, execID, body.Key, cmd) - text := strings.TrimSpace(string(out)) + // Keep the notification path: it is the only channel for anyone + // who has not enabled the phone's in-app output card. + text = strings.TrimSpace(text) if len(text) == 0 { - if err != nil { - text = fmt.Sprintf("Error: %v", err) - } else { - // No output and no error — do not send a notification. - return - } - } else if err != nil { - text = fmt.Sprintf("Error: %v\n\n%s", err, text) + return } // Do not send notifications for massive outputs (e.g. log dumps) @@ -151,8 +183,6 @@ func (p *RunCommandPlugin) Handle(ctx context.Context, dev device.Sender, pkt *p // Send notification back to the phone. // The Android app uses 'appName' as the title and 'ticker' as the body. // It ignores 'title' and 'text'. - // Note: the Android ReceiveNotificationsPlugin is disabled by default. - // Users must enable "Receive notifications" in the device's plugin settings. notifBody := map[string]interface{}{ "id": fmt.Sprintf("%d", time.Now().UnixNano()), "appName": fmt.Sprintf("Run: %s", body.Key), @@ -177,7 +207,24 @@ func (p *RunCommandPlugin) Handle(ctx context.Context, dev device.Sender, pkt *p return nil } +// sendStarted announces a beginning execution. The phone registers the id +// against a display row here, so it must precede any output for that id. +func (p *RunCommandPlugin) sendStarted(dev device.Sender, id int32, command string) error { + pkt, err := protocol.NewPacket(PacketTypeOutput, map[string]any{ + "commandStarted": true, + "id": id, + "command": command, + }) + if err != nil { + return err + } + return dev.Send(pkt) +} + func (p *RunCommandPlugin) OnConnect(dev device.Sender) {} +// OnDisconnect cancels anything still running for the device, so a command +// does not outlive the connection that asked for it. func (p *RunCommandPlugin) OnDisconnect(dev device.Sender) { + p.stopAll(dev.ID()) } diff --git a/internal/plugins/sms/arming_test.go b/internal/plugins/sms/arming_test.go new file mode 100644 index 0000000..14364b2 --- /dev/null +++ b/internal/plugins/sms/arming_test.go @@ -0,0 +1,292 @@ +package sms + +import ( + "context" + "strings" + "sync" + "testing" + "time" + + "github.com/bethropolis/kcd/internal/config" + "github.com/bethropolis/kcd/internal/events" + "github.com/bethropolis/kcd/internal/log" + "github.com/bethropolis/kcd/internal/protocol" +) + +// syncCaptureSender is thread-safe: arming sends from a goroutine so the +// send races the test goroutine's read of the recorded packets. +type syncCaptureSender struct { + captureSender + mu sync.Mutex + sent []*protocol.Packet +} + +func (s *syncCaptureSender) Send(p *protocol.Packet) error { + s.mu.Lock() + defer s.mu.Unlock() + s.sent = append(s.sent, p) + return nil +} + +func (s *syncCaptureSender) types() []string { + s.mu.Lock() + defer s.mu.Unlock() + out := make([]string, 0, len(s.sent)) + for _, p := range s.sent { + out = append(out, p.Type) + } + return out +} + +func (s *syncCaptureSender) count(t string) int { + n := 0 + for _, got := range s.types() { + if got == t { + n++ + } + } + return n +} + +// armingSettled waits for the arming goroutine to drain. +func armingSettled(t *testing.T, dev *syncCaptureSender) { + t.Helper() + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if dev.count(PacketTypeSMSRequestConvs) > 0 { + return + } + time.Sleep(5 * time.Millisecond) + } +} + +func newArmingPlugin(t *testing.T, cfg config.SMSConfig) (*SMSPlugin, *events.Bus) { + t.Helper() + bus := events.NewBus(log.Nop()) + return NewSMSPlugin(cfg, bus, nil, log.Nop()), bus +} + +// Opting in arms the phone on connect. +func TestOnConnectArmsWhenAlwaysArmSet(t *testing.T) { + cfg := config.SMSConfig{AlwaysArm: true} + p, _ := newArmingPlugin(t, cfg) + dev := &syncCaptureSender{} + + p.OnConnect(dev) + armingSettled(t, dev) + + if got := dev.count(PacketTypeSMSRequestConvs); got != 1 { + t.Fatalf("request_conversations sent %d times, want 1", got) + } +} + +// The shipped default must not ask the phone for anything: an armed phone +// cannot be un-armed, so opting in has to be deliberate. +func TestDefaultsDoNotArm(t *testing.T) { + cfg := config.SMSConfig{} + cfg.Defaults() + if cfg.AlwaysArm { + t.Fatal("SMSConfig.Defaults sets AlwaysArm; an armed phone cannot be un-armed") + } + + p, _ := newArmingPlugin(t, cfg) + dev := &syncCaptureSender{} + p.OnConnect(dev) + time.Sleep(100 * time.Millisecond) + + if got := len(dev.types()); got != 0 { + t.Fatalf("default config sent %v, want nothing", dev.types()) + } +} + +// With always_arm off and nobody listening, the phone must not be asked. +func TestNoArmWithoutSubscribers(t *testing.T) { + p, _ := newArmingPlugin(t, config.SMSConfig{AlwaysArm: false}) + dev := &syncCaptureSender{} + + p.OnConnect(dev) + time.Sleep(100 * time.Millisecond) + + if got := len(dev.types()); got != 0 { + t.Fatalf("sent %v with no subscribers, want none", dev.types()) + } +} + +// With always_arm off, a client watching sms.incoming is the opt-in. +func TestArmsWhenSubscribed(t *testing.T) { + p, bus := newArmingPlugin(t, config.SMSConfig{AlwaysArm: false}) + dev := &syncCaptureSender{} + p.OnConnect(dev) + time.Sleep(50 * time.Millisecond) + + sub := bus.Subscribe(1, events.TypeSMSIncoming) + defer sub.Close() + armingSettled(t, dev) + + if got := dev.count(PacketTypeSMSRequestConvs); got != 1 { + t.Fatalf("request_conversations sent %d times after subscribe, want 1", got) + } +} + +// The bus hook fires on every subscribe of any event type. Arming is +// idempotent per connection so a reconnecting watcher cannot re-trigger a +// conversation-head burst each time. +func TestArmsOnlyOncePerConnection(t *testing.T) { + cfg := config.SMSConfig{AlwaysArm: true} + p, bus := newArmingPlugin(t, cfg) + dev := &syncCaptureSender{} + + p.OnConnect(dev) + armingSettled(t, dev) + + for range 5 { + s := bus.Subscribe(1, events.TypeBatteryUpdate) + s.Close() + } + time.Sleep(100 * time.Millisecond) + + if got := dev.count(PacketTypeSMSRequestConvs); got != 1 { + t.Fatalf("request_conversations sent %d times, want 1", got) + } +} + +func TestOnDisconnectForgetsDevice(t *testing.T) { + cfg := config.SMSConfig{AlwaysArm: true} + p, _ := newArmingPlugin(t, cfg) + dev := &syncCaptureSender{} + + p.OnConnect(dev) + armingSettled(t, dev) + p.OnDisconnect(dev) + + p.mu.Lock() + _, stillArmed := p.armedAt[dev.ID()] + p.mu.Unlock() + if stillArmed { + t.Error("device still marked armed after disconnect") + } +} + +// The phone's content observer has no empty guard, so empty batches are +// routine once armed and must not produce events. +func TestEmptyBatchPublishesNothing(t *testing.T) { + cfg := config.SMSConfig{NotifyIncoming: true} + p, bus := newArmingPlugin(t, cfg) + sub := bus.Subscribe(4, events.TypeSMSIncoming) + defer sub.Close() + + pkt, err := protocol.NewPacket(PacketTypeSMSMessages, map[string]any{ + "version": 2, + "messages": []SMSMessage{}, + }) + if err != nil { + t.Fatalf("NewPacket: %v", err) + } + if err := p.Handle(context.Background(), &syncCaptureSender{}, pkt); err != nil { + t.Fatalf("Handle: %v", err) + } + + select { + case ev := <-sub.C: + t.Fatalf("unexpected event for empty batch: %+v", ev) + default: + } +} + +// Messages that predate the arm are the reply burst, not new mail. +func TestShouldNotifyDropsArmingBurst(t *testing.T) { + cfg := config.SMSConfig{NotifyIncoming: true} + p, _ := newArmingPlugin(t, cfg) + + armAt := time.Now() + armAtMs := armAt.UnixMilli() + + old := SMSMessage{Body: "history", Type: 1, Date: armAtMs - 60_000} + if p.shouldNotify(old, true, armAtMs) { + t.Error("notified for a message older than the arm time") + } + + fresh := SMSMessage{Body: "hello", Type: 1, Date: armAtMs + 1} + if !p.shouldNotify(fresh, true, armAtMs) { + t.Error("did not notify for a message newer than the arm time") + } +} + +// The phone cannot be un-armed, so packets keep arriving after every +// client leaves. Notifications must stop regardless. +func TestShouldNotifySilentWhenUnarmed(t *testing.T) { + cfg := config.SMSConfig{NotifyIncoming: true} + p, _ := newArmingPlugin(t, cfg) + + fresh := SMSMessage{Body: "hello", Type: 1, Date: time.Now().UnixMilli() + 1} + if p.shouldNotify(fresh, false, 0) { + t.Error("notified while no client is watching") + } +} + +// Outbound messages come back in the same batch and must not notify. +func TestShouldNotifyIgnoresOutbound(t *testing.T) { + cfg := config.SMSConfig{NotifyIncoming: true} + p, _ := newArmingPlugin(t, cfg) + + outbound := SMSMessage{Body: "sent", Type: 2, Date: time.Now().UnixMilli() + 1} + if p.shouldNotify(outbound, true, 0) { + t.Error("notified for an outbound message") + } +} + +func TestShouldNotifyRespectsConfig(t *testing.T) { + p, _ := newArmingPlugin(t, config.SMSConfig{AlwaysArm: true, NotifyIncoming: false}) + if p.shouldNotify(SMSMessage{Type: 1, Date: 1 << 40}, true, 0) { + t.Error("notified with notify_incoming disabled") + } +} + +// Journals get collected and shipped off-box, so a message body must never +// reach the logger. Only the bus and notify-send may carry it. +func TestMessageBodyNeverLogged(t *testing.T) { + const secret = "my bank code is 1234" + logger, snapshot := log.Observe() + bus := events.NewBus(log.Nop()) + p := NewSMSPlugin(config.SMSConfig{AlwaysArm: true, NotifyIncoming: false}, bus, nil, logger) + sub := bus.Subscribe(4, events.TypeSMSIncoming) + defer sub.Close() + + pkt, err := protocol.NewPacket(PacketTypeSMSMessages, map[string]any{ + "version": 2, + "messages": []SMSMessage{{ + Body: secret, + Type: 1, + Date: time.Now().UnixMilli(), + ThreadID: 7, + Addresses: []SMSAddress{{Address: "+15550100"}}, + }}, + }) + if err != nil { + t.Fatalf("NewPacket: %v", err) + } + if err := p.Handle(context.Background(), &syncCaptureSender{}, pkt); err != nil { + t.Fatalf("Handle: %v", err) + } + + entries := snapshot() + if len(entries) == 0 { + t.Fatal("no log entries captured; the assertion would be vacuous") + } + for _, e := range entries { + if strings.Contains(e, secret) { + t.Errorf("message body reached the log: %q", e) + } + } + + // The body must still reach the event, or the feature is broken. + select { + case ev := <-sub.C: + payload, _ := ev.Payload.(map[string]any) + if payload["body"] != secret { + t.Errorf("event body = %v, want the message", payload["body"]) + } + default: + t.Fatal("no sms.incoming event published") + } +} diff --git a/internal/plugins/sms/handle.go b/internal/plugins/sms/handle.go index 2d24205..54fcda2 100644 --- a/internal/plugins/sms/handle.go +++ b/internal/plugins/sms/handle.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "fmt" + "time" "github.com/bethropolis/kcd/internal/device" "github.com/bethropolis/kcd/internal/events" @@ -36,10 +37,22 @@ func (p *SMSPlugin) handleMessages(_ context.Context, dev device.Sender, pkt *pr return fmt.Errorf("sms: unmarshal messages batch: %w", err) } + // The phone's content observer fires on any SMS database change and has + // no empty guard of its own, so empty batches are routine once armed. + if len(batch.Messages) == 0 { + return nil + } + if len(batch.Messages) > maxSMSMessages { return fmt.Errorf("sms: messages batch too large: %d (max %d)", len(batch.Messages), maxSMSMessages) } + // Resolved once per batch: the arm time and whether we are still + // wanted. See shouldNotify for why both matter. + armAt, isArmed := p.armedAtFor(dev.ID()) + armAtMs := armAt.UnixMilli() + stillArmed := p.armed() + for _, msg := range batch.Messages { if msg.Body == "" { continue @@ -52,9 +65,10 @@ func (p *SMSPlugin) handleMessages(_ context.Context, dev device.Sender, pkt *pr sender = msg.Addresses[0].Address } + // Never log the body: journals are routinely collected and shipped + // off-box, and a message is the most sensitive thing we handle. p.logger.Debug("sms: message received", log.String("from", sender), - log.String("body", msg.Body), log.Int64("thread_id", msg.ThreadID), ) @@ -76,7 +90,9 @@ func (p *SMSPlugin) handleMessages(_ context.Context, dev device.Sender, pkt *pr p.bus.Publish(events.TypeSMSIncoming, dev.ID(), payload) } - if p.cfg.NotifyIncoming { + // Type 1 is an inbound message; type 2 is one we sent, and the + // phone echoes those back in the same batch. + if p.shouldNotify(msg, isArmed && stillArmed, armAtMs) { msgText := msg.Body if len(msgText) > 120 { msgText = msgText[:120] + "…" @@ -93,3 +109,33 @@ func (p *SMSPlugin) handleMessages(_ context.Context, dev device.Sender, pkt *pr return nil } + +// shouldNotify decides whether a message deserves a desktop notification. +// +// Two independent gates, each closing a different hole. The arm time drops +// the conversation-head burst the phone sends in reply to the arming +// request, which is history rather than news. The armed check stops +// notifications once every client has gone away, which matters because the +// phone cannot be un-armed: it keeps pushing for the rest of its lifetime, +// and without this a single transient `kcd watch` would silently turn into +// a permanent notifier. +func (p *SMSPlugin) shouldNotify(msg SMSMessage, isArmed bool, armAtMs int64) bool { + if !p.cfg.NotifyIncoming || !isArmed { + return false + } + // Type 1 is inbound; type 2 is a message we sent, echoed back in the + // same batch. + if msg.Type != 1 { + return false + } + return msg.Date >= armAtMs +} + +// armedAtFor returns when the device was armed and whether it is armed at +// all. Messages older than the arm time are the arming burst, not new SMS. +func (p *SMSPlugin) armedAtFor(deviceID string) (time.Time, bool) { + p.mu.RLock() + defer p.mu.RUnlock() + t, ok := p.armedAt[deviceID] + return t, ok +} diff --git a/internal/plugins/sms/send.go b/internal/plugins/sms/send.go index e4f0388..54cafa6 100644 --- a/internal/plugins/sms/send.go +++ b/internal/plugins/sms/send.go @@ -1,7 +1,11 @@ package sms import ( + "time" + "github.com/bethropolis/kcd/internal/device" + "github.com/bethropolis/kcd/internal/events" + "github.com/bethropolis/kcd/internal/log" "github.com/bethropolis/kcd/internal/protocol" ) @@ -69,7 +73,80 @@ func (p *SMSPlugin) RequestAttachment(dev device.Sender, partID int64, uniqueIde return dev.Send(pkt) } -// --- Lifecycle ------------------------------------------------------------- +// --- Arming ------------------------------------------------------------------ -func (p *SMSPlugin) OnConnect(dev device.Sender) {} -func (p *SMSPlugin) OnDisconnect(dev device.Sender) {} +// armed reports whether the phone should be streaming new SMS to us. The +// phone suppresses every push until it has been asked once, so this is the +// single gate. It is off by default: an armed phone cannot be un-armed. +func (p *SMSPlugin) armed() bool { + if p.cfg.AlwaysArm { + return true + } + return p.bus != nil && p.bus.HasSubscribers(events.TypeSMSIncoming) +} + +// syncArming arms every connected device that is not armed yet. It runs +// from the bus subscriber-change hook, which must return quickly, so the +// device sends happen in a goroutine: dev.Send blocks for up to its write +// timeout when the peer's send channel is full, and a hook that blocked +// would stall bus.Subscribe for every caller. +func (p *SMSPlugin) syncArming() { + if !p.armed() { + return + } + + p.mu.Lock() + targets := make([]device.Sender, 0, len(p.devices)) + now := time.Now() + for id, dev := range p.devices { + if _, ok := p.armedAt[id]; ok { + continue + } + p.armedAt[id] = now + targets = append(targets, dev) + } + p.mu.Unlock() + + if len(targets) == 0 { + return + } + + go func() { + for _, dev := range targets { + if !dev.IsConnected() { + continue + } + if !dev.HasCapability(PacketTypeSMSRequestConvs) { + continue + } + if err := p.RequestConversations(dev); err != nil { + p.logger.Warn("sms: failed to arm push", log.Error(err), + log.String("device_id", dev.ID())) + continue + } + p.logger.Debug("sms: armed push", log.String("device_id", dev.ID())) + } + }() +} + +// --- Lifecycle --------------------------------------------------------------- + +// OnConnect tracks the device and re-arms it. The phone's arm flag survives +// a reconnect, but not a restart of its own app, so re-arming per connection +// is the only recovery available over this protocol. +func (p *SMSPlugin) OnConnect(dev device.Sender) { + p.mu.Lock() + p.devices[dev.ID()] = dev + delete(p.armedAt, dev.ID()) + p.mu.Unlock() + + p.syncArming() +} + +// OnDisconnect forgets the device so the next connect re-arms it. +func (p *SMSPlugin) OnDisconnect(dev device.Sender) { + p.mu.Lock() + delete(p.devices, dev.ID()) + delete(p.armedAt, dev.ID()) + p.mu.Unlock() +} diff --git a/internal/plugins/sms/types.go b/internal/plugins/sms/types.go index da8e1fa..164140a 100644 --- a/internal/plugins/sms/types.go +++ b/internal/plugins/sms/types.go @@ -4,9 +4,11 @@ import ( "crypto/tls" "os" "path/filepath" + "sync" "time" "github.com/bethropolis/kcd/internal/config" + "github.com/bethropolis/kcd/internal/device" "github.com/bethropolis/kcd/internal/events" "github.com/bethropolis/kcd/internal/log" "github.com/bethropolis/kcd/internal/protocol" @@ -39,6 +41,16 @@ type SMSPlugin struct { tlsConfig *tls.Config logger log.Logger cacheDir string + + // mu guards devices and armedAt. + mu sync.RWMutex + // devices holds the currently connected devices, keyed by ID, so the + // bus hook can arm them without reaching for a registry. + devices map[string]device.Sender + // armedAt records when each device was last armed, and doubles as the + // set of already-armed devices: arming twice on one connection would + // re-trigger a full conversation-head burst for no benefit. + armedAt map[string]time.Time } // Options customizes storage, network timeouts, and desktop notification identity. @@ -59,7 +71,7 @@ func NewSMSPlugin(cfg config.SMSConfig, bus *events.Bus, tlsConfig *tls.Config, } _ = os.MkdirAll(cacheDir, 0700) - return &SMSPlugin{ + p := &SMSPlugin{ sidechannel: opts.Sidechannel, notifications: opts.Notifications, cfg: cfg, @@ -67,7 +79,17 @@ func NewSMSPlugin(cfg config.SMSConfig, bus *events.Bus, tlsConfig *tls.Config, tlsConfig: tlsConfig, logger: logger.With(log.String("plugin", "sms")), cacheDir: cacheDir, + devices: make(map[string]device.Sender), + armedAt: make(map[string]time.Time), + } + + // Arming is owner-gated: the phone is only asked to push while a client + // is watching sms.incoming, or when always_arm is set. + if bus != nil { + bus.OnSubscriberChange(p.syncArming) } + + return p } func (p *SMSPlugin) Name() string { return "SMS" } diff --git a/mise.toml b/mise.toml new file mode 100644 index 0000000..c1acd20 --- /dev/null +++ b/mise.toml @@ -0,0 +1,8 @@ +[tools] +hk = "latest" +# Pinned to match .github/workflows/ci.yml (golangci-lint-action v9, v2.11). +# Bump both together; a local version ahead of CI means local checks can pass +# while CI fails. +golangci-lint = "2.11" +# go.mod declares go 1.25.0 and CI uses go-version-file, so track that. +go = "1.25" diff --git a/packaging/kcd.example.toml b/packaging/kcd.example.toml index f3a4c46..8cdd73b 100644 --- a/packaging/kcd.example.toml +++ b/packaging/kcd.example.toml @@ -193,13 +193,12 @@ remotesystemvolume = true # The local D-Bus watcher is event-driven; this section tunes the position # poller that keeps the phone's now-playing display exact while music plays. # poll_while_playing = true # false = pure event-driven (position extrapolates - # from posAnchorMs; no poller and no watchdog, so a - # dropped D-Bus signal can leave the display stale) + # from posAnchorMs; no poller at all, so a dropped + # D-Bus signal can leave the display stale) # position_interval = "2s" # D-Bus re-read cadence while playing only. - # While a player is tracked but paused, a 10s - # watchdog still runs so a missed play signal cannot - # strand the poller; with no player tracked, zero - # timers remain. + # D-Bus signals arm the poller on playback and it + # stops itself on a confirmed pause, so a paused + # or absent player costs zero wakeups. [cache] # Empty values preserve existing storage paths; use absolute paths to override. @@ -308,6 +307,15 @@ suspend = "systemctl suspend" [sms] # notify_incoming = true # show a desktop notification on incoming SMS +# always_arm = false # ask the phone to push new SMS on every + # connect, so notifications work with no + # client attached. Off by default because an + # armed phone cannot be un-armed; without it, + # `kcd watch --events sms.incoming` receives + # messages instead. + # The phone applies its blocked-numbers list + # only to the deprecated telephony push, so + # blocked senders can still arrive here. # ─── Mousepad: remote input settings ───────────────────────────────────────── diff --git a/scripts/preremove.sh b/scripts/preremove.sh index ca4041b..7729324 100755 --- a/scripts/preremove.sh +++ b/scripts/preremove.sh @@ -4,7 +4,7 @@ set -e # Package pre-remove hook for kcd # Try to stop and disable all instances of the kcd template service -# that might be running. +# that might be running. if command -v systemctl >/dev/null 2>&1; then echo "Stopping any running kcd services..." # System-level kcd services (template or plain)