Files
run/runtime/file_artifact_transfer.go

130 lines
3.3 KiB
Go

package runtime
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"io"
"os"
"browser.local/run/protocol"
)
func (worker *Worker) uploadFileArtifact(ctx context.Context, assignment protocol.RunJobAssignment, artifactID string, targetPath string, sizeBytes int64, checksum string) error {
state, err := worker.registeredState()
if err != nil {
return err
}
opened, err := worker.client.OpenArtifactTransfer(ctx, protocol.ArtifactTransferOpenRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
ArtifactID: artifactID,
Direction: "upload",
OwnerKind: "job",
OwnerID: assignment.JobID,
SizeBytes: sizeBytes,
ChunkSizeBytes: fileArtifactChunkSize,
Checksum: checksum,
IdempotencyKey: "file-read:" + assignment.JobID,
})
if err != nil {
return err
}
received := map[int]bool{}
for _, index := range opened.ReceivedChunkIndexes {
received[index] = true
}
file, err := os.Open(targetPath)
if err != nil {
return err
}
defer file.Close()
buffer := make([]byte, fileArtifactChunkSize)
for index, offset := 0, int64(0); offset < sizeBytes; index, offset = index+1, offset+int64(fileArtifactChunkSize) {
length := fileArtifactChunkSize
if remaining := sizeBytes - offset; remaining < int64(length) {
length = int(remaining)
}
if received[index] {
continue
}
if err := ctx.Err(); err != nil {
return err
}
read, err := file.ReadAt(buffer[:length], offset)
if err != nil && !(errors.Is(err, io.EOF) && read == length) {
return err
}
if read != length {
return fmt.Errorf("file artifact chunk is shorter than expected")
}
chunk := buffer[:length]
state, err := worker.registeredState()
if err != nil {
return err
}
if _, err := worker.client.UploadArtifactChunk(ctx, protocol.ArtifactChunkUploadRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
TransferID: opened.TransferID,
ArtifactID: artifactID,
ChunkIndex: index,
Offset: offset,
SizeBytes: length,
Checksum: bytesChecksum(chunk),
Payload: chunk,
}); err != nil {
return err
}
}
state, err = worker.registeredState()
if err != nil {
return err
}
completed, err := worker.client.CompleteArtifactTransfer(ctx, protocol.ArtifactTransferCompleteRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
TransferID: opened.TransferID,
ArtifactID: artifactID,
Checksum: checksum,
SizeBytes: sizeBytes,
})
if err != nil {
return err
}
if !completed.Completed || completed.Artifact.State != "available" {
return fmt.Errorf("artifact transfer did not complete")
}
return nil
}
func checksumServerFileForArtifact(ctx context.Context, targetPath string) (string, error) {
file, err := os.Open(targetPath)
if err != nil {
return "", err
}
defer file.Close()
hash := sha256.New()
buffer := make([]byte, fileArtifactChunkSize)
for {
if err := ctx.Err(); err != nil {
return "", err
}
read, err := file.Read(buffer)
if read > 0 {
if _, writeErr := hash.Write(buffer[:read]); writeErr != nil {
return "", writeErr
}
}
if errors.Is(err, io.EOF) {
break
}
if err != nil {
return "", err
}
}
return "sha256:" + hex.EncodeToString(hash.Sum(nil)), nil
}