mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
chore!: send modules archive over the proto messages (#21398)
# What this does Dynamic parameters caches the `./terraform/modules` directory for parameter usage. What this PR does is send over this archive to the provisioner when building workspaces. This allow terraform to skip downloading modules from their registries, a step that takes seconds. <img width="1223" height="429" alt="Screenshot From 2025-12-29 12-57-52" src="https://github.com/user-attachments/assets/16066e0a-ac79-4296-819d-924f4b0418dc" /> # Wire protocol The wire protocol reuses the same mechanism used to download the modules `provisoner -> coder`. It splits up large archives into multiple protobuf messages so larger archives can be sent under the message size limit. # 🚨 Behavior Change (Breaking Change) 🚨 **Before this PR** modules were downloaded on every workspace build. This means unpinned modules always fetched the latest version **After this PR** modules are cached at template import time, and their versions are effectively pinned for all subsequent workspace builds.
This commit is contained in:
@@ -509,6 +509,18 @@ func (s *server) acquireProtoJob(ctx context.Context, job database.ProvisionerJo
|
||||
if err != nil {
|
||||
return nil, failJob(fmt.Sprintf("get owner: %s", err))
|
||||
}
|
||||
|
||||
// Fetch the file id of the cached module files if it exists.
|
||||
versionModulesFile := ""
|
||||
tfvals, err := s.Database.GetTemplateVersionTerraformValues(ctx, templateVersion.ID)
|
||||
if err != nil && !xerrors.Is(err, sql.ErrNoRows) {
|
||||
// Older templates (before dynamic parameters) will not have cached module files.
|
||||
return nil, failJob(fmt.Sprintf("get template version terraform values: %s", err))
|
||||
}
|
||||
if err == nil && tfvals.CachedModuleFiles.Valid {
|
||||
versionModulesFile = tfvals.CachedModuleFiles.UUID.String()
|
||||
}
|
||||
|
||||
var ownerSSHPublicKey, ownerSSHPrivateKey string
|
||||
if ownerSSHKey, err := s.Database.GetGitSSHKey(ctx, owner.ID); err != nil {
|
||||
if !xerrors.Is(err, sql.ErrNoRows) {
|
||||
@@ -732,6 +744,7 @@ func (s *server) acquireProtoJob(ctx context.Context, job database.ProvisionerJo
|
||||
PrebuiltWorkspaceBuildStage: input.PrebuiltWorkspaceBuildStage,
|
||||
TaskId: task.ID.String(),
|
||||
TaskPrompt: task.Prompt,
|
||||
TemplateVersionModulesFile: versionModulesFile,
|
||||
},
|
||||
LogLevel: input.LogLevel,
|
||||
},
|
||||
@@ -1423,54 +1436,12 @@ func (s *server) prepareForNotifyWorkspaceManualBuildFailed(ctx context.Context,
|
||||
func (s *server) UploadFile(stream proto.DRPCProvisionerDaemon_UploadFileStream) error {
|
||||
var file *sdkproto.DataBuilder
|
||||
// Always terminate the stream with an empty response.
|
||||
//nolint:errcheck // We can't do much about send errors here.
|
||||
defer stream.SendAndClose(&proto.Empty{})
|
||||
|
||||
UploadFileStream:
|
||||
for {
|
||||
msg, err := stream.Recv()
|
||||
if err != nil {
|
||||
return xerrors.Errorf("receive complete job with files: %w", err)
|
||||
}
|
||||
|
||||
switch typed := msg.Type.(type) {
|
||||
case *sdkproto.FileUpload_DataUpload:
|
||||
if file != nil {
|
||||
return xerrors.New("unexpected file upload while waiting for file completion")
|
||||
}
|
||||
|
||||
file, err = sdkproto.NewDataBuilder(&sdkproto.DataUpload{
|
||||
UploadType: typed.DataUpload.UploadType,
|
||||
DataHash: typed.DataUpload.DataHash,
|
||||
FileSize: typed.DataUpload.FileSize,
|
||||
Chunks: typed.DataUpload.Chunks,
|
||||
})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("unable to create file upload: %w", err)
|
||||
}
|
||||
|
||||
if file.IsDone() {
|
||||
// If a file is 0 bytes, we can consider it done immediately.
|
||||
// This should never really happen in practice, but we handle it gracefully.
|
||||
break UploadFileStream
|
||||
}
|
||||
case *sdkproto.FileUpload_ChunkPiece:
|
||||
if file == nil {
|
||||
return xerrors.New("unexpected chunk piece while waiting for file upload")
|
||||
}
|
||||
|
||||
done, err := file.Add(&sdkproto.ChunkPiece{
|
||||
Data: typed.ChunkPiece.Data,
|
||||
FullDataHash: typed.ChunkPiece.FullDataHash,
|
||||
PieceIndex: typed.ChunkPiece.PieceIndex,
|
||||
})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("unable to add chunk piece: %w", err)
|
||||
}
|
||||
|
||||
if done {
|
||||
break UploadFileStream
|
||||
}
|
||||
}
|
||||
file, err := provisionersdk.HandleReceivingDataUpload(stream)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
fileData, err := file.Complete()
|
||||
@@ -1518,9 +1489,71 @@ UploadFileStream:
|
||||
return nil
|
||||
}
|
||||
|
||||
func (*server) DownloadFile(_ *proto.FileRequest, _ proto.DRPCProvisionerDaemon_DownloadFileStream) error {
|
||||
// TODO implemented in follow up PR
|
||||
panic("implement me")
|
||||
// DownloadFile pulls the requested file from the database and sends it over the protobuf stream in chunks.
|
||||
func (s *server) DownloadFile(request *proto.FileRequest, stream proto.DRPCProvisionerDaemon_DownloadFileStream) error {
|
||||
//nolint:errcheck
|
||||
defer stream.CloseSend()
|
||||
//nolint:gocritic // Provisionerd is the actor here.
|
||||
ctx := dbauthz.AsProvisionerd(stream.Context())
|
||||
|
||||
// A graceful error message will help debugging.
|
||||
fail := func(err error) error {
|
||||
_ = stream.Send(&sdkproto.FileUpload{
|
||||
Type: &sdkproto.FileUpload_Error{
|
||||
Error: &sdkproto.FailedFile{
|
||||
Error: err.Error(),
|
||||
},
|
||||
},
|
||||
})
|
||||
return err
|
||||
}
|
||||
if request.FileId == "" || request.FileId == uuid.Nil.String() {
|
||||
return fail(xerrors.New("file id is required"))
|
||||
}
|
||||
|
||||
fid, err := uuid.Parse(request.FileId)
|
||||
if err != nil {
|
||||
return fail(xerrors.Errorf("invalid file id: %w", err))
|
||||
}
|
||||
|
||||
file, err := s.Database.GetFileByID(ctx, fid)
|
||||
if err != nil {
|
||||
return fail(xerrors.Errorf("get file: %w", err))
|
||||
}
|
||||
|
||||
switch request.UploadType {
|
||||
case sdkproto.DataUploadType_UPLOAD_TYPE_MODULE_FILES:
|
||||
// This check is not perfect. If these conditions are not true, then the file is not a modules file.
|
||||
if file.CreatedBy != uuid.Nil || file.Mimetype != tarMimeType {
|
||||
return fail(xerrors.Errorf("file %s is not a modules file", fid))
|
||||
}
|
||||
default:
|
||||
return fail(xerrors.Errorf("unsupported file upload type: %s", request.UploadType))
|
||||
}
|
||||
|
||||
upload, chunks := sdkproto.BytesToDataUpload(sdkproto.DataUploadType_UPLOAD_TYPE_MODULE_FILES, file.Data)
|
||||
|
||||
err = stream.Send(&sdkproto.FileUpload{
|
||||
Type: &sdkproto.FileUpload_DataUpload{DataUpload: upload},
|
||||
})
|
||||
if err != nil {
|
||||
return fail(xerrors.Errorf("send file upload: %w", err))
|
||||
}
|
||||
|
||||
for i, c := range chunks {
|
||||
if ctx.Err() != nil {
|
||||
return fail(ctx.Err())
|
||||
}
|
||||
|
||||
err = stream.Send(&sdkproto.FileUpload{
|
||||
Type: &sdkproto.FileUpload_ChunkPiece{ChunkPiece: c},
|
||||
})
|
||||
if err != nil {
|
||||
return fail(xerrors.Errorf("send chunk piece %d: %w", i, err))
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// CompleteJob is triggered by a provision daemon to mark a provisioner job as completed.
|
||||
|
||||
@@ -120,7 +120,7 @@ func TestUploadFileErrorScenarios(t *testing.T) {
|
||||
stream.messages <- up
|
||||
|
||||
err := server.UploadFile(stream)
|
||||
require.ErrorContains(t, err, "unexpected file upload while waiting for file completion")
|
||||
require.ErrorContains(t, err, "unexpected file download while waiting for file completion")
|
||||
require.True(t, stream.isDone(), "stream should be done after error")
|
||||
})
|
||||
|
||||
|
||||
Reference in New Issue
Block a user