198 lines
6.3 KiB
Go
198 lines
6.3 KiB
Go
package agent
|
|
|
|
import (
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
|
|
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
|
)
|
|
|
|
// UploadAttempt records no signed URLs or credentials. An attempted PUT whose
|
|
// result is unknown must never be repeated by restart recovery.
|
|
type UploadAttempt struct {
|
|
RequestID string `json:"request_id"`
|
|
UploadID string `json:"upload_id"`
|
|
Identity string `json:"identity"`
|
|
State string `json:"state"`
|
|
ObjectKey string `json:"object_key"`
|
|
Result UploadResult `json:"result"`
|
|
Binding *agentv1.ExecutionBinding `json:"binding"`
|
|
Asset *agentv1.AssetDescriptor `json:"asset"`
|
|
FailureFact *agentv1.ExecutionFact `json:"failure_fact,omitempty"`
|
|
FailureDelivered bool `json:"failure_delivered,omitempty"`
|
|
}
|
|
|
|
func (s *Spool) uploadAttemptDir(id string) string { return filepath.Join(s.root, ".uploads", id) }
|
|
|
|
func (s *Spool) ClaimUpload(record UploadAttempt) error {
|
|
if err := validateName(record.UploadID); err != nil {
|
|
return err
|
|
}
|
|
if record.RequestID != "" {
|
|
if err := validateName(record.RequestID); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if record.Identity == "" || record.State != "attempted" || record.FailureFact != nil || record.FailureDelivered {
|
|
return errors.New("upload attempt must start without a failure fact")
|
|
}
|
|
root := filepath.Join(s.root, ".uploads")
|
|
if err := os.MkdirAll(root, 0700); err != nil {
|
|
return err
|
|
}
|
|
// Exclusive directory creation arbitrates across processes, not just goroutines.
|
|
dir := s.uploadAttemptDir(record.UploadID)
|
|
if err := os.Mkdir(dir, 0700); err != nil {
|
|
return err
|
|
}
|
|
if err := syncDirectory(s.root); err != nil {
|
|
return err
|
|
}
|
|
if err := syncDirectory(root); err != nil {
|
|
return err
|
|
}
|
|
if record.RequestID != "" {
|
|
requests := filepath.Join(dir, "requests")
|
|
if err := os.Mkdir(requests, 0700); err != nil {
|
|
return err
|
|
}
|
|
if err := os.Mkdir(filepath.Join(requests, record.RequestID), 0700); err != nil {
|
|
return err
|
|
}
|
|
if err := syncDirectory(requests); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return writeJSONAtomic(filepath.Join(dir, "state.json"), record)
|
|
}
|
|
|
|
// ReserveUploadRetry is used only for an explicit new request, while holding
|
|
// LockUpload. A consumed request identity is never made reusable after a crash.
|
|
func (s *Spool) ReserveUploadRetry(id, requestID string) error {
|
|
if err := validateName(requestID); err != nil {
|
|
return err
|
|
}
|
|
record, err := s.LoadUploadAttempt(id)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if record.State != "attempted" || record.FailureFact != nil || record.RequestID == "" || requestID == record.RequestID {
|
|
return errors.New("only an unsuccessful, non-terminal attempt can use an explicit new request")
|
|
}
|
|
requests := filepath.Join(s.uploadAttemptDir(id), "requests")
|
|
if err := os.Mkdir(filepath.Join(requests, requestID), 0700); err != nil {
|
|
return err
|
|
}
|
|
if err := syncDirectory(requests); err != nil {
|
|
return err
|
|
}
|
|
record.RequestID = requestID
|
|
record.Result = UploadResult{}
|
|
return writeJSONAtomic(filepath.Join(s.uploadAttemptDir(id), "state.json"), record)
|
|
}
|
|
|
|
func (s *Spool) LoadUploadAttempt(id string) (UploadAttempt, error) {
|
|
if err := validateName(id); err != nil {
|
|
return UploadAttempt{}, err
|
|
}
|
|
dir := s.uploadAttemptDir(id)
|
|
data, err := os.ReadFile(filepath.Join(dir, "state.json"))
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
if _, statErr := os.Stat(dir); statErr == nil {
|
|
return UploadAttempt{}, errors.New("upload attempt exists without durable state; PUT outcome is unknown")
|
|
}
|
|
}
|
|
if err != nil {
|
|
return UploadAttempt{}, err
|
|
}
|
|
var record UploadAttempt
|
|
if err := json.Unmarshal(data, &record); err != nil {
|
|
return record, err
|
|
}
|
|
if record.UploadID != id || record.Identity == "" {
|
|
return record, errors.New("invalid persisted upload identity")
|
|
}
|
|
switch record.State {
|
|
case "attempted", "uploaded", "completed":
|
|
default:
|
|
return record, fmt.Errorf("invalid persisted upload state %q", record.State)
|
|
}
|
|
if record.FailureDelivered && record.FailureFact == nil || record.FailureFact != nil && record.State != "attempted" {
|
|
return record, errors.New("persisted upload failure contradicts attempt state")
|
|
}
|
|
return record, nil
|
|
}
|
|
|
|
func (s *Spool) RecordUploadResult(id string, result UploadResult) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
record, err := s.LoadUploadAttempt(id)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if record.State != "attempted" || record.FailureFact != nil {
|
|
return errors.New("only an attempted upload without a terminal failure may record its PUT result")
|
|
}
|
|
if result.SizeBytes <= 0 || result.SHA256 == "" || result.StatusCode < 200 || result.StatusCode >= 300 {
|
|
return errors.New("successful upload result is required")
|
|
}
|
|
record.State, record.Result = "uploaded", result
|
|
return writeJSONAtomic(filepath.Join(s.uploadAttemptDir(id), "state.json"), record)
|
|
}
|
|
|
|
func (s *Spool) CompleteUploadNotification(id string) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
record, err := s.LoadUploadAttempt(id)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if record.State != "uploaded" && record.State != "completed" {
|
|
return errors.New("upload result must precede notification completion")
|
|
}
|
|
record.State = "completed"
|
|
return writeJSONAtomic(filepath.Join(s.uploadAttemptDir(id), "state.json"), record)
|
|
}
|
|
|
|
// PendingUploadNotifications enumerates only successful PUTs. Unknown attempts
|
|
// remain reserved and are never returned as work to retry.
|
|
func (s *Spool) PendingUploadNotifications() ([]UploadAttempt, error) {
|
|
entries, err := os.ReadDir(filepath.Join(s.root, ".uploads"))
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return nil, nil
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var pending []UploadAttempt
|
|
for _, entry := range entries {
|
|
if !entry.IsDir() {
|
|
return nil, fmt.Errorf("unexpected upload journal entry %q", entry.Name())
|
|
}
|
|
record, err := s.LoadUploadAttempt(entry.Name())
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if record.State != "uploaded" {
|
|
continue
|
|
}
|
|
if record.Binding == nil || record.Asset == nil {
|
|
return nil, fmt.Errorf("upload %q lacks notification metadata", record.UploadID)
|
|
}
|
|
pending = append(pending, record)
|
|
}
|
|
return pending, nil
|
|
}
|
|
|
|
func syncDirectory(path string) error {
|
|
dir, err := os.Open(path)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
syncErr := dir.Sync()
|
|
return firstError(syncErr, dir.Close())
|
|
}
|