Files
at-container-registry/docs/XRPC_BLOB_MIGRATION.md
T

26 KiB

XRPC Blob Upload Migration

This document describes how to migrate from separate legacy multipart upload endpoints to a unified com.atproto.repo.uploadBlob endpoint that supports both standard single-blob uploads and OCI container layer multipart uploads.

Current State

Legacy HTTP Endpoints (cmd/hold/main.go)

// Unified presigned URL endpoint (handles upload AND download)
mux.HandleFunc("/presigned-url", service.HandlePresignedURL)

// Internal move operation (used by multipart complete)
mux.HandleFunc("/move", service.HandleMove)

// Multipart upload endpoints
mux.HandleFunc("/start-multipart", service.HandleStartMultipart)
mux.HandleFunc("/part-presigned-url", service.HandleGetPartURL)
mux.HandleFunc("/complete-multipart", service.HandleCompleteMultipart)
mux.HandleFunc("/abort-multipart", service.HandleAbortMultipart)

// Buffered part upload (when presigned URLs unavailable)
mux.HandleFunc("/multipart-parts/", func(w http.ResponseWriter, r *http.Request) {
    // Parse URL: /multipart-parts/{uploadID}/{partNumber}
    // ...
    service.HandleMultipartPartUpload(w, r, uploadID, partNumber, did, service.MultipartMgr)
})

Existing XRPC Endpoint (pkg/hold/pds/xrpc.go)

// Current implementation - redirects to presigned URL
func (h *XRPCHandler) HandleUploadBlob(w http.ResponseWriter, r *http.Request) {
    digest := r.URL.Query().Get("digest")
    uploadURL, err := h.blobStore.GetPresignedUploadURL(digest)
    http.Redirect(w, r, uploadURL, http.StatusFound)
}

Supporting Code

pkg/hold/multipart.go:

  • MultipartManager - Tracks upload sessions
  • MultipartSession - State for each upload (parts, mode, etc.)
  • Modes: S3Native (presigned URLs), Buffered (proxy uploads)

pkg/hold/blobstore_adapter.go:

  • HoldServiceBlobStore - Adapter wrapping HoldService for XRPC handlers
  • Implements presigned URL generation
  • Currently not used by XRPC handlers

pkg/hold/handlers.go:

  • HandlePresignedURL() - Unified endpoint for GET/HEAD/PUT presigned URLs
  • HandleMove() - Moves blob from temp to final location (internal operation)
  • HandleStartMultipart() - Starts upload, returns uploadID
  • HandleGetPartURL() - Returns presigned URL for part
  • HandleCompleteMultipart() - Finalizes upload, assembles parts (calls Move internally)
  • HandleAbortMultipart() - Cancels upload
  • HandleMultipartPartUpload() - Buffered part upload fallback

Legacy Endpoint Mapping

/presigned-url → Multiple XRPC Operations

The legacy /presigned-url endpoint is a unified endpoint that handles both upload and download operations based on the operation field in the JSON body:

Legacy format:

POST /presigned-url
Content-Type: application/json

{
  "operation": "GET",    // or "HEAD" or "PUT"
  "did": "did:plc:alice123",
  "digest": "sha256:abc123...",
  "size": 1234567890     // Only for PUT operations
}

Response:
{
  "url": "https://s3.amazonaws.com/...",
  "expires_at": "2025-10-16T..."
}

XRPC mapping:

  • operation: "GET"GET /xrpc/com.atproto.sync.getBlob?did=...&cid=sha256:abc...
  • operation: "HEAD"HEAD /xrpc/com.atproto.sync.getBlob?did=...&cid=sha256:abc...
  • operation: "PUT"com.atproto.repo.uploadBlob (single upload via presigned URL)

Note: For GET/HEAD operations, AppView passes OCI digest directly as cid parameter. Hold detects sha256: prefix and uses digest directly (no CID conversion needed).

/move → Internal to Multipart Complete

The legacy /move endpoint moves a blob from temporary location to final digest-based location:

Legacy format:

POST /move?from=uploads/temp-123&to=sha256:abc123...&did=did:plc:alice123

Response: 200 OK

Purpose: Server-side S3 copy after multipart assembly. Used in this flow:

  1. Multipart parts uploaded → uploads/temp-{uploadID}/part-1, part-2, etc.
  2. Complete multipart → S3 assembles parts at uploads/temp-{uploadID}
  3. Move operation → S3 copy from uploads/temp-{uploadID}blobs/sha256/ab/abc123...

XRPC mapping:

  • Not a separate endpoint - becomes internal operation in uploadBlob?action=complete
  • The complete action automatically handles the move after multipart assembly
  • AppView doesn't need to call move explicitly in XRPC flow

New Unified Design

Single Endpoint: com.atproto.repo.uploadBlob

Content-Type discrimination determines operation:

  • application/octet-stream → Standard blob upload (profile images, small media)
  • application/json → Multipart operations (large OCI layers)

Complementary Endpoint: com.atproto.sync.getBlob

For blob downloads (maps from legacy /presigned-url with operation=GET/HEAD):

Standard ATProto blobs (CID):

GET /xrpc/com.atproto.sync.getBlob?did={holdDID}&cid=bafyreib...

Response: 307 Temporary Redirect
Location: https://s3.amazonaws.com/bucket/...?presigned-params

OCI container layers (digest):

GET /xrpc/com.atproto.sync.getBlob?did={holdDID}&cid=sha256:abc123...

Response: 307 Temporary Redirect
Location: https://s3.amazonaws.com/bucket/...?presigned-params

Implementation - Flexible CID parameter:

func (h *XRPCHandler) HandleGetBlob(w http.ResponseWriter, r *http.Request) {
    cidOrDigest := r.URL.Query().Get("cid")

    var digest string
    if strings.HasPrefix(cidOrDigest, "sha256:") {
        // OCI digest - use directly (no conversion needed)
        digest = cidOrDigest
    } else {
        // Standard CID - convert to digest
        c, _ := cid.Decode(cidOrDigest)
        digest = cidToDigest(c) // bafyreib... → sha256:abc...
    }

    // Generate presigned URL for S3
    url := h.blobStore.GetPresignedDownloadURL(digest)
    http.Redirect(w, r, url, http.StatusTemporaryRedirect)
}

Key insight: The cid parameter accepts both formats. Hold service checks prefix and handles accordingly. This keeps the endpoint spec-compliant (GET with query params) while supporting OCI digests natively.

API Specification

Standard Single Upload (ATProto Spec Compliant)

POST /xrpc/com.atproto.repo.uploadBlob
Content-Type: application/octet-stream

[raw blob bytes]

Response (200 OK):
{
  "blob": {
    "$type": "blob",
    "ref": {
      "$link": "bafyreib..."  // CID
    },
    "mimeType": "application/octet-stream",
    "size": 12345
  }
}

Use case: Profile images, small media (< 10MB), standard ATProto blobs

Multipart Start (ATCR Extension)

POST /xrpc/com.atproto.repo.uploadBlob
Content-Type: application/json

{
  "action": "start",
  "digest": "sha256:abc123...",
  "size": 1234567890  // Optional hint for storage allocation
}

Response (200 OK):
{
  "uploadId": "upload-1634567890",
  "expiresAt": "2025-10-16T12:00:00Z",
  "mode": "s3-native"  // or "buffered"
}

Implementation:

  • Calls service.StartMultipartUploadWithManager(ctx, digest, multipartMgr)
  • Returns uploadID and mode from MultipartSession

Multipart Get Part URL (ATCR Extension)

POST /xrpc/com.atproto.repo.uploadBlob
Content-Type: application/json

{
  "action": "part",
  "uploadId": "upload-1634567890",
  "partNumber": 1,
  "digest": "sha256:abc123..."
}

Response (200 OK):
{
  "url": "https://s3.amazonaws.com/bucket/...?X-Amz-...",
  "expiresAt": "2025-10-16T12:15:00Z",
  "method": "PUT"
}

// OR for buffered mode:
{
  "url": "https://hold01.atcr.io/xrpc/com.atproto.repo.uploadBlob",
  "method": "PUT",
  "headers": {
    "X-Upload-Id": "upload-1634567890",
    "X-Part-Number": "1"
  },
  "expiresAt": "2025-10-16T12:15:00Z"
}

Implementation:

  • Retrieve session: multipartMgr.GetSession(uploadID)
  • S3Native mode: Call service.GetPartUploadURL(ctx, session, partNumber, did)
  • Buffered mode: Return self-referential URL with headers

Multipart Upload Part (Buffered Mode)

PUT /xrpc/com.atproto.repo.uploadBlob
Content-Type: application/octet-stream
X-Upload-Id: upload-1634567890
X-Part-Number: 1

[part data bytes]

Response (200 OK):
{
  "etag": "abc123def456",
  "partNumber": 1
}

Implementation:

  • Extract headers: X-Upload-Id, X-Part-Number
  • Call service.HandleMultipartPartUpload(w, r, uploadID, partNumber, did, multipartMgr)
  • Return ETag for completion

Multipart Complete (ATCR Extension)

POST /xrpc/com.atproto.repo.uploadBlob
Content-Type: application/json

{
  "action": "complete",
  "uploadId": "upload-1634567890",
  "digest": "sha256:abc123...",
  "parts": [
    { "partNumber": 1, "etag": "abc123" },
    { "partNumber": 2, "etag": "def456" }
  ]
}

Response (200 OK):
{
  "status": "completed",
  "blob": {
    "$type": "blob",
    "ref": {
      "$link": "bafyreib..."  // CID computed from digest
    },
    "mimeType": "application/octet-stream",
    "size": 1234567890
  }
}

Implementation:

  • Retrieve session: multipartMgr.GetSession(uploadID)
  • For S3Native: Record parts via session.RecordS3Part()
  • Call service.CompleteMultipartUploadWithManager(ctx, session, multipartMgr)
    • This internally calls S3 CompleteMultipartUpload to assemble parts
    • Then performs server-side S3 copy from temp location to final digest location
    • Equivalent to legacy /move endpoint operation
  • Convert digest to CID for response

Multipart Abort (ATCR Extension)

POST /xrpc/com.atproto.repo.uploadBlob
Content-Type: application/json

{
  "action": "abort",
  "uploadId": "upload-1634567890",
  "digest": "sha256:abc123..."
}

Response (200 OK):
{
  "status": "aborted"
}

Implementation:

  • Retrieve session: multipartMgr.GetSession(uploadID)
  • Call service.AbortMultipartUploadWithManager(ctx, session, multipartMgr)

Implementation Strategy

Phase 1: Add Unified Handler (Keep Legacy Endpoints)

File: pkg/hold/pds/xrpc.go

// HandleUploadBlob unified handler supporting both single and multipart uploads
func (h *XRPCHandler) HandleUploadBlob(w http.ResponseWriter, r *http.Request) {
    if r.Method != http.MethodPost && r.Method != http.MethodPut {
        http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
        return
    }

    contentType := r.Header.Get("Content-Type")

    // Buffered multipart part upload (PUT with headers)
    if r.Method == http.MethodPut && r.Header.Get("X-Upload-Id") != "" {
        h.handleBufferedPartUpload(w, r)
        return
    }

    // Multipart operations (JSON body)
    if strings.Contains(contentType, "application/json") {
        h.handleMultipartOperation(w, r)
        return
    }

    // Standard single blob upload (raw bytes)
    h.handleSingleBlobUpload(w, r)
}

func (h *XRPCHandler) handleMultipartOperation(w http.ResponseWriter, r *http.Request) {
    var req struct {
        Action     string `json:"action"`
        Digest     string `json:"digest,omitempty"`
        Size       int64  `json:"size,omitempty"`
        UploadID   string `json:"uploadId,omitempty"`
        PartNumber int    `json:"partNumber,omitempty"`
        Parts      []struct {
            PartNumber int    `json:"partNumber"`
            ETag       string `json:"etag"`
        } `json:"parts,omitempty"`
    }

    if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
        http.Error(w, fmt.Sprintf("invalid JSON: %v", err), http.StatusBadRequest)
        return
    }

    // TODO: Add authentication check
    // user, err := ValidateDPoPRequest(r)

    ctx := r.Context()

    switch req.Action {
    case "start":
        h.handleMultipartStart(w, r, req.Digest, req.Size)
    case "part":
        h.handleMultipartPart(w, r, req.UploadID, req.PartNumber, req.Digest)
    case "complete":
        h.handleMultipartComplete(w, r, req.UploadID, req.Digest, req.Parts)
    case "abort":
        h.handleMultipartAbort(w, r, req.UploadID, req.Digest)
    default:
        http.Error(w, "invalid action", http.StatusBadRequest)
    }
}

func (h *XRPCHandler) handleMultipartStart(w http.ResponseWriter, r *http.Request, digest string, size int64) {
    ctx := r.Context()

    // Use HoldService multipart manager
    // Note: h.blobStore is HoldServiceBlobStore which wraps the service
    uploadID, mode, err := h.blobStore.StartMultipart(ctx, digest, size)
    if err != nil {
        http.Error(w, fmt.Sprintf("failed to start upload: %v", err), http.StatusInternalServerError)
        return
    }

    response := map[string]any{
        "uploadId":  uploadID,
        "expiresAt": time.Now().Add(24 * time.Hour),
        "mode":      mode, // "s3-native" or "buffered"
    }

    w.Header().Set("Content-Type", "application/json")
    json.NewEncoder(w).Encode(response)
}

func (h *XRPCHandler) handleMultipartPart(w http.ResponseWriter, r *http.Request, uploadID string, partNumber int, digest string) {
    ctx := r.Context()

    // Get part upload URL (presigned S3 or buffered endpoint)
    partURL, err := h.blobStore.GetPartUploadURL(ctx, uploadID, partNumber, digest)
    if err != nil {
        http.Error(w, fmt.Sprintf("failed to get part URL: %v", err), http.StatusInternalServerError)
        return
    }

    response := map[string]any{
        "url":       partURL,
        "expiresAt": time.Now().Add(15 * time.Minute),
        "method":    "PUT",
    }

    w.Header().Set("Content-Type", "application/json")
    json.NewEncoder(w).Encode(response)
}

func (h *XRPCHandler) handleMultipartComplete(w http.ResponseWriter, r *http.Request, uploadID string, digest string, parts []struct{ PartNumber int; ETag string }) {
    ctx := r.Context()

    // Convert parts format
    completedParts := make([]hold.CompletedPart, len(parts))
    for i, p := range parts {
        completedParts[i] = hold.CompletedPart{
            PartNumber: p.PartNumber,
            ETag:       p.ETag,
        }
    }

    // Complete upload
    if err := h.blobStore.CompleteMultipart(ctx, uploadID, digest, completedParts); err != nil {
        http.Error(w, fmt.Sprintf("failed to complete upload: %v", err), http.StatusInternalServerError)
        return
    }

    // Convert digest to CID for ATProto response format
    cid, err := digestToCID(digest)
    if err != nil {
        http.Error(w, fmt.Sprintf("failed to generate CID: %v", err), http.StatusInternalServerError)
        return
    }

    response := map[string]any{
        "status": "completed",
        "blob": map[string]any{
            "$type": "blob",
            "ref": map[string]any{
                "$link": cid.String(),
            },
            "mimeType": "application/octet-stream",
            // Size would need to be tracked in session
        },
    }

    w.Header().Set("Content-Type", "application/json")
    json.NewEncoder(w).Encode(response)
}

func (h *XRPCHandler) handleMultipartAbort(w http.ResponseWriter, r *http.Request, uploadID string, digest string) {
    ctx := r.Context()

    if err := h.blobStore.AbortMultipart(ctx, uploadID, digest); err != nil {
        http.Error(w, fmt.Sprintf("failed to abort upload: %v", err), http.StatusInternalServerError)
        return
    }

    response := map[string]any{
        "status": "aborted",
    }

    w.Header().Set("Content-Type", "application/json")
    json.NewEncoder(w).Encode(response)
}

func (h *XRPCHandler) handleBufferedPartUpload(w http.ResponseWriter, r *http.Request) {
    uploadID := r.Header.Get("X-Upload-Id")
    partNumberStr := r.Header.Get("X-Part-Number")

    partNumber, err := strconv.Atoi(partNumberStr)
    if err != nil {
        http.Error(w, "invalid part number", http.StatusBadRequest)
        return
    }

    // Stream part data to storage
    etag, err := h.blobStore.UploadPart(r.Context(), uploadID, partNumber, r.Body)
    if err != nil {
        http.Error(w, fmt.Sprintf("failed to upload part: %v", err), http.StatusInternalServerError)
        return
    }

    response := map[string]any{
        "etag":       etag,
        "partNumber": partNumber,
    }

    w.Header().Set("Content-Type", "application/json")
    json.NewEncoder(w).Encode(response)
}

func (h *XRPCHandler) handleSingleBlobUpload(w http.ResponseWriter, r *http.Request) {
    // Standard ATProto uploadBlob behavior
    // Read blob data
    data, err := io.ReadAll(r.Body)
    if err != nil {
        http.Error(w, "failed to read blob", http.StatusInternalServerError)
        return
    }

    // Upload to storage (single operation)
    cid, size, err := h.blobStore.UploadBlob(r.Context(), bytes.NewReader(data))
    if err != nil {
        http.Error(w, fmt.Sprintf("failed to upload blob: %v", err), http.StatusInternalServerError)
        return
    }

    // Standard ATProto blob response format
    response := map[string]any{
        "blob": map[string]any{
            "$type": "blob",
            "ref": map[string]any{
                "$link": cid.String(),
            },
            "mimeType": "application/octet-stream",
            "size":     size,
        },
    }

    w.Header().Set("Content-Type", "application/json")
    json.NewEncoder(w).Encode(response)
}

// digestToCID converts OCI digest (sha256:abc...) to ATProto CID
func digestToCID(digest string) (cid.Cid, error) {
    // Implementation in pkg/hold/cid.go or similar
    // Strip "sha256:" prefix, decode hex, construct CIDv1 with sha256 multihash
    return cid.Undef, fmt.Errorf("not implemented")
}

Phase 2: Extend HoldServiceBlobStore (pkg/hold/blobstore_adapter.go)

The HoldServiceBlobStore currently wraps HoldService for presigned URLs. Extend it to support multipart operations:

// Add multipart methods to HoldServiceBlobStore

func (h *HoldServiceBlobStore) StartMultipart(ctx context.Context, digest string, size int64) (uploadID string, mode string, err error) {
    uploadID, uploadMode, err := h.service.StartMultipartUploadWithManager(ctx, digest, h.service.MultipartMgr)
    if err != nil {
        return "", "", err
    }

    modeStr := "s3-native"
    if uploadMode == hold.Buffered {
        modeStr = "buffered"
    }

    return uploadID, modeStr, nil
}

func (h *HoldServiceBlobStore) GetPartUploadURL(ctx context.Context, uploadID string, partNumber int, digest string) (string, error) {
    session, err := h.service.MultipartMgr.GetSession(uploadID)
    if err != nil {
        return "", err
    }

    // For S3Native: return presigned URL
    // For Buffered: return self-referential URL with upload instructions
    if session.Mode == hold.S3Native {
        return h.service.GetPartUploadURL(ctx, session, partNumber, h.holdDID)
    }

    // Buffered mode: client will PUT to uploadBlob with headers
    return fmt.Sprintf("%s/xrpc/com.atproto.repo.uploadBlob", h.publicURL), nil
}

func (h *HoldServiceBlobStore) UploadPart(ctx context.Context, uploadID string, partNumber int, data io.Reader) (string, error) {
    // Buffered part upload - streams data to storage
    // Used when client PUTs to uploadBlob with X-Upload-Id header
    session, err := h.service.MultipartMgr.GetSession(uploadID)
    if err != nil {
        return "", err
    }

    // Stream to storage, return ETag
    // This wraps HandleMultipartPartUpload logic
    etag, err := h.service.UploadPartBuffered(ctx, session, partNumber, data)
    return etag, err
}

func (h *HoldServiceBlobStore) CompleteMultipart(ctx context.Context, uploadID string, digest string, parts []hold.CompletedPart) error {
    session, err := h.service.MultipartMgr.GetSession(uploadID)
    if err != nil {
        return err
    }

    // For S3Native: record parts ETags
    if session.Mode == hold.S3Native {
        for _, p := range parts {
            session.RecordS3Part(p.PartNumber, p.ETag, 0)
        }
    }

    return h.service.CompleteMultipartUploadWithManager(ctx, session, h.service.MultipartMgr)
}

func (h *HoldServiceBlobStore) AbortMultipart(ctx context.Context, uploadID string, digest string) error {
    session, err := h.service.MultipartMgr.GetSession(uploadID)
    if err != nil {
        return err
    }

    return h.service.AbortMultipartUploadWithManager(ctx, session, h.service.MultipartMgr)
}

func (h *HoldServiceBlobStore) UploadBlob(ctx context.Context, data io.Reader) (cid.Cid, int64, error) {
    // Single blob upload for standard ATProto use case
    // Compute digest, store via service driver
    // Return CID and size
    // Implementation TBD
    return cid.Undef, 0, fmt.Errorf("not implemented")
}

Phase 3: Update AppView Client (pkg/appview/storage/)

Create new XRPC client or update ProxyBlobStore to use unified endpoint:

Download (GET/HEAD):

func (p *ProxyBlobStore) ServeBlob(ctx context.Context, w http.ResponseWriter, r *http.Request, dgst digest.Digest) error {
    // Pass digest directly as cid parameter (no conversion)
    url := fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob?did=%s&cid=%s",
        p.storageEndpoint, p.holdDID, dgst.String()) // cid=sha256:abc...

    http.Redirect(w, r, url, http.StatusTemporaryRedirect)
    return nil
}

Multipart Upload:

func (p *ProxyBlobStore) startMultipartUpload(ctx context.Context, digest string) (string, error) {
    reqBody := map[string]any{
        "action": "start",
        "digest": digest,
    }

    body, _ := json.Marshal(reqBody)
    url := fmt.Sprintf("%s/xrpc/com.atproto.repo.uploadBlob", p.storageEndpoint)
    req, _ := http.NewRequestWithContext(ctx, "POST", url, bytes.NewReader(body))
    req.Header.Set("Content-Type", "application/json")

    resp, err := p.httpClient.Do(req)
    // ... parse response, return uploadID
}

func (p *ProxyBlobStore) getPartPresignedURL(ctx context.Context, digest, uploadID string, partNumber int) (string, error) {
    reqBody := map[string]any{
        "action":     "part",
        "uploadId":   uploadID,
        "partNumber": partNumber,
        "digest":     digest,
    }

    body, _ := json.Marshal(reqBody)
    url := fmt.Sprintf("%s/xrpc/com.atproto.repo.uploadBlob", p.storageEndpoint)
    req, _ := http.NewRequestWithContext(ctx, "POST", url, bytes.NewReader(body))
    req.Header.Set("Content-Type", "application/json")

    resp, err := p.httpClient.Do(req)
    // ... parse response, return presigned URL
}

// Similar for complete, abort

Phase 4: Testing Period

During transition:

  • Both legacy HTTP endpoints AND new XRPC endpoint active
  • AppView can use either based on configuration/feature flag
  • New deployments use XRPC
  • Old deployments continue with legacy

Detection logic:

func (r *RoutingRepository) Blobs(ctx context.Context) distribution.BlobStore {
    // Try XRPC first (check for /.well-known/did.json)
    if supportsXRPC(storageEndpoint) {
        return NewXRPCBlobStore(storageEndpoint, ...)
    }
    // Fallback to legacy
    return NewProxyBlobStore(storageEndpoint, ...)
}

Phase 5: Remove Legacy Endpoints

Once all holds migrated and tested:

cmd/hold/main.go - Remove:

// DELETE these lines
mux.HandleFunc("/presigned-url", service.HandlePresignedURL)
mux.HandleFunc("/move", service.HandleMove)
mux.HandleFunc("/start-multipart", service.HandleStartMultipart)
mux.HandleFunc("/part-presigned-url", service.HandleGetPartURL)
mux.HandleFunc("/complete-multipart", service.HandleCompleteMultipart)
mux.HandleFunc("/abort-multipart", service.HandleAbortMultipart)
mux.HandleFunc("/multipart-parts/", ...)

pkg/hold/handlers.go - Remove HTTP handler wrappers:

// DELETE these functions:
// - HandlePresignedURL() - replaced by uploadBlob + getBlob XRPC endpoints
// - HandleMove() - now internal operation in CompleteMultipartUploadWithManager()
// - HandleStartMultipart() - replaced by uploadBlob?action=start
// - HandleGetPartURL() - replaced by uploadBlob?action=part
// - HandleCompleteMultipart() - replaced by uploadBlob?action=complete
// - HandleAbortMultipart() - replaced by uploadBlob?action=abort
// - HandleMultipartPartUpload() - replaced by uploadBlob PUT with headers

// KEEP internal service methods:
// - s.getPresignedURL() - still used by blobstore_adapter
// - s.driver.Move() - still used for temp→final move
// - s.StartMultipartUploadWithManager() - core multipart logic
// - s.GetPartUploadURL() - presigned URL generation
// - s.CompleteMultipartUploadWithManager() - includes move operation
// - s.AbortMultipartUploadWithManager() - cleanup logic

Key Design Decisions

  1. Content-Type discrimination: Natural way to distinguish single vs multipart uploads
  2. JSON bodies for all parameters: Follows XRPC conventions (like putRecord, deleteRecord)
    • No query parameters - all operation details in request body
    • Makes requests more inspectable and debuggable
    • Easier to extend with new fields
  3. Preserve standard uploadBlob: Raw bytes still work for profile images, small media
  4. Reuse existing code: HoldService multipart logic unchanged, just new HTTP layer
  5. Backward compatibility: Both endpoints active during transition
  6. Action-based routing: Clear, extensible JSON structure
  7. Move is internal: /move endpoint logic absorbed into multipart complete operation
    • No separate XRPC endpoint needed
    • Simplifies AppView client code
  8. Unified presigned URL handling: Single uploadBlob/getBlob pair replaces operation-based routing
  9. Flexible CID parameter: getBlob accepts both standard CIDs and OCI digests via prefix detection
    • Keeps endpoint spec-compliant (GET with query params)
    • No conversion overhead on AppView side
    • Hold does simple prefix check: sha256: → use directly, else → convert CID

Benefits

  • Single endpoint for all blob operations
  • Standard ATProto uploadBlob preserved
  • XRPC-like JSON request/response
  • Reuses existing multipart.go logic
  • Gradual migration path
  • Less endpoints to maintain
  • Cleaner AppView client code

Testing Checklist

  • Single blob upload (< 10MB, raw bytes)
  • Multipart start → part → complete flow
  • S3Native mode (presigned URLs)
  • Buffered mode (proxy uploads)
  • Multipart abort
  • Large blob upload (> 5GB, many parts)
  • Concurrent uploads
  • Upload resume after network failure
  • Legacy endpoint backward compatibility
  • AppView XRPC client integration
  • Performance comparison (XRPC vs legacy)

Migration Timeline

  1. Week 1: Implement unified uploadBlob handler (Phase 1-2)
  2. Week 2: Update AppView client, feature flag (Phase 3)
  3. Week 3: Deploy to dev/staging, test both paths (Phase 4)
  4. Week 4: Roll out to production (gradual)
  5. Week 5-6: Monitor, verify all holds migrated
  6. Week 7: Remove legacy endpoints (Phase 5)

References