Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ jobs:
run: buf generate

- name: Unit tests
run: go test ./src/...
run: go test ./src/... ./pkg/...

- name: Plugin tests
run: go test ./plugins/calendar/... ./plugins/devicemonitor/... ./plugins/docker/... ./plugins/mqtt/... ./plugins/timer/... ./plugins/webhook/...
Expand Down
2 changes: 1 addition & 1 deletion Taskfile.yml
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ tasks:
lint:go:
desc: Run golangci-lint over the Go module
cmds:
- '{{.GOPATH}}/bin/golangci-lint run ./src/... ./plugins/calendar/... ./plugins/devicemonitor/... ./plugins/mqtt/... ./plugins/timer/... ./plugins/webhook/...'
- '{{.GOPATH}}/bin/golangci-lint run ./src/... ./pkg/... ./plugins/calendar/... ./plugins/devicemonitor/... ./plugins/mqtt/... ./plugins/timer/... ./plugins/webhook/...'

lint:proto:
desc: Lint protobuf definitions
Expand Down
5 changes: 2 additions & 3 deletions pkg/config/overlay_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,9 +32,8 @@ type sliceMapConfig struct {
Labels map[string]string
}

func intPtr(v int) *int { return &v }
func boolPtr(v bool) *bool { return &v }
func strPtr(v string) *string { return &v }
func intPtr(v int) *int { return &v }
func boolPtr(v bool) *bool { return &v }

func TestOverlay_NoLayers(t *testing.T) {
base := simpleConfig{Name: "base", Count: 5, Flag: true}
Expand Down
55 changes: 55 additions & 0 deletions pkg/plugin-sdk/plugin-sdk.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
// Package pluginsdk is the stable, supported API for building recur trigger
// plugins. External plugin repositories import this package instead of reaching
// into core internals under src/infra, so the plugin contract can evolve behind
// a curated surface without breaking out-of-tree plugins.
//
// A trigger plugin is a subprocess that recur launches. It reads its trigger
// configuration as JSON from stdin, dials the daemon at the Unix socket named by
// the RECUR_SOCKET environment variable, and reports fired events for the
// trigger identified by RECUR_TRIGGER_ID:
//
// client, err := pluginsdk.Connect(os.Getenv("RECUR_SOCKET"))
// if err != nil {
// log.Fatal(err)
// }
// defer client.Close()
//
// resp, err := client.Service.ReportTriggerEvent(ctx, &pluginsdk.ReportTriggerEventRequest{
// TriggerId: os.Getenv("RECUR_TRIGGER_ID"),
// Context: map[string]string{"Key": "value"},
// })
package pluginsdk

import (
clientgrpc "github.com/directedbits/recur/src/infra/grpc/client"
recurv1 "github.com/directedbits/recur/src/infra/grpc/recur/v1"
)

// Client is a connection to the recur daemon over its gRPC socket. Report fired
// events through the embedded Service and close the connection on shutdown.
type Client = clientgrpc.Client

// RecurClient is the daemon RPC surface reachable via Client.Service.
type RecurClient = recurv1.RecurServiceClient

// ReportTriggerEventRequest reports that a trigger fired. TriggerId identifies
// the trigger (the RECUR_TRIGGER_ID environment variable) and Context carries
// the trigger's context variables.
type ReportTriggerEventRequest = recurv1.ReportTriggerEventRequest

// ReportTriggerEventResponse is the daemon's reply to a reported event. Accepted
// reports whether the daemon acted on the event; Error explains a rejection.
type ReportTriggerEventResponse = recurv1.ReportTriggerEventResponse

// Connect dials the recur daemon at the given Unix socket path (typically the
// value of RECUR_SOCKET). It returns an error if the daemon is not reachable
// within a few seconds.
func Connect(socketPath string) (*Client, error) {
return clientgrpc.Connect(socketPath)
}

// ConnectOrNil dials the daemon like Connect but returns nil (no error) when the
// daemon is unreachable, for callers that can operate without it.
func ConnectOrNil(socketPath string) *Client {
return clientgrpc.ConnectOrNil(socketPath)
}
31 changes: 31 additions & 0 deletions pkg/plugin-sdk/plugin-sdk_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
package pluginsdk_test

import (
"testing"

pluginsdk "github.com/directedbits/recur/pkg/plugin-sdk"
)

func TestConnectUnreachableSocketReturnsError(t *testing.T) {
if _, err := pluginsdk.Connect("/nonexistent/recur.sock"); err == nil {
t.Fatal("expected an error connecting to a nonexistent socket")
}
}

func TestConnectOrNilUnreachableSocketReturnsNil(t *testing.T) {
if c := pluginsdk.ConnectOrNil("/nonexistent/recur.sock"); c != nil {
t.Fatalf("expected nil client for an unreachable socket, got %#v", c)
}
}

// Pins the re-exported message contract so a change to the underlying generated
// type surfaces here rather than only in out-of-tree plugins.
func TestReportTriggerEventRequestFields(t *testing.T) {
req := &pluginsdk.ReportTriggerEventRequest{
TriggerId: "trigger-1",
Context: map[string]string{"Key": "value"},
}
if req.TriggerId != "trigger-1" || req.Context["Key"] != "value" {
t.Fatal("re-exported ReportTriggerEventRequest fields did not round-trip")
}
}
7 changes: 3 additions & 4 deletions plugins/calendar/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,7 @@ import (
"syscall"
"time"

clientgrpc "github.com/directedbits/recur/src/infra/grpc/client"
recurv1 "github.com/directedbits/recur/src/infra/grpc/recur/v1"
sdk "github.com/directedbits/recur/pkg/plugin-sdk"
)

// pluginInput is the JSON payload read from stdin.
Expand Down Expand Up @@ -103,7 +102,7 @@ func main() {
input.TriggerType, source, pollInterval, lookAhead)

// Connect to daemon gRPC socket
client, err := clientgrpc.Connect(socketPath)
client, err := sdk.Connect(socketPath)
if err != nil {
log.Fatalf("connecting to daemon: %v", err)
}
Expand Down Expand Up @@ -135,7 +134,7 @@ func main() {
ctxVars["StartsIn"] = evt.StartsIn.Truncate(time.Second).String()
}

resp, err := client.Service.ReportTriggerEvent(context.Background(), &recurv1.ReportTriggerEventRequest{
resp, err := client.Service.ReportTriggerEvent(context.Background(), &sdk.ReportTriggerEventRequest{
TriggerId: triggerID,
Context: ctxVars,
})
Expand Down
8 changes: 8 additions & 0 deletions plugins/devicemonitor/dbus_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -377,6 +377,7 @@ func TestParseInterfacesAdded_BlockDevice(t *testing.T) {
}
if event == nil {
t.Fatal("expected event, got nil")
return
}
if event.DeviceName != "sdb1" {
t.Errorf("DeviceName = %q, want %q", event.DeviceName, "sdb1")
Expand Down Expand Up @@ -412,6 +413,7 @@ func TestParseInterfacesAdded_Drive(t *testing.T) {
}
if event == nil {
t.Fatal("expected event, got nil")
return
}
if event.DeviceName != "WD_Elements_1234" {
t.Errorf("DeviceName = %q, want %q", event.DeviceName, "WD_Elements_1234")
Expand Down Expand Up @@ -443,6 +445,7 @@ func TestParseInterfacesAdded_USBDrive(t *testing.T) {
}
if event == nil {
t.Fatal("expected event, got nil")
return
}
if event.DeviceType != "drive" {
t.Errorf("DeviceType = %q, want drive (kind-of-object axis)", event.DeviceType)
Expand Down Expand Up @@ -479,6 +482,7 @@ func TestParseInterfacesAdded_PartitionUsesLookupForBus(t *testing.T) {
}
if event == nil {
t.Fatal("expected event")
return
}
if event.DeviceType != "block" {
t.Errorf("DeviceType = %q, want block (no reclassification)", event.DeviceType)
Expand Down Expand Up @@ -517,6 +521,7 @@ func TestParseInterfacesAdded_LoopDeviceHasLoopBusAndEmptyDrivePath(t *testing.T
}
if event == nil {
t.Fatal("expected event")
return
}
if event.DeviceType != "block" {
t.Errorf("DeviceType = %q, want block", event.DeviceType)
Expand Down Expand Up @@ -865,6 +870,7 @@ func TestParseInterfacesAdded_WithMountPoint(t *testing.T) {
}
if event == nil {
t.Fatal("expected event, got nil")
return
}
if event.MountPoint != "/mnt/usb" {
t.Errorf("MountPoint = %q, want %q", event.MountPoint, "/mnt/usb")
Expand Down Expand Up @@ -948,6 +954,7 @@ func TestParseInterfacesRemoved_BlockDevice(t *testing.T) {
}
if event == nil {
t.Fatal("expected event, got nil")
return
}
if event.DeviceName != "sdb1" {
t.Errorf("DeviceName = %q, want %q", event.DeviceName, "sdb1")
Expand Down Expand Up @@ -975,6 +982,7 @@ func TestParseInterfacesRemoved_Drive(t *testing.T) {
}
if event == nil {
t.Fatal("expected event, got nil")
return
}
if event.DeviceType != "drive" {
t.Errorf("DeviceType = %q, want %q", event.DeviceType, "drive")
Expand Down
7 changes: 3 additions & 4 deletions plugins/devicemonitor/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,7 @@ import (
"os/signal"
"syscall"

clientgrpc "github.com/directedbits/recur/src/infra/grpc/client"
recurv1 "github.com/directedbits/recur/src/infra/grpc/recur/v1"
sdk "github.com/directedbits/recur/pkg/plugin-sdk"
)

// pluginInput is the JSON payload read from stdin.
Expand Down Expand Up @@ -94,7 +93,7 @@ func main() {
}

// Connect to daemon gRPC socket
client, err := clientgrpc.Connect(socketPath)
client, err := sdk.Connect(socketPath)
if err != nil {
log.Fatalf("connecting to daemon: %v", err)
}
Expand Down Expand Up @@ -134,7 +133,7 @@ func main() {
ctxVars["MountPoint"] = event.MountPoint
}

resp, err := client.Service.ReportTriggerEvent(context.Background(), &recurv1.ReportTriggerEventRequest{
resp, err := client.Service.ReportTriggerEvent(context.Background(), &sdk.ReportTriggerEventRequest{
TriggerId: triggerID,
Context: ctxVars,
})
Expand Down
7 changes: 3 additions & 4 deletions plugins/docker/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,7 @@ import (
"strings"
"syscall"

clientgrpc "github.com/directedbits/recur/src/infra/grpc/client"
recurv1 "github.com/directedbits/recur/src/infra/grpc/recur/v1"
sdk "github.com/directedbits/recur/pkg/plugin-sdk"
)

// pluginInput is the JSON payload read from stdin.
Expand Down Expand Up @@ -121,7 +120,7 @@ func runTrigger(input *pluginInput) {

log.Printf("watching: host=%s trigger=%s", host, input.TriggerType)

grpcClient, err := clientgrpc.Connect(socketPath)
grpcClient, err := sdk.Connect(socketPath)
if err != nil {
log.Fatalf("connecting to daemon: %v", err)
}
Expand All @@ -145,7 +144,7 @@ func runTrigger(input *pluginInput) {

ctxVars := buildContextVars(input.TriggerType, evt)

resp, err := grpcClient.Service.ReportTriggerEvent(context.Background(), &recurv1.ReportTriggerEventRequest{
resp, err := grpcClient.Service.ReportTriggerEvent(context.Background(), &sdk.ReportTriggerEventRequest{
TriggerId: triggerID,
Context: ctxVars,
})
Expand Down
7 changes: 3 additions & 4 deletions plugins/fileevents/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,7 @@ import (
"syscall"
"time"

clientgrpc "github.com/directedbits/recur/src/infra/grpc/client"
recurv1 "github.com/directedbits/recur/src/infra/grpc/recur/v1"
sdk "github.com/directedbits/recur/pkg/plugin-sdk"
)

// pluginInput is the JSON payload read from stdin.
Expand Down Expand Up @@ -74,7 +73,7 @@ func main() {
log.Fatal(err)
}

client, err := clientgrpc.Connect(socketPath)
client, err := sdk.Connect(socketPath)
if err != nil {
log.Fatalf("connecting to daemon: %v", err)
}
Expand Down Expand Up @@ -127,7 +126,7 @@ func main() {
ctxVars := buildContext(input.TriggerType, action)
ctxVars["TriggeredOn"] = time.Now().UTC().Format(time.RFC3339)

resp, err := client.Service.ReportTriggerEvent(context.Background(), &recurv1.ReportTriggerEventRequest{
resp, err := client.Service.ReportTriggerEvent(context.Background(), &sdk.ReportTriggerEventRequest{
TriggerId: triggerID,
Context: ctxVars,
})
Expand Down
7 changes: 3 additions & 4 deletions plugins/mqtt/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,7 @@ import (
"strconv"
"syscall"

clientgrpc "github.com/directedbits/recur/src/infra/grpc/client"
recurv1 "github.com/directedbits/recur/src/infra/grpc/recur/v1"
sdk "github.com/directedbits/recur/pkg/plugin-sdk"
)

// pluginInput is the JSON payload read from stdin.
Expand Down Expand Up @@ -97,7 +96,7 @@ func runTrigger(input *pluginInput) {

log.Printf("subscribed: broker=%s topic=%s qos=%d", cfg.Broker, cfg.Topic, cfg.QoS)

client, err := clientgrpc.Connect(socketPath)
client, err := sdk.Connect(socketPath)
if err != nil {
log.Fatalf("connecting to daemon: %v", err)
}
Expand All @@ -122,7 +121,7 @@ func runTrigger(input *pluginInput) {
"MessageID": strconv.FormatUint(uint64(msg.MessageID), 10),
}

resp, err := client.Service.ReportTriggerEvent(context.Background(), &recurv1.ReportTriggerEventRequest{
resp, err := client.Service.ReportTriggerEvent(context.Background(), &sdk.ReportTriggerEventRequest{
TriggerId: triggerID,
Context: ctxVars,
})
Expand Down
7 changes: 3 additions & 4 deletions plugins/timer/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,7 @@ import (
"os/signal"
"syscall"

clientgrpc "github.com/directedbits/recur/src/infra/grpc/client"
recurv1 "github.com/directedbits/recur/src/infra/grpc/recur/v1"
sdk "github.com/directedbits/recur/pkg/plugin-sdk"
)

// pluginInput is the JSON payload read from stdin.
Expand Down Expand Up @@ -108,7 +107,7 @@ func main() {
defer stop()

// Connect to daemon gRPC socket
client, err := clientgrpc.Connect(socketPath)
client, err := sdk.Connect(socketPath)
if err != nil {
log.Fatalf("connecting to daemon: %v", err)
}
Expand All @@ -132,7 +131,7 @@ func main() {
"TimeSinceStarted": tick.TimeSinceStarted,
}

resp, err := client.Service.ReportTriggerEvent(context.Background(), &recurv1.ReportTriggerEventRequest{
resp, err := client.Service.ReportTriggerEvent(context.Background(), &sdk.ReportTriggerEventRequest{
TriggerId: triggerID,
Context: ctxVars,
})
Expand Down
27 changes: 13 additions & 14 deletions plugins/webhook/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,7 @@ import (
"strconv"
"syscall"

clientgrpc "github.com/directedbits/recur/src/infra/grpc/client"
recurv1 "github.com/directedbits/recur/src/infra/grpc/recur/v1"
sdk "github.com/directedbits/recur/pkg/plugin-sdk"
)

// pluginInput is the JSON payload read from stdin.
Expand Down Expand Up @@ -137,7 +136,7 @@ func main() {
log.Printf("started: %s port=%s path=%s method=%s max_body_size=%d", proto, port, path, method, maxBodySize)

// Connect to daemon gRPC socket
client, err := clientgrpc.Connect(socketPath)
client, err := sdk.Connect(socketPath)
if err != nil {
log.Fatalf("connecting to daemon: %v", err)
}
Expand All @@ -157,19 +156,19 @@ func main() {
}

ctxVars := map[string]string{
"RequestMethod": evt.Method,
"RequestPath": evt.Path,
"RequestBody": evt.Body,
"QueryString": evt.QueryString,
"RemoteAddr": evt.RemoteAddr,
"ContentType": evt.ContentType,
"Headers": encodeHeaders(evt.Headers),
"UserAgent": evt.UserAgent,
"Referer": evt.Referer,
"XForwardedFor": evt.XForwardedFor,
"RequestMethod": evt.Method,
"RequestPath": evt.Path,
"RequestBody": evt.Body,
"QueryString": evt.QueryString,
"RemoteAddr": evt.RemoteAddr,
"ContentType": evt.ContentType,
"Headers": encodeHeaders(evt.Headers),
"UserAgent": evt.UserAgent,
"Referer": evt.Referer,
"XForwardedFor": evt.XForwardedFor,
}

resp, err := client.Service.ReportTriggerEvent(context.Background(), &recurv1.ReportTriggerEventRequest{
resp, err := client.Service.ReportTriggerEvent(context.Background(), &sdk.ReportTriggerEventRequest{
TriggerId: triggerID,
Context: ctxVars,
})
Expand Down
Loading