From 4fce2f5d9cd545a93d3fd91a0842c3c41a138fe3 Mon Sep 17 00:00:00 2001 From: Lewis Date: Tue, 02 Jun 2026 11:09:18 +0000 Subject: [PATCH] knotserver/xrpc: add/remove knot members directly Lewis: May this revision serve well! --- knotserver/db/known_dids.go | 6 ++++++ knotserver/xrpc/acl_saga.go | 195 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotserver/xrpc/add_member.go | 98 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotserver/xrpc/remove_member.go | 92 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotserver/xrpc/xrpc.go | 2 ++ 5 file(s) changed, 393 insertion(s)(+), 0 deletion(s)(-) diff --git a/knotserver/db/known_dids.go b/knotserver/db/known_dids.go --- a/knotserver/db/known_dids.go +++ b/knotserver/db/known_dids.go @@ -5,6 +5,12 @@ return err } +func IsDidKnown(q DBTX, did string) (bool, error) { + var exists bool + err := q.QueryRow(`select exists (select 1 from known_dids where did = ?)`, did).Scan(&exists) + return exists, err +} + func RemoveDid(q DBTX, did string) error { _, err := q.Exec(`delete from known_dids where did = ?`, did) return err diff --git a/knotserver/xrpc/acl_saga.go b/knotserver/xrpc/acl_saga.go new file mode 100644 --- /dev/null +++ b/knotserver/xrpc/acl_saga.go @@ -0,0 +1,195 @@ +package xrpc + +import ( + "context" + "database/sql" + "log/slog" + "net/http" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/knotserver/db" + "tangled.org/core/knotserver/keys" + xrpcerr "tangled.org/core/xrpc/errors" +) + +type aclGrant struct { + role string + subject syntax.DID + inAcl func() (bool, error) + inTable func() (bool, error) + insertRow func(*sql.Tx) error + deleteRow func() error + grantAcl func() error + emit func() error +} + +type aclRevoke struct { + role string + subject syntax.DID + inAcl func() (bool, error) + inTable func() (bool, error) + removeAcl func() (bool, error) + restoreAcl func() error + deleteRow func(*sql.Tx) error + emit func() error +} + +func (h *Xrpc) applyAclGrant(ctx context.Context, l *slog.Logger, g aclGrant) (int, *xrpcerr.XrpcError) { + fail := func(status int, e xrpcerr.XrpcError) (int, *xrpcerr.XrpcError) { + return status, &e + } + + inAcl, err := g.inAcl() + if err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + inTable, err := g.inTable() + if err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + if inAcl && inTable { + l.Info("subject already granted, no-op", "role", g.role, "subject", g.subject) + return http.StatusOK, nil + } + + didKnown, err := db.IsDidKnown(h.Db, g.subject.String()) + if err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + + tx, err := h.Db.BeginTx(ctx, nil) + if err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + committed := false + defer func() { + if !committed { + tx.Rollback() + } + }() + + if err := db.AddDid(tx, g.subject.String()); err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + if err := g.insertRow(tx); err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + if err := tx.Commit(); err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + committed = true + + if err := g.grantAcl(); err != nil { + if !inTable { + if rbErr := g.deleteRow(); rbErr != nil { + l.Error("failed to roll back row after ACL grant failed", "role", g.role, "subject", g.subject, "error", rbErr) + } + } + if !didKnown { + if rbErr := db.RemoveDid(h.Db, g.subject.String()); rbErr != nil { + l.Error("failed to roll back known_did after ACL grant failed", "role", g.role, "subject", g.subject, "error", rbErr) + } + } + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + + h.Ingester.AddDid(g.subject.String()) + h.fetchKeysAsync(ctx, l, g.subject) + + if g.emit != nil { + if err := g.emit(); err != nil { + l.Error("failed to emit acl grant event, appview reconcile will catch up", "role", g.role, "subject", g.subject, "error", err) + } + } + + l.Info("granted", "role", g.role, "subject", g.subject) + return http.StatusOK, nil +} + +func (h *Xrpc) applyAclRevoke(ctx context.Context, l *slog.Logger, rv aclRevoke) (int, *xrpcerr.XrpcError) { + fail := func(status int, e xrpcerr.XrpcError) (int, *xrpcerr.XrpcError) { + return status, &e + } + + inAcl, err := rv.inAcl() + if err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + inTable, err := rv.inTable() + if err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + if !inAcl && !inTable { + l.Info("subject not granted, no-op", "role", rv.role, "subject", rv.subject) + return http.StatusOK, nil + } + + removed := false + if inAcl { + removed, err = rv.removeAcl() + if err != nil { + return fail(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + } + + failRestoringACL := func(status int, e xrpcerr.XrpcError) (int, *xrpcerr.XrpcError) { + if removed { + if rbErr := rv.restoreAcl(); rbErr != nil { + l.Error("failed to restore ACL after remove rollback", "role", rv.role, "subject", rv.subject, "error", rbErr) + } + } + return fail(status, e) + } + + stillKnown, err := h.Enforcer.HasAnyPolicyForUser(rv.subject.String()) + if err != nil { + return failRestoringACL(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + + tx, err := h.Db.BeginTx(ctx, nil) + if err != nil { + return failRestoringACL(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + committed := false + defer func() { + if !committed { + tx.Rollback() + } + }() + + if err := rv.deleteRow(tx); err != nil { + return failRestoringACL(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + if !stillKnown { + if err := db.RemoveDid(tx, rv.subject.String()); err != nil { + return failRestoringACL(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + } + if err := tx.Commit(); err != nil { + return failRestoringACL(http.StatusInternalServerError, xrpcerr.GenericError(err)) + } + committed = true + + if !stillKnown { + h.Ingester.RemoveDid(rv.subject.String()) + } + + if rv.emit != nil { + if err := rv.emit(); err != nil { + l.Error("failed to emit acl revoke event, appview reconcile will catch up", "role", rv.role, "subject", rv.subject, "error", err) + } + } + + l.Info("revoked", "role", rv.role, "subject", rv.subject, "did_dropped", !stillKnown) + return http.StatusOK, nil +} + +func (h *Xrpc) fetchKeysAsync(ctx context.Context, l *slog.Logger, subject syntax.DID) { + kctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), keyFetchTimeout) + go func() { + defer cancel() + if err := keys.FetchAndStore(kctx, h.Resolver.Directory(), h.Db, subject.String()); err != nil { + l.Warn("failed to fetch subject public keys, continuing", "subject", subject, "error", err) + } + }() +} diff --git a/knotserver/xrpc/add_member.go b/knotserver/xrpc/add_member.go new file mode 100644 --- /dev/null +++ b/knotserver/xrpc/add_member.go @@ -0,0 +1,98 @@ +package xrpc + +import ( + "context" + "database/sql" + "encoding/json" + "log/slog" + "net/http" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/knotserver/db" + "tangled.org/core/rbac" + xrpcerr "tangled.org/core/xrpc/errors" +) + +const keyFetchTimeout = 15 * time.Second + +func (h *Xrpc) AddMember(w http.ResponseWriter, r *http.Request) { + l := h.Logger.With("handler", "AddMember") + fail := func(e xrpcerr.XrpcError, status int) { + l.Error("failed", "kind", e.Tag, "error", e.Message) + writeError(w, e, status) + } + + actorDid, ok := r.Context().Value(ActorDid).(syntax.DID) + if !ok { + fail(xrpcerr.MissingActorDidError, http.StatusForbidden) + return + } + + allowed, err := h.Enforcer.IsKnotInviteAllowed(actorDid.String(), rbac.ThisServer) + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + if !allowed { + fail(xrpcerr.AccessControlError(actorDid.String()), http.StatusForbidden) + return + } + + var data tangled.KnotAddMember_Input + if err := json.NewDecoder(r.Body).Decode(&data); err != nil { + fail(xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + subject, err := syntax.ParseDID(data.Subject) + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + status, xerr := h.addMemberToKnot(r.Context(), l, actorDid, subject) + if xerr != nil { + fail(*xerr, status) + return + } + w.WriteHeader(status) +} + +func (h *Xrpc) addMemberToKnot(ctx context.Context, l *slog.Logger, addedBy, subject syntax.DID) (int, *xrpcerr.XrpcError) { + isOwner, err := h.Enforcer.IsKnotOwner(subject.String(), rbac.ThisServer) + if err != nil { + e := xrpcerr.GenericError(err) + return http.StatusInternalServerError, &e + } + if isOwner { + l.Info("subject is the knot owner, no-op", "subject", subject) + return http.StatusOK, nil + } + + return h.applyAclGrant(ctx, l, aclGrant{ + role: "member", + subject: subject, + inAcl: func() (bool, error) { + return h.Enforcer.IsKnotMember(subject.String(), rbac.ThisServer) + }, + inTable: func() (bool, error) { + n, err := db.CountKnotMembersBySubject(h.Db, subject.String()) + return n > 0, err + }, + insertRow: func(tx *sql.Tx) error { + return db.AddKnotMemberDirect(tx, addedBy, subject) + }, + deleteRow: func() error { + return db.RemoveKnotMemberDirect(h.Db, subject) + }, + grantAcl: func() error { + _, err := h.Enforcer.TryAddKnotMember(rbac.ThisServer, subject.String()) + return err + }, + emit: func() error { + return h.Db.EmitKnotMemberUpdate(h.Notifier, db.AclOpAdd, subject) + }, + }) +} diff --git a/knotserver/xrpc/remove_member.go b/knotserver/xrpc/remove_member.go new file mode 100644 --- /dev/null +++ b/knotserver/xrpc/remove_member.go @@ -0,0 +1,92 @@ +package xrpc + +import ( + "database/sql" + "encoding/json" + "net/http" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/knotserver/db" + "tangled.org/core/rbac" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func (h *Xrpc) RemoveMember(w http.ResponseWriter, r *http.Request) { + l := h.Logger.With("handler", "RemoveMember") + fail := func(e xrpcerr.XrpcError, status int) { + l.Error("failed", "kind", e.Tag, "error", e.Message) + writeError(w, e, status) + } + + actorDid, ok := r.Context().Value(ActorDid).(syntax.DID) + if !ok { + fail(xrpcerr.MissingActorDidError, http.StatusForbidden) + return + } + + allowed, err := h.Enforcer.IsKnotInviteAllowed(actorDid.String(), rbac.ThisServer) + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + if !allowed { + fail(xrpcerr.AccessControlError(actorDid.String()), http.StatusForbidden) + return + } + + var data tangled.KnotRemoveMember_Input + if err := json.NewDecoder(r.Body).Decode(&data); err != nil { + fail(xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + subject, err := syntax.ParseDID(data.Subject) + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusBadRequest) + return + } + + isOwner, err := h.Enforcer.IsKnotOwner(subject.String(), rbac.ThisServer) + if err != nil { + fail(xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + if isOwner { + fail(xrpcerr.NewXrpcError( + xrpcerr.WithTag("InvalidRequest"), + xrpcerr.WithMessage("cannot remove the knot owner"), + ), http.StatusBadRequest) + return + } + + status, xerr := h.applyAclRevoke(r.Context(), l, aclRevoke{ + role: "member", + subject: subject, + inAcl: func() (bool, error) { + return h.Enforcer.IsKnotMember(subject.String(), rbac.ThisServer) + }, + inTable: func() (bool, error) { + n, err := db.CountKnotMembersBySubject(h.Db, subject.String()) + return n > 0, err + }, + removeAcl: func() (bool, error) { + return h.Enforcer.TryRemoveKnotMember(rbac.ThisServer, subject.String()) + }, + restoreAcl: func() error { + _, err := h.Enforcer.TryAddKnotMember(rbac.ThisServer, subject.String()) + return err + }, + deleteRow: func(tx *sql.Tx) error { + return db.RemoveKnotMemberBySubject(tx, subject) + }, + emit: func() error { + return h.Db.EmitKnotMemberUpdate(h.Notifier, db.AclOpRemove, subject) + }, + }) + if xerr != nil { + fail(*xerr, status) + return + } + w.WriteHeader(status) +} diff --git a/knotserver/xrpc/xrpc.go b/knotserver/xrpc/xrpc.go --- a/knotserver/xrpc/xrpc.go +++ b/knotserver/xrpc/xrpc.go @@ -56,6 +56,8 @@ r.Post("/"+tangled.RepoForkSyncNSID, x.ForkSync) r.Post("/"+tangled.RepoHiddenRefNSID, x.HiddenRef) r.Post("/"+tangled.RepoMergeNSID, x.Merge) + r.Post("/"+tangled.KnotAddMemberNSID, x.AddMember) + r.Post("/"+tangled.KnotRemoveMemberNSID, x.RemoveMember) }) // merge check is an open endpoint -- tangled.sh