diff --git a/server/handle_repo_delete_record.go b/server/handle_repo_delete_record.go new file mode 100644 --- /dev/null +++ b/server/handle_repo_delete_record.go @@ -0,0 +1,55 @@ +package server + +import ( + "github.com/haileyok/cocoon/internal/helpers" + "github.com/haileyok/cocoon/models" + "github.com/labstack/echo/v4" +) + +type ComAtprotoRepoDeleteRecordRequest struct { + Repo string `json:"repo" validate:"required,atproto-did"` + Collection string `json:"collection" validate:"required,atproto-nsid"` + Rkey string `json:"rkey" validate:"required,atproto-rkey"` + SwapRecord *string `json:"swapRecord"` + SwapCommit *string `json:"swapCommit"` +} + +func (s *Server) handleDeleteRecord(e echo.Context) error { + repo := e.Get("repo").(*models.RepoActor) + + var req ComAtprotoRepoDeleteRecordRequest + if err := e.Bind(&req); err != nil { + s.logger.Error("error binding", "error", err) + return helpers.ServerError(e, nil) + } + + if err := e.Validate(req); err != nil { + s.logger.Error("error validating", "error", err) + return helpers.InputError(e, nil) + } + + if repo.Repo.Did != req.Repo { + s.logger.Warn("mismatched repo/auth") + return helpers.InputError(e, nil) + } + + results, err := s.repoman.applyWrites(repo.Repo, []Op{ + { + Type: OpTypeDelete, + Collection: req.Collection, + Rkey: &req.Rkey, + SwapRecord: req.SwapRecord, + }, + }, req.SwapCommit) + if err != nil { + s.logger.Error("error applying writes", "error", err) + return helpers.ServerError(e, nil) + } + + results[0].Type = nil + results[0].Uri = nil + results[0].Cid = nil + results[0].ValidationStatus = nil + + return e.JSON(200, results[0]) +} diff --git a/server/repo.go b/server/repo.go --- a/server/repo.go +++ b/server/repo.go @@ -84,10 +84,10 @@ } type ApplyWriteResult struct { Type *string `json:"$type,omitempty"` - Uri string `json:"uri"` - Cid string `json:"cid"` + Uri *string `json:"uri,omitempty"` + Cid *string `json:"cid,omitempty"` Commit *RepoCommit `json:"commit,omitempty"` - ValidationStatus *string `json:"validationStatus"` + ValidationStatus *string `json:"validationStatus,omitempty"` } type RepoCommit struct { @@ -139,15 +139,28 @@ Value: d, }) results = append(results, ApplyWriteResult{ Type: to.StringPtr(OpTypeCreate.String()), - Uri: "at://" + urepo.Did + "/" + op.Collection + "/" + *op.Rkey, - Cid: nc.String(), + Uri: to.StringPtr("at://" + urepo.Did + "/" + op.Collection + "/" + *op.Rkey), + Cid: to.StringPtr(nc.String()), ValidationStatus: to.StringPtr("valid"), // TODO: obviously this might not be true atm lol }) case OpTypeDelete: + var old models.Record + if err := rm.db.Raw("SELECT value FROM records WHERE did = ? AND nsid = ? AND rkey = ?", urepo.Did, op.Collection, op.Rkey).Scan(&old).Error; err != nil { + return nil, err + } + entries = append(entries, models.Record{ + Did: urepo.Did, + Nsid: op.Collection, + Rkey: *op.Rkey, + Value: old.Value, + }) err := r.DeleteRecord(context.TODO(), op.Collection+"/"+*op.Rkey) if err != nil { return nil, err } + results = append(results, ApplyWriteResult{ + Type: to.StringPtr(OpTypeDelete.String()), + }) case OpTypeUpdate: nc, err := r.UpdateRecord(context.TODO(), op.Collection+"/"+*op.Rkey, op.Record) if err != nil { @@ -165,8 +178,8 @@ Value: d, }) results = append(results, ApplyWriteResult{ Type: to.StringPtr(OpTypeUpdate.String()), - Uri: "at://" + urepo.Did + "/" + op.Collection + "/" + *op.Rkey, - Cid: nc.String(), + Uri: to.StringPtr("at://" + urepo.Did + "/" + op.Collection + "/" + *op.Rkey), + Cid: to.StringPtr(nc.String()), ValidationStatus: to.StringPtr("valid"), // TODO: obviously this might not be true atm lol }) } @@ -211,10 +224,12 @@ Cid: &ll, }) case "del": + ll := lexutil.LexLink(op.OldCid) ops = append(ops, &atproto.SyncSubscribeRepos_RepoOp{ Action: "delete", Path: op.Rpath, Cid: nil, + Prev: &ll, }) } @@ -236,17 +251,27 @@ } var blobs []lexutil.LexLink for _, entry := range entries { - if err := rm.s.db.Clauses(clause.OnConflict{ - Columns: []clause.Column{{Name: "did"}, {Name: "nsid"}, {Name: "rkey"}}, - UpdateAll: true, - }).Create(&entry).Error; err != nil { - return nil, err - } + var cids []cid.Cid + if entry.Cid != "" { + if err := rm.s.db.Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "did"}, {Name: "nsid"}, {Name: "rkey"}}, + UpdateAll: true, + }).Create(&entry).Error; err != nil { + return nil, err + } - // we should actually check the type (i.e. delete, create,., update) here but we'll do it later - cids, err := rm.incrementBlobRefs(urepo, entry.Value) - if err != nil { - return nil, err + cids, err = rm.incrementBlobRefs(urepo, entry.Value) + if err != nil { + return nil, err + } + } else { + if err := rm.s.db.Delete(&entry).Error; err != nil { + return nil, err + } + cids, err = rm.decrementBlobRefs(urepo, entry.Value) + if err != nil { + return nil, err + } } for _, c := range cids { @@ -314,6 +339,34 @@ for _, c := range cids { if err := rm.db.Exec("UPDATE blobs SET ref_count = ref_count + 1 WHERE did = ? AND cid = ?", urepo.Did, c.Bytes()).Error; err != nil { return nil, err + } + } + + return cids, nil +} + +func (rm *RepoMan) decrementBlobRefs(urepo models.Repo, cbor []byte) ([]cid.Cid, error) { + cids, err := getBlobCidsFromCbor(cbor) + if err != nil { + return nil, err + } + + for _, c := range cids { + var res struct { + ID uint + Count int + } + if err := rm.db.Raw("UPDATE blobs SET ref_count = ref_count - 1 WHERE did = ? AND cid = ? RETURNING id, ref_count", urepo.Did, c.Bytes()).Scan(&res).Error; err != nil { + return nil, err + } + + if res.Count == 0 { + if err := rm.db.Exec("DELETE FROM blobs WHERE id = ?", res.ID).Error; err != nil { + return nil, err + } + if err := rm.db.Exec("DELETE FROM blob_parts WHERE blob_id = ?", res.ID).Error; err != nil { + return nil, err + } } } diff --git a/server/server.go b/server/server.go --- a/server/server.go +++ b/server/server.go @@ -391,6 +391,7 @@ // repo s.echo.POST("/xrpc/com.atproto.repo.createRecord", s.handleCreateRecord, s.handleSessionMiddleware) s.echo.POST("/xrpc/com.atproto.repo.putRecord", s.handlePutRecord, s.handleSessionMiddleware) + s.echo.POST("/xrpc/com.atproto.repo.deleteRecord", s.handleDeleteRecord, s.handleSessionMiddleware) s.echo.POST("/xrpc/com.atproto.repo.applyWrites", s.handleApplyWrites, s.handleSessionMiddleware) s.echo.POST("/xrpc/com.atproto.repo.uploadBlob", s.handleRepoUploadBlob, s.handleSessionMiddleware)