Files

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())
}