413 lines
16 KiB
Go
413 lines
16 KiB
Go
package service
|
|
|
|
import (
|
|
"bytes"
|
|
"errors"
|
|
"fmt"
|
|
"sort"
|
|
"time"
|
|
|
|
"browser.local/platform/domain"
|
|
"browser.local/platform/repo"
|
|
"browser.local/platform/validator"
|
|
)
|
|
|
|
func (svc *CoreService) OpenArtifactTransfer(open domain.ArtifactTransferOpen) (domain.ArtifactTransferOpenResult, error) {
|
|
open = domain.CopyArtifactTransferOpen(open)
|
|
if err := validator.ValidateArtifactTransferOpen(open); err != nil {
|
|
return domain.ArtifactTransferOpenResult{}, err
|
|
}
|
|
if err := svc.validateRunSession(open.RunEndpointID, open.SessionToken); err != nil {
|
|
return domain.ArtifactTransferOpenResult{}, err
|
|
}
|
|
|
|
stamp := svc.now()
|
|
svc.artifactMu.Lock()
|
|
defer svc.artifactMu.Unlock()
|
|
|
|
if session, exists := svc.findArtifactTransferByIdempotency(open.RunEndpointID, open.IdempotencyKey); exists {
|
|
if err := validateArtifactTransferOpenMatchesSession(open, session); err != nil {
|
|
return domain.ArtifactTransferOpenResult{}, err
|
|
}
|
|
artifact, err := svc.store.Artifacts().Get(session.ArtifactID)
|
|
if err != nil {
|
|
return domain.ArtifactTransferOpenResult{}, err
|
|
}
|
|
return artifactTransferOpenResult(session, artifact, true, stamp), nil
|
|
}
|
|
|
|
if err := svc.validateArtifactTransferOwner(open); err != nil {
|
|
return domain.ArtifactTransferOpenResult{}, err
|
|
}
|
|
|
|
artifact, err := svc.store.Artifacts().Get(open.ArtifactID)
|
|
if err != nil {
|
|
if !errors.Is(err, repo.ErrNotFound) {
|
|
return domain.ArtifactTransferOpenResult{}, err
|
|
}
|
|
artifact = domain.Artifact{
|
|
ID: open.ArtifactID,
|
|
OwnerKind: open.OwnerKind,
|
|
OwnerID: open.OwnerID,
|
|
SizeBytes: open.SizeBytes,
|
|
Checksum: open.Checksum,
|
|
State: domain.ArtifactStateUploading,
|
|
CreatedAt: stamp,
|
|
UpdatedAt: stamp,
|
|
}
|
|
if err := validator.ValidateArtifact(artifact); err != nil {
|
|
return domain.ArtifactTransferOpenResult{}, err
|
|
}
|
|
if err := svc.store.Artifacts().Create(artifact); err != nil {
|
|
return domain.ArtifactTransferOpenResult{}, err
|
|
}
|
|
} else if err := validateArtifactMatchesTransferOpen(artifact, open); err != nil {
|
|
return domain.ArtifactTransferOpenResult{}, err
|
|
}
|
|
|
|
svc.artifactTransferSeq++
|
|
session := domain.ArtifactTransferSession{
|
|
TransferID: fmt.Sprintf("artifact-transfer:%s:%d:%d", open.ArtifactID, stamp.UnixNano(), svc.artifactTransferSeq),
|
|
RunEndpointID: open.RunEndpointID,
|
|
ArtifactID: open.ArtifactID,
|
|
Direction: open.Direction,
|
|
OwnerKind: open.OwnerKind,
|
|
OwnerID: open.OwnerID,
|
|
SizeBytes: open.SizeBytes,
|
|
ChunkSizeBytes: open.ChunkSizeBytes,
|
|
Checksum: open.Checksum,
|
|
IdempotencyKey: open.IdempotencyKey,
|
|
TotalChunks: artifactTotalChunks(open.SizeBytes, open.ChunkSizeBytes),
|
|
ReceivedChunks: map[int]domain.ArtifactChunkRecord{},
|
|
CreatedAt: stamp,
|
|
UpdatedAt: stamp,
|
|
}
|
|
svc.artifactTransfers[session.TransferID] = domain.CopyArtifactTransferSession(session)
|
|
return artifactTransferOpenResult(session, artifact, false, stamp), nil
|
|
}
|
|
|
|
func (svc *CoreService) UploadArtifactChunk(chunk domain.ArtifactChunkUpload) (domain.ArtifactChunkUploadResult, error) {
|
|
chunk = domain.CopyArtifactChunkUpload(chunk)
|
|
if err := validator.ValidateArtifactChunkUpload(chunk); err != nil {
|
|
return domain.ArtifactChunkUploadResult{}, err
|
|
}
|
|
if err := svc.validateRunSession(chunk.RunEndpointID, chunk.SessionToken); err != nil {
|
|
return domain.ArtifactChunkUploadResult{}, err
|
|
}
|
|
|
|
stamp := svc.now()
|
|
svc.artifactMu.Lock()
|
|
defer svc.artifactMu.Unlock()
|
|
|
|
session, err := svc.getArtifactTransferSession(chunk.TransferID)
|
|
if err != nil {
|
|
return domain.ArtifactChunkUploadResult{}, err
|
|
}
|
|
if err := validateArtifactChunkMatchesSession(chunk, session); err != nil {
|
|
return domain.ArtifactChunkUploadResult{}, err
|
|
}
|
|
|
|
if existing, exists := session.ReceivedChunks[chunk.ChunkIndex]; exists {
|
|
if existing.Offset == chunk.Offset && existing.SizeBytes == chunk.SizeBytes && existing.Checksum == chunk.Checksum && bytes.Equal(existing.Payload, chunk.Payload) {
|
|
return artifactChunkUploadResult(session, chunk.ChunkIndex, true, stamp), nil
|
|
}
|
|
return domain.ArtifactChunkUploadResult{}, validationError("artifact chunk conflicts with acknowledged chunk")
|
|
}
|
|
|
|
session.ReceivedChunks[chunk.ChunkIndex] = domain.ArtifactChunkRecord{
|
|
ChunkIndex: chunk.ChunkIndex,
|
|
Offset: chunk.Offset,
|
|
SizeBytes: chunk.SizeBytes,
|
|
Checksum: chunk.Checksum,
|
|
Payload: domain.CopyBytes(chunk.Payload),
|
|
ReceivedAt: stamp,
|
|
}
|
|
session.UpdatedAt = stamp
|
|
svc.artifactTransfers[session.TransferID] = domain.CopyArtifactTransferSession(session)
|
|
return artifactChunkUploadResult(session, chunk.ChunkIndex, false, stamp), nil
|
|
}
|
|
|
|
func (svc *CoreService) QueryArtifactTransferStatus(query domain.ArtifactTransferStatusQuery) (domain.ArtifactTransferStatusResult, error) {
|
|
if err := validator.ValidateArtifactTransferStatusQuery(query); err != nil {
|
|
return domain.ArtifactTransferStatusResult{}, err
|
|
}
|
|
if err := svc.validateRunSession(query.RunEndpointID, query.SessionToken); err != nil {
|
|
return domain.ArtifactTransferStatusResult{}, err
|
|
}
|
|
|
|
stamp := svc.now()
|
|
svc.artifactMu.Lock()
|
|
defer svc.artifactMu.Unlock()
|
|
|
|
session, err := svc.getArtifactTransferSession(query.TransferID)
|
|
if err != nil {
|
|
return domain.ArtifactTransferStatusResult{}, err
|
|
}
|
|
if err := validateArtifactTransferStatusMatchesSession(query, session); err != nil {
|
|
return domain.ArtifactTransferStatusResult{}, err
|
|
}
|
|
return artifactTransferStatusResult(session, stamp), nil
|
|
}
|
|
|
|
func (svc *CoreService) CompleteArtifactTransfer(complete domain.ArtifactTransferComplete) (domain.ArtifactTransferCompleteResult, error) {
|
|
if err := validator.ValidateArtifactTransferComplete(complete); err != nil {
|
|
return domain.ArtifactTransferCompleteResult{}, err
|
|
}
|
|
if err := svc.validateRunSession(complete.RunEndpointID, complete.SessionToken); err != nil {
|
|
return domain.ArtifactTransferCompleteResult{}, err
|
|
}
|
|
|
|
stamp := svc.now()
|
|
svc.artifactMu.Lock()
|
|
defer svc.artifactMu.Unlock()
|
|
|
|
session, err := svc.getArtifactTransferSession(complete.TransferID)
|
|
if err != nil {
|
|
return domain.ArtifactTransferCompleteResult{}, err
|
|
}
|
|
if err := validateArtifactCompleteMatchesSession(complete, session); err != nil {
|
|
return domain.ArtifactTransferCompleteResult{}, err
|
|
}
|
|
artifact, err := svc.store.Artifacts().Get(session.ArtifactID)
|
|
if err != nil {
|
|
return domain.ArtifactTransferCompleteResult{}, err
|
|
}
|
|
if session.Completed {
|
|
return domain.ArtifactTransferCompleteResult{Accepted: true, TransferID: session.TransferID, Artifact: artifact, Completed: true, ServerTime: stamp}, nil
|
|
}
|
|
if len(session.ReceivedChunks) != session.TotalChunks {
|
|
return domain.ArtifactTransferCompleteResult{}, validationError("artifact transfer has missing chunks")
|
|
}
|
|
|
|
payload := make([]byte, 0, int(session.SizeBytes))
|
|
for index := 0; index < session.TotalChunks; index++ {
|
|
record, exists := session.ReceivedChunks[index]
|
|
if !exists {
|
|
return domain.ArtifactTransferCompleteResult{}, validationError("artifact transfer has missing chunks")
|
|
}
|
|
payload = append(payload, record.Payload...)
|
|
}
|
|
if int64(len(payload)) != session.SizeBytes {
|
|
return domain.ArtifactTransferCompleteResult{}, validationError("artifact transfer size does not match metadata")
|
|
}
|
|
if checksum := validator.BytesChecksum(payload); checksum != session.Checksum {
|
|
return domain.ArtifactTransferCompleteResult{}, validationError("artifact transfer checksum does not match metadata")
|
|
}
|
|
|
|
artifact.SizeBytes = session.SizeBytes
|
|
artifact.Checksum = session.Checksum
|
|
artifact.State = domain.ArtifactStateAvailable
|
|
artifact.UpdatedAt = stamp
|
|
if err := validator.ValidateArtifact(artifact); err != nil {
|
|
return domain.ArtifactTransferCompleteResult{}, err
|
|
}
|
|
if err := svc.store.Artifacts().Update(artifact); err != nil {
|
|
return domain.ArtifactTransferCompleteResult{}, err
|
|
}
|
|
session.Completed = true
|
|
session.UpdatedAt = stamp
|
|
svc.artifactTransfers[session.TransferID] = domain.CopyArtifactTransferSession(session)
|
|
return domain.ArtifactTransferCompleteResult{Accepted: true, TransferID: session.TransferID, Artifact: artifact, Completed: true, ServerTime: stamp}, nil
|
|
}
|
|
|
|
func (svc *CoreService) getArtifactTransferSession(transferID string) (domain.ArtifactTransferSession, error) {
|
|
session, exists := svc.artifactTransfers[transferID]
|
|
if !exists {
|
|
return domain.ArtifactTransferSession{}, repo.ErrNotFound
|
|
}
|
|
return domain.CopyArtifactTransferSession(session), nil
|
|
}
|
|
|
|
func (svc *CoreService) findArtifactTransferByIdempotency(runEndpointID string, idempotencyKey string) (domain.ArtifactTransferSession, bool) {
|
|
transferIDs := make([]string, 0, len(svc.artifactTransfers))
|
|
for transferID := range svc.artifactTransfers {
|
|
transferIDs = append(transferIDs, transferID)
|
|
}
|
|
sort.Strings(transferIDs)
|
|
for _, transferID := range transferIDs {
|
|
session := svc.artifactTransfers[transferID]
|
|
if session.RunEndpointID == runEndpointID && session.IdempotencyKey == idempotencyKey {
|
|
return domain.CopyArtifactTransferSession(session), true
|
|
}
|
|
}
|
|
return domain.ArtifactTransferSession{}, false
|
|
}
|
|
|
|
func (svc *CoreService) validateArtifactTransferOwner(open domain.ArtifactTransferOpen) error {
|
|
switch open.OwnerKind {
|
|
case domain.ArtifactOwnerKindJob:
|
|
job, err := svc.store.Jobs().Get(open.OwnerID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if job.RunEndpointID != open.RunEndpointID {
|
|
return validationError("artifact owner job must belong to runEndpointId")
|
|
}
|
|
case domain.ArtifactOwnerKindServerInstance:
|
|
instance, err := svc.store.ServerInstances().Get(open.OwnerID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if instance.RunEndpointID != open.RunEndpointID {
|
|
return validationError("artifact owner server instance must belong to runEndpointId")
|
|
}
|
|
if instance.State == domain.ServerInstanceStateDeleted {
|
|
return validationError("artifact owner server instance must not be deleted")
|
|
}
|
|
default:
|
|
return validationError("ownerKind must be job or server-instance for run uploads")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func validateArtifactMatchesTransferOpen(artifact domain.Artifact, open domain.ArtifactTransferOpen) error {
|
|
if artifact.OwnerKind != open.OwnerKind {
|
|
return validationError("artifact ownerKind must match transfer")
|
|
}
|
|
if artifact.OwnerID != open.OwnerID {
|
|
return validationError("artifact ownerId must match transfer")
|
|
}
|
|
if artifact.SizeBytes != open.SizeBytes {
|
|
return validationError("artifact sizeBytes must match transfer")
|
|
}
|
|
if artifact.Checksum != open.Checksum {
|
|
return validationError("artifact checksum must match transfer")
|
|
}
|
|
if artifact.State != domain.ArtifactStateUploading {
|
|
return validationError("artifact must be uploading")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func validateArtifactTransferOpenMatchesSession(open domain.ArtifactTransferOpen, session domain.ArtifactTransferSession) error {
|
|
if session.ArtifactID != open.ArtifactID || session.Direction != open.Direction || session.OwnerKind != open.OwnerKind || session.OwnerID != open.OwnerID || session.SizeBytes != open.SizeBytes || session.ChunkSizeBytes != open.ChunkSizeBytes || session.Checksum != open.Checksum {
|
|
return validationError("artifact transfer idempotency key conflicts with existing transfer")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func validateArtifactChunkMatchesSession(chunk domain.ArtifactChunkUpload, session domain.ArtifactTransferSession) error {
|
|
if session.Completed {
|
|
return validationError("artifact transfer is already complete")
|
|
}
|
|
if session.RunEndpointID != chunk.RunEndpointID {
|
|
return validationError("runEndpointId must match artifact transfer")
|
|
}
|
|
if session.ArtifactID != chunk.ArtifactID {
|
|
return validationError("artifactId must match artifact transfer")
|
|
}
|
|
if chunk.ChunkIndex >= session.TotalChunks {
|
|
return validationError("chunkIndex exceeds transfer chunk count")
|
|
}
|
|
expectedOffset := int64(chunk.ChunkIndex) * int64(session.ChunkSizeBytes)
|
|
if chunk.Offset != expectedOffset {
|
|
return validationError("offset must match chunk index")
|
|
}
|
|
expectedSize := expectedChunkSize(session, chunk.ChunkIndex)
|
|
if chunk.SizeBytes != expectedSize {
|
|
return validationError("sizeBytes must match expected chunk size")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func validateArtifactTransferStatusMatchesSession(query domain.ArtifactTransferStatusQuery, session domain.ArtifactTransferSession) error {
|
|
if session.RunEndpointID != query.RunEndpointID {
|
|
return validationError("runEndpointId must match artifact transfer")
|
|
}
|
|
if session.ArtifactID != query.ArtifactID {
|
|
return validationError("artifactId must match artifact transfer")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func validateArtifactCompleteMatchesSession(complete domain.ArtifactTransferComplete, session domain.ArtifactTransferSession) error {
|
|
if session.RunEndpointID != complete.RunEndpointID {
|
|
return validationError("runEndpointId must match artifact transfer")
|
|
}
|
|
if session.ArtifactID != complete.ArtifactID {
|
|
return validationError("artifactId must match artifact transfer")
|
|
}
|
|
if session.SizeBytes != complete.SizeBytes {
|
|
return validationError("sizeBytes must match artifact transfer")
|
|
}
|
|
if session.Checksum != complete.Checksum {
|
|
return validationError("checksum must match artifact transfer")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func artifactTransferOpenResult(session domain.ArtifactTransferSession, artifact domain.Artifact, duplicate bool, stamp time.Time) domain.ArtifactTransferOpenResult {
|
|
return domain.ArtifactTransferOpenResult{
|
|
Accepted: true,
|
|
TransferID: session.TransferID,
|
|
Direction: session.Direction,
|
|
Artifact: artifact,
|
|
TotalChunks: session.TotalChunks,
|
|
ChunkSizeBytes: session.ChunkSizeBytes,
|
|
ReceivedChunkIndexes: receivedArtifactChunkIndexes(session),
|
|
NextMissingChunkIndex: nextMissingArtifactChunkIndex(session),
|
|
Completed: session.Completed,
|
|
Duplicate: duplicate,
|
|
ServerTime: stamp,
|
|
}
|
|
}
|
|
|
|
func artifactChunkUploadResult(session domain.ArtifactTransferSession, chunkIndex int, duplicate bool, stamp time.Time) domain.ArtifactChunkUploadResult {
|
|
return domain.ArtifactChunkUploadResult{
|
|
Accepted: true,
|
|
TransferID: session.TransferID,
|
|
ArtifactID: session.ArtifactID,
|
|
ChunkIndex: chunkIndex,
|
|
ReceivedChunkIndexes: receivedArtifactChunkIndexes(session),
|
|
NextMissingChunkIndex: nextMissingArtifactChunkIndex(session),
|
|
Duplicate: duplicate,
|
|
ServerTime: stamp,
|
|
}
|
|
}
|
|
|
|
func artifactTransferStatusResult(session domain.ArtifactTransferSession, stamp time.Time) domain.ArtifactTransferStatusResult {
|
|
return domain.ArtifactTransferStatusResult{
|
|
Accepted: true,
|
|
TransferID: session.TransferID,
|
|
ArtifactID: session.ArtifactID,
|
|
Direction: session.Direction,
|
|
TotalChunks: session.TotalChunks,
|
|
ChunkSizeBytes: session.ChunkSizeBytes,
|
|
ReceivedChunkIndexes: receivedArtifactChunkIndexes(session),
|
|
NextMissingChunkIndex: nextMissingArtifactChunkIndex(session),
|
|
Completed: session.Completed,
|
|
ServerTime: stamp,
|
|
}
|
|
}
|
|
|
|
func artifactTotalChunks(sizeBytes int64, chunkSizeBytes int) int {
|
|
return int((sizeBytes + int64(chunkSizeBytes) - 1) / int64(chunkSizeBytes))
|
|
}
|
|
|
|
func expectedChunkSize(session domain.ArtifactTransferSession, chunkIndex int) int {
|
|
offset := int64(chunkIndex) * int64(session.ChunkSizeBytes)
|
|
remaining := session.SizeBytes - offset
|
|
if remaining < int64(session.ChunkSizeBytes) {
|
|
return int(remaining)
|
|
}
|
|
return session.ChunkSizeBytes
|
|
}
|
|
|
|
func receivedArtifactChunkIndexes(session domain.ArtifactTransferSession) []int {
|
|
indexes := make([]int, 0, len(session.ReceivedChunks))
|
|
for index := range session.ReceivedChunks {
|
|
indexes = append(indexes, index)
|
|
}
|
|
sort.Ints(indexes)
|
|
return indexes
|
|
}
|
|
|
|
func nextMissingArtifactChunkIndex(session domain.ArtifactTransferSession) int {
|
|
for index := 0; index < session.TotalChunks; index++ {
|
|
if _, exists := session.ReceivedChunks[index]; !exists {
|
|
return index
|
|
}
|
|
}
|
|
return session.TotalChunks
|
|
}
|