122 lines
5.1 KiB
Go
122 lines
5.1 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"errors"
|
|
"fmt"
|
|
"net/url"
|
|
"os"
|
|
"strings"
|
|
|
|
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
|
"git.ipao.vip/rogee/go-sip/internal/agent"
|
|
"git.ipao.vip/rogee/go-sip/internal/config"
|
|
"google.golang.org/protobuf/proto"
|
|
)
|
|
|
|
type recordingUploadRPC interface {
|
|
RequestUpload(context.Context, *agentv1.RequestUploadRequest) (*agentv1.RequestUploadResponse, error)
|
|
CompleteUpload(context.Context, *agentv1.CompleteUploadRequest) (*agentv1.CompleteUploadResponse, error)
|
|
}
|
|
|
|
func uploadRecording(ctx context.Context, cfg config.Config, client recordingUploadRPC, uploader agent.UploadClient, spool *agent.Spool, binding *agentv1.ExecutionBinding, asset *agentv1.AssetDescriptor, path string) (returned agent.UploadAttempt, returnErr error) {
|
|
if binding == nil || asset == nil {
|
|
return returned, errors.New("upload binding and asset are required")
|
|
}
|
|
id := stableUploadID(binding, asset)
|
|
lock, err := spool.LockUpload(id)
|
|
if err != nil {
|
|
return returned, err
|
|
}
|
|
defer func() { returnErr = errors.Join(returnErr, lock.Close()) }()
|
|
identityBytes, err := (proto.MarshalOptions{Deterministic: true}).Marshal(&agentv1.RequestUploadRequest{Binding: binding, Asset: asset, UploadId: id})
|
|
if err != nil {
|
|
return agent.UploadAttempt{}, err
|
|
}
|
|
digest := sha256.Sum256(identityBytes)
|
|
identity := hex.EncodeToString(digest[:])
|
|
record, err := spool.LoadUploadAttempt(id)
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
response, err := client.RequestUpload(ctx, &agentv1.RequestUploadRequest{Meta: uploadMeta(cfg, "request", id), Binding: binding, Asset: asset, UploadId: id})
|
|
if err != nil {
|
|
return record, fmt.Errorf("request upload %s: %w", id, err)
|
|
}
|
|
if response.GetReceipt().GetResult() != agentv1.ResultCode_RESULT_CODE_ACCEPTED || response.GetGrant() == nil {
|
|
return record, fmt.Errorf("request upload %s has no accepted grant", id)
|
|
}
|
|
grant := response.Grant
|
|
if grant.UploadId != id {
|
|
return record, errors.New("upload grant identity mismatch")
|
|
}
|
|
parsed, err := url.Parse(grant.TargetUrl)
|
|
if err != nil || parsed.Host == "" {
|
|
return record, errors.New("upload grant has invalid target URL")
|
|
}
|
|
record = agent.UploadAttempt{UploadID: id, Identity: identity, State: "attempted", ObjectKey: grant.ObjectKey, Binding: binding, Asset: asset, RequestID: uploadMeta(cfg, "request", id).OperationId}
|
|
// Durable exclusive claim must precede any network PUT, including an attempt
|
|
// that ends in an ambiguous transport failure.
|
|
if err := spool.ClaimUpload(record); err != nil {
|
|
return record, err
|
|
}
|
|
uploader.AllowedHosts = map[string]struct{}{strings.ToLower(parsed.Host): {}}
|
|
result, err := uploader.UploadFile(ctx, grant, path)
|
|
if err != nil {
|
|
return record, fmt.Errorf("upload %s data plane: %w", id, err)
|
|
}
|
|
if result.SizeBytes != asset.SizeBytes || !strings.EqualFold(result.SHA256, asset.ChecksumSha256) {
|
|
return record, errors.New("upload result does not match asset")
|
|
}
|
|
if err := spool.RecordUploadResult(id, result); err != nil {
|
|
return record, err
|
|
}
|
|
record.State, record.Result = "uploaded", result
|
|
} else if err != nil {
|
|
return record, err
|
|
} else if record.Identity != identity {
|
|
return record, errors.New("persisted upload binding or asset mismatch")
|
|
} else if record.State == "attempted" {
|
|
return record, errors.New("prior PUT outcome unknown or failed; automatic re-upload is forbidden")
|
|
}
|
|
if record.State == "completed" {
|
|
return record, nil
|
|
}
|
|
return notifyUploadedRecordingLocked(ctx, cfg, client, spool, record)
|
|
}
|
|
|
|
func notifyUploadedRecording(ctx context.Context, cfg config.Config, client recordingUploadRPC, spool *agent.Spool, record agent.UploadAttempt) (returned agent.UploadAttempt, returnErr error) {
|
|
lock, err := spool.LockUpload(record.UploadID)
|
|
if err != nil {
|
|
return returned, err
|
|
}
|
|
defer func() { returnErr = errors.Join(returnErr, lock.Close()) }()
|
|
current, err := spool.LoadUploadAttempt(record.UploadID)
|
|
if err != nil {
|
|
return returned, err
|
|
}
|
|
if current.State == "completed" {
|
|
return current, nil
|
|
}
|
|
return notifyUploadedRecordingLocked(ctx, cfg, client, spool, current)
|
|
}
|
|
|
|
func notifyUploadedRecordingLocked(ctx context.Context, cfg config.Config, client recordingUploadRPC, spool *agent.Spool, record agent.UploadAttempt) (agent.UploadAttempt, error) {
|
|
if record.State != "uploaded" || record.Binding == nil || record.Asset == nil {
|
|
return record, errors.New("successful upload and original notification metadata are required")
|
|
}
|
|
id := record.UploadID
|
|
response, err := client.CompleteUpload(ctx, &agentv1.CompleteUploadRequest{Meta: uploadMeta(cfg, "complete", id), Binding: record.Binding, Asset: record.Asset, UploadId: id, UploadedSizeBytes: record.Result.SizeBytes, UploadedChecksumSha256: record.Result.SHA256})
|
|
if err != nil {
|
|
return record, fmt.Errorf("notify upload %s: %w", id, err)
|
|
}
|
|
if response.GetReceipt().GetResult() != agentv1.ResultCode_RESULT_CODE_ACCEPTED || response.GetState() != agentv1.UploadState_UPLOAD_STATE_COMPLETED {
|
|
return record, errors.New("upload notification has not completed MQ delivery")
|
|
}
|
|
if err := spool.CompleteUploadNotification(id); err != nil {
|
|
return record, err
|
|
}
|
|
record.State = "completed"
|
|
return record, nil
|
|
}
|