diff --git a/alertmanager/alerts.go b/alertmanager/alerts.go index fcbf97797..50d883f35 100644 --- a/alertmanager/alerts.go +++ b/alertmanager/alerts.go @@ -173,7 +173,6 @@ var pdpTasks = []string{ tasknames.PDPDelDataSet, tasknames.PDPInitPP, tasknames.PDPProvingPeriod, - tasknames.PDPNotify, tasknames.PDPCommP, tasknames.PDPSaveCache, tasknames.AggregatePDPDeal, @@ -183,7 +182,6 @@ var pdpTasks = []string{ tasknames.PDPv0_SaveCache, tasknames.PDPv0_InitPP, tasknames.PDPv0_ProvPeriod, - tasknames.PDPv0_Notify, } // taskFailureCheckWith is the parameterized core shared by taskFailureCheck diff --git a/cmd/pdptool/main.go b/cmd/pdptool/main.go index da7c5780d..ad7fd7d09 100644 --- a/cmd/pdptool/main.go +++ b/cmd/pdptool/main.go @@ -12,7 +12,6 @@ import ( "encoding/pem" "fmt" "io" - "net" "net/http" "os" "strconv" @@ -375,50 +374,7 @@ var piecePrepareCmd = &cli.Command{ }, } -func startLocalNotifyServer() (string, chan struct{}, error) { - var notifyReceived chan struct{} - var server *http.Server - var ln net.Listener - - notifyReceived = make(chan struct{}) - var err error - ln, err = net.Listen("tcp", "127.0.0.1:0") - if err != nil { - return "", nil, fmt.Errorf("failed to start local HTTP server: %v", err) - } - serverAddr := fmt.Sprintf("http://%s/notify", ln.Addr().String()) - - mux := http.NewServeMux() - mux.HandleFunc("/notify", func(w http.ResponseWriter, r *http.Request) { - fmt.Println("Received notification from server.") - b, err := io.ReadAll(r.Body) - if err != nil { - fmt.Printf("Failed to read notification body: %v\n", err) - w.WriteHeader(http.StatusInternalServerError) - return - } - fmt.Printf("Notification body: %s\n", string(b)) - w.WriteHeader(http.StatusOK) - // Signal that notification was received - close(notifyReceived) - }) - - server = &http.Server{Handler: mux} - - go func() { - if err := server.Serve(ln); err != nil && err != http.ErrServerClosed { - fmt.Printf("HTTP server error: %v\n", err) - } - }() - - defer func() { - _ = server.Close() - _ = ln.Close() - }() - return serverAddr, notifyReceived, nil -} - -func uploadOnePiece(client *http.Client, serviceURL string, reqBody []byte, jwtToken string, r io.ReadSeeker, pieceSize int64, localNotifWait bool, notifyReceived chan struct{}, verbose bool) error { +func uploadOnePiece(client *http.Client, serviceURL string, reqBody []byte, jwtToken string, r io.ReadSeeker, pieceSize int64, verbose bool) error { req, err := http.NewRequest("POST", serviceURL+"/pdp/piece", bytes.NewReader(reqBody)) if err != nil { return fmt.Errorf("failed to create request: %v", err) @@ -491,13 +447,6 @@ func uploadOnePiece(client *http.Client, serviceURL string, reqBody []byte, jwtT body, _ := io.ReadAll(uploadResp.Body) return fmt.Errorf("upload failed with status code %d: %s", uploadResp.StatusCode, string(body)) } - if localNotifWait { - if verbose { - fmt.Println("Waiting for server notification...") - } - <-notifyReceived - } - return nil default: body, _ := io.ReadAll(resp.Body) @@ -523,20 +472,11 @@ var pieceUploadCmd = &cli.Command{ Name: "service-name", Usage: "Service Name to include in the JWT token (used if --jwt-token is not provided)", }, - &cli.StringFlag{ - Name: "notify-url", - Usage: "Notification URL", - Required: false, - }, &cli.StringFlag{ Name: "hash-type", Usage: "Hash type to use for verification (sha256 or commp)", Value: "sha256", }, - &cli.BoolFlag{ - Name: "local-notif-wait", - Usage: "Wait for server notification by spawning a temporary local HTTP server", - }, }, Action: func(cctx *cli.Context) error { inputFile := cctx.Args().Get(0) @@ -546,10 +486,8 @@ var pieceUploadCmd = &cli.Command{ serviceURL := cctx.String("service-url") jwtToken := cctx.String("jwt-token") - notifyURL := cctx.String("notify-url") serviceName := cctx.String("service-name") hashType := cctx.String("hash-type") - localNotifWait := cctx.Bool("local-notif-wait") if jwtToken == "" { if serviceName == "" { @@ -566,20 +504,8 @@ var pieceUploadCmd = &cli.Command{ return fmt.Errorf("invalid hash type: %s", hashType) } - if localNotifWait && notifyURL != "" { - return fmt.Errorf("cannot specify both --notify-url and --local-notif-wait") - } - - var notifyReceived chan struct{} var err error - if localNotifWait { - notifyURL, notifyReceived, err = startLocalNotifyServer() - if err != nil { - return fmt.Errorf("failed to start local HTTP server: %v", err) - } - } - // Open input file file, err := os.Open(inputFile) if err != nil { @@ -626,15 +552,12 @@ var pieceUploadCmd = &cli.Command{ return fmt.Errorf("unsupported hash type: %s", hashType) } - if notifyURL != "" { - reqData["notify"] = notifyURL - } reqBody, err = json.Marshal(reqData) if err != nil { return fmt.Errorf("failed to marshal request data: %v", err) } client := &http.Client{} - if err := uploadOnePiece(client, serviceURL, reqBody, jwtToken, file, pieceSize, localNotifWait, notifyReceived, true); err != nil { + if err := uploadOnePiece(client, serviceURL, reqBody, jwtToken, file, pieceSize, true); err != nil { return fmt.Errorf("failed to upload piece: %v", err) } @@ -662,20 +585,11 @@ var uploadFileCmd = &cli.Command{ Name: "service-name", Usage: "Service Name to include in the JWT token (used if --jwt-token is not provided)", }, - &cli.StringFlag{ - Name: "notify-url", - Usage: "Notification URL", - Required: false, - }, &cli.StringFlag{ Name: "hash-type", Usage: "Hash type to use for verification (sha256 or commp)", Value: "sha256", }, - &cli.BoolFlag{ - Name: "local-notif-wait", - Usage: "Wait for server notification by spawning a temporary local HTTP server", - }, &cli.BoolFlag{ Name: "verbose", Usage: "Verbose output", @@ -702,8 +616,6 @@ var uploadFileCmd = &cli.Command{ jwtToken := cctx.String("jwt-token") serviceName := cctx.String("service-name") hashType := cctx.String("hash-type") - localNotifWait := cctx.Bool("local-notif-wait") - notifyURL := cctx.String("notify-url") verbose := cctx.Bool("verbose") dryRun := cctx.Bool("dry-run") chunkFileName := cctx.String("chunk-file") @@ -754,15 +666,6 @@ var uploadFileCmd = &cli.Command{ bar = progressbar.NewOptions(int(fileSize/chunkSize), progressbar.OptionSetDescription("Uploading...")) } - // Setup local server if needed - var notifyReceived chan struct{} - if localNotifWait { - notifyURL, notifyReceived, err = startLocalNotifyServer() - if err != nil { - return fmt.Errorf("failed to start local HTTP server: %v", err) - } - } - // group piece aggregations for tracking as onchain pieces into sector size chunks type pieceSetInfo struct { pieces []abi.PieceInfo @@ -820,16 +723,13 @@ var uploadFileCmd = &cli.Command{ return fmt.Errorf("unsupported hash type: %s", hashType) } - if notifyURL != "" { - reqData["notify"] = notifyURL - } reqBody, err = json.Marshal(reqData) if err != nil { return fmt.Errorf("failed to marshal request data: %v", err) } // Upload the piece - err = uploadOnePiece(client, serviceURL, reqBody, jwtToken, chunkReader, int64(n), localNotifWait, notifyReceived, verbose) + err = uploadOnePiece(client, serviceURL, reqBody, jwtToken, chunkReader, int64(n), verbose) if err != nil { return fmt.Errorf("failed to upload piece: %v", err) } @@ -1530,20 +1430,11 @@ var streamingPieceUploadCmd = &cli.Command{ Name: "service-name", Usage: "Service Name to include in the JWT token (used if --jwt-token is not provided)", }, - &cli.StringFlag{ - Name: "notify-url", - Usage: "Notification URL", - Required: false, - }, &cli.StringFlag{ Name: "hash-type", Usage: "Hash type to use for verification (sha256 or commp)", Value: "commp", }, - &cli.BoolFlag{ - Name: "local-notif-wait", - Usage: "Wait for server notification by spawning a temporary local HTTP server", - }, }, Action: func(cctx *cli.Context) error { inputFile := cctx.Args().Get(0) @@ -1553,10 +1444,8 @@ var streamingPieceUploadCmd = &cli.Command{ serviceURL := cctx.String("service-url") jwtToken := cctx.String("jwt-token") - notifyURL := cctx.String("notify-url") serviceName := cctx.String("service-name") hashType := cctx.String("hash-type") - localNotifWait := cctx.Bool("local-notif-wait") if jwtToken == "" { if serviceName == "" { @@ -1573,20 +1462,8 @@ var streamingPieceUploadCmd = &cli.Command{ return fmt.Errorf("invalid hash type: %s", hashType) } - if localNotifWait && notifyURL != "" { - return fmt.Errorf("cannot specify both --notify-url and --local-notif-wait") - } - - var notifyReceived chan struct{} var err error - if localNotifWait { - notifyURL, notifyReceived, err = startLocalNotifyServer() - if err != nil { - return fmt.Errorf("failed to start local HTTP server: %v", err) - } - } - // Open the input file file, err := os.Open(inputFile) if err != nil { @@ -1715,17 +1592,12 @@ var streamingPieceUploadCmd = &cli.Command{ type finalize struct { PieceCID string `json:"pieceCid"` - Notify string `json:"notify,omitempty"` } bd := finalize{ PieceCID: pcid2.String(), } - if notifyURL != "" { - bd.Notify = notifyURL - } - bodyBytes, err := json.Marshal(bd) if err != nil { return fmt.Errorf("failed to marshal finalize request body: %v", err) @@ -1758,10 +1630,6 @@ var streamingPieceUploadCmd = &cli.Command{ fmt.Printf("Piece CID: %s\n", pcid2.String()) fmt.Println("Piece uploaded successfully.") - if localNotifWait { - fmt.Println("Waiting for server notification...") - <-notifyReceived - } return nil }, } diff --git a/cuhttp/server.go b/cuhttp/server.go index 1a7bcb197..c98431807 100644 --- a/cuhttp/server.go +++ b/cuhttp/server.go @@ -14,6 +14,7 @@ import ( "github.com/filecoin-project/curio/cuhttp/servicedeps" "github.com/filecoin-project/curio/deps" + "github.com/filecoin-project/curio/lib/piecestore" mhttp "github.com/filecoin-project/curio/market/http" "github.com/filecoin-project/curio/market/libp2p" "github.com/filecoin-project/curio/pdp" @@ -97,12 +98,12 @@ func attachRouters(ctx context.Context, r *chi.Mux, d *deps.Deps, sd *ServiceDep if sd.EthSender != nil { if err := pdp.MountRoutes(ctx, r, pdp.MountDeps{ - DB: d.DB, - LocalStore: d.LocalStore, - EthClient: must.One(d.EthClient.Get()), - Chain: d.Chain, - EthSender: sd.EthSender, - AlertTask: sd.AlertTask, + DB: d.DB, + PieceIO: piecestore.New(d.Stor, d.LocalStore, d.Si), + EthClient: must.One(d.EthClient.Get()), + Chain: d.Chain, + EthSender: sd.EthSender, + AlertTask: sd.AlertTask, }, ipp); err != nil { return nil, err } diff --git a/documentation/en/experimental-features/PDPCURIOSPEC.md b/documentation/en/experimental-features/PDPCURIOSPEC.md index 0b08426a9..6ea677ada 100644 --- a/documentation/en/experimental-features/PDPCURIOSPEC.md +++ b/documentation/en/experimental-features/PDPCURIOSPEC.md @@ -20,7 +20,6 @@ TODO: the pdp datasets table etc ## Managing Pieces TODO: add, upload, delete, pull, info TODO: the pdp pieceref table etc -TODO: "finalization" notify task # Storage @@ -143,7 +142,7 @@ Harmony tasks are created through three trigger mechanisms: - **Chain handlers** — Registered via `chainsched.AddHandler`, these callbacks fire on every chain head change. They inspect on-chain state (e.g. transaction receipts, epoch thresholds) and call `AddTask` to insert work into the harmony_task queue when conditions are met. The proving-cycle tasks (InitPP, ProvPeriod, Prove) use this mechanism so they respond immediately to new tipsets. - **IAmBored** — An optional callback in `TaskTypeDetails` that the task engine invokes when a machine has spare capacity and no queued work exists for that task type. A `passcall.Every(duration, ...)` wrapper rate-limits invocations. Tasks like TerminateFWSS (1 min), DeleteDataSet (1 hour), and Settle (12 hours) use IAmBored because they generate work opportunistically rather than in response to chain events. -- **Polling** — Some tasks use a dedicated poller goroutine that periodically queries the database for pending work and calls `AddTask`. The Notify (2s) and PullPiece (10s) tasks use this pattern because their triggers are purely database-driven (new uploads or pull requests) with no chain dependency. +- **Polling** — Some tasks use a dedicated poller goroutine that periodically queries the database for pending work and calls `AddTask`. PullPiece uses this pattern because its trigger is a database-backed pull request rather than a chain event. All three mechanisms funnel through `harmonytask.AddTask()`, which atomically inserts a task record. The main poller loop (every 3s) then discovers unowned tasks and assigns them to machines with available resources. @@ -154,7 +153,6 @@ All three mechanisms funnel through `harmonytask.AddTask()`, which atomically in | `PDPv0_InitPP` | `InitProvingPeriodTask` | `tasks/pdpv0/task_init_pp.go` | Chain handler | | `PDPv0_ProvPeriod` | `NextProvingPeriodTask` | `tasks/pdpv0/task_next_pp.go` | Chain handler | | `PDPv0_Prove` | `ProveTask` | `tasks/pdpv0/task_prove.go` | Chain handler | -| `PDPv0_Notify` | `PDPNotifyTask` | `tasks/pdpv0/notify_task.go` | Polling (2s) | | `PDPv0_PullPiece` | `PDPPullPieceTask` | `tasks/pdpv0/task_pull_piece.go` | Polling (10s) | | `PDPv0_Indexing` | `PDPIndexingTask` | `tasks/indexing/task_pdp_v0_indexing.go` | IAmBored (3s) | | `PDPv0_IPNI` | `PDPIPNITask` | `tasks/indexing/task_pdp_v0_ipni.go` | IAmBored (30s) | @@ -180,11 +178,17 @@ The sections below describe how tasks and watchers connect to form the PDP lifec ## Piece Ingestion -There are two paths for getting pieces into the system: +There are three paths for getting pieces into the system: -**Direct upload path:** -1. Client uploads piece data via HTTP. The data is written to `parked_pieces`. -2. When the upload completes (`parked_pieces.complete = TRUE`), **PDPv0_Notify** picks it up, sends an HTTP callback to `notify_url` if configured, and moves the reference from `pdp_piece_uploads` to `pdp_piecerefs`. +**Known-CID direct upload path:** +1. The client posts the PieceCID. If a complete long-term copy already exists, the handler creates its `parked_piece_refs` and `pdp_piecerefs` rows immediately. +2. Otherwise, the client PUTs the bytes to the returned upload URL. The handler claims a `parked_pieces` row, writes the request body once directly to final piece storage while computing CommP, and validates the declared size and PieceCID. +3. After the write succeeds, one database transaction marks the parked piece complete, creates `pdp_piecerefs`, and deletes the upload row. The legacy `notify` request field is accepted for compatibility but this path does not call the URL or schedule a notify task. + +**Streaming upload path:** +1. The client creates an upload session and PUTs bytes without declaring a PieceCID. The handler claims a provisional `parked_pieces` identity for that session, then streams the request once directly into final piece storage while calculating CommP and the raw size. +2. After the PieceCID is known, one transaction promotes the provisional row to the calculated identity. If a matching complete piece already exists, the session is pointed at that piece and the unreferenced provisional copy is left for normal parked-piece cleanup. +3. The client finalizes with the calculated PieceCID. The handler validates it, creates `pdp_piecerefs`, and deletes the streaming session in one transaction. It does not use scratch space, create an intermediate `pdp_piece_uploads` row, or schedule a notify task. **Pull path:** 1. Client submits a pull request via HTTP, creating a row in `pdp_piece_pull_items`. diff --git a/harmony/harmonydb/downgrade/20260820-pdp-streaming-timestamptz.sql b/harmony/harmonydb/downgrade/20260820-pdp-streaming-timestamptz.sql new file mode 100644 index 000000000..e4f32bf88 --- /dev/null +++ b/harmony/harmonydb/downgrade/20260820-pdp-streaming-timestamptz.sql @@ -0,0 +1,4 @@ +-- The canonical pre-migration schema already used TIMESTAMPTZ for both +-- columns, so retain their types and restore only the previous default. +ALTER TABLE pdp_piece_streaming_uploads + ALTER COLUMN created_at SET DEFAULT TIMEZONE('UTC', NOW()); diff --git a/harmony/harmonydb/sql/20260820-pdp-streaming-timestamptz.sql b/harmony/harmonydb/sql/20260820-pdp-streaming-timestamptz.sql new file mode 100644 index 000000000..4261dcd12 --- /dev/null +++ b/harmony/harmonydb/sql/20260820-pdp-streaming-timestamptz.sql @@ -0,0 +1,40 @@ +/* + Streaming-upload timestamps represent absolute instants. Convert legacy + naive columns, if present, by interpreting their stored values as UTC. + Existing TIMESTAMPTZ columns already store absolute instants and are left + unchanged. +*/ + +DO $$ +BEGIN + IF EXISTS ( + SELECT 1 + FROM information_schema.columns + WHERE table_schema = 'public' + AND table_name = 'pdp_piece_streaming_uploads' + AND column_name = 'created_at' + AND data_type = 'timestamp without time zone' + ) THEN + ALTER TABLE pdp_piece_streaming_uploads + ALTER COLUMN created_at DROP DEFAULT; + ALTER TABLE pdp_piece_streaming_uploads + ALTER COLUMN created_at TYPE TIMESTAMPTZ + USING created_at AT TIME ZONE 'UTC'; + END IF; + + IF EXISTS ( + SELECT 1 + FROM information_schema.columns + WHERE table_schema = 'public' + AND table_name = 'pdp_piece_streaming_uploads' + AND column_name = 'completed_at' + AND data_type = 'timestamp without time zone' + ) THEN + ALTER TABLE pdp_piece_streaming_uploads + ALTER COLUMN completed_at TYPE TIMESTAMPTZ + USING completed_at AT TIME ZONE 'UTC'; + END IF; +END $$; + +ALTER TABLE pdp_piece_streaming_uploads + ALTER COLUMN created_at SET DEFAULT NOW(); diff --git a/lib/piecestore/io.go b/lib/piecestore/io.go index 8a349c2a2..3dfc0b44c 100644 --- a/lib/piecestore/io.go +++ b/lib/piecestore/io.go @@ -123,39 +123,19 @@ func (s *Store) WriteUploadPiece(ctx context.Context, pieceID storiface.PieceNum }() copyStart := time.Now() - - wr := new(commp.Calc) - defer wr.Reset() - writers := io.MultiWriter(wr, destFile) - - n, err := io.CopyBuffer(writers, io.LimitReader(data, size), make([]byte, 8<<20)) + pieceInfo, rawSize, err := writeUploadPieceData(destFile, size, data, verifySize) if err != nil { _ = destFile.Close() - return abi.PieceInfo{}, 0, xerrors.Errorf("copying piece data: %w", err) + return abi.PieceInfo{}, 0, err } if err := destFile.Close(); err != nil { return abi.PieceInfo{}, 0, xerrors.Errorf("closing temp piece file: %w", err) } - if verifySize && n != size { - return abi.PieceInfo{}, 0, xerrors.Errorf("short write: %d", n) - } - - digest, pieceSize, err := wr.Digest() - if err != nil { - return abi.PieceInfo{}, 0, xerrors.Errorf("computing piece digest: %w", err) - } - - pcid, err := commcid.DataCommitmentV1ToCID(digest) - if err != nil { - return abi.PieceInfo{}, 0, xerrors.Errorf("computing piece CID: %w", err) - } - psize := abi.PaddedPieceSize(pieceSize) - copyEnd := time.Now() - log.Infow("wrote piece", "piece", pieceID, "size", n, "duration", copyEnd.Sub(copyStart), "dest", dest, "MiB/s", float64(size)/(1<<20)/copyEnd.Sub(copyStart).Seconds()) + log.Infow("wrote piece", "piece", pieceID, "size", rawSize, "duration", copyEnd.Sub(copyStart), "dest", dest, "MiB/s", float64(rawSize)/(1<<20)/copyEnd.Sub(copyStart).Seconds()) if err := os.Rename(tempDest, dest); err != nil { return abi.PieceInfo{}, 0, xerrors.Errorf("rename temp piece to dest %s -> %s: %w", tempDest, dest, err) @@ -168,5 +148,40 @@ func (s *Store) WriteUploadPiece(ctx context.Context, pieceID storiface.PieceNum return abi.PieceInfo{}, 0, xerrors.Errorf("ensure one copy: %w", err) } - return abi.PieceInfo{PieceCID: pcid, Size: psize}, uint64(n), nil + return pieceInfo, rawSize, nil +} + +func writeUploadPieceData(destFile *os.File, size int64, data io.Reader, verifySize bool) (abi.PieceInfo, uint64, error) { + copyLimit := size + if !verifySize { + // Read one byte beyond the maximum so an oversized stream cannot be + // silently accepted as a truncated piece. + copyLimit++ + } + + wr := new(commp.Calc) + defer wr.Reset() + writers := io.MultiWriter(wr, destFile) + + n, err := io.CopyBuffer(writers, io.LimitReader(data, copyLimit), make([]byte, 8<<20)) + if err != nil { + return abi.PieceInfo{}, 0, xerrors.Errorf("copying piece data: %w", err) + } + if verifySize && n != size { + return abi.PieceInfo{}, 0, xerrors.Errorf("short write: %d", n) + } + if !verifySize && n > size { + return abi.PieceInfo{}, 0, xerrors.Errorf("%w: limit %d bytes", ErrPieceTooLarge, size) + } + + digest, pieceSize, err := wr.Digest() + if err != nil { + return abi.PieceInfo{}, 0, xerrors.Errorf("computing piece digest: %w", err) + } + pcid, err := commcid.DataCommitmentV1ToCID(digest) + if err != nil { + return abi.PieceInfo{}, 0, xerrors.Errorf("computing piece CID: %w", err) + } + + return abi.PieceInfo{PieceCID: pcid, Size: abi.PaddedPieceSize(pieceSize)}, uint64(n), nil } diff --git a/lib/piecestore/io_test.go b/lib/piecestore/io_test.go new file mode 100644 index 000000000..34e3e7d53 --- /dev/null +++ b/lib/piecestore/io_test.go @@ -0,0 +1,52 @@ +package piecestore + +import ( + "bytes" + "os" + "testing" + + "github.com/stretchr/testify/require" + + commcid "github.com/filecoin-project/go-fil-commcid" + commp "github.com/filecoin-project/go-fil-commp-hashhash" +) + +func TestWriteUploadPieceDataUnknownSize(t *testing.T) { + body := bytes.Repeat([]byte{0x5a}, 1024) + const maxSize = int64(4096) + + f, err := os.CreateTemp(t.TempDir(), "piece-") + require.NoError(t, err) + + pieceInfo, rawSize, err := writeUploadPieceData(f, maxSize, bytes.NewReader(body), false) + require.NoError(t, err) + require.NoError(t, f.Close()) + require.Equal(t, uint64(len(body)), rawSize) + + stored, err := os.ReadFile(f.Name()) + require.NoError(t, err) + require.Equal(t, body, stored) + + calc := &commp.Calc{} + t.Cleanup(calc.Reset) + _, err = calc.Write(body) + require.NoError(t, err) + digest, paddedSize, err := calc.Digest() + require.NoError(t, err) + expectedCID, err := commcid.DataCommitmentV1ToCID(digest) + require.NoError(t, err) + require.True(t, expectedCID.Equals(pieceInfo.PieceCID)) + require.Equal(t, paddedSize, uint64(pieceInfo.Size)) +} + +func TestWriteUploadPieceDataRejectsOversizedStream(t *testing.T) { + const maxSize = int64(1024) + body := bytes.Repeat([]byte{0x6b}, int(maxSize)+1) + + f, err := os.CreateTemp(t.TempDir(), "piece-") + require.NoError(t, err) + t.Cleanup(func() { _ = f.Close() }) + + _, _, err = writeUploadPieceData(f, maxSize, bytes.NewReader(body), false) + require.ErrorIs(t, err, ErrPieceTooLarge) +} diff --git a/lib/piecestore/pieceio.go b/lib/piecestore/pieceio.go index 5dbe5d087..06a5c212c 100644 --- a/lib/piecestore/pieceio.go +++ b/lib/piecestore/pieceio.go @@ -2,6 +2,7 @@ package piecestore import ( "context" + "errors" "io" "github.com/filecoin-project/go-state-types/abi" @@ -10,9 +11,15 @@ import ( "github.com/filecoin-project/curio/lib/storiface" ) +// ErrPieceTooLarge reports that an unknown-size upload exceeded its hard limit. +var ErrPieceTooLarge = errors.New("piece data exceeds the maximum size") + // PieceIO provides FFI-free piece storage operations. type PieceIO interface { WritePiece(ctx context.Context, taskID *harmonytask.TaskID, pieceID storiface.PieceNumber, size int64, data io.Reader, storageType storiface.PathType) error + // WriteUploadPiece writes directly to the requested storage class while + // calculating CommP. When verifySize is true, short input is rejected; + // otherwise size is a hard maximum. WriteUploadPiece(ctx context.Context, pieceID storiface.PieceNumber, size int64, data io.Reader, storageType storiface.PathType, verifySize bool) (abi.PieceInfo, uint64, error) PieceReader(ctx context.Context, id storiface.PieceNumber) (io.ReadCloser, error) RemovePiece(ctx context.Context, id storiface.PieceNumber) error diff --git a/market/mk20/mk20_utils.go b/market/mk20/mk20_utils.go index 640aa8d53..5c1b1697a 100644 --- a/market/mk20/mk20_utils.go +++ b/market/mk20/mk20_utils.go @@ -315,7 +315,7 @@ func NewTimeoutLimitReader(r io.Reader, timeout time.Duration) *TimeoutLimitRead } } -const UploadSizeLimit = int64(1 * 1024 * 1024 * 1024) +const UploadSizeLimit = int64(((1 * 1024 * 1024 * 1024) * 127) / 128) func (t *TimeoutLimitReader) Read(p []byte) (int, error) { deadline := time.Now().Add(t.timeout) diff --git a/pdp/README.md b/pdp/README.md index 492782d7b..5b4743774 100644 --- a/pdp/README.md +++ b/pdp/README.md @@ -57,14 +57,12 @@ All endpoints are rooted at `/pdp`. ```json { - "pieceCid": "", - "notify": "" + "pieceCid": "" } ``` - **Fields:** - `pieceCid`: The Piece CID in CommP **v2** format (CIDv1 with `fil-commitment-unsealed` codec and raw size encoded). This uniquely identifies the piece and encodes size information. - - `notify`: *(Optional)* A URL to be notified when the piece has been processed successfully. #### Responses @@ -197,7 +195,7 @@ The streaming upload API provides a way to upload large pieces in a streaming fa 2. Stream the data via `PUT`. 3. Finalize the upload with the pieceCid to link and validate. -> **Note:** Each streaming upload chunk is limited to **1 GiB** (unpadded). The server computes the CommP on-the-fly. +> **Note:** Each streaming upload is limited to **1,065,353,216 raw bytes** (1 GiB padded). The server writes it once directly to piece storage while computing CommP on-the-fly. #### 3.1. Create Streaming Upload Session @@ -236,6 +234,7 @@ The streaming upload API provides a way to upload large pieces in a streaming fa - `400 Bad Request`: Invalid UUID. - `401 Unauthorized`: Missing or invalid JWT token. - `404 Not Found`: Upload session not found. +- `409 Conflict`: The upload session or calculated piece is already being written. - `413 Payload Too Large`: Data exceeds size limit. --- @@ -251,8 +250,7 @@ The streaming upload API provides a way to upload large pieces in a streaming fa ```json { - "pieceCid": "", - "notify": "" + "pieceCid": "" } ``` @@ -265,50 +263,13 @@ The streaming upload API provides a way to upload large pieces in a streaming fa - `400 Bad Request`: Invalid pieceCid, size mismatch, or CID does not match the uploaded data. - `401 Unauthorized`: Missing or invalid JWT token. - `404 Not Found`: Upload session not found. +- `409 Conflict`: The PUT has not completed final storage yet. --- -### 4. Notifications +### 4. Upload Completion -When you initiate an upload with the `notify` field specified, the PDP Service will send a notification to the provided URL once the piece has been successfully processed and stored. - -#### 4.1. Notification Request - -- **Method:** `POST` -- **URL:** The `notify` URL provided during the upload initiation (`POST /pdp/piece`). -- **Headers:** - - `Content-Type`: `application/json` -- **Request Body:** - -```json -{ - "id": "", - "service": "", - "pieceCID": "", - "notify_url": "", - "check_hash_codec": "", - "check_hash": "" -} -``` - -- **Fields:** - - `id`: The upload ID. - - `service`: The service name. - - `pieceCID`: The Piece CID of the stored piece (may be `null` if not applicable). - - `notify_url`: The original notification URL provided. - - `check_hash_codec`: The hash function used (e.g., `"sha2-256-trunc254-padded"`). - - `check_hash`: The byte array of the original hash provided in the upload initiation. - -#### 4.2. Expected Response from Your Server - -- **Status Code:** `200 OK` to acknowledge receipt. -- **Response Body:** (Optional) Can be empty or contain a message. - -#### 4.3. Notes - -- The PDP Service may retry the notification if it fails. -- Ensure that your server is accessible from the PDP Service and can handle incoming POST requests. -- The notification does not include the piece data; it confirms that the piece has been successfully stored. +Uploads complete synchronously. A `204 No Content` response from the known-CID PUT means the piece is stored and registered for PDP. For streaming uploads, the PUT stores the piece and a `200 OK` response from finalize registers it for PDP. The service does not send upload callback requests. --- @@ -1008,8 +969,7 @@ Error responses typically include an error message in the response body. Content-Type: application/json { - "pieceCid": "", - "notify": "https://example.com/notify" + "pieceCid": "" } ``` @@ -1052,31 +1012,6 @@ Error responses typically include an error message in the response body. HTTP/1.1 204 No Content ``` -3. **Receive Notification (if `notify` was provided):** - - **Server's Notification Request:** - - ```http - POST /notify HTTP/1.1 - Host: example.com - Content-Type: application/json - - { - "id": "", - "service": "", - "pieceCID": "", - "notify_url": "https://example.com/notify", - "check_hash_codec": "sha2-256-trunc254-padded", - "check_hash": "" - } - ``` - - **Your Response:** - - ```http - HTTP/1.1 200 OK - ``` - ### Uploading a Piece (Streaming) 1. **Create Upload Session:** diff --git a/pdp/handlers.go b/pdp/handlers.go index 7f5767501..c0aee6de2 100644 --- a/pdp/handlers.go +++ b/pdp/handlers.go @@ -25,7 +25,7 @@ import ( "github.com/filecoin-project/curio/api" "github.com/filecoin-project/curio/harmony/harmonydb" "github.com/filecoin-project/curio/lib/ethchain" - "github.com/filecoin-project/curio/lib/paths" + "github.com/filecoin-project/curio/lib/piecestore" ipni_provider "github.com/filecoin-project/curio/market/ipni/ipni-provider" "github.com/filecoin-project/curio/pdp/contract" "github.com/filecoin-project/curio/tasks/indexing" @@ -74,7 +74,7 @@ type ETHTxSender interface { type PDPService struct { Auth db *harmonydb.DB - storage paths.StashStore + pieceIO piecestore.PieceIO sender ETHTxSender ethClient ethchain.EthClient @@ -97,7 +97,7 @@ type PDPServiceNodeApi interface { func NewPDPService( ctx context.Context, db *harmonydb.DB, - stor paths.StashStore, + pieceIO piecestore.PieceIO, ec ethchain.EthClient, fc PDPServiceNodeApi, sn ETHTxSender, @@ -110,7 +110,7 @@ func NewPDPService( p := &PDPService{ Auth: auth, db: db, - storage: stor, + pieceIO: pieceIO, sender: sn, ethClient: ec, @@ -1312,21 +1312,32 @@ func (p *PDPService) handleGetDataSetPiece(w http.ResponseWriter, r *http.Reques func (p *PDPService) cleanup(ctx context.Context) { rm := func(ctx context.Context, db *harmonydb.DB) { - var RefIDs []int64 - - err := db.QueryRow(ctx, `SELECT COALESCE(array_agg(piece_ref), '{}') AS ref_ids - FROM pdp_piece_streaming_uploads - WHERE complete = TRUE - AND completed_at <= TIMEZONE('UTC', NOW()) - INTERVAL '60 minutes';`).Scan(&RefIDs) - if err != nil { - log.Errorw("failed to get non-finalized uploads", "error", err) + if err := p.cleanupExpiredDirectUploadClaims(ctx); err != nil { + log.Errorw("failed to clean up expired direct upload claims", "error", err) + } + if err := p.cleanupExpiredStreamingUploadClaims(ctx); err != nil { + log.Errorw("failed to clean up expired streaming upload claims", "error", err) } - if len(RefIDs) > 0 { - _, err := db.Exec(ctx, `DELETE FROM parked_piece_refs WHERE ref_id = ANY($1);`, RefIDs) - if err != nil { - log.Errorw("failed to delete non-finalized uploads", "error", err) - } + _, err := db.Exec(ctx, ` + WITH expired AS ( + DELETE FROM pdp_piece_streaming_uploads su + WHERE su.complete = TRUE + AND su.completed_at <= NOW() - INTERVAL '60 minutes' + AND NOT EXISTS ( + SELECT 1 FROM pdp_piecerefs pr WHERE pr.piece_ref = su.piece_ref + ) + RETURNING su.piece_ref + ) + DELETE FROM parked_piece_refs ppr + USING expired + WHERE ppr.ref_id = expired.piece_ref + AND NOT EXISTS ( + SELECT 1 FROM pdp_piecerefs pr WHERE pr.piece_ref = ppr.ref_id + ) + `) + if err != nil { + log.Errorw("failed to delete non-finalized uploads", "error", err) } // Clean up old piece pull records (older than 5 days). Pull items only diff --git a/pdp/handlers_upload.go b/pdp/handlers_upload.go index fdfd41905..6029c9439 100644 --- a/pdp/handlers_upload.go +++ b/pdp/handlers_upload.go @@ -2,14 +2,13 @@ package pdp import ( "bytes" + "context" "database/sql" - "encoding/hex" "encoding/json" "errors" "fmt" "io" "net/http" - "os" "path" "time" @@ -20,15 +19,17 @@ import ( "github.com/multiformats/go-multicodec" "github.com/multiformats/go-multihash" "github.com/yugabyte/pgx/v5" + "github.com/yugabyte/pgx/v5/pgconn" commcid "github.com/filecoin-project/go-fil-commcid" commp "github.com/filecoin-project/go-fil-commp-hashhash" "github.com/filecoin-project/go-state-types/abi" "github.com/filecoin-project/curio/harmony/harmonydb" - "github.com/filecoin-project/curio/lib/dealdata" "github.com/filecoin-project/curio/lib/parkpiece" + "github.com/filecoin-project/curio/lib/piecestore" "github.com/filecoin-project/curio/lib/proof" + "github.com/filecoin-project/curio/lib/storiface" ) var log = logging.Logger("pdpv0") @@ -43,8 +44,566 @@ var ( ErrPieceTooSmall = fmt.Errorf("piece data is below the minimum allowed size (%d bytes)", PieceSizeMinLimit) ErrPieceTooLarge = fmt.Errorf("piece data exceeds the maximum allowed size (%d bytes)", PieceSizeMaxLimit) ErrExceedsDeclaredSize = fmt.Errorf("piece data exceeds the declared piece size") + errPieceTooShort = errors.New("piece data is shorter than the declared piece size") + errUploadInProgress = errors.New("another upload is already writing this piece") + errUploadClaimed = errors.New("upload UUID has already been claimed") ) +const minPaddedPieceSizeForCache = int64(32 * 1024 * 1024) + +type exactSizeReader struct { + r io.Reader + remaining int64 +} + +func (r *exactSizeReader) Read(p []byte) (int, error) { + if r.remaining == 0 { + return 0, io.EOF + } + if int64(len(p)) > r.remaining { + p = p[:r.remaining] + } + + n, err := r.r.Read(p) + r.remaining -= int64(n) + if errors.Is(err, io.EOF) && r.remaining > 0 { + return n, errPieceTooShort + } + return n, err +} + +func readHasExtraByte(r io.Reader) (bool, error) { + var buf [1]byte + for { + n, err := r.Read(buf[:]) + if n > 0 { + return true, nil + } + if err != nil { + // TimeoutLimitReader reports an over-limit byte as an error rather + // than returning it. It still proves that the body exceeded the + // exact size declared by PieceCIDv2. + if errors.Is(err, ErrPieceTooLarge) { + return true, nil + } + if errors.Is(err, io.EOF) { + return false, nil + } + return false, err + } + } +} + +func needsSaveCache(rawSize int64) bool { + return PadPieceSize(rawSize) >= minPaddedPieceSizeForCache +} + +func insertPDPReference(tx *harmonydb.Tx, service, pieceCID string, pieceRef, rawSize int64) error { + n, err := tx.Exec(` + INSERT INTO pdp_piecerefs (service, piece_cid, piece_ref, created_at, needs_save_cache) + VALUES ($1, $2, $3, NOW(), $4) + `, service, pieceCID, pieceRef, needsSaveCache(rawSize)) + if err != nil { + return fmt.Errorf("failed to insert pdp_piecerefs: %w", err) + } + if n != 1 { + return fmt.Errorf("failed to insert pdp_piecerefs: expected 1 row, got %d", n) + } + return nil +} + +func deleteClaimedUpload(tx *harmonydb.Tx, uploadID string, pieceRef int64) error { + n, err := tx.Exec(`DELETE FROM pdp_piece_uploads WHERE id = $1 AND piece_ref = $2`, uploadID, pieceRef) + if err != nil { + return fmt.Errorf("failed to delete pdp_piece_uploads row: %w", err) + } + if n != 1 { + return fmt.Errorf("failed to delete pdp_piece_uploads row: expected 1 row, got %d", n) + } + return nil +} + +type parkedPieceClaim struct { + parkedPieceID int64 + pieceRefID int64 + created bool + complete bool +} + +// claimParkedPiece returns a per-upload ref to the active long-term parked +// piece. A newly inserted skip=true row is owned by the caller and may be +// written directly. Existing incomplete rows remain owned by their current +// writer; existing complete rows can be reused without writing any bytes. +func claimParkedPiece(tx *harmonydb.Tx, pieceCID string, rawSize, paddedSize int64) (parkedPieceClaim, error) { + var claim parkedPieceClaim + var err error + claim.parkedPieceID, claim.created, err = parkpiece.UpsertSkipWithInserted(tx, pieceCID, paddedSize, rawSize, true, true) + if err != nil { + return parkedPieceClaim{}, fmt.Errorf("failed to claim parked piece: %w", err) + } + + if !claim.created { + err = tx.QueryRow(`SELECT complete FROM parked_pieces WHERE id = $1`, claim.parkedPieceID).Scan(&claim.complete) + if err != nil { + return parkedPieceClaim{}, fmt.Errorf("failed to inspect existing parked piece: %w", err) + } + if !claim.complete { + return parkedPieceClaim{}, errUploadInProgress + } + } + + err = tx.QueryRow(` + INSERT INTO parked_piece_refs (piece_id, long_term) + VALUES ($1, TRUE) + RETURNING ref_id + `, claim.parkedPieceID).Scan(&claim.pieceRefID) + if err != nil { + return parkedPieceClaim{}, fmt.Errorf("failed to create parked piece ref: %w", err) + } + + return claim, nil +} + +// claimDirectUpload atomically claims an unclaimed upload intent. For a new +// piece, it binds the intent to a skip=true parked-piece ref so this handler +// owns the direct write. If the piece became complete after POST, it publishes +// the PDP ref and consumes the intent without reading the request body. +func (p *PDPService) claimDirectUpload(ctx context.Context, uploadID, service, pieceCID string, rawSize, paddedSize int64) (parkedPieceClaim, error) { + var claim parkedPieceClaim + committed, err := p.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { + var err error + claim, err = claimParkedPiece(tx, pieceCID, rawSize, paddedSize) + if err != nil { + return false, err + } + + if claim.complete { + n, err := tx.Exec(`DELETE FROM pdp_piece_uploads WHERE id = $1 AND piece_ref IS NULL`, uploadID) + if err != nil { + return false, fmt.Errorf("failed to consume upload UUID: %w", err) + } + if n != 1 { + return false, errUploadClaimed + } + if err := insertPDPReference(tx, service, pieceCID, claim.pieceRefID, rawSize); err != nil { + return false, err + } + return true, nil + } + + n, err := tx.Exec(` + UPDATE pdp_piece_uploads + SET piece_ref = $1, piece_cid = $2, created_at = NOW() + WHERE id = $3 AND piece_ref IS NULL + `, claim.pieceRefID, pieceCID, uploadID) + if err != nil { + return false, fmt.Errorf("failed to claim upload UUID: %w", err) + } + if n != 1 { + return false, errUploadClaimed + } + + return true, nil + }, harmonydb.OptionRetry()) + if err != nil { + return parkedPieceClaim{}, err + } + if !committed { + return parkedPieceClaim{}, errors.New("failed to commit direct upload claim") + } + return claim, nil +} + +// claimStreamingUpload binds the session to an upload-owned provisional piece +// before any request bytes are read. The provisional identity is replaced by +// the computed PieceCID after the one-pass final-storage write. +func (p *PDPService) claimStreamingUpload(ctx context.Context, uploadID, service string) (parkedPieceClaim, error) { + calc := &commp.Calc{} + defer calc.Reset() + temporaryRawSize, err := io.WriteString(calc, uploadID+":"+uuid.NewString()) + if err != nil { + return parkedPieceClaim{}, fmt.Errorf("failed to generate provisional piece identity: %w", err) + } + temporaryDigest, temporaryPaddedSize, err := calc.Digest() + if err != nil { + return parkedPieceClaim{}, fmt.Errorf("failed to generate provisional piece identity: %w", err) + } + temporaryCID, err := commcid.DataCommitmentV1ToCID(temporaryDigest) + if err != nil { + return parkedPieceClaim{}, fmt.Errorf("failed to generate provisional piece CID: %w", err) + } + + var claim parkedPieceClaim + committed, err := p.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { + var err error + claim, err = claimParkedPiece(tx, temporaryCID.String(), int64(temporaryRawSize), int64(temporaryPaddedSize)) + if err != nil { + return false, err + } + if !claim.created { + return false, errUploadInProgress + } + + n, err := tx.Exec(` + UPDATE pdp_piece_streaming_uploads + SET piece_ref = $1, + created_at = NOW() + WHERE id = $2 + AND service = $3 + AND piece_ref IS NULL + AND COALESCE(complete, FALSE) = FALSE + `, claim.pieceRefID, uploadID, service) + if err != nil { + return false, fmt.Errorf("failed to claim streaming upload UUID: %w", err) + } + if n != 1 { + return false, errUploadClaimed + } + + return true, nil + }, harmonydb.OptionRetry()) + if err != nil { + return parkedPieceClaim{}, err + } + if !committed { + return parkedPieceClaim{}, errors.New("failed to commit streaming upload claim") + } + return claim, nil +} + +func (p *PDPService) completeStreamingUpload(ctx context.Context, uploadID, service string, claim parkedPieceClaim, pieceInfo abi.PieceInfo, rawSize int64) error { + committed, err := p.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { + var existingPieceID int64 + var existingRawSize int64 + var existingComplete bool + err := tx.QueryRow(` + SELECT id, piece_raw_size, complete + FROM parked_pieces + WHERE piece_cid = $1 + AND piece_padded_size = $2 + AND long_term = TRUE + AND cleanup_task_id IS NULL + AND id != $3 + ORDER BY id + LIMIT 1 + `, pieceInfo.PieceCID.String(), int64(pieceInfo.Size), claim.parkedPieceID).Scan(&existingPieceID, &existingRawSize, &existingComplete) + switch { + case err == nil && (!existingComplete || existingRawSize != rawSize): + return false, errUploadInProgress + case err == nil: + n, err := tx.Exec(` + UPDATE parked_piece_refs + SET piece_id = $1 + WHERE ref_id = $2 AND piece_id = $3 + `, existingPieceID, claim.pieceRefID, claim.parkedPieceID) + if err != nil { + return false, fmt.Errorf("failed to reuse completed parked piece: %w", err) + } + if n != 1 { + return false, fmt.Errorf("failed to reuse completed parked piece: expected 1 row, got %d", n) + } + case errors.Is(err, pgx.ErrNoRows): + n, err := tx.Exec(` + UPDATE parked_pieces + SET piece_cid = $1, + piece_padded_size = $2, + piece_raw_size = $3, + complete = TRUE + WHERE id = $4 + AND complete = FALSE + AND skip = TRUE + AND cleanup_task_id IS NULL + `, pieceInfo.PieceCID.String(), int64(pieceInfo.Size), rawSize, claim.parkedPieceID) + if err != nil { + var pgErr *pgconn.PgError + if errors.As(err, &pgErr) && pgErr.Code == "23505" { + return false, errUploadInProgress + } + return false, fmt.Errorf("failed to promote provisional parked piece: %w", err) + } + if n != 1 { + return false, fmt.Errorf("failed to promote provisional parked piece: expected 1 row, got %d", n) + } + default: + return false, fmt.Errorf("failed to inspect completed parked piece: %w", err) + } + + n, err := tx.Exec(` + UPDATE pdp_piece_streaming_uploads + SET piece_cid = $1, + piece_size = $2, + raw_size = $3, + complete = TRUE, + completed_at = NOW() + WHERE id = $4 + AND service = $5 + AND piece_ref = $6 + AND COALESCE(complete, FALSE) = FALSE + `, pieceInfo.PieceCID.String(), int64(pieceInfo.Size), rawSize, uploadID, service, claim.pieceRefID) + if err != nil { + return false, fmt.Errorf("failed to mark streaming upload complete: %w", err) + } + if n != 1 { + return false, fmt.Errorf("failed to mark streaming upload complete: expected 1 row, got %d", n) + } + + return true, nil + }, harmonydb.OptionRetry()) + if err != nil { + return err + } + if !committed { + return errors.New("failed to commit streaming upload completion") + } + return nil +} + +func (p *PDPService) finalizeDirectUpload(ctx context.Context, uploadID, service, pieceCID string, claim parkedPieceClaim, rawSize int64) error { + committed, err := p.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { + n, err := tx.Exec(` + UPDATE parked_pieces + SET complete = TRUE + WHERE id = $1 AND complete = FALSE AND skip = TRUE AND cleanup_task_id IS NULL + `, claim.parkedPieceID) + if err != nil { + return false, fmt.Errorf("failed to mark parked piece complete: %w", err) + } + if n != 1 { + return false, fmt.Errorf("failed to mark parked piece complete: expected 1 row, got %d", n) + } + + if err := insertPDPReference(tx, service, pieceCID, claim.pieceRefID, rawSize); err != nil { + return false, err + } + if err := deleteClaimedUpload(tx, uploadID, claim.pieceRefID); err != nil { + return false, err + } + return true, nil + }, harmonydb.OptionRetry()) + if err != nil { + return err + } + if !committed { + return errors.New("failed to commit direct upload completion") + } + return nil +} + +// releaseDirectUploadClaim makes a failed or expired upload retryable. It only +// removes the parked row when no other subsystem attached a ref to it. +func (p *PDPService) releaseDirectUploadClaim(ctx context.Context, uploadID string, pieceRefID, parkedPieceID int64) error { + removePiece := false + committed, err := p.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { + removePiece = false + n, err := tx.Exec(` + UPDATE pdp_piece_uploads + SET piece_ref = NULL, created_at = NOW() + WHERE id = $1 AND piece_ref = $2 + `, uploadID, pieceRefID) + if err != nil { + return false, fmt.Errorf("failed to release upload claim: %w", err) + } + if n == 0 { + var finalized bool + err = tx.QueryRow(`SELECT EXISTS(SELECT 1 FROM pdp_piecerefs WHERE piece_ref = $1)`, pieceRefID).Scan(&finalized) + if err != nil { + return false, fmt.Errorf("failed to check direct upload finalization: %w", err) + } + if finalized { + return true, nil + } + } + if n > 1 { + return false, fmt.Errorf("failed to release upload claim: expected 1 row, got %d", n) + } + + _, err = tx.Exec(`DELETE FROM parked_piece_refs WHERE ref_id = $1 AND piece_id = $2`, pieceRefID, parkedPieceID) + if err != nil { + return false, fmt.Errorf("failed to delete direct upload ref: %w", err) + } + + n, err = tx.Exec(` + DELETE FROM parked_pieces pp + WHERE pp.id = $1 + AND pp.complete = FALSE + AND pp.skip = TRUE + AND NOT EXISTS ( + SELECT 1 FROM parked_piece_refs ppr WHERE ppr.piece_id = pp.id + ) + `, parkedPieceID) + if err != nil { + return false, fmt.Errorf("failed to delete abandoned parked piece: %w", err) + } + removePiece = n == 1 + + // Pull or another subsystem may have attached a usable source while the + // direct upload was running. Let StorePiece take over in that case. + if !removePiece { + _, err = tx.Exec(` + UPDATE parked_pieces pp + SET skip = FALSE + WHERE pp.id = $1 + AND pp.complete = FALSE + AND pp.skip = TRUE + AND EXISTS ( + SELECT 1 FROM parked_piece_refs ppr + WHERE ppr.piece_id = pp.id AND ppr.data_url IS NOT NULL + ) + `, parkedPieceID) + if err != nil { + return false, fmt.Errorf("failed to release parked piece to StorePiece: %w", err) + } + } + return true, nil + }, harmonydb.OptionRetry()) + if err != nil { + return err + } + if !committed { + return errors.New("failed to commit direct upload claim release") + } + if removePiece { + if err := p.pieceIO.RemovePiece(ctx, storiface.PieceNumber(parkedPieceID)); err != nil { + return fmt.Errorf("failed to remove abandoned piece %d: %w", parkedPieceID, err) + } + } + return nil +} + +// releaseStreamingUploadClaim resets a failed streaming session and drops its +// ref. The zero-ref provisional parked piece is left for normal piece cleanup. +func (p *PDPService) releaseStreamingUploadClaim(ctx context.Context, uploadID, service string, claim parkedPieceClaim) error { + committed, err := p.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { + n, err := tx.Exec(` + UPDATE pdp_piece_streaming_uploads + SET piece_ref = NULL, + piece_cid = NULL, + piece_size = NULL, + raw_size = NULL, + complete = NULL, + completed_at = NULL + WHERE id = $1 + AND service = $2 + AND piece_ref = $3 + AND COALESCE(complete, FALSE) = FALSE + `, uploadID, service, claim.pieceRefID) + if err != nil { + return false, fmt.Errorf("failed to release streaming upload claim: %w", err) + } + if n == 0 { + var retained bool + err = tx.QueryRow(` + SELECT EXISTS( + SELECT 1 + FROM pdp_piece_streaming_uploads + WHERE id = $1 AND service = $2 AND piece_ref = $3 AND complete = TRUE + UNION ALL + SELECT 1 FROM pdp_piecerefs WHERE piece_ref = $3 + ) + `, uploadID, service, claim.pieceRefID).Scan(&retained) + if err != nil { + return false, fmt.Errorf("failed to check streaming upload completion: %w", err) + } + if retained { + return true, nil + } + } + if n > 1 { + return false, fmt.Errorf("failed to release streaming upload claim: expected 1 row, got %d", n) + } + + _, err = tx.Exec(`DELETE FROM parked_piece_refs WHERE ref_id = $1 AND piece_id = $2`, claim.pieceRefID, claim.parkedPieceID) + if err != nil { + return false, fmt.Errorf("failed to delete streaming upload ref: %w", err) + } + + return true, nil + }, harmonydb.OptionRetry()) + if err != nil { + return err + } + if !committed { + return errors.New("failed to commit streaming upload claim release") + } + return nil +} + +func (p *PDPService) cleanupExpiredDirectUploadClaims(ctx context.Context) error { + var claims []struct { + UploadID string `db:"upload_id"` + PieceRefID int64 `db:"piece_ref_id"` + ParkedPieceID int64 `db:"parked_piece_id"` + } + err := p.db.Select(ctx, &claims, ` + SELECT pu.id::TEXT AS upload_id, + pu.piece_ref AS piece_ref_id, + pp.id AS parked_piece_id + FROM pdp_piece_uploads pu + JOIN parked_piece_refs ppr ON ppr.ref_id = pu.piece_ref + JOIN parked_pieces pp ON pp.id = ppr.piece_id + WHERE pu.piece_ref IS NOT NULL + AND pu.created_at <= NOW() - INTERVAL '1 hour' + AND ppr.data_url IS NULL + AND pp.complete = FALSE + AND pp.skip = TRUE + ORDER BY pu.created_at, pu.id + LIMIT 256 + `) + if err != nil { + return fmt.Errorf("select expired direct upload claims: %w", err) + } + + var cleanupErr error + for _, claim := range claims { + if err := p.releaseDirectUploadClaim(ctx, claim.UploadID, claim.PieceRefID, claim.ParkedPieceID); err != nil { + cleanupErr = errors.Join(cleanupErr, fmt.Errorf("release expired upload %s: %w", claim.UploadID, err)) + } + } + return cleanupErr +} + +func (p *PDPService) cleanupExpiredStreamingUploadClaims(ctx context.Context) error { + var claims []struct { + UploadID string `db:"upload_id"` + Service string `db:"service"` + PieceRefID int64 `db:"piece_ref_id"` + ParkedPieceID int64 `db:"parked_piece_id"` + } + err := p.db.Select(ctx, &claims, ` + SELECT su.id::TEXT AS upload_id, + su.service, + su.piece_ref AS piece_ref_id, + pp.id AS parked_piece_id + FROM pdp_piece_streaming_uploads su + JOIN parked_piece_refs ppr ON ppr.ref_id = su.piece_ref + JOIN parked_pieces pp ON pp.id = ppr.piece_id + WHERE su.piece_ref IS NOT NULL + AND COALESCE(su.complete, FALSE) = FALSE + AND su.created_at <= NOW() - INTERVAL '1 hour' + AND ppr.data_url IS NULL + AND pp.complete = FALSE + AND pp.skip = TRUE + ORDER BY su.created_at, su.id + LIMIT 256 + `) + if err != nil { + return fmt.Errorf("select expired streaming upload claims: %w", err) + } + + var cleanupErr error + for _, stale := range claims { + claim := parkedPieceClaim{ + parkedPieceID: stale.ParkedPieceID, + pieceRefID: stale.PieceRefID, + created: true, + } + if err := p.releaseStreamingUploadClaim(ctx, stale.UploadID, stale.Service, claim); err != nil { + cleanupErr = errors.Join(cleanupErr, fmt.Errorf("release expired streaming upload %s: %w", stale.UploadID, err)) + } + } + return cleanupErr +} + func (p *PDPService) handlePiecePost(w http.ResponseWriter, r *http.Request) { // Verify that the request is authorized using ECDSA JWT serviceID, err := p.AuthService(r) @@ -116,16 +675,10 @@ func (p *PDPService) handlePiecePost(w http.ResponseWriter, r *http.Request) { } log.Debugw("[handlePiecePost] -- new parked piece ref", "parkedPieceRefID", parkedPieceRefID, "pieceCidV1", pieceCidV1) - // Create a new 'pdp_piece_uploads' entry pointing to the 'parked_piece_refs' entry - uploadUUID = uuid.New() - _, err = tx.Exec(` - INSERT INTO pdp_piece_uploads (id, service, piece_cid, notify_url, piece_ref, check_hash_codec, check_hash, check_size) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8) - `, uploadUUID.String(), serviceID, pieceCidV1.String(), req.Notify, parkedPieceRefID, multicodec.Sha2_256Trunc254Padded.String(), dmh.Digest, size) - if err != nil { - return false, fmt.Errorf("failed to insert into pdp_piece_uploads: %w", err) + if err := insertPDPReference(tx, serviceID, pieceCidV1.String(), parkedPieceRefID, int64(size)); err != nil { + return false, err } - log.Debugw("[handlePiecePost] -- new pdp_piece_uploads", "uploadUUID", uploadUUID, "pieceCidV1", pieceCidV1) + log.Debugw("[handlePiecePost] -- new pdp_piecerefs", "parkedPieceRefID", parkedPieceRefID, "pieceCidV1", pieceCidV1) responseStatus = http.StatusOK return true, nil // Commit the transaction @@ -189,15 +742,16 @@ func (p *PDPService) handlePieceUpload(w http.ResponseWriter, r *http.Request) { log.Debugw("[handlePieceUpload] -- upload started", "uploadUUID", uploadUUID) ctx := r.Context() - // Lookup the expected pieceCID, notify_url, and piece_ref from the database using uploadUUID - var pieceCIDStr *string - var notifyURL string + // Lookup the expected piece and current claim from the database. + var serviceID string + var pieceCIDStr string var checkSize int64 - var pieceRef sql.NullInt64 err = p.db.QueryRow(ctx, ` - SELECT piece_cid, notify_url, piece_ref, check_size FROM pdp_piece_uploads WHERE id = $1 - `, uploadUUID.String()).Scan(&pieceCIDStr, ¬ifyURL, &pieceRef, &checkSize) + SELECT service, piece_cid, piece_ref, check_size + FROM pdp_piece_uploads + WHERE id = $1 + `, uploadUUID.String()).Scan(&serviceID, &pieceCIDStr, &pieceRef, &checkSize) if err != nil { if errors.Is(err, pgx.ErrNoRows) { httpServerError(w, http.StatusNotFound, "Upload UUID not found", err) @@ -207,156 +761,100 @@ func (p *PDPService) handlePieceUpload(w http.ResponseWriter, r *http.Request) { return } log.Debugw("[handlePieceUpload] -- upload lookup done", "uploadUUID", uploadUUID) - // Check that piece_ref is null; non-null means data was already uploaded + // A non-null ref is an active direct-write claim (or a completed legacy upload). if pieceRef.Valid { httpServerError(w, http.StatusConflict, "Data has already been uploaded", err) return } - pieceCidV1, err := cid.Parse(*pieceCIDStr) + pieceCidV1, err := cid.Parse(pieceCIDStr) if err != nil { httpServerError(w, http.StatusInternalServerError, "Failed to convert piece CID (v1): "+err.Error(), err) return } - dmh, err := multihash.Decode(pieceCidV1.Hash()) + paddedSize := PadPieceSize(checkSize) + claim, err := p.claimDirectUpload(ctx, uploadUUID.String(), serviceID, pieceCidV1.String(), checkSize, paddedSize) if err != nil { - httpServerError(w, http.StatusInternalServerError, "Failed to decode piece CID: "+err.Error(), err) + switch { + case errors.Is(err, errUploadInProgress): + httpServerError(w, http.StatusConflict, "This piece is already being uploaded", err) + case errors.Is(err, errUploadClaimed): + httpServerError(w, http.StatusConflict, "Data has already been uploaded", err) + default: + httpServerError(w, http.StatusInternalServerError, "Failed to claim piece upload", err) + } + return + } + if claim.complete { + w.WriteHeader(http.StatusNoContent) return } - declaredPieceSize := checkSize - - // Create a commp.Calc instance for calculating commP - cp := &commp.Calc{} - defer cp.Reset() - readSize := int64(0) - - // Function to write data into StashStore and calculate commP - writeFunc := func(f *os.File) error { - limitedReader := io.LimitReader(r.Body, declaredPieceSize+1) // +1 to detect exceeding the limit - multiWriter := io.MultiWriter(cp, f) - - // Copy data from limitedReader to multiWriter - n, err := io.Copy(multiWriter, limitedReader) - if err != nil { - return fmt.Errorf("failed to read and write piece data: %w", err) + cleanupClaim := true + defer func() { + if !cleanupClaim { + return } - - if n > declaredPieceSize { - return fmt.Errorf("read %d bytes, declared %d: %w", n, declaredPieceSize, ErrExceedsDeclaredSize) + cleanupCtx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + if err := p.releaseDirectUploadClaim(cleanupCtx, uploadUUID.String(), claim.pieceRefID, claim.parkedPieceID); err != nil { + log.Errorw("failed to release direct upload claim", "uploadUUID", uploadUUID, "piece", claim.parkedPieceID, "error", err) } - - readSize = n - - return nil - } - // Upload into StashStore - stashID, err := p.storage.StashCreate(ctx, declaredPieceSize, writeFunc) + }() + + bodyReader := NewTimeoutLimitReader(r.Body, 5*time.Second) + exactReader := &exactSizeReader{r: bodyReader, remaining: checkSize} + pieceInfo, readSize, err := p.pieceIO.WriteUploadPiece( + ctx, + storiface.PieceNumber(claim.parkedPieceID), + checkSize, + exactReader, + storiface.PathStorage, + true, + ) if err != nil { - if errors.Is(err, ErrExceedsDeclaredSize) { - msg := fmt.Sprintf("piece data exceeds the size declared in pieceCid (%d bytes)", declaredPieceSize) - httpServerError(w, http.StatusBadRequest, msg, err) - return + if errors.Is(err, errPieceTooShort) { + httpServerError(w, http.StatusBadRequest, "Piece size does not match the expected size", err) } else { - log.Errorw("Failed to store piece data in StashStore", "error", err) + log.Errorw("failed to write uploaded piece directly to storage", "uploadUUID", uploadUUID, "error", err) httpServerError(w, http.StatusInternalServerError, "Failed to store piece data", err) - return } + return } - log.Debugw("[handlePieceUpload] -- uploaded into StashStore", "uploadUUID", uploadUUID) - - // Finalize the commP calculation - digest, paddedPieceSize, err := cp.Digest() - if err != nil { - // Remove the stash file as the data is invalid - _ = p.storage.StashRemove(ctx, stashID) - httpServerError(w, http.StatusInternalServerError, "Failed to finalize commP calculation", err) + if readSize != uint64(checkSize) { + httpServerError(w, http.StatusBadRequest, "Piece size does not match the expected size", nil) return } - if readSize != checkSize { - log.Debugw("[handlePieceUpload] -- piece size does not match the expected size removing from stash store", "uploadUUID", uploadUUID) - _ = p.storage.StashRemove(ctx, stashID) - httpServerError(w, http.StatusBadRequest, "Piece size does not match the expected size", err) + hasExtra, err := readHasExtraByte(bodyReader) + if err != nil { + httpServerError(w, http.StatusInternalServerError, "Failed to verify uploaded piece size", err) return } - - outHash := digest - - if !bytes.Equal(outHash, dmh.Digest) { - log.Debugw("[handlePieceUpload] -- computed hash does not match expected hash removing from stash store", "uploadUUID", uploadUUID) - // Remove the stash file as the data is invalid - _ = p.storage.StashRemove(ctx, stashID) - log.Warnw("Computed hash does not match expected hash", "computed", hex.EncodeToString(outHash), "expected", hex.EncodeToString(dmh.Digest), "pieceCid", pieceCidV1.String()) - httpServerError(w, http.StatusBadRequest, "Computed hash does not match expected hash", err) + if hasExtra { + msg := fmt.Sprintf("piece data exceeds the size declared in pieceCid (%d bytes)", checkSize) + httpServerError(w, http.StatusBadRequest, msg, ErrExceedsDeclaredSize) return } - - // Convert commP digest into a piece CID - pieceCIDComputed, err := commcid.DataCommitmentV1ToCID(digest) - if err != nil { - // Remove the stash file as the data is invalid - _ = p.storage.StashRemove(ctx, stashID) - httpServerError(w, http.StatusInternalServerError, "Failed to convert commP to CID", err) + if !pieceInfo.PieceCID.Equals(pieceCidV1) { + log.Warnw("computed piece CID does not match expected piece CID", "computed", pieceInfo.PieceCID, "expected", pieceCidV1, "uploadUUID", uploadUUID) + httpServerError(w, http.StatusBadRequest, "Computed piece CID does not match expected piece CID", nil) return } - - // Compare the computed piece CID with the expected one from the database - if pieceCIDStr != nil && pieceCIDComputed.String() != *pieceCIDStr { - log.Debugw("[handlePieceUpload] -- computed piece CID does not match expected piece CID removing from stash store", "uploadUUID", uploadUUID) - // Remove the stash file as the data is invalid - _ = p.storage.StashRemove(ctx, stashID) - httpServerError(w, http.StatusBadRequest, "Computed piece CID does not match expected piece CID", err) + if pieceInfo.Size != abi.PaddedPieceSize(paddedSize) { + httpServerError(w, http.StatusBadRequest, "Computed padded piece size does not match expected piece size", nil) return } - didCommit, err := p.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { - // 1. Create a long-term parked piece entry - parkedPieceID, err := parkpiece.Upsert(tx, pieceCIDComputed.String(), int64(paddedPieceSize), readSize, true) - if err != nil { - return false, fmt.Errorf("failed to create parked_pieces entry: %w", err) - } - log.Debugw("[handlePieceUpload] -- parked pieces entry created", "uploadUUID", uploadUUID) - // 2. Create a piece ref with data_url being "stashstore://" - // Get StashURL - stashURL, err := p.storage.StashURL(stashID) - if err != nil { - return false, fmt.Errorf("failed to get stash URL: %w", err) - } - - // Change scheme to "custore" - stashURL.Scheme = dealdata.CustoreScheme - dataURL := stashURL.String() - - var pieceRefID int64 - err = tx.QueryRow(` - INSERT INTO parked_piece_refs (piece_id, data_url, long_term) - VALUES ($1, $2, TRUE) RETURNING ref_id - `, parkedPieceID, dataURL).Scan(&pieceRefID) - if err != nil { - return false, fmt.Errorf("failed to create parked_piece_refs entry: %w", err) - } - log.Debugw("[handlePieceUpload] -- parked piece ref created", "uploadUUID", uploadUUID) - // 3. Update the pdp_piece_uploads entry to contain the created piece_ref - _, err = tx.Exec(` - UPDATE pdp_piece_uploads SET piece_ref = $1, piece_cid = $2 WHERE id = $3 - `, pieceRefID, pieceCIDComputed.String(), uploadUUID.String()) - if err != nil { - return false, fmt.Errorf("failed to update pdp_piece_uploads: %w", err) - } - log.Debugw("[handlePieceUpload] -- pdp_piece_uploads entry updated", "uploadUUID", uploadUUID) - return true, nil // Commit the transaction - }, harmonydb.OptionRetry()) - - if err != nil || !didCommit { - // Remove the stash file as the transaction failed - _ = p.storage.StashRemove(ctx, stashID) - httpServerError(w, http.StatusInternalServerError, "Failed to process piece upload", err) + finalizeCtx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + if err := p.finalizeDirectUpload(finalizeCtx, uploadUUID.String(), serviceID, pieceCidV1.String(), claim, checkSize); err != nil { + httpServerError(w, http.StatusInternalServerError, "Failed to finalize piece upload", err) return } + cleanupClaim = false log.Debugw("[handlePieceUpload] -- piece upload done, writing response", "uploadUUID", uploadUUID) - // Respond with 204 No Content w.WriteHeader(http.StatusNoContent) } @@ -475,116 +973,73 @@ func (p *PDPService) handleStreamingUpload(w http.ResponseWriter, r *http.Reques return } - reader := NewTimeoutLimitReader(r.Body, 5*time.Second) - cp := &commp.Calc{} - defer cp.Reset() - readSize := int64(0) - - // Function to write data into StashStore and calculate commP - writeFunc := func(f *os.File) error { - multiWriter := io.MultiWriter(cp, f) - - // Copy data from limitedReader to multiWriter - n, err := io.Copy(multiWriter, reader) - if err != nil { - return fmt.Errorf("failed to read and write piece data: %w", err) - } - - // already limited the maximum read size in TimeoutLimitReader - if n < int64(PieceSizeMinLimit) { - return ErrPieceTooSmall - } - - readSize = n - - return nil - } - - // Upload into StashStore - stashID, err := p.storage.StashCreate(ctx, int64(PieceSizeMaxLimit), writeFunc) - if err != nil { - if errors.Is(err, ErrPieceTooLarge) { - httpServerError(w, http.StatusRequestEntityTooLarge, ErrPieceTooLarge.Error(), err) - return - } else if errors.Is(err, ErrPieceTooSmall) { - httpServerError(w, http.StatusBadRequest, ErrPieceTooSmall.Error(), err) - return - } else { - log.Errorw("Failed to store piece data in StashStore", "error", err) - httpServerError(w, http.StatusInternalServerError, "Failed to store piece data", err) + bodyReader := NewTimeoutLimitReader(r.Body, 5*time.Second) + prefix := make([]byte, int(PieceSizeMinLimit)) + if _, err := io.ReadFull(bodyReader, prefix); err != nil { + if errors.Is(err, io.EOF) || errors.Is(err, io.ErrUnexpectedEOF) { + httpServerError(w, http.StatusBadRequest, ErrPieceTooSmall.Error(), ErrPieceTooSmall) return } - } - - // Finalize the commP calculation - digest, paddedPieceSize, err := cp.Digest() - if err != nil { - log.Errorw("Failed to finalize commP calculation", "error", err) - // Remove the stash file as the data is invalid - _ = p.storage.StashRemove(ctx, stashID) - httpServerError(w, http.StatusInternalServerError, "Failed to finalize commP calculation", err) + httpServerError(w, http.StatusInternalServerError, "Failed to read piece data", err) return } - pcid, err := commcid.DataCommitmentV1ToCID(digest) + claim, err := p.claimStreamingUpload(ctx, uploadUUID.String(), serviceID) if err != nil { - log.Errorw("Failed to calculate PieceCIDV2", "error", err) - _ = p.storage.StashRemove(ctx, stashID) - httpServerError(w, http.StatusInternalServerError, "Failed to calculate PieceCIDV2", err) + switch { + case errors.Is(err, errUploadInProgress): + httpServerError(w, http.StatusConflict, "This piece is already being uploaded", err) + case errors.Is(err, errUploadClaimed): + httpServerError(w, http.StatusConflict, "Data has already been uploaded", err) + default: + log.Errorw("Failed to claim streaming upload", "uploadUUID", uploadUUID, "error", err) + httpServerError(w, http.StatusInternalServerError, "Failed to claim streaming upload", err) + } return } - didCommit, err := p.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { - // 1. Create a long-term parked piece entry - parkedPieceID, err := parkpiece.Upsert(tx, pcid.String(), int64(paddedPieceSize), readSize, true) - if err != nil { - return false, fmt.Errorf("failed to create parked_pieces entry: %w", err) - } - - // 2. Create a piece ref with data_url being "stashstore://" - // Get StashURL - stashURL, err := p.storage.StashURL(stashID) - if err != nil { - return false, fmt.Errorf("failed to get stash URL: %w", err) + cleanupClaim := true + defer func() { + if !cleanupClaim { + return } - - // Change scheme to "custore" - stashURL.Scheme = dealdata.CustoreScheme - dataURL := stashURL.String() - - var pieceRefID int64 - err = tx.QueryRow(` - INSERT INTO parked_piece_refs (piece_id, data_url, long_term) - VALUES ($1, $2, TRUE) RETURNING ref_id - `, parkedPieceID, dataURL).Scan(&pieceRefID) - if err != nil { - return false, fmt.Errorf("failed to create parked_piece_refs entry: %w", err) + cleanupCtx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + if err := p.releaseStreamingUploadClaim(cleanupCtx, uploadUUID.String(), serviceID, claim); err != nil { + log.Errorw("failed to release streaming upload claim", "uploadUUID", uploadUUID, "piece", claim.parkedPieceID, "error", err) } - - // 3. Update the pdp_piece_streaming_uploads entry - _, err = tx.Exec(` - UPDATE pdp_piece_streaming_uploads SET piece_ref = $1, piece_cid = $2, piece_size = $3, raw_size = $4, complete = TRUE, completed_at = NOW() AT TIME ZONE 'UTC' WHERE id = $5 and service = $6 - `, pieceRefID, pcid.String(), paddedPieceSize, readSize, uploadUUID.String(), serviceID) - if err != nil { - return false, fmt.Errorf("failed to update pdp_piece_streaming_uploads: %w", err) + }() + + pieceInfo, readSize, err := p.pieceIO.WriteUploadPiece( + ctx, + storiface.PieceNumber(claim.parkedPieceID), + int64(PieceSizeMaxLimit), + io.MultiReader(bytes.NewReader(prefix), bodyReader), + storiface.PathStorage, + false, + ) + if err != nil { + if errors.Is(err, ErrPieceTooLarge) || errors.Is(err, piecestore.ErrPieceTooLarge) { + httpServerError(w, http.StatusRequestEntityTooLarge, ErrPieceTooLarge.Error(), err) + return } + log.Errorw("Failed to write streaming upload directly to storage", "uploadUUID", uploadUUID, "error", err) + httpServerError(w, http.StatusInternalServerError, "Failed to store piece data", err) + return + } - return true, nil // Commit the transaction - }, harmonydb.OptionRetry()) - - if err != nil || !didCommit { - // Remove the stash file as the transaction failed - if err != nil { - log.Errorw("Failed to process piece upload", "error", err) - } else { - log.Errorw("Failed to process piece upload", "error", "failed to commit transaction") + completeCtx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + if err := p.completeStreamingUpload(completeCtx, uploadUUID.String(), serviceID, claim, pieceInfo, int64(readSize)); err != nil { + if errors.Is(err, errUploadInProgress) { + httpServerError(w, http.StatusConflict, "This piece is already being uploaded or conflicts with an existing piece", err) + return } - _ = p.storage.StashRemove(ctx, stashID) - httpServerError(w, http.StatusInternalServerError, "Failed to process piece upload", err) + httpServerError(w, http.StatusInternalServerError, "Failed to complete streaming upload", err) return } + cleanupClaim = false - // Respond with 204 No Content w.WriteHeader(http.StatusNoContent) } @@ -640,20 +1095,26 @@ func (p *PDPService) handleFinalizeStreamingUpload(w http.ResponseWriter, r *htt } pieceCidV1 := pieceInfo.CidV1 - // Get digest for insertion - digest, err := commcid.CIDToDataCommitmentV1(pieceCidV1) - if err != nil { - httpServerError(w, http.StatusBadRequest, "Invalid request body: invalid pieceCid", err) - return - } - // Query database for stored piece info var dPcidStr string var pref int64 var rawSize uint64 - err = p.db.QueryRow(ctx, `SELECT piece_cid, piece_ref, raw_size FROM pdp_piece_streaming_uploads WHERE id = $1 AND service = $2 AND complete = TRUE`, uploadUUID.String(), serviceID).Scan(&dPcidStr, &pref, &rawSize) + err = p.db.QueryRow(ctx, ` + SELECT su.piece_cid, su.piece_ref, su.raw_size + FROM pdp_piece_streaming_uploads su + JOIN parked_piece_refs ppr ON ppr.ref_id = su.piece_ref + JOIN parked_pieces pp ON pp.id = ppr.piece_id + WHERE su.id = $1 + AND su.service = $2 + AND su.complete = TRUE + AND pp.complete = TRUE + `, uploadUUID.String(), serviceID).Scan(&dPcidStr, &pref, &rawSize) if err != nil { + if errors.Is(err, pgx.ErrNoRows) { + httpServerError(w, http.StatusConflict, "Streaming upload is not complete", err) + return + } log.Errorw("Failed to query pdp_piece_streaming_uploads", "error", err) httpServerError(w, http.StatusInternalServerError, "Database error", err) return @@ -681,20 +1142,34 @@ func (p *PDPService) handleFinalizeStreamingUpload(w http.ResponseWriter, r *htt comm, err := p.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (commit bool, err error) { n, err := tx.Exec(` - INSERT INTO pdp_piece_uploads (id, service, piece_cid, notify_url, check_hash_codec, check_hash, check_size, piece_ref) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8) - `, uploadUUID.String(), serviceID, pieceCidV1.String(), req.Notify, multicodec.Sha2_256Trunc254Padded.String(), digest, pieceInfo.RawSize, pref) + INSERT INTO pdp_piecerefs (service, piece_cid, piece_ref, created_at, needs_save_cache) + SELECT su.service, su.piece_cid, su.piece_ref, NOW(), $4 + FROM pdp_piece_streaming_uploads su + JOIN parked_piece_refs ppr ON ppr.ref_id = su.piece_ref + JOIN parked_pieces pp ON pp.id = ppr.piece_id + WHERE su.id = $1 + AND su.service = $2 + AND su.piece_ref = $3 + AND su.complete = TRUE + AND pp.complete = TRUE + `, uploadUUID.String(), serviceID, pref, needsSaveCache(int64(rawSize))) if err != nil { - return false, fmt.Errorf("failed to store upload request in database: %w", err) + return false, fmt.Errorf("failed to create PDP piece reference: %w", err) } if n != 1 { - return false, fmt.Errorf("failed to store upload request in database: expected 1 row but got %d", n) + return false, fmt.Errorf("failed to create PDP piece reference: expected 1 row but got %d", n) } - _, err = tx.Exec(`DELETE FROM pdp_piece_streaming_uploads WHERE id = $1 AND service = $2 AND complete = TRUE`, uploadUUID.String(), serviceID) + n, err = tx.Exec(` + DELETE FROM pdp_piece_streaming_uploads + WHERE id = $1 AND service = $2 AND piece_ref = $3 AND complete = TRUE + `, uploadUUID.String(), serviceID, pref) if err != nil { return false, fmt.Errorf("failed to delete pdp_piece_streaming_uploads entry: %w", err) } + if n != 1 { + return false, fmt.Errorf("failed to delete pdp_piece_streaming_uploads entry: expected 1 row but got %d", n) + } return true, nil }, harmonydb.OptionRetry()) if err != nil { diff --git a/pdp/handlers_upload_test.go b/pdp/handlers_upload_test.go new file mode 100644 index 000000000..404d7c5bb --- /dev/null +++ b/pdp/handlers_upload_test.go @@ -0,0 +1,621 @@ +package pdp + +import ( + "bytes" + "context" + "database/sql" + "errors" + "io" + "net/http" + "net/http/httptest" + "path" + "testing" + "time" + + "github.com/go-chi/chi/v5" + "github.com/stretchr/testify/require" + + commcid "github.com/filecoin-project/go-fil-commcid" + commp "github.com/filecoin-project/go-fil-commp-hashhash" + "github.com/filecoin-project/go-state-types/abi" + + "github.com/filecoin-project/curio/harmony/harmonydb" + "github.com/filecoin-project/curio/harmony/harmonytask" + "github.com/filecoin-project/curio/lib/storiface" +) + +type uploadTestPieceIO struct { + writes int + removes int + declaredSize int64 + verifySize bool + storageType storiface.PathType + pieceID storiface.PieceNumber + writtenPieceData []byte + writeUploadErr error +} + +func (m *uploadTestPieceIO) WritePiece(context.Context, *harmonytask.TaskID, storiface.PieceNumber, int64, io.Reader, storiface.PathType) error { + return errors.New("unexpected WritePiece call") +} + +func (m *uploadTestPieceIO) WriteUploadPiece(_ context.Context, pieceID storiface.PieceNumber, size int64, data io.Reader, storageType storiface.PathType, verifySize bool) (abi.PieceInfo, uint64, error) { + m.writes++ + m.declaredSize = size + m.verifySize = verifySize + m.storageType = storageType + m.pieceID = pieceID + + body, err := io.ReadAll(data) + if err != nil { + return abi.PieceInfo{}, 0, err + } + m.writtenPieceData = append([]byte(nil), body...) + if m.writeUploadErr != nil { + return abi.PieceInfo{}, 0, m.writeUploadErr + } + calc := &commp.Calc{} + defer calc.Reset() + if _, err := calc.Write(body); err != nil { + return abi.PieceInfo{}, 0, err + } + digest, paddedSize, err := calc.Digest() + if err != nil { + return abi.PieceInfo{}, 0, err + } + pieceCID, err := commcid.DataCommitmentV1ToCID(digest) + if err != nil { + return abi.PieceInfo{}, 0, err + } + return abi.PieceInfo{PieceCID: pieceCID, Size: abi.PaddedPieceSize(paddedSize)}, uint64(len(body)), nil +} + +func (m *uploadTestPieceIO) PieceReader(context.Context, storiface.PieceNumber) (io.ReadCloser, error) { + return nil, errors.New("unexpected PieceReader call") +} + +func (m *uploadTestPieceIO) RemovePiece(_ context.Context, pieceID storiface.PieceNumber) error { + m.removes++ + m.pieceID = pieceID + return nil +} + +func testPieceCIDs(t *testing.T, body []byte) (string, string, int64) { + t.Helper() + calc := &commp.Calc{} + defer calc.Reset() + _, err := calc.Write(body) + require.NoError(t, err) + digest, paddedSize, err := calc.Digest() + require.NoError(t, err) + pieceCIDV1, err := commcid.DataCommitmentV1ToCID(digest) + require.NoError(t, err) + pieceCIDV2, err := commcid.DataCommitmentToPieceCidv2(digest, uint64(len(body))) + require.NoError(t, err) + return pieceCIDV1.String(), pieceCIDV2.String(), int64(paddedSize) +} + +func createClassicUpload(t *testing.T, service *PDPService, pieceCIDV2 string) string { + t.Helper() + req := httptest.NewRequest(http.MethodPost, "/pdp/piece", bytes.NewBufferString(`{"pieceCid":"`+pieceCIDV2+`"}`)) + rec := httptest.NewRecorder() + service.handlePiecePost(rec, req) + require.Equal(t, http.StatusCreated, rec.Code, rec.Body.String()) + return path.Base(rec.Header().Get("Location")) +} + +func putClassicUpload(service *PDPService, uploadID string, body []byte) *httptest.ResponseRecorder { + req := httptest.NewRequest(http.MethodPut, "/pdp/piece/upload/"+uploadID, bytes.NewReader(body)) + routeCtx := chi.NewRouteContext() + routeCtx.URLParams.Add("uploadUUID", uploadID) + req = req.WithContext(context.WithValue(req.Context(), chi.RouteCtxKey, routeCtx)) + rec := httptest.NewRecorder() + service.handlePieceUpload(rec, req) + return rec +} + +func createStreamingUpload(t *testing.T, service *PDPService) string { + t.Helper() + req := httptest.NewRequest(http.MethodPost, "/pdp/piece/uploads", nil) + rec := httptest.NewRecorder() + service.handleStreamingUploadURL(rec, req) + require.Equal(t, http.StatusCreated, rec.Code, rec.Body.String()) + return path.Base(rec.Header().Get("Location")) +} + +func putStreamingUpload(service *PDPService, uploadID string, body []byte) *httptest.ResponseRecorder { + req := httptest.NewRequest(http.MethodPut, "/pdp/piece/uploads/"+uploadID, bytes.NewReader(body)) + routeCtx := chi.NewRouteContext() + routeCtx.URLParams.Add("uploadUUID", uploadID) + req = req.WithContext(context.WithValue(req.Context(), chi.RouteCtxKey, routeCtx)) + rec := httptest.NewRecorder() + service.handleStreamingUpload(rec, req) + return rec +} + +func finalizeStreamingUpload(service *PDPService, uploadID, pieceCIDV2 string) *httptest.ResponseRecorder { + req := httptest.NewRequest(http.MethodPost, "/pdp/piece/uploads/"+uploadID, bytes.NewBufferString(`{"pieceCid":"`+pieceCIDV2+`"}`)) + routeCtx := chi.NewRouteContext() + routeCtx.URLParams.Add("uploadUUID", uploadID) + req = req.WithContext(context.WithValue(req.Context(), chi.RouteCtxKey, routeCtx)) + rec := httptest.NewRecorder() + service.handleFinalizeStreamingUpload(rec, req) + return rec +} + +func TestExactSizeReader(t *testing.T) { + t.Run("exact", func(t *testing.T) { + reader := &exactSizeReader{r: bytes.NewReader([]byte("abcd")), remaining: 4} + body, err := io.ReadAll(reader) + require.NoError(t, err) + require.Equal(t, []byte("abcd"), body) + }) + + t.Run("short", func(t *testing.T) { + reader := &exactSizeReader{r: bytes.NewReader([]byte("abc")), remaining: 4} + body, err := io.ReadAll(reader) + require.ErrorIs(t, err, errPieceTooShort) + require.Equal(t, []byte("abc"), body) + }) + + t.Run("leaves trailing byte", func(t *testing.T) { + body := NewTimeoutLimitReader(bytes.NewReader([]byte("abcde")), time.Second) + reader := &exactSizeReader{r: body, remaining: 4} + read, err := io.ReadAll(reader) + require.NoError(t, err) + require.Equal(t, []byte("abcd"), read) + hasExtra, err := readHasExtraByte(body) + require.NoError(t, err) + require.True(t, hasExtra) + }) + + t.Run("timeout limit overrun is a trailing byte", func(t *testing.T) { + body := NewTimeoutLimitReader(bytes.NewReader([]byte("x")), time.Second) + body.totalBytes = int64(PieceSizeMaxLimit) + hasExtra, err := readHasExtraByte(body) + require.NoError(t, err) + require.True(t, hasExtra) + }) +} + +func TestNeedsSaveCacheBoundary(t *testing.T) { + maxRawBelowThreshold := minPaddedPieceSizeForCache * 127 / 256 + require.False(t, needsSaveCache(maxRawBelowThreshold)) + require.True(t, needsSaveCache(maxRawBelowThreshold+1)) +} + +func TestHandlePieceUploadWritesDirectlyAndPublishes(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0xab}, 1024) + pieceCIDV1, pieceCIDV2, paddedSize := testPieceCIDs(t, body) + pio := &uploadTestPieceIO{} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + uploadID := createClassicUpload(t, service, pieceCIDV2) + + rec := putClassicUpload(service, uploadID, body) + require.Equal(t, http.StatusNoContent, rec.Code, rec.Body.String()) + require.Equal(t, 1, pio.writes) + require.Equal(t, int64(len(body)), pio.declaredSize) + require.True(t, pio.verifySize) + require.Equal(t, storiface.PathStorage, pio.storageType) + + var parkedPieceID, pieceRefID int64 + var complete, skip, cache bool + var dataURL sql.NullString + err = db.QueryRow(t.Context(), ` + SELECT pp.id, ppr.ref_id, pp.complete, pp.skip, ppr.data_url, pr.needs_save_cache + FROM pdp_piecerefs pr + JOIN parked_piece_refs ppr ON ppr.ref_id = pr.piece_ref + JOIN parked_pieces pp ON pp.id = ppr.piece_id + WHERE pr.service = 'public' AND pr.piece_cid = $1 + `, pieceCIDV1).Scan(&parkedPieceID, &pieceRefID, &complete, &skip, &dataURL, &cache) + require.NoError(t, err) + require.Equal(t, int64(pio.pieceID), parkedPieceID) + require.True(t, complete) + require.True(t, skip) + require.False(t, dataURL.Valid) + require.False(t, cache) + + var padded int64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT piece_padded_size FROM parked_pieces WHERE id = $1`, parkedPieceID).Scan(&padded)) + require.Equal(t, paddedSize, padded) + + var uploadExists bool + require.NoError(t, db.QueryRow(t.Context(), `SELECT EXISTS(SELECT 1 FROM pdp_piece_uploads WHERE id = $1)`, uploadID).Scan(&uploadExists)) + require.False(t, uploadExists) +} + +func TestStreamingUploadWritesDirectlyAndPublishesOnFinalize(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0x6a}, 1024) + pieceCIDV1, pieceCIDV2, paddedSize := testPieceCIDs(t, body) + pio := &uploadTestPieceIO{} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + uploadID := createStreamingUpload(t, service) + + earlyFinalize := finalizeStreamingUpload(service, uploadID, pieceCIDV2) + require.Equal(t, http.StatusConflict, earlyFinalize.Code, earlyFinalize.Body.String()) + + put := putStreamingUpload(service, uploadID, body) + require.Equal(t, http.StatusNoContent, put.Code, put.Body.String()) + require.Equal(t, 1, pio.writes) + require.Equal(t, int64(PieceSizeMaxLimit), pio.declaredSize) + require.False(t, pio.verifySize) + require.Equal(t, storiface.PathStorage, pio.storageType) + require.Equal(t, body, pio.writtenPieceData) + + var parkedPieceID, pieceRefID int64 + var streamingComplete, parkedComplete, skip bool + var dataURL sql.NullString + err = db.QueryRow(t.Context(), ` + SELECT pp.id, ppr.ref_id, su.complete, pp.complete, pp.skip, ppr.data_url + FROM pdp_piece_streaming_uploads su + JOIN parked_piece_refs ppr ON ppr.ref_id = su.piece_ref + JOIN parked_pieces pp ON pp.id = ppr.piece_id + WHERE su.id = $1 AND su.service = 'public' + `, uploadID).Scan(&parkedPieceID, &pieceRefID, &streamingComplete, &parkedComplete, &skip, &dataURL) + require.NoError(t, err) + require.Equal(t, int64(pio.pieceID), parkedPieceID) + require.True(t, streamingComplete) + require.True(t, parkedComplete) + require.True(t, skip) + require.False(t, dataURL.Valid) + + var storedPaddedSize int64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT piece_padded_size FROM parked_pieces WHERE id = $1`, parkedPieceID).Scan(&storedPaddedSize)) + require.Equal(t, paddedSize, storedPaddedSize) + + var refs int64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT COUNT(*) FROM pdp_piecerefs WHERE piece_ref = $1`, pieceRefID).Scan(&refs)) + require.Zero(t, refs) + + finalize := finalizeStreamingUpload(service, uploadID, pieceCIDV2) + require.Equal(t, http.StatusOK, finalize.Code, finalize.Body.String()) + + var serviceID, storedPieceCID string + var cache bool + require.NoError(t, db.QueryRow(t.Context(), ` + SELECT service, piece_cid, needs_save_cache + FROM pdp_piecerefs + WHERE piece_ref = $1 + `, pieceRefID).Scan(&serviceID, &storedPieceCID, &cache)) + require.Equal(t, "public", serviceID) + require.Equal(t, pieceCIDV1, storedPieceCID) + require.False(t, cache) + + var streamingExists, legacyUploadExists bool + require.NoError(t, db.QueryRow(t.Context(), `SELECT EXISTS(SELECT 1 FROM pdp_piece_streaming_uploads WHERE id = $1)`, uploadID).Scan(&streamingExists)) + require.NoError(t, db.QueryRow(t.Context(), `SELECT EXISTS(SELECT 1 FROM pdp_piece_uploads WHERE id = $1)`, uploadID).Scan(&legacyUploadExists)) + require.False(t, streamingExists) + require.False(t, legacyUploadExists) +} + +func TestStreamingUploadWriteFailureReleasesClaimForRetry(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0x7b}, 1024) + pieceCIDV1, pieceCIDV2, _ := testPieceCIDs(t, body) + pio := &uploadTestPieceIO{writeUploadErr: errors.New("injected write failure")} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + uploadID := createStreamingUpload(t, service) + + failed := putStreamingUpload(service, uploadID, body) + require.Equal(t, http.StatusInternalServerError, failed.Code, failed.Body.String()) + require.Equal(t, 1, pio.writes) + require.Equal(t, int64(PieceSizeMaxLimit), pio.declaredSize) + require.False(t, pio.verifySize) + require.Equal(t, storiface.PathStorage, pio.storageType) + require.Equal(t, body, pio.writtenPieceData) + require.Zero(t, pio.removes) + failedPieceID := pio.pieceID + + var pieceRef sql.NullInt64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT piece_ref FROM pdp_piece_streaming_uploads WHERE id = $1`, uploadID).Scan(&pieceRef)) + require.False(t, pieceRef.Valid) + + var parkedExists, parkedComplete bool + var refCount int + require.NoError(t, db.QueryRow(t.Context(), ` + SELECT EXISTS(SELECT 1 FROM parked_pieces WHERE id = $1), + COALESCE((SELECT complete FROM parked_pieces WHERE id = $1), FALSE), + COALESCE((SELECT ref_count FROM parked_pieces WHERE id = $1), -1) + `, int64(failedPieceID)).Scan(&parkedExists, &parkedComplete, &refCount)) + require.True(t, parkedExists) + require.False(t, parkedComplete) + require.Zero(t, refCount) + + pio.writeUploadErr = nil + retry := putStreamingUpload(service, uploadID, body) + require.Equal(t, http.StatusNoContent, retry.Code, retry.Body.String()) + require.Equal(t, 2, pio.writes) + require.Zero(t, pio.removes) + require.NotEqual(t, failedPieceID, pio.pieceID) + + finalize := finalizeStreamingUpload(service, uploadID, pieceCIDV2) + require.Equal(t, http.StatusOK, finalize.Code, finalize.Body.String()) + + var refs int64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT COUNT(*) FROM pdp_piecerefs WHERE piece_cid = $1`, pieceCIDV1).Scan(&refs)) + require.Equal(t, int64(1), refs) +} + +func TestStreamingUploadReusesCompletedPiece(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0x8c}, 1024) + pieceCIDV1, pieceCIDV2, paddedSize := testPieceCIDs(t, body) + var parkedPieceID int64 + err = db.QueryRow(t.Context(), ` + INSERT INTO parked_pieces (piece_cid, piece_padded_size, piece_raw_size, complete, long_term) + VALUES ($1, $2, $3, TRUE, TRUE) + RETURNING id + `, pieceCIDV1, paddedSize, len(body)).Scan(&parkedPieceID) + require.NoError(t, err) + + pio := &uploadTestPieceIO{} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + uploadID := createStreamingUpload(t, service) + + put := putStreamingUpload(service, uploadID, body) + require.Equal(t, http.StatusNoContent, put.Code, put.Body.String()) + require.Equal(t, 1, pio.writes) + require.Equal(t, int64(PieceSizeMaxLimit), pio.declaredSize) + require.False(t, pio.verifySize) + require.Equal(t, storiface.PathStorage, pio.storageType) + require.Equal(t, body, pio.writtenPieceData) + require.Zero(t, pio.removes) + require.NotEqual(t, storiface.PieceNumber(parkedPieceID), pio.pieceID) + + var claimedParkedPieceID, pieceRefID int64 + var complete bool + require.NoError(t, db.QueryRow(t.Context(), ` + SELECT ppr.piece_id, su.piece_ref, su.complete + FROM pdp_piece_streaming_uploads su + JOIN parked_piece_refs ppr ON ppr.ref_id = su.piece_ref + WHERE su.id = $1 + `, uploadID).Scan(&claimedParkedPieceID, &pieceRefID, &complete)) + require.Equal(t, parkedPieceID, claimedParkedPieceID) + require.True(t, complete) + + finalize := finalizeStreamingUpload(service, uploadID, pieceCIDV2) + require.Equal(t, http.StatusOK, finalize.Code, finalize.Body.String()) + + var publishedPieceRefID int64 + require.NoError(t, db.QueryRow(t.Context(), ` + SELECT piece_ref FROM pdp_piecerefs + WHERE service = 'public' AND piece_cid = $1 + `, pieceCIDV1).Scan(&publishedPieceRefID)) + require.Equal(t, pieceRefID, publishedPieceRefID) +} + +func TestCleanupExpiredStreamingUploadClaims(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0x9d}, 1024) + _, pieceCIDV2, _ := testPieceCIDs(t, body) + pio := &uploadTestPieceIO{} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + uploadID := createStreamingUpload(t, service) + + claim, err := service.claimStreamingUpload(t.Context(), uploadID, "public") + require.NoError(t, err) + require.True(t, claim.created) + _, err = db.Exec(t.Context(), `UPDATE pdp_piece_streaming_uploads SET created_at = NOW() - INTERVAL '2 hours' WHERE id = $1`, uploadID) + require.NoError(t, err) + + require.NoError(t, service.cleanupExpiredStreamingUploadClaims(t.Context())) + require.Zero(t, pio.removes) + + var pieceRef sql.NullInt64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT piece_ref FROM pdp_piece_streaming_uploads WHERE id = $1`, uploadID).Scan(&pieceRef)) + require.False(t, pieceRef.Valid) + + var refExists, parkedExists bool + var refCount int + require.NoError(t, db.QueryRow(t.Context(), `SELECT EXISTS(SELECT 1 FROM parked_piece_refs WHERE ref_id = $1)`, claim.pieceRefID).Scan(&refExists)) + require.NoError(t, db.QueryRow(t.Context(), ` + SELECT EXISTS(SELECT 1 FROM parked_pieces WHERE id = $1), + COALESCE((SELECT ref_count FROM parked_pieces WHERE id = $1), -1) + `, claim.parkedPieceID).Scan(&parkedExists, &refCount)) + require.False(t, refExists) + require.True(t, parkedExists) + require.Zero(t, refCount) + + retry := putStreamingUpload(service, uploadID, body) + require.Equal(t, http.StatusNoContent, retry.Code, retry.Body.String()) + require.Equal(t, 1, pio.writes) + require.Equal(t, int64(PieceSizeMaxLimit), pio.declaredSize) + require.False(t, pio.verifySize) + require.Equal(t, storiface.PathStorage, pio.storageType) + finalize := finalizeStreamingUpload(service, uploadID, pieceCIDV2) + require.Equal(t, http.StatusOK, finalize.Code, finalize.Body.String()) +} + +func TestHandlePieceUploadRejectsInvalidLengthsAndReleasesClaim(t *testing.T) { + tests := []struct { + name string + body func([]byte) []byte + }{ + {name: "short", body: func(body []byte) []byte { return body[:len(body)-1] }}, + {name: "oversized", body: func(body []byte) []byte { return append(append([]byte{}, body...), 0xff) }}, + {name: "cid mismatch", body: func(body []byte) []byte { + wrong := append([]byte{}, body...) + wrong[0] ^= 0xff + return wrong + }}, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + expected := bytes.Repeat([]byte{0xcd}, 1024) + _, pieceCIDV2, _ := testPieceCIDs(t, expected) + pio := &uploadTestPieceIO{} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + uploadID := createClassicUpload(t, service, pieceCIDV2) + + rec := putClassicUpload(service, uploadID, test.body(expected)) + require.Equal(t, http.StatusBadRequest, rec.Code, rec.Body.String()) + + var pieceRef sql.NullInt64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT piece_ref FROM pdp_piece_uploads WHERE id = $1`, uploadID).Scan(&pieceRef)) + require.False(t, pieceRef.Valid) + require.Equal(t, 1, pio.removes) + + retry := putClassicUpload(service, uploadID, expected) + require.Equal(t, http.StatusNoContent, retry.Code, retry.Body.String()) + }) + } +} + +func TestHandlePiecePostPublishesAlreadyStoredPiece(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0xef}, 1024) + pieceCIDV1, pieceCIDV2, paddedSize := testPieceCIDs(t, body) + _, err = db.Exec(t.Context(), ` + INSERT INTO parked_pieces (piece_cid, piece_padded_size, piece_raw_size, complete, long_term, skip) + VALUES ($1, $2, $3, TRUE, TRUE, TRUE) + `, pieceCIDV1, paddedSize, len(body)) + require.NoError(t, err) + + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: &uploadTestPieceIO{}} + req := httptest.NewRequest(http.MethodPost, "/pdp/piece", bytes.NewBufferString(`{"pieceCid":"`+pieceCIDV2+`"}`)) + rec := httptest.NewRecorder() + service.handlePiecePost(rec, req) + require.Equal(t, http.StatusOK, rec.Code, rec.Body.String()) + + var refs, uploads int64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT COUNT(*) FROM pdp_piecerefs WHERE piece_cid = $1`, pieceCIDV1).Scan(&refs)) + require.NoError(t, db.QueryRow(t.Context(), `SELECT COUNT(*) FROM pdp_piece_uploads WHERE piece_cid = $1`, pieceCIDV1).Scan(&uploads)) + require.Equal(t, int64(1), refs) + require.Zero(t, uploads) +} + +func TestSecondPendingUploadReusesCompletedPiece(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0x17}, 1024) + pieceCIDV1, pieceCIDV2, _ := testPieceCIDs(t, body) + pio := &uploadTestPieceIO{} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + firstUpload := createClassicUpload(t, service, pieceCIDV2) + secondUpload := createClassicUpload(t, service, pieceCIDV2) + + first := putClassicUpload(service, firstUpload, body) + require.Equal(t, http.StatusNoContent, first.Code, first.Body.String()) + second := putClassicUpload(service, secondUpload, body) + require.Equal(t, http.StatusNoContent, second.Code, second.Body.String()) + require.Equal(t, 1, pio.writes) + + var refs, uploads int64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT COUNT(*) FROM pdp_piecerefs WHERE piece_cid = $1`, pieceCIDV1).Scan(&refs)) + require.NoError(t, db.QueryRow(t.Context(), `SELECT COUNT(*) FROM pdp_piece_uploads WHERE piece_cid = $1`, pieceCIDV1).Scan(&uploads)) + require.Equal(t, int64(2), refs) + require.Zero(t, uploads) +} + +func TestActiveSameCIDUploadRejectsSecondClaim(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0x31}, 1024) + pieceCIDV1, pieceCIDV2, paddedSize := testPieceCIDs(t, body) + pio := &uploadTestPieceIO{} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + firstUpload := createClassicUpload(t, service, pieceCIDV2) + secondUpload := createClassicUpload(t, service, pieceCIDV2) + + firstClaim, err := service.claimDirectUpload(t.Context(), firstUpload, "public", pieceCIDV1, int64(len(body)), paddedSize) + require.NoError(t, err) + require.True(t, firstClaim.created) + + second := putClassicUpload(service, secondUpload, body) + require.Equal(t, http.StatusConflict, second.Code, second.Body.String()) + require.Zero(t, pio.writes) + require.Zero(t, pio.removes) + + var firstRef, secondRef sql.NullInt64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT piece_ref FROM pdp_piece_uploads WHERE id = $1`, firstUpload).Scan(&firstRef)) + require.NoError(t, db.QueryRow(t.Context(), `SELECT piece_ref FROM pdp_piece_uploads WHERE id = $1`, secondUpload).Scan(&secondRef)) + require.True(t, firstRef.Valid) + require.Equal(t, firstClaim.pieceRefID, firstRef.Int64) + require.False(t, secondRef.Valid) + + var complete bool + require.NoError(t, db.QueryRow(t.Context(), `SELECT complete FROM parked_pieces WHERE id = $1`, firstClaim.parkedPieceID).Scan(&complete)) + require.False(t, complete) +} + +func TestCleanupExpiredDirectUploadClaims(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0x42}, 1024) + pieceCIDV1, pieceCIDV2, paddedSize := testPieceCIDs(t, body) + pio := &uploadTestPieceIO{} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + uploadID := createClassicUpload(t, service, pieceCIDV2) + claim, err := service.claimDirectUpload(t.Context(), uploadID, "public", pieceCIDV1, int64(len(body)), paddedSize) + require.NoError(t, err) + require.True(t, claim.created) + + _, err = db.Exec(t.Context(), `UPDATE pdp_piece_uploads SET created_at = NOW() - INTERVAL '2 hours' WHERE id = $1`, uploadID) + require.NoError(t, err) + require.NoError(t, service.cleanupExpiredDirectUploadClaims(t.Context())) + + var pieceRef sql.NullInt64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT piece_ref FROM pdp_piece_uploads WHERE id = $1`, uploadID).Scan(&pieceRef)) + require.False(t, pieceRef.Valid) + + var refExists bool + require.NoError(t, db.QueryRow(t.Context(), `SELECT EXISTS(SELECT 1 FROM parked_piece_refs WHERE ref_id = $1)`, claim.pieceRefID).Scan(&refExists)) + require.False(t, refExists) + + var parkedExists bool + require.NoError(t, db.QueryRow(t.Context(), `SELECT EXISTS(SELECT 1 FROM parked_pieces WHERE id = $1)`, claim.parkedPieceID).Scan(&parkedExists)) + require.False(t, parkedExists) + require.Equal(t, 1, pio.removes) + require.Equal(t, storiface.PieceNumber(claim.parkedPieceID), pio.pieceID) + + retry := putClassicUpload(service, uploadID, body) + require.Equal(t, http.StatusNoContent, retry.Code, retry.Body.String()) +} + +func TestCleanupExpiredDirectUploadClaimsLeavesFreshClaim(t *testing.T) { + db, err := harmonydb.NewFromConfigWithITestID(t) + require.NoError(t, err) + + body := bytes.Repeat([]byte{0x53}, 1024) + pieceCIDV1, pieceCIDV2, paddedSize := testPieceCIDs(t, body) + pio := &uploadTestPieceIO{} + service := &PDPService{Auth: &NullAuth{}, db: db, pieceIO: pio} + uploadID := createClassicUpload(t, service, pieceCIDV2) + claim, err := service.claimDirectUpload(t.Context(), uploadID, "public", pieceCIDV1, int64(len(body)), paddedSize) + require.NoError(t, err) + + require.NoError(t, service.cleanupExpiredDirectUploadClaims(t.Context())) + + var pieceRef sql.NullInt64 + require.NoError(t, db.QueryRow(t.Context(), `SELECT piece_ref FROM pdp_piece_uploads WHERE id = $1`, uploadID).Scan(&pieceRef)) + require.True(t, pieceRef.Valid) + require.Equal(t, claim.pieceRefID, pieceRef.Int64) + + var refExists, parkedExists bool + require.NoError(t, db.QueryRow(t.Context(), `SELECT EXISTS(SELECT 1 FROM parked_piece_refs WHERE ref_id = $1)`, claim.pieceRefID).Scan(&refExists)) + require.NoError(t, db.QueryRow(t.Context(), `SELECT EXISTS(SELECT 1 FROM parked_pieces WHERE id = $1)`, claim.parkedPieceID).Scan(&parkedExists)) + require.True(t, refExists) + require.True(t, parkedExists) + require.Zero(t, pio.removes) +} diff --git a/pdp/mount.go b/pdp/mount.go index 902500ce2..dd9bc044f 100644 --- a/pdp/mount.go +++ b/pdp/mount.go @@ -10,18 +10,18 @@ import ( "github.com/filecoin-project/curio/api" "github.com/filecoin-project/curio/harmony/harmonydb" "github.com/filecoin-project/curio/lib/ethchain" - "github.com/filecoin-project/curio/lib/paths" + "github.com/filecoin-project/curio/lib/piecestore" ipni_provider "github.com/filecoin-project/curio/market/ipni/ipni-provider" ) // MountDeps holds dependencies for mounting PDP HTTP routes. type MountDeps struct { - DB *harmonydb.DB - LocalStore paths.StashStore - EthClient ethchain.EthClient - Chain api.Chain - EthSender ETHTxSender - AlertTask *alertmanager.AlertTask + DB *harmonydb.DB + PieceIO piecestore.PieceIO + EthClient ethchain.EthClient + Chain api.Chain + EthSender ETHTxSender + AlertTask *alertmanager.AlertTask } // MountRoutes registers PDP HTTP routes on an existing router. @@ -29,8 +29,11 @@ func MountRoutes(ctx context.Context, r chi.Router, d MountDeps, ipp *ipni_provi if d.EthSender == nil { return xerrors.Errorf("eth sender required for PDP routes") } + if d.PieceIO == nil { + return xerrors.Errorf("piece IO required for PDP routes") + } - pdsvc := NewPDPService(ctx, d.DB, d.LocalStore, d.EthClient, d.Chain, d.EthSender, d.AlertTask, ipp) + pdsvc := NewPDPService(ctx, d.DB, d.PieceIO, d.EthClient, d.Chain, d.EthSender, d.AlertTask, ipp) Routes(r, pdsvc) return nil } diff --git a/pdpnode/routes.go b/pdpnode/routes.go index 4f37bf21c..48c66b5ad 100644 --- a/pdpnode/routes.go +++ b/pdpnode/routes.go @@ -16,12 +16,12 @@ import ( // MountPDPRoutes attaches PDP HTTP routes using an existing IPNI provider. func MountPDPRoutes(ctx context.Context, r chi.Router, d *Deps, sd *servicedeps.Deps, ipp *ipni_provider.Provider) error { return pdp.MountRoutes(ctx, r, pdp.MountDeps{ - DB: d.DB, - LocalStore: d.LocalStore, - EthClient: must.One(d.EthClient.Val()), - Chain: d.Chain, - EthSender: sd.EthSender, - AlertTask: sd.AlertTask, + DB: d.DB, + PieceIO: d.PieceIO, + EthClient: must.One(d.EthClient.Val()), + Chain: d.Chain, + EthSender: sd.EthSender, + AlertTask: sd.AlertTask, }, ipp) } diff --git a/pdpnode/skiff_lifecycle_integration_test.go b/pdpnode/skiff_lifecycle_integration_test.go index f0b115c7c..119996ab2 100644 --- a/pdpnode/skiff_lifecycle_integration_test.go +++ b/pdpnode/skiff_lifecycle_integration_test.go @@ -35,6 +35,7 @@ import ( "github.com/filecoin-project/curio/lib/ethchain" "github.com/filecoin-project/curio/lib/paths" "github.com/filecoin-project/curio/lib/pieceprovider" + "github.com/filecoin-project/curio/lib/piecestore" "github.com/filecoin-project/curio/lib/storiface" "github.com/filecoin-project/curio/pdp" "github.com/filecoin-project/curio/pdp/contract" @@ -123,10 +124,10 @@ func TestSkiffCreateAddRetrieveLifecycle(t *testing.T) { svcCtx, svcCancel := context.WithCancel(ctx) t.Cleanup(svcCancel) require.NoError(t, pdp.MountRoutes(svcCtx, mux, pdp.MountDeps{ - DB: db, - LocalStore: localStore, - EthClient: stubEthClient{}, - EthSender: mockSender, + DB: db, + PieceIO: piecestore.New(remote, localStore, index), + EthClient: stubEthClient{}, + EthSender: mockSender, }, nil)) srv := httptest.NewServer(mux) @@ -177,28 +178,20 @@ func TestSkiffCreateAddRetrieveLifecycle(t *testing.T) { require.Equal(t, http.StatusNoContent, putRes.StatusCode, httpBody(t, putRes)) require.NoError(t, putRes.Body.Close()) - var parkedID, pieceRef int64 - var uploadID string + var parkedID, pieceRef, pdpPieceRefID int64 + var parkedComplete bool err = db.QueryRow(ctx, ` - SELECT pu.id, pu.piece_ref, pp.id - FROM pdp_piece_uploads pu - JOIN parked_piece_refs ppr ON ppr.ref_id = pu.piece_ref + SELECT pr.id, pr.piece_ref, pp.id, pp.complete + FROM pdp_piecerefs pr + JOIN parked_piece_refs ppr ON ppr.ref_id = pr.piece_ref JOIN parked_pieces pp ON pp.id = ppr.piece_id - WHERE pu.piece_cid = $1`, pieceCidV1.String()).Scan(&uploadID, &pieceRef, &parkedID) + WHERE pr.piece_cid = $1`, pieceCidV1.String()).Scan(&pdpPieceRefID, &pieceRef, &parkedID, &parkedComplete) require.NoError(t, err) + require.True(t, parkedComplete) - require.NoError(t, writeParkedPieceBytes(storageDir, parkedID, raw)) - _, err = db.Exec(ctx, `UPDATE parked_pieces SET complete = TRUE WHERE id = $1`, parkedID) - require.NoError(t, err) - - var pdpPieceRefID int64 - err = db.QueryRow(ctx, ` - INSERT INTO pdp_piecerefs (service, piece_cid, piece_ref, created_at) - VALUES ('public', $1, $2, NOW()) - RETURNING id`, pieceCidV1.String(), pieceRef).Scan(&pdpPieceRefID) - require.NoError(t, err) - _, err = db.Exec(ctx, `DELETE FROM pdp_piece_uploads WHERE id = $1`, uploadID) - require.NoError(t, err) + var uploadExists bool + require.NoError(t, db.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM pdp_piece_uploads WHERE piece_cid = $1)`, pieceCidV1.String()).Scan(&uploadExists)) + require.False(t, uploadExists) addBody, err := json.Marshal(map[string]any{ "pieces": []map[string]any{{ @@ -266,18 +259,6 @@ func httpBody(t *testing.T, res *http.Response) string { return string(b) } -func writeParkedPieceBytes(storageDir string, pieceID int64, data []byte) error { - path := filepath.Join( - storageDir, - storiface.FTPiece.String(), - storiface.SectorName(storiface.PieceNumber(pieceID).Ref().ID), - ) - if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { - return err - } - return os.WriteFile(path, data, 0o644) -} - type stubEthClient struct{} var _ ethchain.EthClient = stubEthClient{} diff --git a/pdpnode/tasks.go b/pdpnode/tasks.go index ecf530e8e..7c8d139c4 100644 --- a/pdpnode/tasks.go +++ b/pdpnode/tasks.go @@ -61,7 +61,6 @@ func buildPDPTasks(ctx context.Context, d *Deps, chainSched *chainsched.CurioCha pdp.NewWatcherPieceDelete(db, chainSched) tasks = append(tasks, - pdp.NewPDPNotifyTask(db), pdp.NewProveTask(chainSched, db, ethClient, d.Chain, senderEth, d.CachedPieceReader, d.IndexStore), pdp.NewNextProvingPeriodTask(db, ethClient, d.Chain, chainSched, senderEth), pdp.NewInitProvingPeriodTask(db, ethClient, d.Chain, chainSched, senderEth), @@ -86,7 +85,6 @@ func buildPDPTasks(ctx context.Context, d *Deps, chainSched *chainsched.CurioCha pdpv0.NewProveTask(db, ethClient, d.Chain, w, senderEth, d.CachedPieceReader, d.IndexStore), pdpv0.NewNextProvingPeriodTask(db, ethClient, d.Chain, w, senderEth), pdpv0.NewInitProvingPeriodTask(db, ethClient, d.Chain, w, senderEth), - pdpv0.NewPDPNotifyTask(ctx, db), pdpv0.NewPDPPullPieceTask(ctx, db, d.PieceIO, cfg.Subsystems.PDPPullPieceMaxTasks), pdpv0.NewTerminateServiceTask(db, ethClient, senderEth), pdpv0.NewDeleteDataSetTask(db, ethClient, senderEth), diff --git a/tasks/pdp/notify_task.go b/tasks/pdp/notify_task.go deleted file mode 100644 index f243f0c88..000000000 --- a/tasks/pdp/notify_task.go +++ /dev/null @@ -1,181 +0,0 @@ -package pdp - -import ( - "bytes" - "context" - "database/sql" - "encoding/json" - "fmt" - "net/http" - "time" - - "github.com/curiostorage/harmonyquery" - logger "github.com/ipfs/go-log/v2" - "golang.org/x/xerrors" - - "github.com/filecoin-project/curio/harmony/harmonydb" - "github.com/filecoin-project/curio/harmony/harmonytask" - "github.com/filecoin-project/curio/harmony/resources" - "github.com/filecoin-project/curio/harmony/taskhelp" - "github.com/filecoin-project/curio/lib/passcall" - "github.com/filecoin-project/curio/tasks/tasknames" -) - -var log = logger.Logger("pdp") - -type PDPNotifyTask struct { - db *harmonydb.DB - client *http.Client -} - -func NewPDPNotifyTask(db *harmonydb.DB) *PDPNotifyTask { - client := &http.Client{ - Timeout: 15 * time.Second, - Transport: &http.Transport{ - ResponseHeaderTimeout: 10 * time.Second, - IdleConnTimeout: 30 * time.Second, - }, - } - return &PDPNotifyTask{db: db, client: client} -} - -func (t *PDPNotifyTask) Do(ctx context.Context, taskID harmonytask.TaskID, stillOwned func() bool) (done bool, err error) { - - // Fetch the pdp_piece_uploads entry associated with the taskID - var upload struct { - ID string `db:"id" json:"id"` - Service string `db:"service" json:"service"` - PieceCID sql.NullString `db:"piece_cid" json:"piece_cid"` - NotifyURL string `db:"notify_url" json:"notify_url"` - PieceRef int64 `db:"piece_ref" json:"piece_ref"` - CheckHashCodec string `db:"check_hash_codec" json:"check_hash_codec"` - CheckHash []byte `db:"check_hash" json:"check_hash"` - } - err = t.db.QueryRow(ctx, ` - SELECT id, service, piece_cid, notify_url, piece_ref, check_hash_codec, check_hash - FROM pdp_piece_uploads - WHERE notify_task_id = $1`, taskID).Scan( - &upload.ID, &upload.Service, &upload.PieceCID, &upload.NotifyURL, &upload.PieceRef, &upload.CheckHashCodec, &upload.CheckHash) - if err != nil { - return false, fmt.Errorf("failed to query pdp_piece_uploads for task %d: %w", taskID, err) - } - - // Perform HTTP Post request to the notify URL - upJson, err := json.Marshal(upload) - if err != nil { - return false, fmt.Errorf("failed to marshal upload to JSON: %w", err) - } - - log.Infow("PDP notify", "upload", upload, "task_id", taskID) - - if upload.NotifyURL != "" { - - resp, err := t.client.Post(upload.NotifyURL, "application/json", bytes.NewReader(upJson)) - if err != nil { - log.Errorw("HTTP POST request to notify_url failed", "notify_url", upload.NotifyURL, "upload_id", upload.ID, "error", err) - } else { - defer func() { - _ = resp.Body.Close() - }() - // Not reading the body as per requirement - log.Infow("HTTP GET request to notify_url succeeded", "notify_url", upload.NotifyURL, "upload_id", upload.ID) - } - } - - // Move the entry from pdp_piece_uploads to pdp_piecerefs - comm, err := t.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { - _, err := tx.Exec(` - INSERT INTO pdp_piecerefs (service, piece_cid, piece_ref, created_at) - VALUES ($1, $2, $3, NOW())`, - upload.Service, upload.PieceCID, upload.PieceRef) - if err != nil { - return false, fmt.Errorf("failed to insert into pdp_piecerefs: %w", err) - } - - _, err = tx.Exec(`DELETE FROM pdp_piece_uploads WHERE id = $1`, upload.ID) - if err != nil { - return false, fmt.Errorf("failed to delete upload ID %s from pdp_piece_uploads: %w", upload.ID, err) - } - - return true, nil - }, harmonyquery.OptionRetry()) - if err != nil { - return false, fmt.Errorf("failed to move upload to piecerefs: %w", err) - } - if !comm { - return false, fmt.Errorf("transaction to move upload to piecerefs was not committed") - } - - log.Infof("Successfully processed PDP notify task %d for upload ID %s", taskID, upload.ID) - - return true, nil -} - -func (t *PDPNotifyTask) CanAccept(ids []harmonytask.TaskID, engine *harmonytask.TaskEngine) ([]harmonytask.TaskID, error) { - if len(ids) == 0 { - return []harmonytask.TaskID{}, nil - } - return ids, nil -} - -func (t *PDPNotifyTask) TypeDetails() harmonytask.TaskTypeDetails { - return harmonytask.TaskTypeDetails{ - Name: tasknames.PDPNotify, - MayFollow: []string{tasknames.PDPAddPiece}, - Cost: resources.Resources{ - Cpu: 0, - Ram: 128 << 20, // 128MB - }, - MaxFailures: 14, - RetryWait: taskhelp.RetryWaitExp(5*time.Second, 2), - IAmBored: passcall.Every(time.Second, func(taskFunc harmonytask.AddTaskFunc) error { - return t.schedule(context.Background(), taskFunc) - }), - } -} - -func (t *PDPNotifyTask) schedule(ctx context.Context, taskFunc harmonytask.AddTaskFunc) error { - for { - stop := true - taskFunc(func(id harmonytask.TaskID, tx *harmonydb.Tx) (shouldCommit bool, seriousError error) { - n, err := tx.Exec(` - WITH pending AS ( - SELECT pu.id - FROM pdp_piece_uploads pu - JOIN parked_piece_refs pr ON pr.ref_id = pu.piece_ref - JOIN parked_pieces pp ON pp.id = pr.piece_id - WHERE pu.piece_ref IS NOT NULL - AND pp.complete = TRUE - AND pu.notify_task_id IS NULL - LIMIT 1 - ) - UPDATE pdp_piece_uploads pu - SET notify_task_id = $1 - FROM pending - WHERE pu.id = pending.id - AND pu.notify_task_id IS NULL - `, id) - if err != nil { - return false, xerrors.Errorf("updating notify_task_id: %w", err) - } - if n == 0 { - return false, nil - } - if n != 1 { - return false, xerrors.Errorf("updated %d rows assigning pdp notify task", n) - } - - stop = false // Continue scheduling as there might be more tasks - return true, nil // Commit the transaction - }) - if stop { - return nil - } - } -} - -func (t *PDPNotifyTask) Adder(taskFunc harmonytask.AddTaskFunc) { -} - -var _ = harmonytask.Reg(&PDPNotifyTask{}) -var _ harmonytask.TaskInterface = &PDPNotifyTask{} diff --git a/tasks/pdp/task_init_pp.go b/tasks/pdp/task_init_pp.go index 541bd7404..07a737498 100644 --- a/tasks/pdp/task_init_pp.go +++ b/tasks/pdp/task_init_pp.go @@ -8,6 +8,7 @@ import ( "github.com/ethereum/go-ethereum/accounts/abi/bind" "github.com/ethereum/go-ethereum/core/types" + logging "github.com/ipfs/go-log/v2" "github.com/yugabyte/pgx/v5" "golang.org/x/xerrors" @@ -24,6 +25,8 @@ import ( chainTypes "github.com/filecoin-project/lotus/chain/types" ) +var log = logging.Logger("pdp") + type InitProvingPeriodTask struct { db *harmonydb.DB ethClient ethchain.EthClient diff --git a/tasks/pdpv0/notify_task.go b/tasks/pdpv0/notify_task.go deleted file mode 100644 index 30d294301..000000000 --- a/tasks/pdpv0/notify_task.go +++ /dev/null @@ -1,212 +0,0 @@ -package pdpv0 - -import ( - "bytes" - "context" - "encoding/json" - "fmt" - "net/http" - "time" - - logging "github.com/ipfs/go-log/v2" - "golang.org/x/xerrors" - - "github.com/filecoin-project/go-padreader" - "github.com/filecoin-project/go-state-types/abi" - - "github.com/filecoin-project/curio/harmony/harmonydb" - "github.com/filecoin-project/curio/harmony/harmonytask" - "github.com/filecoin-project/curio/harmony/resources" - "github.com/filecoin-project/curio/harmony/taskhelp" - "github.com/filecoin-project/curio/lib/promise" - "github.com/filecoin-project/curio/tasks/tasknames" -) - -var log = logging.Logger("pdpv0") - -// NotifyPollInterval is how often to poll for uploads ready to finalize. -var NotifyPollInterval = 5 * time.Second - -// PDPNotifyTask finalizes completed piece uploads. -// -// When piece data finishes uploading (parked_pieces.complete = TRUE), this task: -// 1. Sends an optional HTTP POST callback to notify_url if configured -// 2. Creates a permanent reference in pdp_piecerefs linking piece_cid to piece_ref -// 3. Removes the temporary upload record from pdp_piece_uploads -// -// The poll goroutine watches for uploads where the underlying piece is complete -// but no finalization task has been assigned yet. -type PDPNotifyTask struct { - db *harmonydb.DB - TF promise.Promise[harmonytask.AddTaskFunc] - client *http.Client -} - -func NewPDPNotifyTask(ctx context.Context, db *harmonydb.DB) *PDPNotifyTask { - client := &http.Client{ - Timeout: 15 * time.Second, - Transport: &http.Transport{ - ResponseHeaderTimeout: 10 * time.Second, - IdleConnTimeout: 30 * time.Second, - }, - } - n := &PDPNotifyTask{db: db, client: client} - go n.poll(ctx) - return n -} - -func (t *PDPNotifyTask) poll(ctx context.Context) { - ticker := time.NewTicker(NotifyPollInterval) - defer ticker.Stop() - - for { - select { - case <-ctx.Done(): - return - case <-ticker.C: - } - - var uploads []struct { - ID string `db:"id"` - } - - err := t.db.Select(ctx, &uploads, ` - SELECT pu.id - FROM pdp_piece_uploads pu - JOIN parked_piece_refs pr ON pr.ref_id = pu.piece_ref - JOIN parked_pieces pp ON pp.id = pr.piece_id - WHERE - pu.piece_ref IS NOT NULL - AND pp.complete = TRUE - AND pu.notify_task_id IS NULL LIMIT 10`) - if err != nil { - log.Errorf("getting uploads to notify: %s", err) - continue - } - - if len(uploads) == 0 { - continue - } - - for _, upload := range uploads { - failed := false - - t.TF.Val(ctx)(func(id harmonytask.TaskID, tx *harmonydb.Tx) (shouldCommit bool, err error) { - n, err := tx.Exec(` - UPDATE pdp_piece_uploads - SET notify_task_id = $1 - WHERE id = $2 AND notify_task_id IS NULL`, id, upload.ID) - if err != nil { - failed = true - return false, xerrors.Errorf("updating notify_task_id: %w", err) - } - return n > 0, nil - }) - if failed { - break - } - } - } -} - -func (t *PDPNotifyTask) Do(ctx context.Context, taskID harmonytask.TaskID, stillOwned func() bool) (done bool, err error) { - - // Fetch the pdp_piece_uploads entry associated with the taskID - var upload struct { - ID string `db:"id" json:"id"` - Service string `db:"service" json:"service"` - PieceCID *string `db:"piece_cid" json:"piece_cid"` - NotifyURL string `db:"notify_url" json:"notify_url"` - PieceRef int64 `db:"piece_ref" json:"piece_ref"` - CheckHashCodec string `db:"check_hash_codec" json:"check_hash_codec"` - CheckHash []byte `db:"check_hash" json:"check_hash"` - PieceRawSize uint64 `db:"piece_raw_size" json:"piece_raw_size"` - } - err = t.db.QueryRow(ctx, ` - SELECT pu.id, pu.service, pu.piece_cid, pu.notify_url, pu.piece_ref, pu.check_hash_codec, pu.check_hash, - pp.piece_raw_size - FROM pdp_piece_uploads pu - JOIN parked_piece_refs ppr ON ppr.ref_id = pu.piece_ref - JOIN parked_pieces pp ON pp.id = ppr.piece_id - WHERE pu.notify_task_id = $1`, taskID).Scan( - &upload.ID, &upload.Service, &upload.PieceCID, &upload.NotifyURL, &upload.PieceRef, &upload.CheckHashCodec, &upload.CheckHash, &upload.PieceRawSize) - if err != nil { - return false, fmt.Errorf("failed to query pdp_piece_uploads for task %d: %w", taskID, err) - } - - // Perform HTTP Post request to the notify URL - upJson, err := json.Marshal(upload) - if err != nil { - return false, fmt.Errorf("failed to marshal upload to JSON: %w", err) - } - - log.Infow("PDP notify", "upload", upload, "task_id", taskID) - - if upload.NotifyURL != "" { - - resp, err := t.client.Post(upload.NotifyURL, "application/json", bytes.NewReader(upJson)) - if err != nil { - log.Errorw("HTTP POST request to notify_url failed", "notify_url", upload.NotifyURL, "upload_id", upload.ID, "error", err) - } else { - defer func() { - _ = resp.Body.Close() - }() - // Not reading the body as per requirement - log.Infow("HTTP GET request to notify_url succeeded", "notify_url", upload.NotifyURL, "upload_id", upload.ID) - } - } - - comm, err := t.db.BeginTransaction(ctx, func(tx *harmonydb.Tx) (bool, error) { - // Move the entry from pdp_piece_uploads to pdp_piecerefs - // Insert into pdp_piecerefs - // Set needs_save_cache=TRUE for large pieces to enable proactive caching - needsSaveCache := padreader.PaddedSize(upload.PieceRawSize).Padded() >= abi.PaddedPieceSize(MinSizeForCache) - _, err = tx.Exec(` - INSERT INTO pdp_piecerefs (service, piece_cid, piece_ref, created_at, needs_save_cache) - VALUES ($1, $2, $3, NOW(), $4)`, - upload.Service, upload.PieceCID, upload.PieceRef, needsSaveCache) - if err != nil { - return false, fmt.Errorf("failed to insert into pdp_piecerefs: %w", err) - } - - _, err = tx.Exec(`DELETE FROM pdp_piece_uploads WHERE id = $1`, upload.ID) - if err != nil { - return false, fmt.Errorf("failed to delete upload ID %s from pdp_piece_uploads: %w", upload.ID, err) - } - - return true, nil - }, harmonydb.OptionRetry()) - if err != nil { - return false, fmt.Errorf("failed to move upload to piecerefs: %w", err) - } - if !comm { - return false, fmt.Errorf("transaction to move upload to piecerefs was not committed") - } - - log.Infof("Successfully processed PDP notify task %d for upload ID %s", taskID, upload.ID) - - return true, nil -} - -func (t *PDPNotifyTask) CanAccept(ids []harmonytask.TaskID, engine *harmonytask.TaskEngine) ([]harmonytask.TaskID, error) { - return ids, nil -} - -func (t *PDPNotifyTask) TypeDetails() harmonytask.TaskTypeDetails { - return harmonytask.TaskTypeDetails{ - Name: tasknames.PDPv0_Notify, - Cost: resources.Resources{ - Cpu: 0, - Ram: 128 << 20, // 128MB - }, - MaxFailures: 14, - RetryWait: taskhelp.RetryWaitExp(5*time.Second, 2), - } -} - -func (t *PDPNotifyTask) Adder(taskFunc harmonytask.AddTaskFunc) { - t.TF.Set(taskFunc) -} - -var _ = harmonytask.Reg(&PDPNotifyTask{}) -var _ harmonytask.TaskInterface = &PDPNotifyTask{} diff --git a/tasks/pdpv0/task_init_pp.go b/tasks/pdpv0/task_init_pp.go index ae9c3c341..de886da10 100644 --- a/tasks/pdpv0/task_init_pp.go +++ b/tasks/pdpv0/task_init_pp.go @@ -9,6 +9,7 @@ import ( "time" "github.com/ethereum/go-ethereum/core/types" + logging "github.com/ipfs/go-log/v2" "github.com/yugabyte/pgx/v5" "golang.org/x/xerrors" @@ -26,6 +27,8 @@ import ( chainTypes "github.com/filecoin-project/lotus/chain/types" ) +var log = logging.Logger("pdpv0") + const alertNameInitPP = "InitProvingPeriod" type InitProvingPeriodTask struct { @@ -280,7 +283,7 @@ func (ipp *InitProvingPeriodTask) CanAccept(ids []harmonytask.TaskID, engine *ha func (ipp *InitProvingPeriodTask) TypeDetails() harmonytask.TaskTypeDetails { return harmonytask.TaskTypeDetails{ Name: tasknames.PDPv0_InitPP, - // Handoff from data onboarding (PDPv0_Notify → PDPv0_PullPiece → PDPv0_SaveCache). + // Handoff from data onboarding (PDPv0_PullPiece → PDPv0_SaveCache). // InitPP checks on-chain leaf count before the first challenge request; proving // continues PDPv0_InitPP → PDPv0_Prove. MayFollow: []string{tasknames.PDPv0_SaveCache}, diff --git a/tasks/pdpv0/task_pull_piece.go b/tasks/pdpv0/task_pull_piece.go index 543a95186..762d70656 100644 --- a/tasks/pdpv0/task_pull_piece.go +++ b/tasks/pdpv0/task_pull_piece.go @@ -1347,9 +1347,8 @@ func (t *PDPPullPieceTask) CanAccept(ids []harmonytask.TaskID, engine *harmonyta func (t *PDPPullPieceTask) TypeDetails() harmonytask.TaskTypeDetails { return harmonytask.TaskTypeDetails{ - Name: tasknames.PDPv0_PullPiece, - Max: taskhelp.Max(t.max), - MayFollow: []string{tasknames.PDPv0_Notify}, + Name: tasknames.PDPv0_PullPiece, + Max: taskhelp.Max(t.max), Cost: resources.Resources{ Cpu: 0, Gpu: 0, diff --git a/tasks/tasknames/names.go b/tasks/tasknames/names.go index 5fb100622..950c41180 100644 --- a/tasks/tasknames/names.go +++ b/tasks/tasknames/names.go @@ -89,7 +89,6 @@ const ( PDPDelDataSet = "PDPDelDataSet" PDPInitPP = "PDPInitPP" PDPProvingPeriod = "PDPProvingPeriod" - PDPNotify = "PDPNotify" PDPCommP = "PDPCommP" PDPSaveCache = "PDPSaveCache" AggregatePDPDeal = "AggregatePDPDeal" @@ -101,7 +100,6 @@ const ( PDPv0_SaveCache = "PDPv0_SaveCache" PDPv0_InitPP = "PDPv0_InitPP" PDPv0_ProvPeriod = "PDPv0_ProvPeriod" - PDPv0_Notify = "PDPv0_Notify" PDPv0_DelDataSet = "PDPv0_DelDataSet" PDPv0_Cleanup = "PDPv0_Cleanup" PDPv0_PieceGC = "PDPv0_PieceGC"