Files

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
}