Files
2026-08-26 22:00:23 +08:00

426 lines
16 KiB
Go

package service
import (
"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,
}
if err := svc.artifactStore.SaveTransfer(session); err != nil {
return domain.ArtifactTransferOpenResult{}, err
}
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 {
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
if err := svc.artifactStore.SaveTransfer(session); err != nil {
return domain.ArtifactChunkUploadResult{}, err
}
svc.artifactTransfers[session.TransferID] = svc.artifactTransferSessionForMemory(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")
}
for index := 0; index < session.TotalChunks; index++ {
if _, exists := session.ReceivedChunks[index]; !exists {
return domain.ArtifactTransferCompleteResult{}, validationError("artifact transfer has missing chunks")
}
}
if err := svc.artifactStore.CommitTransferPayload(session); err != nil {
return domain.ArtifactTransferCompleteResult{}, err
}
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
if err := svc.artifactStore.SaveTransfer(session); err != nil {
return domain.ArtifactTransferCompleteResult{}, err
}
svc.artifactTransfers[session.TransferID] = svc.artifactTransferSessionForMemory(session)
return domain.ArtifactTransferCompleteResult{Accepted: true, TransferID: session.TransferID, Artifact: artifact, Completed: true, ServerTime: stamp}, nil
}
func (svc *CoreService) artifactTransferSessionForMemory(session domain.ArtifactTransferSession) domain.ArtifactTransferSession {
cached := domain.CopyArtifactTransferSession(session)
if _, durableFileStore := svc.artifactStore.(*FileArtifactBodyStore); durableFileStore {
for index, record := range cached.ReceivedChunks {
record.Payload = nil
cached.ReceivedChunks[index] = record
}
}
return cached
}
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
}