Skip to content
Open
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
3 changes: 2 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
@@ -1 +1,2 @@
tmp
tmp
.env
9 changes: 8 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,8 @@ A Go implementation of a Temporal Codec server that provides encryption and decr
## Features

- `/encode` endpoint for encrypting payloads using AES-GCM
- `/decoder` endpoint for encdecryptingypting payloads using AES-GCM
- `/decode` endpoint for decrypting payloads
- Configurable timeout simulation for testing
- CORS support for Temporal Web UI integration
- Health check endpoint
- Key rotation support via key IDs
Expand Down Expand Up @@ -34,6 +34,13 @@ For detailed implementation examples and usage instructions, see [codec/README.m

## Running the Server

### Docker Compose (includes Temporal)

1. docker compose up
2. docker compose down

### Standalone

1. Install dependencies:
```bash
go mod download
Expand Down
11 changes: 11 additions & 0 deletions codec/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,17 @@ func main() {
}
```

## Run steps

1) Run the following command to start the worker
```
go run codec/worker/main.go
```
2) Run the following command to start the example
```
go run codec/starter/main.go
```

## Important Notes

1. The codec server must be running before starting your worker
Expand Down
25 changes: 5 additions & 20 deletions codec/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package codec

import (
"bytes"
"codec-server/models"
"encoding/json"
"fmt"
"io"
Expand All @@ -16,25 +17,9 @@ type HTTPCodec struct {
CodecServerURL string
}

// Payload represents the structure for codec server communication
type Payload struct {
Metadata map[string]string `json:"metadata"`
Data string `json:"data"`
}

// CodecRequest represents the request to the codec server
type CodecRequest struct {
Payloads []*commonpb.Payload `json:"payloads"`
}

// CodecResponse represents the response from the codec server
type CodecResponse struct {
Payloads []*commonpb.Payload `json:"payloads"`
}

func (c *HTTPCodec) Encode(payloads []*commonpb.Payload) ([]*commonpb.Payload, error) {
// Send request to codec server
reqBody, err := json.Marshal(CodecRequest{Payloads: payloads})
reqBody, err := json.Marshal(models.CodecRequest{Payloads: payloads})
if err != nil {
return nil, fmt.Errorf("failed to marshal request: %v", err)
}
Expand All @@ -51,7 +36,7 @@ func (c *HTTPCodec) Encode(payloads []*commonpb.Payload) ([]*commonpb.Payload, e
return nil, fmt.Errorf("codec server returned error: %s (URL: %s, Body: %s)", resp.Status, c.CodecServerURL, string(bodyBytes))
}

var codecResp CodecResponse
var codecResp models.CodecResponse
if err := json.NewDecoder(resp.Body).Decode(&codecResp); err != nil {
return nil, fmt.Errorf("failed to decode response: %v", err)
}
Expand All @@ -61,7 +46,7 @@ func (c *HTTPCodec) Encode(payloads []*commonpb.Payload) ([]*commonpb.Payload, e

func (c *HTTPCodec) Decode(payloads []*commonpb.Payload) ([]*commonpb.Payload, error) {
// Send request to codec server
reqBody, err := json.Marshal(CodecRequest{Payloads: payloads})
reqBody, err := json.Marshal(models.CodecRequest{Payloads: payloads})
if err != nil {
return nil, fmt.Errorf("failed to marshal request: %v", err)
}
Expand All @@ -78,7 +63,7 @@ func (c *HTTPCodec) Decode(payloads []*commonpb.Payload) ([]*commonpb.Payload, e
return nil, fmt.Errorf("codec server returned error: %s", resp.Status)
}

var codecResp CodecResponse
var codecResp models.CodecResponse
if err := json.NewDecoder(resp.Body).Decode(&codecResp); err != nil {
return nil, fmt.Errorf("failed to decode response: %v", err)
}
Expand Down
69 changes: 69 additions & 0 deletions codec/starter/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
package main

import (
codec "codec-server/codec"
"codec-server/models"
"context"
"log"

"go.temporal.io/sdk/client"
"go.temporal.io/sdk/converter"
"go.temporal.io/sdk/workflow"
)

func main() {

// Create the codec instance
cInstance := codec.New("http://localhost:8081")

// Create a data converter that uses our codec
dataConverter := converter.NewCodecDataConverter(
converter.GetDefaultDataConverter(),
cInstance,
)

// Create a data converter without deadlock detection
dataConverterWithoutDeadlock := workflow.DataConverterWithoutDeadlockDetection(dataConverter)

c, err := client.Dial(client.Options{
HostPort: "localhost:7233",
Namespace: "default",
DataConverter: dataConverterWithoutDeadlock,
})

if err != nil {
log.Fatalln("Unable to create client", err)
}
defer c.Close()

workflowOptions := client.StartWorkflowOptions{
ID: "codecserver_workflowID",
TaskQueue: "codecserver",
}

payload := models.PayloadData{
Data: "Plain text input",
ActivityType: "SimpleActivity",
Timeout: 9,
}
// The workflow input "My Compressed Friend" will be encoded by the codec before being sent to Temporal
we, err := c.ExecuteWorkflow(
context.Background(),
workflowOptions,
codec.Workflow,
payload,
)
if err != nil {
log.Fatalln("Unable to execute workflow", err)
}

log.Println("Started workflow", "WorkflowID", we.GetID(), "RunID", we.GetRunID())

// Synchronously wait for the workflow completion.
var result string
err = we.Get(context.Background(), &result)
if err != nil {
log.Fatalln("Unable get workflow result", err)
}
log.Println("Workflow result:", result)
}
48 changes: 48 additions & 0 deletions codec/worker/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
package main

import (
codec "codec-server/codec"
"log"

"go.temporal.io/sdk/client"
"go.temporal.io/sdk/converter"
"go.temporal.io/sdk/worker"
"go.temporal.io/sdk/workflow"
)

func main() {

// Create the codec instance
cInstance := codec.New("http://localhost:8081")

// Create a data converter that uses our codec
dataConverter := converter.NewCodecDataConverter(
converter.GetDefaultDataConverter(),
cInstance,
)

// Create a data converter without deadlock detection
dataConverterWithoutDeadlock := workflow.DataConverterWithoutDeadlockDetection(dataConverter)

c, err := client.Dial(client.Options{
HostPort: "localhost:7233",
Namespace: "default",
DataConverter: dataConverterWithoutDeadlock,
})

if err != nil {
log.Fatalln("Unable to create client", err)
}
defer c.Close()

w := worker.New(c, "codecserver", worker.Options{})

w.RegisterWorkflow(codec.Workflow)
w.RegisterActivity(codec.Activity)
w.RegisterActivity(codec.TimeoutActivity)

err = w.Run(worker.InterruptCh())
if err != nil {
log.Fatalln("Unable to start worker", err)
}
}
92 changes: 92 additions & 0 deletions codec/workflow.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
package codec

import (
"codec-server/models"
"context"
"errors"
"fmt"
"io"
"net/http"
"time"

"github.com/google/uuid"
"go.temporal.io/sdk/activity"
"go.temporal.io/sdk/workflow"
)

// Workflow is a standard workflow definition.
// Note that the Workflow and Activity don't need to care that
// their inputs/results are being encoded.
func Workflow(ctx workflow.Context, input models.PayloadData) (string, error) {
ao := workflow.ActivityOptions{
StartToCloseTimeout: 5 * time.Second,
}
lao := workflow.LocalActivityOptions{
StartToCloseTimeout: 5 * time.Second,
}
ctx = workflow.WithActivityOptions(ctx, ao)
ctx = workflow.WithLocalActivityOptions(ctx, lao)

logger := workflow.GetLogger(ctx)
logger.Info("Codec Server workflow started", "input", input)

var result string

err := workflow.SideEffect(ctx, func(ctx workflow.Context) interface{} {
return uuid.New()
}).Get(&result)
if err != nil {
logger.Error("SideEffect failed.", "Error", err)
return "", err
}

err = workflow.ExecuteLocalActivity(ctx, Activity, input).Get(ctx, &result)
if err != nil {
logger.Error("Local Activity failed.", "Error", err)
return "", err
}

err = workflow.ExecuteActivity(ctx, Activity, input).Get(ctx, &result)
if err != nil {
logger.Error("Activity failed.", "Error", err)
return "", err
}

logger.Info("About to execute next activity...")
err = workflow.ExecuteActivity(ctx, Activity, input).Get(ctx, &result)
if err != nil {
logger.Error("Activity failed.", "Error", err)
return "", err
}

logger.Info("Codec Server workflow completed.", "result", result)

return result, nil
}

func Activity(ctx context.Context, input models.PayloadData) (string, error) {
logger := activity.GetLogger(ctx)
logger.Info("Activity", "input", input)

return fmt.Sprintf("Received %s", input.Data), nil
}

func TimeoutActivity(ctx context.Context, input models.PayloadData) error {

resp, err := http.Get("http://localhost:5173/timeout")
if err != nil {
return err
}
body, err := io.ReadAll(resp.Body)
_ = resp.Body.Close()
if err != nil {
return err
}

if string(body) == "SUCCEED" {
activity.GetLogger(ctx).Info("Activity unexpectedly succeeded.", "input", input)
return nil
}

return errors.New(string(body))
}
12 changes: 12 additions & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ toolchain go1.23.4
require (
github.com/gin-contrib/cors v1.5.0
github.com/gin-gonic/gin v1.9.1
github.com/google/uuid v1.6.0
go.temporal.io/api v1.49.1
go.temporal.io/sdk v1.34.0
)
Expand All @@ -15,28 +16,39 @@ require (
github.com/bytedance/sonic v1.10.1 // indirect
github.com/chenzhuoyu/base64x v0.0.0-20230717121745-296ad89f973d // indirect
github.com/chenzhuoyu/iasm v0.9.0 // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/facebookgo/clock v0.0.0-20150410010913-600d898af40a // indirect
github.com/gabriel-vasile/mimetype v1.4.2 // indirect
github.com/gin-contrib/sse v0.1.0 // indirect
github.com/go-playground/locales v0.14.1 // indirect
github.com/go-playground/universal-translator v0.18.1 // indirect
github.com/go-playground/validator/v10 v10.15.5 // indirect
github.com/goccy/go-json v0.10.2 // indirect
github.com/gogo/protobuf v1.3.2 // indirect
github.com/golang/mock v1.6.0 // indirect
github.com/grpc-ecosystem/go-grpc-middleware v1.4.0 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.22.0 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/klauspost/cpuid/v2 v2.2.5 // indirect
github.com/leodido/go-urn v1.2.4 // indirect
github.com/mattn/go-isatty v0.0.19 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.2 // indirect
github.com/nexus-rpc/sdk-go v0.3.0 // indirect
github.com/pelletier/go-toml/v2 v2.1.0 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
github.com/robfig/cron v1.2.0 // indirect
github.com/stretchr/objx v0.5.2 // indirect
github.com/stretchr/testify v1.10.0 // indirect
github.com/twitchyliquid64/golang-asm v0.15.1 // indirect
github.com/ugorji/go/codec v1.2.11 // indirect
golang.org/x/arch v0.5.0 // indirect
golang.org/x/crypto v0.35.0 // indirect
golang.org/x/net v0.36.0 // indirect
golang.org/x/sync v0.11.0 // indirect
golang.org/x/sys v0.30.0 // indirect
golang.org/x/text v0.22.0 // indirect
golang.org/x/time v0.3.0 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20240827150818-7e3bb234dfed // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20240827150818-7e3bb234dfed // indirect
google.golang.org/grpc v1.66.0 // indirect
Expand Down
Loading