feat: Add echo provisioner (#162)

This replaces the cdr-basic provisioner type with
"echo". It reads binary data from the directory
and returns the responses in order.

This is used to test project and workspace job logic.
This commit is contained in:
Kyle Carberry
2022-02-04 16:51:54 -06:00
committed by GitHub
parent 65de6eef9c
commit 682238d384
15 changed files with 366 additions and 129 deletions
+138
View File
@@ -0,0 +1,138 @@
package echo
import (
"archive/tar"
"bytes"
"context"
"fmt"
"os"
"path/filepath"
"golang.org/x/xerrors"
protobuf "google.golang.org/protobuf/proto"
"github.com/coder/coder/provisionersdk"
"github.com/coder/coder/provisionersdk/proto"
)
var (
// ParseComplete is a helper to indicate an empty parse completion.
ParseComplete = []*proto.Parse_Response{{
Type: &proto.Parse_Response_Complete{
Complete: &proto.Parse_Complete{},
},
}}
// ProvisionComplete is a helper to indicate an empty provision completion.
ProvisionComplete = []*proto.Provision_Response{{
Type: &proto.Provision_Response_Complete{
Complete: &proto.Provision_Complete{},
},
}}
)
// Serve starts the echo provisioner.
func Serve(ctx context.Context, options *provisionersdk.ServeOptions) error {
return provisionersdk.Serve(ctx, &echo{}, options)
}
// The echo provisioner serves as a dummy provisioner primarily
// used for testing. It echos responses from JSON files in the
// format %d.protobuf. It's used for testing.
type echo struct {
}
// Parse reads requests from the provided directory to stream responses.
func (*echo) Parse(request *proto.Parse_Request, stream proto.DRPCProvisioner_ParseStream) error {
for index := 0; ; index++ {
path := filepath.Join(request.Directory, fmt.Sprintf("%d.parse.protobuf", index))
_, err := os.Stat(path)
if err != nil {
break
}
data, err := os.ReadFile(path)
if err != nil {
return xerrors.Errorf("read file %q: %w", path, err)
}
var response proto.Parse_Response
err = protobuf.Unmarshal(data, &response)
if err != nil {
return xerrors.Errorf("unmarshal: %w", err)
}
err = stream.Send(&response)
if err != nil {
return err
}
}
return nil
}
// Provision reads requests from the provided directory to stream responses.
func (*echo) Provision(request *proto.Provision_Request, stream proto.DRPCProvisioner_ProvisionStream) error {
for index := 0; ; index++ {
path := filepath.Join(request.Directory, fmt.Sprintf("%d.provision.protobuf", index))
_, err := os.Stat(path)
if err != nil {
break
}
data, err := os.ReadFile(path)
if err != nil {
return xerrors.Errorf("read file %q: %w", path, err)
}
var response proto.Provision_Response
err = protobuf.Unmarshal(data, &response)
if err != nil {
return xerrors.Errorf("unmarshal: %w", err)
}
err = stream.Send(&response)
if err != nil {
return err
}
}
return nil
}
// Tar returns a tar archive of responses to provisioner operations.
func Tar(parseResponses []*proto.Parse_Response, provisionResponses []*proto.Provision_Response) ([]byte, error) {
var buffer bytes.Buffer
writer := tar.NewWriter(&buffer)
for index, response := range parseResponses {
data, err := protobuf.Marshal(response)
if err != nil {
return nil, err
}
err = writer.WriteHeader(&tar.Header{
Name: fmt.Sprintf("%d.parse.protobuf", index),
Size: int64(len(data)),
})
if err != nil {
return nil, err
}
_, err = writer.Write(data)
if err != nil {
return nil, err
}
}
for index, response := range provisionResponses {
data, err := protobuf.Marshal(response)
if err != nil {
return nil, err
}
err = writer.WriteHeader(&tar.Header{
Name: fmt.Sprintf("%d.provision.protobuf", index),
Size: int64(len(data)),
})
if err != nil {
return nil, err
}
_, err = writer.Write(data)
if err != nil {
return nil, err
}
}
err := writer.Flush()
if err != nil {
return nil, err
}
return buffer.Bytes(), nil
}
+123
View File
@@ -0,0 +1,123 @@
package echo_test
import (
"archive/tar"
"bytes"
"context"
"io"
"os"
"path/filepath"
"testing"
"github.com/stretchr/testify/require"
"github.com/coder/coder/provisioner/echo"
"github.com/coder/coder/provisionersdk"
"github.com/coder/coder/provisionersdk/proto"
)
func TestEcho(t *testing.T) {
t.Parallel()
// Create an in-memory provisioner to communicate with.
client, server := provisionersdk.TransportPipe()
ctx, cancelFunc := context.WithCancel(context.Background())
t.Cleanup(func() {
_ = client.Close()
_ = server.Close()
cancelFunc()
})
go func() {
err := echo.Serve(ctx, &provisionersdk.ServeOptions{
Listener: server,
})
require.NoError(t, err)
}()
api := proto.NewDRPCProvisionerClient(provisionersdk.Conn(client))
t.Run("Parse", func(t *testing.T) {
t.Parallel()
responses := []*proto.Parse_Response{{
Type: &proto.Parse_Response_Log{
Log: &proto.Log{
Output: "log-output",
},
},
}, {
Type: &proto.Parse_Response_Complete{
Complete: &proto.Parse_Complete{
ParameterSchemas: []*proto.ParameterSchema{{
Name: "parameter-schema",
}},
},
},
}}
data, err := echo.Tar(responses, nil)
require.NoError(t, err)
client, err := api.Parse(ctx, &proto.Parse_Request{
Directory: unpackTar(t, data),
})
require.NoError(t, err)
log, err := client.Recv()
require.NoError(t, err)
require.Equal(t, responses[0].GetLog().Output, log.GetLog().Output)
complete, err := client.Recv()
require.NoError(t, err)
require.Equal(t, responses[1].GetComplete().ParameterSchemas[0].Name,
complete.GetComplete().ParameterSchemas[0].Name)
})
t.Run("Provision", func(t *testing.T) {
t.Parallel()
responses := []*proto.Provision_Response{{
Type: &proto.Provision_Response_Log{
Log: &proto.Log{
Output: "log-output",
},
},
}, {
Type: &proto.Provision_Response_Complete{
Complete: &proto.Provision_Complete{
Resources: []*proto.Resource{{
Name: "resource",
}},
},
},
}}
data, err := echo.Tar(nil, responses)
require.NoError(t, err)
client, err := api.Provision(ctx, &proto.Provision_Request{
Directory: unpackTar(t, data),
})
require.NoError(t, err)
log, err := client.Recv()
require.NoError(t, err)
require.Equal(t, responses[0].GetLog().Output, log.GetLog().Output)
complete, err := client.Recv()
require.NoError(t, err)
require.Equal(t, responses[1].GetComplete().Resources[0].Name,
complete.GetComplete().Resources[0].Name)
})
}
func unpackTar(t *testing.T, data []byte) string {
directory := t.TempDir()
reader := tar.NewReader(bytes.NewReader(data))
for {
header, err := reader.Next()
if err != nil {
break
}
// #nosec
path := filepath.Join(directory, header.Name)
file, err := os.OpenFile(path, os.O_CREATE|os.O_RDWR, 0600)
require.NoError(t, err)
_, err = io.CopyN(file, reader, 1<<20)
require.ErrorIs(t, err, io.EOF)
err = file.Close()
require.NoError(t, err)
}
return directory
}