Quick Start Guide

Install the Go SDK and configure the [R]DP instance to use:

go mod init rdp-quickstart
go get github.com/raft-tech/[email protected]

export RDP_SERVER_URL=https://rdp-example.com
export RDP_API_KEY=my-rdp-api-key

Create a Transformer

package main

import (
	"context"
	"fmt"
	"log"

	rdpsdk "github.com/raft-tech/rdp-sdk-go"
	transformersapi "github.com/raft-tech/rdp-sdk-go/transformers"
)

func main() {
	cfg, err := rdpsdk.LoadConfig()
	if err != nil {
		log.Fatal(err)
	}
	client, err := rdpsdk.NewFromConfig(cfg)
	if err != nil {
		log.Fatal(err)
	}

	inputName := "Incoming messages"
	inputType := "INTERNAL_STREAMING"
	outputName := "Outgoing files"
	outputType := "INTERNAL_OBJECT_STORE"

	transformer := transformersapi.Transformer{
		Uid:              "new-transformer-abc",
		Name:             "My New Transformer",
		Description:      "Transforms data.",
		Status:           "available",
		SecurityMarkings: "UNCLASSIFIED",
		Types:            []transformersapi.TransformerType{},
		Inputs: map[string]transformersapi.DataConnTpl{
			"INPUT_STREAMING_TOPIC": {
				DisplayName: &inputName,
				ConnType:    &inputType,
			},
		},
		Outputs: map[string]transformersapi.DataConnTpl{
			"OUTPUT_OBJECT_BUCKET": {
				DisplayName: &outputName,
				ConnType:    &outputType,
			},
		},
		Configuration: transformersapi.ConfigurationTpl{
			"environment": []map[string]any{
				{
					"name":        "LOG_LEVEL",
					"required":    true,
					"description": "The log level used at startup.",
				},
			},
		},
		Instantiation: transformersapi.CreationConfig{
			"job_image": map[string]any{
				"image":            "my-container-image",
				"pull_policy":      "IfNotPresent",
				"default_replicas": 1,
			},
		},
	}

	resp, err := client.Transformers().Create(context.Background(), transformer)
	if err != nil {
		log.Fatal(err)
	}
	if resp.JSON201 == nil {
		log.Fatalf("create transformer failed: %s", resp.Status())
	}

	fmt.Printf("Created Transformer with UID %q\n", resp.JSON201.Uid)
}

Retrieve a Transformer

package main

import (
	"context"
	"fmt"
	"log"

	rdpsdk "github.com/raft-tech/rdp-sdk-go"
)

func main() {
	cfg, err := rdpsdk.LoadConfig()
	if err != nil {
		log.Fatal(err)
	}
	client, err := rdpsdk.NewFromConfig(cfg)
	if err != nil {
		log.Fatal(err)
	}

	resp, err := client.Transformers().ID("new-transformer-abc").Get(context.Background())
	if err != nil {
		log.Fatal(err)
	}
	if resp.JSON200 == nil {
		log.Fatalf("get transformer failed: %s", resp.Status())
	}

	fmt.Printf("Found transformer: %+v\n", *resp.JSON200)
}

List all Pipelines

package main

import (
	"context"
	"fmt"
	"log"

	rdpsdk "github.com/raft-tech/rdp-sdk-go"
)

func main() {
	cfg, err := rdpsdk.LoadConfig()
	if err != nil {
		log.Fatal(err)
	}
	client, err := rdpsdk.NewFromConfig(cfg)
	if err != nil {
		log.Fatal(err)
	}

	resp, err := client.Pipelines().Instances().List(context.Background(), nil)
	if err != nil {
		log.Fatal(err)
	}
	if resp.JSON200 == nil {
		log.Fatalf("list pipelines failed: %s", resp.Status())
	}

	fmt.Printf("Found %d pipelines.\n", len(*resp.JSON200))
}

Get the current status of a Pipeline

package main

import (
	"context"
	"encoding/json"
	"fmt"
	"log"

	rdpsdk "github.com/raft-tech/rdp-sdk-go"
)

func main() {
	cfg, err := rdpsdk.LoadConfig()
	if err != nil {
		log.Fatal(err)
	}
	client, err := rdpsdk.NewFromConfig(cfg)
	if err != nil {
		log.Fatal(err)
	}

	uid := "my-pipeline-uid"
	resp, err := client.Pipelines().Instances().ID(uid).Status(context.Background())
	if err != nil {
		log.Fatal(err)
	}
	if resp.JSON200 == nil {
		log.Fatalf("get pipeline status failed: %s", resp.Status())
	}

	status, err := json.MarshalIndent(resp.JSON200, "", "  ")
	if err != nil {
		log.Fatal(err)
	}
	fmt.Printf("Status for pipeline %s:\n%s\n", uid, status)
}

Create a Pipeline

package main

import (
	"context"
	"fmt"
	"log"

	rdpsdk "github.com/raft-tech/rdp-sdk-go"
	pipelinesapi "github.com/raft-tech/rdp-sdk-go/pipelines"
)

func main() {
	cfg, err := rdpsdk.LoadConfig()
	if err != nil {
		log.Fatal(err)
	}
	client, err := rdpsdk.NewFromConfig(cfg)
	if err != nil {
		log.Fatal(err)
	}

	name := "example-pipeline"
	description := "Example pipeline"
	securityMarkings := "UNCLASSIFIED"
	environment := map[string]string{"LOG_LEVEL": "info"}
	transformers := []pipelinesapi.TransformerInstancePost{
		{
			TemplateUid: "new-transformer-abc",
			Uid:         "transformer",
			Configuration: &pipelinesapi.Configuration{
				Environment: &environment,
			},
		},
	}
	datasets := []pipelinesapi.DatasetRefPost{
		{
			Uid: "input-dataset",
			Configuration: &pipelinesapi.DatasetConfiguration{
				Ref:      "quickstart-input",
				ConnType: pipelinesapi.INTERNALSTREAMING,
			},
		},
		{
			Uid: "output-dataset",
			Configuration: &pipelinesapi.DatasetConfiguration{
				Ref:      "quickstart-output",
				ConnType: pipelinesapi.INTERNALOBJECTSTORE,
			},
		},
	}
	connections := []pipelinesapi.Connection{
		{
			From: pipelinesapi.ConnectionEndpoint{
				Conn: "",
				Id:   "input-dataset",
				Type: pipelinesapi.Dataset,
			},
			To: pipelinesapi.ConnectionEndpoint{
				Conn: "INPUT_STREAMING_TOPIC",
				Id:   "transformer",
				Type: pipelinesapi.Transformer,
			},
		},
		{
			From: pipelinesapi.ConnectionEndpoint{
				Conn: "OUTPUT_OBJECT_BUCKET",
				Id:   "transformer",
				Type: pipelinesapi.Transformer,
			},
			To: pipelinesapi.ConnectionEndpoint{
				Conn: "",
				Id:   "output-dataset",
				Type: pipelinesapi.Dataset,
			},
		},
	}

	pipeline := pipelinesapi.Pipeline{
		Name:             &name,
		Description:      &description,
		SecurityMarkings: &securityMarkings,
		Transformers:     &transformers,
		Datasets:         &datasets,
		Connections:      &connections,
	}
	resp, err := client.Pipelines().Instances().Create(context.Background(), nil, pipeline)
	if err != nil {
		log.Fatal(err)
	}
	if resp.JSON201 == nil || resp.JSON201.Uid == nil {
		log.Fatalf("create pipeline failed: %s", resp.Status())
	}

	fmt.Printf("Created Pipeline with UID %q\n", *resp.JSON201.Uid)
}

Manage a Pipeline’s state

package main

import (
	"context"
	"log"

	rdpsdk "github.com/raft-tech/rdp-sdk-go"
)

func main() {
	cfg, err := rdpsdk.LoadConfig()
	if err != nil {
		log.Fatal(err)
	}
	client, err := rdpsdk.NewFromConfig(cfg)
	if err != nil {
		log.Fatal(err)
	}

	ctx := context.Background()
	pipeline := client.Pipelines().Instances().ID("my-pipeline-uid")

	start, err := pipeline.Start(ctx)
	if err != nil {
		log.Fatal(err)
	}
	if start.StatusCode() >= 300 {
		log.Fatalf("start pipeline failed: %s", start.Status())
	}

	stop, err := pipeline.Stop(ctx)
	if err != nil {
		log.Fatal(err)
	}
	if stop.StatusCode() >= 300 {
		log.Fatalf("stop pipeline failed: %s", stop.Status())
	}

	restart, err := pipeline.Restart(ctx)
	if err != nil {
		log.Fatal(err)
	}
	if restart.StatusCode() >= 300 {
		log.Fatalf("restart pipeline failed: %s", restart.Status())
	}
}