From 23d1f195396e98fcb61e134cf6f618977ee7b723 Mon Sep 17 00:00:00 2001 From: Alan Shaw Date: Tue, 24 Feb 2026 15:25:32 +0000 Subject: [PATCH] wip: adding repair cmd --- cmd/blob/repair.go | 152 +++++++++++++++++++++++++++++++++++++++++++ cmd/blob/root.go | 1 + pkg/client/client.go | 5 ++ 3 files changed, 158 insertions(+) create mode 100644 cmd/blob/repair.go diff --git a/cmd/blob/repair.go b/cmd/blob/repair.go new file mode 100644 index 00000000..abb1645d --- /dev/null +++ b/cmd/blob/repair.go @@ -0,0 +1,152 @@ +package blob + +import ( + "errors" + "fmt" + "math/rand/v2" + + "github.com/multiformats/go-multihash" + "github.com/spf13/cobra" + "github.com/storacha/go-libstoracha/capabilities/assert" + spaceblobcap "github.com/storacha/go-libstoracha/capabilities/space/blob" + "github.com/storacha/go-libstoracha/digestutil" + "github.com/storacha/go-ucanto/core/dag/blockstore" + "github.com/storacha/go-ucanto/core/delegation" + "github.com/storacha/go-ucanto/core/result" + fdm "github.com/storacha/go-ucanto/core/result/failure/datamodel" + "github.com/storacha/go-ucanto/did" + "github.com/storacha/go-ucanto/validator" + "github.com/storacha/guppy/internal/cmdutil" + "github.com/storacha/guppy/pkg/config" + "github.com/storacha/guppy/pkg/receipt" + indexer_types "github.com/storacha/indexing-service/pkg/types" +) + +const minReplicas = 3 + +var repairFlags struct { + ignore []string +} + +func init() { + repairCmd.Flags().StringSliceVarP(&repairFlags.ignore, "ignore", "i", []string{}, "DIDs of storage providers to ignore in results.") +} + +var repairCmd = &cobra.Command{ + Use: "repair ", + Short: "Repair blobs in a space with insufficient replicas", + Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + cfg, err := config.Load[config.Config]() + if err != nil { + return err + } + c := cmdutil.MustGetClient(cfg.Repo.Dir) + + indexer, _ := cmdutil.MustGetIndexClient() + + space, err := cmdutil.ResolveSpace(c, args[0]) + if err != nil { + return err + } + + ignore := map[did.DID]struct{}{} + for _, didStr := range repairFlags.ignore { + did, err := did.Parse(didStr) + cobra.CheckErr(err) + ignore[did] = struct{}{} + } + + var i int + var cursor *string + size := uint64(pageSize) + for { + listOk, err := c.SpaceBlobList( + cmd.Context(), + space, + spaceblobcap.ListCaveats{Cursor: cursor, Size: &size}) + if err != nil { + return err + } + + for _, r := range listOk.Results { + i++ + queryResult, err := indexer.QueryClaims(cmd.Context(), indexer_types.Query{ + Type: indexer_types.QueryTypeLocation, + Hashes: []multihash.Multihash{r.Blob.Digest}, + }) + cobra.CheckErr(err) + + bs, err := blockstore.NewBlockReader(blockstore.WithBlocksIterator(queryResult.Blocks())) + cobra.CheckErr(err) + + var commitments []delegation.Delegation + for _, claimID := range queryResult.Claims() { + claim, err := delegation.NewDelegationView(claimID, bs) + cobra.CheckErr(err) + _, err = assert.Location.Match(validator.NewSource(claim.Capabilities()[0], claim)) + if err != nil { + continue + } + commitments = append(commitments, claim) + } + + var filteredCommitments []delegation.Delegation + for _, commL := range commitments { + if _, ok := ignore[commL.Issuer().DID()]; !ok { + filteredCommitments = append(filteredCommitments, commL) + } + } + + needRepair := len(commitments) < minReplicas || len(filteredCommitments) < minReplicas + replicas := minReplicas + + if len(filteredCommitments) < minReplicas { + replicas = minReplicas + (len(commitments) - len(filteredCommitments)) + } + + if !needRepair { + cmd.Println(fmt.Sprintf("%d: %s does not need repair", i, digestutil.Format(r.Blob.Digest))) + continue + } + + cmd.Println(fmt.Sprintf("%d: %s", i, digestutil.Format(r.Blob.Digest))) + res, _, err := c.SpaceBlobReplicate( + cmd.Context(), + space, + r.Blob, + uint(replicas), + commitments[rand.IntN(len(commitments))], + ) + if err != nil { + cmd.PrintErrln(fmt.Sprintf(" - error invoking replication: %s", err)) + continue + } + for _, site := range res.Site { + cmd.PrintErrln(fmt.Sprintf(" - transfer task: %s", site.UcanAwait.Link)) + + r, err := c.Receipts().Fetch(cmd.Context(), site.UcanAwait.Link) + if err != nil { + if errors.Is(err, receipt.ErrNotFound) { + continue + } + cobra.CheckErr(err) + } + + _, x := result.Unwrap(r.Out()) + if x != nil { + failure := fdm.Bind(x) + cmd.PrintErrln(fmt.Sprintf(" - transfer failed: %s", failure.Message)) + } + } + } + + if listOk.Cursor == nil { + break + } + cursor = listOk.Cursor + } + + return nil + }, +} diff --git a/cmd/blob/root.go b/cmd/blob/root.go index bf6c2d3e..d37ff937 100644 --- a/cmd/blob/root.go +++ b/cmd/blob/root.go @@ -12,5 +12,6 @@ var Cmd = &cobra.Command{ func init() { Cmd.AddCommand( lsCmd, + repairCmd, ) } diff --git a/pkg/client/client.go b/pkg/client/client.go index 08378ea7..dee70dbf 100644 --- a/pkg/client/client.go +++ b/pkg/client/client.go @@ -153,6 +153,11 @@ func (c *Client) AddProofs(delegations ...delegation.Delegation) error { return c.store.AddDelegations(delegations...) } +// Receipts returns a client for fetching invocation receipts. +func (c *Client) Receipts() *receiptclient.Client { + return c.receiptsClient +} + // Reset clears all delegations from the store while preserving the principal. func (c *Client) Reset() error { return c.store.Reset()