diff --git a/accountlimit/accountlimit.go b/accountlimit/accountlimit.go index 76927ef..ea799fe 100644 --- a/accountlimit/accountlimit.go +++ b/accountlimit/accountlimit.go @@ -42,6 +42,11 @@ type SpaceLimits struct { // FileStorageLimitBytes is the default account file-storage pool for // identities without an explicit accountLimit document (files v2). FileStorageLimitBytes uint64 `yaml:"fileStorageLimitBytes" bson:"fileStorageLimitBytes"` + // ExternalSeatsLimit is the default external-seat pool (nested spaces) for + // identities without an explicit accountLimit document. Production networks + // keep it 0 (seats are granted per identity by the payment node); self-hosted + // networks set it to open external admissions up. + ExternalSeatsLimit uint32 `yaml:"externalSeatsLimit" bson:"externalSeatsLimit"` } type Limits struct { @@ -51,7 +56,10 @@ type Limits struct { SpaceMembersRead uint32 `bson:"spaceMembersRead"` SpaceMembersWrite uint32 `bson:"spaceMembersWrite"` SharedSpacesLimit uint32 `bson:"sharedSpacesLimit"` - UpdatedTime time.Time `bson:"updatedTime"` + // ExternalSeatsLimit caps the distinct non-org identities admitted across all + // child spaces of the identity's spaces (nested spaces / paid external seats) + ExternalSeatsLimit uint32 `bson:"externalSeatsLimit"` + UpdatedTime time.Time `bson:"updatedTime"` } type AccountLimit interface { @@ -125,12 +133,13 @@ func (al *accountLimit) GetLimits(ctx context.Context, identity string) (limits } // default limit return Limits{ - Identity: identity, - SpaceMembersRead: al.defaultLimits.SpaceMembersRead, - SpaceMembersWrite: al.defaultLimits.SpaceMembersWrite, - SharedSpacesLimit: al.defaultLimits.SharedSpacesLimit, - FileStorageBytes: al.defaultLimits.FileStorageLimitBytes, - UpdatedTime: time.Now(), + Identity: identity, + SpaceMembersRead: al.defaultLimits.SpaceMembersRead, + SpaceMembersWrite: al.defaultLimits.SpaceMembersWrite, + SharedSpacesLimit: al.defaultLimits.SharedSpacesLimit, + FileStorageBytes: al.defaultLimits.FileStorageLimitBytes, + ExternalSeatsLimit: al.defaultLimits.ExternalSeatsLimit, + UpdatedTime: time.Now(), }, nil } diff --git a/coordinator/coordinator.go b/coordinator/coordinator.go index 692ee12..02718e1 100644 --- a/coordinator/coordinator.go +++ b/coordinator/coordinator.go @@ -4,6 +4,7 @@ import ( "context" "errors" "slices" + "sync" "time" "github.com/anyproto/any-sync/accountservice" @@ -11,6 +12,7 @@ import ( "github.com/anyproto/any-sync/app" "github.com/anyproto/any-sync/app/logger" "github.com/anyproto/any-sync/commonspace/object/accountdata" + "github.com/anyproto/any-sync/commonspace/object/acl/aclrecordproto" "github.com/anyproto/any-sync/commonspace/object/acl/list" "github.com/anyproto/any-sync/commonspace/spacepayloads" "github.com/anyproto/any-sync/commonspace/spacesyncproto" @@ -21,14 +23,18 @@ import ( "github.com/anyproto/any-sync/net/pool" "github.com/anyproto/any-sync/net/rpc/server" "github.com/anyproto/any-sync/nodeconf" + "github.com/anyproto/any-sync/util/cidutil" "github.com/anyproto/any-sync/util/crypto" "go.uber.org/zap" "storj.io/drpc" + "go.mongodb.org/mongo-driver/mongo" + "github.com/anyproto/any-sync-coordinator/accountlimit" "github.com/anyproto/any-sync-coordinator/acleventlog" "github.com/anyproto/any-sync-coordinator/config" "github.com/anyproto/any-sync-coordinator/coordinatorlog" + "github.com/anyproto/any-sync-coordinator/db" "github.com/anyproto/any-sync-coordinator/deletionlog" "github.com/anyproto/any-sync-coordinator/fileusage" "github.com/anyproto/any-sync-coordinator/inbox" @@ -70,6 +76,9 @@ type coordinator struct { accountLimit accountlimit.AccountLimit fileUsage fileusage.FileUsage acl acl.AclService + externalSeats *mongo.Collection + orgSeatLocks [orgSeatLockStripes]sync.Mutex + orgSeatLease *db.Lease inbox inbox.InboxService subscribe subscribe.SubscribeService drpcHandler *rpcHandler @@ -88,6 +97,11 @@ func (c *coordinator) Init(a *app.App) (err error) { c.subscribe = a.MustComponent(subscribe.CName).(subscribe.SubscribeService) c.deletionLog = app.MustComponent[deletionlog.DeletionLog](a) c.acl = app.MustComponent[acl.AclService](a) + // optional so unit fixtures without mongo still assemble; production registers db + if dbComp, _ := a.Component(db.CName).(db.Database); dbComp != nil { + c.externalSeats = dbComp.Db().Collection(externalSeatsCollName) + c.orgSeatLease = db.NewLease(dbComp.Db().Collection(orgSeatLeaseCollName), orgSeatLeaseTTL, orgSeatLeasePoll) + } c.inbox = app.MustComponent[inbox.InboxService](a) c.accountLimit = app.MustComponent[accountlimit.AccountLimit](a) c.fileUsage = app.MustComponent[fileusage.FileUsage](a) @@ -153,11 +167,28 @@ func (c *coordinator) SpaceDelete(ctx context.Context, spaceId string, deletionD if err != nil { return } + var authorizedByLegalOwner bool + entry, err := c.spaceStatus.Status(ctx, spaceId) + if err != nil { + return + } + if entry.ParentSpaceId != "" && entry.Identity != accountPubKey.Account() { + // legalOwner delete: the parent's CURRENT owner may delete a child it cannot read + parentOwner, ownerErr := c.acl.OwnerPubKey(ctx, entry.ParentSpaceId) + if ownerErr != nil { + return 0, ownerErr + } + if !parentOwner.Equals(accountPubKey) { + return 0, coordinatorproto.ErrForbidden + } + authorizedByLegalOwner = true + } return c.spaceStatus.SpaceDelete(ctx, spacestatus.SpaceDeletion{ - DeletionPayload: payload, - DeletionPayloadId: payloadId, - SpaceId: spaceId, - DeletionPeriod: time.Duration(deletionDurationSecs) * time.Minute, + DeletionPayload: payload, + DeletionPayloadId: payloadId, + SpaceId: spaceId, + DeletionPeriod: time.Duration(deletionDurationSecs) * time.Minute, + AuthorizedByLegalOwner: authorizedByLegalOwner, AccountInfo: spacestatus.AccountInfo{ Identity: accountPubKey, PeerId: peerId, @@ -195,7 +226,7 @@ func (c *coordinator) StatusChange(ctx context.Context, spaceId string, deletion }) } -func (c *coordinator) SpaceSign(ctx context.Context, spaceId string, spaceHeader []byte, force bool) (signedReceipt *coordinatorproto.SpaceReceiptWithSignature, err error) { +func (c *coordinator) SpaceSign(ctx context.Context, spaceId string, spaceHeader []byte, force bool, parentAclRecordId string) (signedReceipt *coordinatorproto.SpaceReceiptWithSignature, err error) { // TODO: Think about how to make it more evident that account.SignKey is actually a network key // on a coordinator level networkKey := c.account.SignKey @@ -215,10 +246,53 @@ func (c *coordinator) SpaceSign(ctx context.Context, spaceId string, spaceHeader if err != nil { return } - err = c.spaceStatus.NewStatus(ctx, spaceId, accountPubKey, spaceType, force) + header, err := unmarshalSpaceHeader(spaceHeader) if err != nil { return } + if header.ParentSpaceId != "" { + if spaceType != spacestatus.SpaceTypeRegular { + // allowlist: a child must be a regular space. tech/1-1 would be permanently + // undeletable (SpaceDelete refuses those types); personal (timestamp==0) is an + // unintended state. Chat headers already map to SpaceTypeRegular. + return nil, coordinatorproto.ErrForbidden + } + if header.Version != spacesyncproto.SpaceHeaderVersion_SpaceHeaderVersion1 { + // nested spaces require the V1 header so the acl root is bound into the signed + // header — otherwise the coordinator could sign against one acl root while a + // different one is pushed to nodes (defense-in-depth; nodes also enforce this) + return nil, coordinatorproto.ErrForbidden + } + // the pinned legalOwner is the parent owner at the child's genesis; it must equal the + // parent's CURRENT owner only when the child is first created (anchoring the induction + // chain). On re-sign (receipt renewal) the parent may have transferred ownership since, + // so we authorize/bill against the current owner instead of failing forever. + _, statusErr := c.spaceStatus.Status(ctx, spaceId) + isNew := statusErr == coordinatorproto.ErrSpaceNotExists + if statusErr != nil && !isNew { + return nil, statusErr + } + legalOwner, err := c.verifyNestedSpace(ctx, spaceId, header, parentAclRecordId, isNew) + if err != nil { + return nil, err + } + limits, err := c.accountLimit.GetLimits(ctx, legalOwner.Account()) + if err != nil { + return nil, err + } + err = c.spaceStatus.NewChildStatus(ctx, spaceId, accountPubKey, header.ParentSpaceId, legalOwner.Account(), limits.SharedSpacesLimit, spaceType, force) + if err != nil { + return nil, err + } + } else { + if parentAclRecordId != "" { + return nil, coordinatorproto.ErrForbidden + } + err = c.spaceStatus.NewStatus(ctx, spaceId, accountPubKey, spaceType, force) + if err != nil { + return + } + } signedReceipt, err = coordinatorproto.PrepareSpaceReceipt(spaceId, peerId, spaceReceiptValidPeriod, accountPubKey, networkKey) if err != nil { return @@ -272,7 +346,8 @@ func (c *coordinator) AccountLimitsSet(ctx context.Context, req *coordinatorprot FileStorageBytes: req.FileStorageLimitBytes, SpaceMembersRead: req.SpaceMembersRead, SpaceMembersWrite: req.SpaceMembersWrite, - SharedSpacesLimit: req.SharedSpacesLimit, + SharedSpacesLimit: req.SharedSpacesLimit, + ExternalSeatsLimit: req.ExternalSeatsLimit, }) } @@ -303,6 +378,35 @@ func (c *coordinator) AclAddRecord(ctx context.Context, spaceId string, payload return } + if err = c.verifyKeylessGovernanceRecord(ctx, rec, statusEntry); err != nil { + return + } + + var admittedExternals []crypto.PubKey + if statusEntry.ParentSpaceId != "" { + if admitted := admittedIdentities(rec); len(admitted) > 0 { + // held until this record is committed (or rejected): the seat count is read + // live from the org's child acls, so check→commit must not interleave with a + // concurrent admission into a sibling child. The striped mutex serializes + // this instance; the mongo lease extends the same window across coordinator + // replicas (sibling children are different consensus logs, so nothing else + // orders two admissions against each other). + defer c.lockOrgSeats(statusEntry.ParentSpaceId)() + if c.orgSeatLease != nil { + release, lerr := c.orgSeatLease.Acquire(ctx, statusEntry.ParentSpaceId) + if lerr != nil { + err = lerr + return + } + defer release() + } + admittedExternals, err = c.verifyExternalSeats(ctx, statusEntry, admitted) + if err != nil { + return + } + } + } + limits, err := c.accountLimit.GetLimitsBySpace(ctx, spaceId) if err != nil { return nil, err @@ -321,6 +425,17 @@ func (c *coordinator) AclAddRecord(ctx context.Context, spaceId string, payload log.Debug("ACL change ID:", zap.String("rawRecordWithId.Id", rawRecordWithId.Id)) + if statusEntry.ParentSpaceId != "" { + if len(admittedExternals) > 0 { + if orgOwner, ownerErr := c.acl.OwnerPubKey(ctx, statusEntry.ParentSpaceId); ownerErr == nil { + c.recordExternalSeats(ctx, spaceId, orgOwner.Account(), admittedExternals) + } + } + if removed := removedIdentities(rec); len(removed) > 0 { + c.dropExternalSeats(ctx, spaceId, removed) + } + } + err = c.aclEventLog.AddLog(ctx, acleventlog.AclEventLogEntry{ SpaceId: spaceId, PeerId: peerId, @@ -377,12 +492,15 @@ func (c *coordinator) MakeSpaceShareable(ctx context.Context, spaceId string) (e if err != nil { return } - if statusEntry.Identity != pubKey.Account() { - return coordinatorproto.ErrForbidden - } + // already-shareable is a no-op for ANY caller: a non-owner member legitimately + // ensures shareability before writing to the acl (e.g. an org admin registering + // a child space) — only the actual flip is owner-gated if statusEntry.IsShareable { return nil } + if statusEntry.Identity != pubKey.Account() { + return coordinatorproto.ErrForbidden + } limits, err := c.accountLimit.GetLimitsBySpace(ctx, spaceId) if err != nil { @@ -462,3 +580,145 @@ func (c *coordinator) InboxAddMessage(ctx context.Context, message *inbox.InboxM func (c *coordinator) AddStream(eventType coordinatorproto.NotifyEventType, accountId, peerId string, stream coordinatorproto.DRPCCoordinator_NotifySubscribeStream) error { return c.subscribe.AddStream(eventType, accountId, peerId, stream) } + +func unmarshalSpaceHeader(spaceHeader []byte) (header *spacesyncproto.SpaceHeader, err error) { + rawHeader := &spacesyncproto.RawSpaceHeader{} + if err = rawHeader.UnmarshalVT(spaceHeader); err != nil { + return + } + header = &spacesyncproto.SpaceHeader{} + err = header.UnmarshalVT(rawHeader.SpaceHeader) + return +} + +// verifyNestedSpace is the trust gate for child (nested) spaces at SpaceSign time: nothing else +// enforces the parent's authority. It resolves the legalOwner as the parent's CURRENT owner and +// validates the registration record in the parent acl. The coordinator's acl cache stays current +// for records it accepted itself (all acl writes flow through this coordinator's acl service). +func (c *coordinator) verifyNestedSpace(ctx context.Context, spaceId string, header *spacesyncproto.SpaceHeader, parentAclRecordId string, requirePinnedMatch bool) (legalOwner crypto.PubKey, err error) { + if parentAclRecordId == "" { + return nil, coordinatorproto.ErrForbidden + } + parentEntry, err := c.spaceStatus.Status(ctx, header.ParentSpaceId) + if err != nil { + return nil, err + } + if parentEntry.Status != spacestatus.SpaceStatusCreated { + return nil, coordinatorproto.ErrSpaceIsDeleted + } + if parentEntry.Type == spacestatus.SpaceTypeTech || parentEntry.Type == spacestatus.SpaceTypeOneToOne { + return nil, coordinatorproto.ErrForbidden + } + if parentEntry.BilledIdentity != "" { + // v1: single-level nesting — a child cannot itself be a parent + return nil, coordinatorproto.ErrForbidden + } + if len(header.AclPayload) == 0 { + // nested spaces require the V1 header carrying the acl root + return nil, coordinatorproto.ErrForbidden + } + childAclRootId, err := cidutil.NewCidFromBytes(header.AclPayload) + if err != nil { + return nil, err + } + var rawAclRoot consensusproto.RawRecord + if err = rawAclRoot.UnmarshalVT(header.AclPayload); err != nil { + return nil, err + } + var childAclRoot aclrecordproto.AclRoot + if err = childAclRoot.UnmarshalVT(rawAclRoot.Payload); err != nil { + return nil, err + } + if childAclRoot.ParentSpaceId != header.ParentSpaceId || len(childAclRoot.LegalOwner) == 0 { + return nil, coordinatorproto.ErrForbidden + } + pinnedLegalOwner, err := crypto.UnmarshalEd25519PublicKeyProto(childAclRoot.LegalOwner) + if err != nil { + return nil, err + } + var currentOwner crypto.PubKey + err = c.acl.ReadState(ctx, header.ParentSpaceId, func(st *list.AclState) error { + owner, err := st.OwnerPubKey() + if err != nil { + return err + } + currentOwner = owner + // the pinned parentAclRootId must be the parent's ACTUAL acl root — else a registering + // admin could pin a binding scope they control, weakening the legalOwner-proof binding + if childAclRoot.ParentAclRootId != st.Id() { + return coordinatorproto.ErrForbidden + } + if requirePinnedMatch && !owner.Equals(pinnedLegalOwner) { + // at creation the pinned legalOwner must be the parent's current owner (anchors + // the induction chain); on re-sign the parent may have transferred ownership + return coordinatorproto.ErrForbidden + } + if requirePinnedMatch { + // the "members may create children" toggle gates CREATION only; enforcing it on + // re-sign would starve existing children of receipt renewals when it is turned off + if opts := st.CurrentOptions(); opts != nil && opts.ChildrenCreationDisallowed { + return coordinatorproto.ErrForbidden + } + } + reg, ok := st.ChildRegistration(spaceId) + if !ok || reg.Revoked || reg.RecordId != parentAclRecordId || reg.ChildAclRootId != childAclRootId { + return coordinatorproto.ErrForbidden + } + perms, err := st.PermissionsAtRecord(reg.RecordId, reg.Author) + if err != nil { + return err + } + if !perms.CanManageAccounts() { + return coordinatorproto.ErrForbidden + } + return nil + }) + if err != nil { + return nil, err + } + // bill/authorize against the parent's CURRENT owner (hybrid derive-fresh) + return currentOwner, nil +} + +// verifyKeylessGovernanceRecord is the fresh-derivation half of the hybrid legalOwner model: +// a keyless-governance record (AclAccountRemoveNoRotate / AclLegalOwnerUpdate) must be authored +// by the parent space's CURRENT owner. Clients validate against the child's stored legalOwner key, +// which may lag a parent ownership transfer; this check keeps the window fail-closed — an ex-owner +// passes stale clients but is rejected here at submit. +func (c *coordinator) verifyKeylessGovernanceRecord(ctx context.Context, rec *consensusproto.RawRecord, statusEntry spacestatus.StatusEntry) (err error) { + // content sniffing only: anything that doesn't parse falls through to the + // full validation inside acl.AddRecord + var aclRec consensusproto.Record + if uErr := aclRec.UnmarshalVT(rec.Payload); uErr != nil { + return nil + } + var aclData aclrecordproto.AclData + if uErr := aclData.UnmarshalVT(aclRec.Data); uErr != nil { + return nil + } + var keyless bool + for _, content := range aclData.GetAclContent() { + if content.GetAccountRemoveNoRotate() != nil || content.GetLegalOwnerUpdate() != nil { + keyless = true + break + } + } + if !keyless { + return nil + } + if statusEntry.ParentSpaceId == "" { + return coordinatorproto.ErrForbidden + } + author, err := crypto.UnmarshalEd25519PublicKeyProto(aclRec.Identity) + if err != nil { + return + } + parentOwner, err := c.acl.OwnerPubKey(ctx, statusEntry.ParentSpaceId) + if err != nil { + return + } + if !parentOwner.Equals(author) { + return coordinatorproto.ErrForbidden + } + return nil +} diff --git a/coordinator/externalseats.go b/coordinator/externalseats.go new file mode 100644 index 0000000..1889a08 --- /dev/null +++ b/coordinator/externalseats.go @@ -0,0 +1,283 @@ +package coordinator + +import ( + "context" + "hash/fnv" + "time" + + "github.com/anyproto/any-sync/commonspace/object/acl/aclrecordproto" + "github.com/anyproto/any-sync/commonspace/object/acl/list" + "github.com/anyproto/any-sync/consensus/consensusproto" + "github.com/anyproto/any-sync/coordinator/coordinatorproto" + "github.com/anyproto/any-sync/net/peer" + "github.com/anyproto/any-sync/util/crypto" + "go.mongodb.org/mongo-driver/bson" + "go.mongodb.org/mongo-driver/mongo/options" + "go.uber.org/zap" + + "github.com/anyproto/any-sync-coordinator/spacestatus" +) + +// External seats (nested spaces, docs/16 phase 5): a compartment may admit an +// identity that is NOT a member of the parent (org) space only within the +// legalOwner's externalSeatsLimit — one per-org pool of DISTINCT external +// identities across all of the org's children (grooming decision, Q9). +// +// The gate runs live off the acl states (no persistent counter to drift); the +// externalSeats collection is a best-effort registry that backs the +// ExternalCompartments discovery rpc for the externals themselves, who cannot +// read the parent's registrations. + +const externalSeatsCollName = "externalSeats" + +// orgSeatLockStripes sizes the striped lock set serializing admissions per org. +const orgSeatLockStripes = 256 + +// orgSeatLeaseCollName backs the cross-instance org lease (db.Lease): the striped +// mutex below only serializes ONE coordinator process, and with several replicas +// two admissions into sibling children can land on different instances. The TTL +// bounds how long a crashed holder blocks its org; it must sit well above the +// worst-case check→commit (one consensus write + a few acl reads). +const ( + orgSeatLeaseCollName = "orgSeatLeases" + orgSeatLeaseTTL = 30 * time.Second + orgSeatLeasePoll = 100 * time.Millisecond +) + +// lockOrgSeats serializes the external-seat check→commit window per org (parent +// space id) WITHIN this instance: the pool is counted live across ALL of the +// org's children, and two concurrent admissions into sibling children are +// different consensus logs, so nothing else orders the read against the other +// admission's commit. Striped by hash, so unrelated orgs rarely contend; the +// same-instance fast path in front of the cross-instance mongo lease. +func (c *coordinator) lockOrgSeats(parentSpaceId string) (unlock func()) { + h := fnv.New32a() + h.Write([]byte(parentSpaceId)) + mu := &c.orgSeatLocks[h.Sum32()%orgSeatLockStripes] + mu.Lock() + return mu.Unlock +} + +type externalSeatEntry struct { + Id string `bson:"_id"` // spaceId + "/" + identity + Identity string `bson:"identity"` + SpaceId string `bson:"spaceId"` + OrgOwner string `bson:"orgOwner"` +} + +// admittedIdentities extracts the identities an acl record admits as members +// (direct adds, request accepts, invite joins) or elevates to a non-None role +// (permission changes) — an elevation from None occupies a seat exactly like an +// admission, so it must route through the same gate. Content sniffing only — +// anything unparseable is handled by the full validation in acl.AddRecord. +func admittedIdentities(rec *consensusproto.RawRecord) (admitted []crypto.PubKey) { + var aclRec consensusproto.Record + if err := aclRec.UnmarshalVT(rec.Payload); err != nil { + return nil + } + var aclData aclrecordproto.AclData + if err := aclData.UnmarshalVT(aclRec.Data); err != nil { + return nil + } + addKey := func(raw []byte) { + if key, err := crypto.UnmarshalEd25519PublicKeyProto(raw); err == nil { + admitted = append(admitted, key) + } + } + for _, content := range aclData.GetAclContent() { + switch { + case content.GetAccountsAdd() != nil: + for _, add := range content.GetAccountsAdd().GetAdditions() { + addKey(add.Identity) + } + case content.GetRequestAccept() != nil: + addKey(content.GetRequestAccept().Identity) + case content.GetInviteJoin() != nil: + addKey(content.GetInviteJoin().Identity) + case content.GetPermissionChange() != nil: + if pc := content.GetPermissionChange(); pc.Permissions != aclrecordproto.AclUserPermissions_None { + addKey(pc.Identity) + } + case content.GetPermissionChanges() != nil: + for _, pc := range content.GetPermissionChanges().GetChanges() { + if pc.Permissions != aclrecordproto.AclUserPermissions_None { + addKey(pc.Identity) + } + } + } + } + return admitted +} + +// removedIdentities extracts the identities an acl record removes. +func removedIdentities(rec *consensusproto.RawRecord) (removed []crypto.PubKey) { + var aclRec consensusproto.Record + if err := aclRec.UnmarshalVT(rec.Payload); err != nil { + return nil + } + var aclData aclrecordproto.AclData + if err := aclData.UnmarshalVT(aclRec.Data); err != nil { + return nil + } + addKey := func(raw []byte) { + if key, err := crypto.UnmarshalEd25519PublicKeyProto(raw); err == nil { + removed = append(removed, key) + } + } + for _, content := range aclData.GetAclContent() { + switch { + case content.GetAccountRemove() != nil: + for _, raw := range content.GetAccountRemove().GetIdentities() { + addKey(raw) + } + case content.GetAccountRemoveNoRotate() != nil: + for _, raw := range content.GetAccountRemoveNoRotate().GetIdentities() { + addKey(raw) + } + } + } + return removed +} + +// verifyExternalSeats gates external admissions into a child space against the +// legalOwner's per-org pool and returns the externals among the admitted set. +func (c *coordinator) verifyExternalSeats(ctx context.Context, statusEntry spacestatus.StatusEntry, admitted []crypto.PubKey) (externals []crypto.PubKey, err error) { + // which of the admitted are NOT parent members + err = c.acl.ReadState(ctx, statusEntry.ParentSpaceId, func(st *list.AclState) error { + for _, identity := range admitted { + if st.Permissions(identity).NoPermissions() { + externals = append(externals, identity) + } + } + return nil + }) + if err != nil { + return nil, err + } + if len(externals) == 0 { + return nil, nil + } + orgOwner, err := c.acl.OwnerPubKey(ctx, statusEntry.ParentSpaceId) + if err != nil { + return nil, err + } + limits, err := c.accountLimit.GetLimits(ctx, orgOwner.Account()) + if err != nil { + return nil, err + } + // distinct externals across all of the org's children, live from the acls + var ( + parentMembers = map[string]struct{}{} + childIds []string + ) + err = c.acl.ReadState(ctx, statusEntry.ParentSpaceId, func(st *list.AclState) error { + for _, acc := range st.CurrentAccounts() { + if !acc.Permissions.NoPermissions() { + parentMembers[acc.PubKey.Account()] = struct{}{} + } + } + for _, reg := range st.ChildRegistrations() { + if !reg.Revoked { + childIds = append(childIds, reg.ChildSpaceId) + } + } + return nil + }) + if err != nil { + return nil, err + } + seats := map[string]struct{}{} + for _, childId := range childIds { + err = c.acl.ReadState(ctx, childId, func(st *list.AclState) error { + for _, acc := range st.CurrentAccounts() { + if acc.Permissions.NoPermissions() || acc.Status != list.StatusActive { + continue + } + account := acc.PubKey.Account() + if _, isMember := parentMembers[account]; !isMember { + seats[account] = struct{}{} + } + } + return nil + }) + if err != nil { + // fail closed: a child ACL we cannot count would UNDER-count the pool and let an + // over-limit admission through, so reject rather than skip + log.Warn("external-seat count: child acl unavailable, rejecting admission", zap.String("childId", childId), zap.Error(err)) + return nil, err + } + } + for _, identity := range externals { + seats[identity.Account()] = struct{}{} + } + if uint32(len(seats)) > limits.ExternalSeatsLimit { + return nil, coordinatorproto.ErrSpaceLimitReached + } + return externals, nil +} + +// recordExternalSeats maintains the discovery registry, best-effort. +func (c *coordinator) recordExternalSeats(ctx context.Context, spaceId, orgOwner string, externals []crypto.PubKey) { + if c.externalSeats == nil { + return + } + for _, identity := range externals { + account := identity.Account() + _, err := c.externalSeats.UpdateOne(ctx, + bson.M{"_id": spaceId + "/" + account}, + bson.M{"$set": externalSeatEntry{ + Id: spaceId + "/" + account, + Identity: account, + SpaceId: spaceId, + OrgOwner: orgOwner, + }}, + options.Update().SetUpsert(true), + ) + if err != nil { + log.Warn("external-seat registry upsert failed", zap.String("spaceId", spaceId), zap.Error(err)) + } + } +} + +// dropExternalSeats removes registry rows for identities removed from spaceId. +func (c *coordinator) dropExternalSeats(ctx context.Context, spaceId string, removed []crypto.PubKey) { + if c.externalSeats == nil { + return + } + for _, identity := range removed { + if _, err := c.externalSeats.DeleteOne(ctx, bson.M{"_id": spaceId + "/" + identity.Account()}); err != nil { + log.Debug("external-seat registry delete failed", zap.String("spaceId", spaceId), zap.Error(err)) + } + } +} + +// ExternalCompartments lists the child spaces the calling identity holds an +// external seat in. Rows whose space is gone are filtered out lazily. +func (c *coordinator) ExternalCompartments(ctx context.Context) (spaceIds []string, err error) { + accountPubKey, err := peer.CtxPubKey(ctx) + if err != nil { + return nil, err + } + if c.externalSeats == nil { + return nil, coordinatorproto.ErrUnexpected + } + cur, err := c.externalSeats.Find(ctx, bson.M{"identity": accountPubKey.Account()}) + if err != nil { + return nil, err + } + defer cur.Close(ctx) + for cur.Next(ctx) { + var entry externalSeatEntry + if err := cur.Decode(&entry); err != nil { + continue + } + if entry.SpaceId == "" { + continue + } + if st, err := c.spaceStatus.Status(ctx, entry.SpaceId); err != nil || st.Status != spacestatus.SpaceStatusCreated { + continue + } + spaceIds = append(spaceIds, entry.SpaceId) + } + return spaceIds, nil +} diff --git a/coordinator/nestedspaces_test.go b/coordinator/nestedspaces_test.go new file mode 100644 index 0000000..a64e776 --- /dev/null +++ b/coordinator/nestedspaces_test.go @@ -0,0 +1,658 @@ +package coordinator + +import ( + "testing" + + "github.com/anyproto/any-sync/commonspace/object/accountdata" + "github.com/anyproto/any-sync/commonspace/object/acl/aclrecordproto" + "github.com/anyproto/any-sync/commonspace/object/acl/list" + "github.com/anyproto/any-sync/commonspace/object/acl/list/listtest" + "github.com/anyproto/any-sync/commonspace/spacepayloads" + "github.com/anyproto/any-sync/consensus/consensusproto" + "github.com/anyproto/any-sync/coordinator/coordinatorproto" + "github.com/anyproto/any-sync/net/peer" + "github.com/anyproto/any-sync/nodeconf" + "github.com/anyproto/any-sync/util/crypto" + "github.com/stretchr/testify/require" + "go.uber.org/mock/gomock" + + "github.com/anyproto/any-sync-coordinator/accountlimit" + "github.com/anyproto/any-sync-coordinator/spacestatus" +) + +func buildNestedSetup(t *testing.T) (creatorKeys *accountdata.AccountKeys, parentAcl list.AclList, childId string, childHeader []byte, childAclRootId string) { + parentEx := list.NewAclExecutor("parent.id") + require.NoError(t, parentEx.Execute("a.init::a")) + parentAcl = parentEx.ActualAccounts()["a"].Acl + + creatorKeys, err := accountdata.NewRandom() + require.NoError(t, err) + master, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + metaKey, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + readKey, _ := crypto.NewRandomAES() + out, err := spacepayloads.StoragePayloadForSpaceCreateV1(spacepayloads.SpaceCreatePayload{ + SigningKey: creatorKeys.SignKey, + SpaceType: "anytype.space", + ReplicationKey: 42, + MasterKey: master, + ReadKey: readKey, + MetadataKey: metaKey, + Metadata: []byte("md"), + ParentSpaceId: "parent.id", + LegalOwner: parentAcl.AclState().Identity(), + ParentAclRootId: parentAcl.Id(), + }) + require.NoError(t, err) + return creatorKeys, parentAcl, out.SpaceHeaderWithId.Id, out.SpaceHeaderWithId.RawHeader, out.AclWithId.Id +} + +func registerChild(t *testing.T, parentAcl list.AclList, childId, childAclRootId string) (parentRecId string) { + reg, err := parentAcl.RecordBuilder().BuildChildRegister(list.ChildRegisterPayload{ + ChildSpaceId: childId, + ChildAclRootId: childAclRootId, + }) + require.NoError(t, err) + require.NoError(t, parentAcl.AddRawRecord(listtest.WrapAclRecord(reg))) + return parentAcl.Head().Id +} + +func TestCoordinator_SpaceSignNested(t *testing.T) { + creatorKeys, parentAcl, childId, childHeader, childAclRootId := buildNestedSetup(t) + pubKeyData, err := creatorKeys.SignKey.GetPublic().Marshall() + require.NoError(t, err) + signCtx := peer.CtxWithPeerId(peer.CtxWithIdentity(ctx, pubKeyData), "peer.id") + parentOwner := parentAcl.AclState().Identity() + + expectParentReads := func(fx *fixture, entry spacestatus.StatusEntry) { + fx.spaceStatus.EXPECT().Status(gomock.Any(), "parent.id").Return(entry, nil) + } + activeParent := spacestatus.StatusEntry{ + SpaceId: "parent.id", + Identity: parentOwner.Account(), + Status: spacestatus.SpaceStatusCreated, + Type: spacestatus.SpaceTypeRegular, + } + expectReadState := func(fx *fixture) { + fx.acl.EXPECT().ReadState(gomock.Any(), "parent.id", gomock.Any()).DoAndReturn( + func(_ interface{ Done() <-chan struct{} }, _ string, f func(s *list.AclState) error) error { + return f(parentAcl.AclState()) + }) + } + // the nested branch checks the child's own status first (new vs re-sign) + expectChildNew := func(fx *fixture) { + fx.spaceStatus.EXPECT().Status(gomock.Any(), childId).Return(spacestatus.StatusEntry{}, coordinatorproto.ErrSpaceNotExists) + } + + t.Run("success", func(t *testing.T) { + parentRecId := registerChild(t, parentAcl, childId, childAclRootId) + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + expectChildNew(fx) + expectParentReads(fx, activeParent) + expectReadState(fx) + fx.accountLimit.EXPECT().GetLimits(gomock.Any(), parentOwner.Account()).Return(accountlimit.Limits{SharedSpacesLimit: 5}, nil) + fx.spaceStatus.EXPECT().NewChildStatus(gomock.Any(), childId, gomock.Any(), "parent.id", parentOwner.Account(), uint32(5), spacestatus.SpaceTypeRegular, false).Return(nil) + fx.coordLog.EXPECT().SpaceReceipt(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + fx.aclEventLog.EXPECT().AddLog(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + + receipt, err := fx.SpaceSign(signCtx, childId, childHeader, false, parentRecId) + require.NoError(t, err) + require.NotNil(t, receipt) + }) + + t.Run("missing parentAclRecordId", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + expectChildNew(fx) + _, err := fx.SpaceSign(signCtx, childId, childHeader, false, "") + require.ErrorIs(t, err, coordinatorproto.ErrForbidden) + }) + + t.Run("parent deleted", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + expectChildNew(fx) + expectParentReads(fx, spacestatus.StatusEntry{ + SpaceId: "parent.id", + Status: spacestatus.SpaceStatusDeleted, + Type: spacestatus.SpaceTypeRegular, + }) + _, err := fx.SpaceSign(signCtx, childId, childHeader, false, "some.rec") + require.ErrorIs(t, err, coordinatorproto.ErrSpaceIsDeleted) + }) + + t.Run("parent is itself a child", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + expectChildNew(fx) + nested := activeParent + nested.BilledIdentity = "some.org" + expectParentReads(fx, nested) + _, err := fx.SpaceSign(signCtx, childId, childHeader, false, "some.rec") + require.ErrorIs(t, err, coordinatorproto.ErrForbidden) + }) + + t.Run("wrong registration record id", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + expectChildNew(fx) + expectParentReads(fx, activeParent) + expectReadState(fx) + _, err := fx.SpaceSign(signCtx, childId, childHeader, false, "not.the.registration") + require.ErrorIs(t, err, coordinatorproto.ErrForbidden) + }) + + t.Run("mismatched parentAclRootId rejected", func(t *testing.T) { + // a child that pins a binding scope other than the real parent acl root id + badCreator, err := accountdata.NewRandom() + require.NoError(t, err) + master, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + metaKey, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + rk, _ := crypto.NewRandomAES() + out, err := spacepayloads.StoragePayloadForSpaceCreateV1(spacepayloads.SpaceCreatePayload{ + SigningKey: badCreator.SignKey, + SpaceType: "anytype.space", + ReplicationKey: 99, + MasterKey: master, + ReadKey: rk, + MetadataKey: metaKey, + Metadata: []byte("md"), + ParentSpaceId: "parent.id", + LegalOwner: parentOwner, + ParentAclRootId: "attacker-controlled-root", + }) + require.NoError(t, err) + badKeyData, err := badCreator.SignKey.GetPublic().Marshall() + require.NoError(t, err) + badCtx := peer.CtxWithPeerId(peer.CtxWithIdentity(ctx, badKeyData), "peer.id") + + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + fx.spaceStatus.EXPECT().Status(gomock.Any(), out.SpaceHeaderWithId.Id).Return(spacestatus.StatusEntry{}, coordinatorproto.ErrSpaceNotExists) + expectParentReads(fx, activeParent) + expectReadState(fx) + _, err = fx.SpaceSign(badCtx, out.SpaceHeaderWithId.Id, out.SpaceHeaderWithId.RawHeader, false, "some.rec") + require.ErrorIs(t, err, coordinatorproto.ErrForbidden) + }) + + t.Run("top-level sign with parentAclRecordId is rejected", func(t *testing.T) { + topKeys, err := accountdata.NewRandom() + require.NoError(t, err) + master, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + metaKey, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + readKey, _ := crypto.NewRandomAES() + out, err := spacepayloads.StoragePayloadForSpaceCreateV1(spacepayloads.SpaceCreatePayload{ + SigningKey: topKeys.SignKey, + SpaceType: "anytype.space", + ReplicationKey: 43, + MasterKey: master, + ReadKey: readKey, + MetadataKey: metaKey, + Metadata: []byte("md"), + }) + require.NoError(t, err) + topKeyData, err := topKeys.SignKey.GetPublic().Marshall() + require.NoError(t, err) + topCtx := peer.CtxWithPeerId(peer.CtxWithIdentity(ctx, topKeyData), "peer.id") + + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + _, err = fx.SpaceSign(topCtx, out.SpaceHeaderWithId.Id, out.SpaceHeaderWithId.RawHeader, false, "some.rec") + require.ErrorIs(t, err, coordinatorproto.ErrForbidden) + }) + + t.Run("children disallowed by parent options", func(t *testing.T) { + // separate parent with the toggle set + parentEx := list.NewAclExecutor("parent.id") + require.NoError(t, parentEx.Execute("a.init::a")) + lockedParent := parentEx.ActualAccounts()["a"].Acl + optsChange, err := lockedParent.RecordBuilder().BuildSpaceOptionsChange(&aclrecordproto.AclSpaceOptions{ + ChildrenCreationDisallowed: true, + }) + require.NoError(t, err) + require.NoError(t, lockedParent.AddRawRecord(listtest.WrapAclRecord(optsChange))) + + lockedCreator, err := accountdata.NewRandom() + require.NoError(t, err) + master, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + metaKey, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + readKey, _ := crypto.NewRandomAES() + out, err := spacepayloads.StoragePayloadForSpaceCreateV1(spacepayloads.SpaceCreatePayload{ + SigningKey: lockedCreator.SignKey, + SpaceType: "anytype.space", + ReplicationKey: 44, + MasterKey: master, + ReadKey: readKey, + MetadataKey: metaKey, + Metadata: []byte("md"), + ParentSpaceId: "parent.id", + LegalOwner: lockedParent.AclState().Identity(), + ParentAclRootId: lockedParent.Id(), + }) + require.NoError(t, err) + lockedKeyData, err := lockedCreator.SignKey.GetPublic().Marshall() + require.NoError(t, err) + lockedCtx := peer.CtxWithPeerId(peer.CtxWithIdentity(ctx, lockedKeyData), "peer.id") + + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + fx.spaceStatus.EXPECT().Status(gomock.Any(), out.SpaceHeaderWithId.Id).Return(spacestatus.StatusEntry{}, coordinatorproto.ErrSpaceNotExists) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "parent.id").Return(activeParent, nil) + fx.acl.EXPECT().ReadState(gomock.Any(), "parent.id", gomock.Any()).DoAndReturn( + func(_ interface{ Done() <-chan struct{} }, _ string, f func(s *list.AclState) error) error { + return f(lockedParent.AclState()) + }) + _, err = fx.SpaceSign(lockedCtx, out.SpaceHeaderWithId.Id, out.SpaceHeaderWithId.RawHeader, false, "some.rec") + require.ErrorIs(t, err, coordinatorproto.ErrForbidden) + }) +} + +func TestCoordinator_SpaceSignNested_ResignAfterOwnershipTransfer(t *testing.T) { + // build a parent whose ownership will move a -> b, and a child pinned to a (the genesis owner) + parentEx := list.NewAclExecutor("parent.id") + for _, cmd := range []string{ + "a.init::a", + "a.invite::inv", + "b.join::inv", + "a.approve::b,adm", + } { + require.NoError(t, parentEx.Execute(cmd)) + } + parentAcl := parentEx.ActualAccounts()["a"].Acl + genesisOwner := parentAcl.AclState().Identity() // account a + + creatorKeys, err := accountdata.NewRandom() + require.NoError(t, err) + master, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + metaKey, _, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + readKey, _ := crypto.NewRandomAES() + out, err := spacepayloads.StoragePayloadForSpaceCreateV1(spacepayloads.SpaceCreatePayload{ + SigningKey: creatorKeys.SignKey, + SpaceType: "anytype.space", + ReplicationKey: 77, + MasterKey: master, + ReadKey: readKey, + MetadataKey: metaKey, + Metadata: []byte("md"), + ParentSpaceId: "parent.id", + LegalOwner: genesisOwner, // pinned to a + ParentAclRootId: parentAcl.Id(), + }) + require.NoError(t, err) + childId := out.SpaceHeaderWithId.Id + + // transfer parent ownership a -> b FIRST (via the executor, which tracks the head), + // then register the child directly as the last acl write + require.NoError(t, parentEx.Execute("a.ownership_change::b,adm")) + currentOwner, err := parentAcl.AclState().OwnerPubKey() + require.NoError(t, err) + require.False(t, currentOwner.Equals(genesisOwner), "ownership must have moved") + registerChild(t, parentAcl, childId, out.AclWithId.Id) + + pubKeyData, err := creatorKeys.SignKey.GetPublic().Marshall() + require.NoError(t, err) + signCtx := peer.CtxWithPeerId(peer.CtxWithIdentity(ctx, pubKeyData), "peer.id") + + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + // child already exists → re-sign path (isNew = false) + fx.spaceStatus.EXPECT().Status(gomock.Any(), childId).Return(spacestatus.StatusEntry{ + SpaceId: childId, + Status: spacestatus.SpaceStatusCreated, + Type: spacestatus.SpaceTypeRegular, + ParentSpaceId: "parent.id", + }, nil) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "parent.id").Return(spacestatus.StatusEntry{ + SpaceId: "parent.id", + Identity: currentOwner.Account(), + Status: spacestatus.SpaceStatusCreated, + Type: spacestatus.SpaceTypeRegular, + }, nil) + fx.acl.EXPECT().ReadState(gomock.Any(), "parent.id", gomock.Any()).DoAndReturn( + func(_ interface{ Done() <-chan struct{} }, _ string, f func(s *list.AclState) error) error { + return f(parentAcl.AclState()) + }) + // billing must target the CURRENT owner, not the pinned genesis owner + fx.accountLimit.EXPECT().GetLimits(gomock.Any(), currentOwner.Account()).Return(accountlimit.Limits{SharedSpacesLimit: 5}, nil) + fx.spaceStatus.EXPECT().NewChildStatus(gomock.Any(), childId, gomock.Any(), "parent.id", currentOwner.Account(), uint32(5), spacestatus.SpaceTypeRegular, false).Return(nil) + fx.coordLog.EXPECT().SpaceReceipt(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + fx.aclEventLog.EXPECT().AddLog(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + + // re-sign: header pins the OLD owner but the space is already Created, so isNew=false + // skips the pinned==current check and succeeds, billing the new owner + receipt, err := fx.SpaceSign(signCtx, childId, out.SpaceHeaderWithId.RawHeader, false, parentAcl.Head().Id) + require.NoError(t, err) + require.NotNil(t, receipt) +} + +// networkKeys works around the fixture's registration order: the coordinator's Init reads the +// account before the test account service generated it, so SpaceSign needs the key set directly. +func networkKeys(t *testing.T) *accountdata.AccountKeys { + keys, err := accountdata.NewRandom() + require.NoError(t, err) + return keys +} + +func TestCoordinator_SpaceDeleteAsLegalOwner(t *testing.T) { + legalOwnerKeys, err := accountdata.NewRandom() + require.NoError(t, err) + strangerKeys, err := accountdata.NewRandom() + require.NoError(t, err) + legalOwnerKeyData, err := legalOwnerKeys.SignKey.GetPublic().Marshall() + require.NoError(t, err) + strangerKeyData, err := strangerKeys.SignKey.GetPublic().Marshall() + require.NoError(t, err) + legalOwnerCtx := peer.CtxWithPeerId(peer.CtxWithIdentity(ctx, legalOwnerKeyData), "peer.id") + strangerCtx := peer.CtxWithPeerId(peer.CtxWithIdentity(ctx, strangerKeyData), "peer.id") + + childEntry := spacestatus.StatusEntry{ + SpaceId: "child.id", + Identity: "creator.identity", + Status: spacestatus.SpaceStatusCreated, + Type: spacestatus.SpaceTypeRegular, + ParentSpaceId: "parent.id", + } + + t.Run("legal owner deletes a child it does not own", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(childEntry, nil) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "parent.id").Return(legalOwnerKeys.SignKey.GetPublic(), nil) + fx.nodeConf.EXPECT().Configuration().Return(nodeconf.Configuration{NetworkId: "net"}) + fx.spaceStatus.EXPECT().SpaceDelete(gomock.Any(), gomock.Any()).DoAndReturn( + func(_ interface{ Done() <-chan struct{} }, payload spacestatus.SpaceDeletion) (int64, error) { + require.True(t, payload.AuthorizedByLegalOwner) + require.Equal(t, "child.id", payload.SpaceId) + return 42, nil + }) + toBeDeleted, err := fx.SpaceDelete(legalOwnerCtx, "child.id", 60, []byte("payload"), "payload.id") + require.NoError(t, err) + require.Equal(t, int64(42), toBeDeleted) + }) + + t.Run("non-owner non-legal-owner is rejected", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(childEntry, nil) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "parent.id").Return(legalOwnerKeys.SignKey.GetPublic(), nil) + _, err := fx.SpaceDelete(strangerCtx, "child.id", 60, []byte("payload"), "payload.id") + require.ErrorIs(t, err, coordinatorproto.ErrForbidden) + }) + + t.Run("top-level space keeps the owner-only path", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + fx.account = networkKeys(t) + topEntry := childEntry + topEntry.ParentSpaceId = "" + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(topEntry, nil) + fx.nodeConf.EXPECT().Configuration().Return(nodeconf.Configuration{NetworkId: "net"}) + fx.spaceStatus.EXPECT().SpaceDelete(gomock.Any(), gomock.Any()).DoAndReturn( + func(_ interface{ Done() <-chan struct{} }, payload spacestatus.SpaceDeletion) (int64, error) { + require.False(t, payload.AuthorizedByLegalOwner) + return 7, nil + }) + _, err := fx.SpaceDelete(strangerCtx, "child.id", 60, []byte("payload"), "payload.id") + require.NoError(t, err) + }) +} + +func keylessRemoveRecord(t *testing.T, signer *accountdata.AccountKeys, target crypto.PubKey) []byte { + targetProto, err := target.Marshall() + require.NoError(t, err) + data := &aclrecordproto.AclData{AclContent: []*aclrecordproto.AclContentValue{{ + Value: &aclrecordproto.AclContentValue_AccountRemoveNoRotate{ + AccountRemoveNoRotate: &aclrecordproto.AclAccountRemoveNoRotate{Identities: [][]byte{targetProto}}, + }, + }}} + marshalledData, err := data.MarshalVT() + require.NoError(t, err) + signerProto, err := signer.SignKey.GetPublic().Marshall() + require.NoError(t, err) + rec := &consensusproto.Record{PrevId: "prev", Identity: signerProto, Data: marshalledData} + marshalledRec, err := rec.MarshalVT() + require.NoError(t, err) + sig, err := signer.SignKey.Sign(marshalledRec) + require.NoError(t, err) + raw, err := (&consensusproto.RawRecord{Payload: marshalledRec, Signature: sig}).MarshalVT() + require.NoError(t, err) + return raw +} + +func TestCoordinator_AclAddRecordKeylessGate(t *testing.T) { + legalOwnerKeys, err := accountdata.NewRandom() + require.NoError(t, err) + exOwnerKeys, err := accountdata.NewRandom() + require.NoError(t, err) + targetKeys, err := accountdata.NewRandom() + require.NoError(t, err) + callerKeyData, err := legalOwnerKeys.SignKey.GetPublic().Marshall() + require.NoError(t, err) + callCtx := peer.CtxWithPeerId(peer.CtxWithIdentity(ctx, callerKeyData), "peer.id") + + creatorKeys, err := accountdata.NewRandom() + require.NoError(t, err) + childEntry := spacestatus.StatusEntry{ + SpaceId: "child.id", + Identity: creatorKeys.SignKey.GetPublic().Account(), + Status: spacestatus.SpaceStatusCreated, + Type: spacestatus.SpaceTypeRegular, + IsShareable: true, + ParentSpaceId: "parent.id", + } + + t.Run("current parent owner passes", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + payload := keylessRemoveRecord(t, legalOwnerKeys, targetKeys.SignKey.GetPublic()) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(childEntry, nil) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "parent.id").Return(legalOwnerKeys.SignKey.GetPublic(), nil) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "child.id").Return(creatorKeys.SignKey.GetPublic(), nil) + fx.accountLimit.EXPECT().GetLimitsBySpace(gomock.Any(), "child.id").Return(accountlimit.SpaceLimits{SpaceMembersRead: 10, SpaceMembersWrite: 10}, nil) + fx.acl.EXPECT().AddRecord(gomock.Any(), "child.id", gomock.Any(), gomock.Any()).Return(&consensusproto.RawRecordWithId{Id: "rec.id"}, nil) + fx.aclEventLog.EXPECT().AddLog(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + _, err := fx.AclAddRecord(callCtx, "child.id", payload) + require.NoError(t, err) + }) + + t.Run("ex-owner is rejected at submit (fail-closed window)", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + payload := keylessRemoveRecord(t, exOwnerKeys, targetKeys.SignKey.GetPublic()) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(childEntry, nil) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "parent.id").Return(legalOwnerKeys.SignKey.GetPublic(), nil) + _, err := fx.AclAddRecord(callCtx, "child.id", payload) + require.ErrorIs(t, err, coordinatorproto.ErrForbidden) + }) + + t.Run("keyless record on a top-level space is rejected", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + payload := keylessRemoveRecord(t, legalOwnerKeys, targetKeys.SignKey.GetPublic()) + topEntry := childEntry + topEntry.ParentSpaceId = "" + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(topEntry, nil) + _, err := fx.AclAddRecord(callCtx, "child.id", payload) + require.ErrorIs(t, err, coordinatorproto.ErrForbidden) + }) +} + +func accountsAddRecord(t *testing.T, signer *accountdata.AccountKeys, target crypto.PubKey) []byte { + targetProto, err := target.Marshall() + require.NoError(t, err) + data := &aclrecordproto.AclData{AclContent: []*aclrecordproto.AclContentValue{{ + Value: &aclrecordproto.AclContentValue_AccountsAdd{ + AccountsAdd: &aclrecordproto.AclAccountsAdd{Additions: []*aclrecordproto.AclAccountAdd{{ + Identity: targetProto, + Permissions: aclrecordproto.AclUserPermissions_Writer, + }}}, + }, + }}} + marshalledData, err := data.MarshalVT() + require.NoError(t, err) + signerProto, err := signer.SignKey.GetPublic().Marshall() + require.NoError(t, err) + rec := &consensusproto.Record{PrevId: "prev", Identity: signerProto, Data: marshalledData} + marshalledRec, err := rec.MarshalVT() + require.NoError(t, err) + sig, err := signer.SignKey.Sign(marshalledRec) + require.NoError(t, err) + raw, err := (&consensusproto.RawRecord{Payload: marshalledRec, Signature: sig}).MarshalVT() + require.NoError(t, err) + return raw +} + +// permissionChangeRecord builds a raw acl record moving target to the given role. +func permissionChangeRecord(t *testing.T, signer *accountdata.AccountKeys, target crypto.PubKey, perms aclrecordproto.AclUserPermissions) []byte { + targetProto, err := target.Marshall() + require.NoError(t, err) + data := &aclrecordproto.AclData{AclContent: []*aclrecordproto.AclContentValue{{ + Value: &aclrecordproto.AclContentValue_PermissionChanges{ + PermissionChanges: &aclrecordproto.AclAccountPermissionChanges{Changes: []*aclrecordproto.AclAccountPermissionChange{{ + Identity: targetProto, + Permissions: perms, + }}}, + }, + }}} + marshalledData, err := data.MarshalVT() + require.NoError(t, err) + signerProto, err := signer.SignKey.GetPublic().Marshall() + require.NoError(t, err) + rec := &consensusproto.Record{PrevId: "prev", Identity: signerProto, Data: marshalledData} + marshalledRec, err := rec.MarshalVT() + require.NoError(t, err) + sig, err := signer.SignKey.Sign(marshalledRec) + require.NoError(t, err) + raw, err := (&consensusproto.RawRecord{Payload: marshalledRec, Signature: sig}).MarshalVT() + require.NoError(t, err) + return raw +} + +func TestCoordinator_AclAddRecordExternalSeats(t *testing.T) { + parentEx := list.NewAclExecutor("parent.id") + for _, cmd := range []string{ + "a.init::a", + "a.invite::inv", + "b.join::inv", + "a.approve::b,rw", + } { + require.NoError(t, parentEx.Execute(cmd)) + } + var ( + parentAcl = parentEx.ActualAccounts()["a"].Acl + orgOwnerKey = parentEx.ActualAccounts()["a"].Keys + orgMember = parentEx.ActualAccounts()["b"].Keys + ) + external, err := accountdata.NewRandom() + require.NoError(t, err) + creatorKeys, err := accountdata.NewRandom() + require.NoError(t, err) + callerKeyData, err := creatorKeys.SignKey.GetPublic().Marshall() + require.NoError(t, err) + callCtx := peer.CtxWithPeerId(peer.CtxWithIdentity(ctx, callerKeyData), "peer.id") + + childEntry := spacestatus.StatusEntry{ + SpaceId: "child.id", + Identity: creatorKeys.SignKey.GetPublic().Account(), + Status: spacestatus.SpaceStatusCreated, + Type: spacestatus.SpaceTypeRegular, + IsShareable: true, + ParentSpaceId: "parent.id", + } + expectReadState := func(fx *fixture) { + fx.acl.EXPECT().ReadState(gomock.Any(), "parent.id", gomock.Any()).DoAndReturn( + func(_ interface{ Done() <-chan struct{} }, _ string, f func(s *list.AclState) error) error { + return f(parentAcl.AclState()) + }).AnyTimes() + } + + t.Run("org member admission needs no seats", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + payload := accountsAddRecord(t, creatorKeys, orgMember.SignKey.GetPublic()) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(childEntry, nil) + expectReadState(fx) + fx.accountLimit.EXPECT().GetLimitsBySpace(gomock.Any(), "child.id").Return(accountlimit.SpaceLimits{SpaceMembersRead: 10, SpaceMembersWrite: 10}, nil) + fx.acl.EXPECT().AddRecord(gomock.Any(), "child.id", gomock.Any(), gomock.Any()).Return(&consensusproto.RawRecordWithId{Id: "rec.id"}, nil) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "child.id").Return(creatorKeys.SignKey.GetPublic(), nil) + fx.aclEventLog.EXPECT().AddLog(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + _, err := fx.AclAddRecord(callCtx, "child.id", payload) + require.NoError(t, err) + }) + + t.Run("external rejected when the pool is exhausted", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + payload := accountsAddRecord(t, creatorKeys, external.SignKey.GetPublic()) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(childEntry, nil) + expectReadState(fx) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "parent.id").Return(orgOwnerKey.SignKey.GetPublic(), nil) + fx.accountLimit.EXPECT().GetLimits(gomock.Any(), orgOwnerKey.SignKey.GetPublic().Account()).Return(accountlimit.Limits{ExternalSeatsLimit: 0}, nil) + _, err := fx.AclAddRecord(callCtx, "child.id", payload) + require.ErrorIs(t, err, coordinatorproto.ErrSpaceLimitReached) + }) + + t.Run("external elevation from None consumes a seat", func(t *testing.T) { + // admitting at None then elevating must hit the same gate as a direct admission + fx := newFixture(t) + defer fx.finish(t) + payload := permissionChangeRecord(t, creatorKeys, external.SignKey.GetPublic(), aclrecordproto.AclUserPermissions_Writer) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(childEntry, nil) + expectReadState(fx) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "parent.id").Return(orgOwnerKey.SignKey.GetPublic(), nil) + fx.accountLimit.EXPECT().GetLimits(gomock.Any(), orgOwnerKey.SignKey.GetPublic().Account()).Return(accountlimit.Limits{ExternalSeatsLimit: 0}, nil) + _, err := fx.AclAddRecord(callCtx, "child.id", payload) + require.ErrorIs(t, err, coordinatorproto.ErrSpaceLimitReached) + }) + + t.Run("demotion to None does not consume a seat", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + payload := permissionChangeRecord(t, creatorKeys, external.SignKey.GetPublic(), aclrecordproto.AclUserPermissions_None) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(childEntry, nil) + expectReadState(fx) + fx.accountLimit.EXPECT().GetLimitsBySpace(gomock.Any(), "child.id").Return(accountlimit.SpaceLimits{SpaceMembersRead: 10, SpaceMembersWrite: 10}, nil) + fx.acl.EXPECT().AddRecord(gomock.Any(), "child.id", gomock.Any(), gomock.Any()).Return(&consensusproto.RawRecordWithId{Id: "rec.id"}, nil) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "child.id").Return(creatorKeys.SignKey.GetPublic(), nil) + fx.aclEventLog.EXPECT().AddLog(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + _, err := fx.AclAddRecord(callCtx, "child.id", payload) + require.NoError(t, err) + }) + + t.Run("external admitted within the pool", func(t *testing.T) { + fx := newFixture(t) + defer fx.finish(t) + payload := accountsAddRecord(t, creatorKeys, external.SignKey.GetPublic()) + fx.spaceStatus.EXPECT().Status(gomock.Any(), "child.id").Return(childEntry, nil) + expectReadState(fx) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "parent.id").Return(orgOwnerKey.SignKey.GetPublic(), nil).AnyTimes() + fx.accountLimit.EXPECT().GetLimits(gomock.Any(), orgOwnerKey.SignKey.GetPublic().Account()).Return(accountlimit.Limits{ExternalSeatsLimit: 3}, nil) + fx.accountLimit.EXPECT().GetLimitsBySpace(gomock.Any(), "child.id").Return(accountlimit.SpaceLimits{SpaceMembersRead: 10, SpaceMembersWrite: 10}, nil) + fx.acl.EXPECT().AddRecord(gomock.Any(), "child.id", gomock.Any(), gomock.Any()).Return(&consensusproto.RawRecordWithId{Id: "rec.id"}, nil) + fx.acl.EXPECT().OwnerPubKey(gomock.Any(), "child.id").Return(creatorKeys.SignKey.GetPublic(), nil) + fx.aclEventLog.EXPECT().AddLog(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + _, err := fx.AclAddRecord(callCtx, "child.id", payload) + require.NoError(t, err) + }) +} diff --git a/coordinator/rpchandler.go b/coordinator/rpchandler.go index 631ecbb..d86f7e8 100644 --- a/coordinator/rpchandler.go +++ b/coordinator/rpchandler.go @@ -217,7 +217,7 @@ func (r *rpcHandler) SpaceSign(ctx context.Context, req *coordinatorproto.SpaceS ) }() - receipt, err := r.c.SpaceSign(ctx, req.SpaceId, req.Header, req.ForceRequest) + receipt, err := r.c.SpaceSign(ctx, req.SpaceId, req.Header, req.ForceRequest, req.ParentAclRecordId) if err != nil { return nil, err } @@ -602,3 +602,19 @@ func (r *rpcHandler) FileUsageReport(ctx context.Context, req *coordinatorproto. } return &coordinatorproto.FileUsageReportResponse{}, nil } + +func (r *rpcHandler) ExternalCompartments(ctx context.Context, req *coordinatorproto.ExternalCompartmentsRequest) (resp *coordinatorproto.ExternalCompartmentsResponse, err error) { + st := time.Now() + defer func() { + r.c.metric.RequestLog(ctx, "coordinator.externalCompartments", + metric.TotalDur(time.Since(st)), + zap.String("addr", peer.CtxPeerAddr(ctx)), + zap.Error(err), + ) + }() + spaceIds, err := r.c.ExternalCompartments(ctx) + if err != nil { + return nil, err + } + return &coordinatorproto.ExternalCompartmentsResponse{SpaceIds: spaceIds}, nil +} diff --git a/db/lease.go b/db/lease.go new file mode 100644 index 0000000..33cd223 --- /dev/null +++ b/db/lease.go @@ -0,0 +1,102 @@ +package db + +import ( + "context" + "crypto/rand" + "encoding/hex" + "math/big" + "time" + + "go.mongodb.org/mongo-driver/bson" + "go.mongodb.org/mongo-driver/mongo" + "go.uber.org/zap" +) + +// Lease is a coarse cross-instance mutex backed by a mongo collection: one +// document per key, held for a TTL and taken over once expired. It serializes +// short critical sections across coordinator replicas — an in-process mutex +// only covers one instance. +// +// It is NOT a fencing lock: a holder stalled past the TTL loses exclusivity +// silently while its critical section may still be running. Size the TTL well +// above the worst-case critical section and treat the expiry as a bounded +// residual race, not an impossibility. +type Lease struct { + coll *mongo.Collection + ttl time.Duration + poll time.Duration +} + +// NewLease wraps coll as a lease registry. ttl bounds how long a crashed +// holder blocks the key; poll is the base wait between takeover attempts +// (jittered up to 2x). +func NewLease(coll *mongo.Collection, ttl, poll time.Duration) *Lease { + return &Lease{coll: coll, ttl: ttl, poll: poll} +} + +type leaseDoc struct { + Id string `bson:"_id"` + Holder string `bson:"holder"` + ExpiresAt time.Time `bson:"expiresAt"` +} + +// Acquire blocks until it holds key's lease or ctx is done. The returned +// release is safe to call exactly once and works even when ctx is already +// canceled (it uses its own short deadline), so `defer release()` after a +// failed request still frees the key for other instances. +func (l *Lease) Acquire(ctx context.Context, key string) (release func(), err error) { + tokenBytes := make([]byte, 16) + if _, err = rand.Read(tokenBytes); err != nil { + return nil, err + } + token := hex.EncodeToString(tokenBytes) + release = func() { + relCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) + defer cancel() + if _, derr := l.coll.DeleteOne(relCtx, bson.M{"_id": key, "holder": token}); derr != nil { + // the TTL reclaims the key eventually; log so a systematic release + // failure (mongo outage) is visible rather than silent serialization + log.Warn("lease release failed, ttl will reclaim", zap.String("key", key), zap.Error(derr)) + } + } + for { + now := time.Now() + // take over a free (expired) lease + res, uerr := l.coll.UpdateOne(ctx, + bson.M{"_id": key, "expiresAt": bson.M{"$lt": now}}, + bson.M{"$set": bson.M{"holder": token, "expiresAt": now.Add(l.ttl)}}, + ) + if uerr != nil { + return nil, uerr + } + if res.MatchedCount == 1 { + return release, nil + } + // no expired doc: the key is either absent (claim it) or actively held (wait) + _, ierr := l.coll.InsertOne(ctx, leaseDoc{Id: key, Holder: token, ExpiresAt: now.Add(l.ttl)}) + if ierr == nil { + return release, nil + } + if !mongo.IsDuplicateKeyError(ierr) { + return nil, ierr + } + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-time.After(l.poll + jitter(l.poll)): + } + } +} + +// jitter returns a uniform random duration in [0, max) to de-synchronize +// competing pollers. +func jitter(max time.Duration) time.Duration { + if max <= 0 { + return 0 + } + n, err := rand.Int(rand.Reader, big.NewInt(int64(max))) + if err != nil { + return 0 + } + return time.Duration(n.Int64()) +} diff --git a/db/lease_test.go b/db/lease_test.go new file mode 100644 index 0000000..1011774 --- /dev/null +++ b/db/lease_test.go @@ -0,0 +1,108 @@ +package db + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.mongodb.org/mongo-driver/mongo" + "go.mongodb.org/mongo-driver/mongo/options" +) + +func newTestLeaseColl(t *testing.T) *mongo.Collection { + t.Helper() + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + client, err := mongo.Connect(ctx, options.Client().ApplyURI("mongodb://localhost:27017")) + require.NoError(t, err) + coll := client.Database("coordinator_unittest_lease").Collection("leases") + require.NoError(t, coll.Drop(ctx)) + t.Cleanup(func() { + cctx, ccancel := context.WithTimeout(context.Background(), 5*time.Second) + defer ccancel() + _ = coll.Drop(cctx) + _ = client.Disconnect(cctx) + }) + return coll +} + +func TestLease_MutualExclusion(t *testing.T) { + lease := NewLease(newTestLeaseColl(t), 30*time.Second, 10*time.Millisecond) + ctx := context.Background() + + release1, err := lease.Acquire(ctx, "org.1") + require.NoError(t, err) + + // an unrelated key is not blocked + releaseOther, err := lease.Acquire(ctx, "org.2") + require.NoError(t, err) + releaseOther() + + acquired := make(chan struct{}) + go func() { + release2, err := lease.Acquire(ctx, "org.1") + assert.NoError(t, err) + close(acquired) + release2() + }() + + select { + case <-acquired: + t.Fatal("second acquire must block while the lease is held") + case <-time.After(150 * time.Millisecond): + } + + release1() + select { + case <-acquired: + case <-time.After(5 * time.Second): + t.Fatal("second acquire must proceed after release") + } +} + +func TestLease_ExpiredLeaseIsTakenOver(t *testing.T) { + lease := NewLease(newTestLeaseColl(t), 100*time.Millisecond, 10*time.Millisecond) + ctx := context.Background() + + // acquired and never released — a crashed holder + _, err := lease.Acquire(ctx, "org.1") + require.NoError(t, err) + + start := time.Now() + release2, err := lease.Acquire(ctx, "org.1") + require.NoError(t, err) + release2() + assert.GreaterOrEqual(t, time.Since(start), 50*time.Millisecond, "takeover must wait for expiry") +} + +func TestLease_ReleaseSurvivesCanceledContext(t *testing.T) { + lease := NewLease(newTestLeaseColl(t), 30*time.Second, 10*time.Millisecond) + + reqCtx, cancel := context.WithCancel(context.Background()) + release, err := lease.Acquire(reqCtx, "org.1") + require.NoError(t, err) + cancel() // request failed/canceled after the lease was taken + release() + + // the key is free immediately, not after the 30s ttl + ctx, tcancel := context.WithTimeout(context.Background(), 2*time.Second) + defer tcancel() + release2, err := lease.Acquire(ctx, "org.1") + require.NoError(t, err) + release2() +} + +func TestLease_AcquireHonorsContext(t *testing.T) { + lease := NewLease(newTestLeaseColl(t), 30*time.Second, 10*time.Millisecond) + + release, err := lease.Acquire(context.Background(), "org.1") + require.NoError(t, err) + defer release() + + ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond) + defer cancel() + _, err = lease.Acquire(ctx, "org.1") + require.ErrorIs(t, err, context.DeadlineExceeded) +} diff --git a/go.mod b/go.mod index f795f87..a072844 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,7 @@ go 1.25.7 require ( github.com/ahmetb/govvv v0.3.0 - github.com/anyproto/any-sync v0.13.0-alpha.5 + github.com/anyproto/any-sync v0.13.0-alpha.5.0.20260707151954-0499748aa6a1 github.com/anyproto/go-chash v0.1.0 github.com/gogo/protobuf v1.3.2 github.com/stretchr/testify v1.11.1 diff --git a/go.sum b/go.sum index 66f1892..5409cf4 100644 --- a/go.sum +++ b/go.sum @@ -7,8 +7,8 @@ github.com/ahmetb/govvv v0.3.0 h1:YGLGwEyiUwHFy5eh/RUhdupbuaCGBYn5T5GWXp+WJB0= github.com/ahmetb/govvv v0.3.0/go.mod h1:4WRFpdWtc/YtKgPFwa1dr5+9hiRY5uKAL08bOlxOR6s= github.com/anyproto/any-store v0.4.7 h1:329NWY/xUzGdwKSqFgjAFaPAVkTJbR+WdJirPrvTm/w= github.com/anyproto/any-store v0.4.7/go.mod h1:8cqb52gjZSaYnlybugqpSqSG1RQygY96D2vWMbSJsLo= -github.com/anyproto/any-sync v0.13.0-alpha.5 h1:r+osXFVeC/SGae05TiNJY9prH+uOK2Q4kOOiidSp7yg= -github.com/anyproto/any-sync v0.13.0-alpha.5/go.mod h1:HObtLnsHG20r8rYfwDkpHF7FHzLMjQ/USGhEDLPBVc8= +github.com/anyproto/any-sync v0.13.0-alpha.5.0.20260707151954-0499748aa6a1 h1:kWEoaOCZb8/SKrOz4QipLAlJ1STPdVVPQseR6059hU0= +github.com/anyproto/any-sync v0.13.0-alpha.5.0.20260707151954-0499748aa6a1/go.mod h1:HObtLnsHG20r8rYfwDkpHF7FHzLMjQ/USGhEDLPBVc8= github.com/anyproto/go-bip39 v1.0.0 h1:T6/7WowKYDeyuX/QyXtt98ZX0XXaoOh17M/LFF2M5yk= github.com/anyproto/go-bip39 v1.0.0/go.mod h1:l0rcxmXRyiWAYzE1noMAc4qbeNrbhUwxM3rqSO9ILwo= github.com/anyproto/go-chash v0.1.0 h1:I9meTPjXFRfXZHRJzjOHC/XF7Q5vzysKkiT/grsogXY= diff --git a/spacestatus/mock_spacestatus/mock_spacestatus.go b/spacestatus/mock_spacestatus/mock_spacestatus.go index 2a2a8c5..deac475 100644 --- a/spacestatus/mock_spacestatus/mock_spacestatus.go +++ b/spacestatus/mock_spacestatus/mock_spacestatus.go @@ -172,6 +172,20 @@ func (mr *MockSpaceStatusMockRecorder) Name() *gomock.Call { return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Name", reflect.TypeOf((*MockSpaceStatus)(nil).Name)) } +// NewChildStatus mocks base method. +func (m *MockSpaceStatus) NewChildStatus(ctx context.Context, spaceId string, identity crypto.PubKey, parentSpaceId, billedIdentity string, billedLimit uint32, spaceType spacestatus.SpaceType, force bool) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "NewChildStatus", ctx, spaceId, identity, parentSpaceId, billedIdentity, billedLimit, spaceType, force) + ret0, _ := ret[0].(error) + return ret0 +} + +// NewChildStatus indicates an expected call of NewChildStatus. +func (mr *MockSpaceStatusMockRecorder) NewChildStatus(ctx, spaceId, identity, parentSpaceId, billedIdentity, billedLimit, spaceType, force any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "NewChildStatus", reflect.TypeOf((*MockSpaceStatus)(nil).NewChildStatus), ctx, spaceId, identity, parentSpaceId, billedIdentity, billedLimit, spaceType, force) +} + // NewStatus mocks base method. func (m *MockSpaceStatus) NewStatus(ctx context.Context, spaceId string, identity crypto.PubKey, spaceType spacestatus.SpaceType, force bool) error { m.ctrl.T.Helper() diff --git a/spacestatus/spacedeleter.go b/spacestatus/spacedeleter.go index 535fd7a..3524608 100644 --- a/spacestatus/spacedeleter.go +++ b/spacestatus/spacedeleter.go @@ -66,6 +66,10 @@ type StatusEntry struct { Status int `bson:"status"` Type SpaceType `bson:"type"` IsShareable bool `bson:"isShareable"` + // ParentSpaceId + BilledIdentity are set for child (nested) spaces: the declared parent and + // the identity whose limits the space consumes + ParentSpaceId string `bson:"parentSpaceId,omitempty"` + BilledIdentity string `bson:"billedIdentity,omitempty"` } type spaceDeleter struct { diff --git a/spacestatus/spacestatus.go b/spacestatus/spacestatus.go index 54a2b93..0dcec74 100644 --- a/spacestatus/spacestatus.go +++ b/spacestatus/spacestatus.go @@ -55,6 +55,9 @@ type SpaceDeletion struct { DeletionPayloadId string SpaceId string DeletionPeriod time.Duration + // AuthorizedByLegalOwner marks a deletion already authorized by the coordinator as the + // parent space's current owner (nested spaces): the doc-identity predicate is skipped + AuthorizedByLegalOwner bool AccountInfo } @@ -87,6 +90,9 @@ type configProvider interface { type SpaceStatus interface { NewStatus(ctx context.Context, spaceId string, identity crypto.PubKey, spaceType SpaceType, force bool) (err error) + // NewChildStatus registers a child (nested) space: the creator stays the space identity, while + // the space is billed to the legalOwner (the parent space's owner) against its shared-spaces limit + NewChildStatus(ctx context.Context, spaceId string, identity crypto.PubKey, parentSpaceId, billedIdentity string, billedLimit uint32, spaceType SpaceType, force bool) (err error) // ChangeStatus is deprecated, use only for backwards compatibility ChangeStatus(ctx context.Context, change StatusChange) (entry StatusEntry, err error) ChangeOwner(ctx context.Context, spaceId, newOwnerId string) (err error) @@ -163,6 +169,10 @@ type insertNewSpaceOp struct { Type SpaceType `bson:"type"` IsShareable bool `bson:"isShareable"` SpaceId string `bson:"_id"` + // ParentSpaceId + BilledIdentity are set for child (nested) spaces: the declared parent and + // the identity whose limits the space consumes + ParentSpaceId string `bson:"parentSpaceId,omitempty"` + BilledIdentity string `bson:"billedIdentity,omitempty"` } func (s *spaceStatus) AccountDelete(ctx context.Context, payload AccountDeletion) (toBeDeleted int64, err error) { @@ -293,11 +303,17 @@ func (s *spaceStatus) SpaceDelete(ctx context.Context, payload SpaceDeletion) (t log.Debug("cannot delete tech space", zap.Error(err), zap.String("spaceId", payload.SpaceId)) return coordinatorproto.ErrUnexpected } + deleterIdentity := payload.Identity + if payload.AuthorizedByLegalOwner { + // the coordinator verified the caller is the parent's current owner; + // nil identity skips the doc-owner predicate in modifyStatus + deleterIdentity = nil + } change := StatusChange{ DeletionPayloadType: coordinatorproto.DeletionPayloadType_Confirm, DeletionPayload: payload.DeletionPayload, DeletionPayloadId: payload.DeletionPayloadId, - Identity: payload.Identity, + Identity: deleterIdentity, PeerId: payload.PeerId, NetworkId: payload.NetworkId, Status: SpaceStatusDeletionPending, @@ -508,6 +524,77 @@ func (s *spaceStatus) NewStatus(ctx context.Context, spaceId string, identity cr }) } +// NewChildStatus mirrors NewStatus for child (nested) spaces: the doc carries billedIdentity and +// the insert is limited by the billed identity's shared-spaces allowance instead of only the +// creator's global space cap. +func (s *spaceStatus) NewChildStatus(ctx context.Context, spaceId string, identity crypto.PubKey, parentSpaceId, billedIdentity string, billedLimit uint32, spaceType SpaceType, force bool) error { + return s.db.Tx(ctx, func(txCtx mongo.SessionContext) error { + if s.accountStatusFindTx(txCtx, identity.Account(), SpaceStatusDeletionPending) { + return coordinatorproto.ErrAccountIsDeleted + } + entry, err := s.Status(txCtx, spaceId) + notFound := err == coordinatorproto.ErrSpaceNotExists + if err != nil && !notFound { + return err + } + if entry.Status == SpaceStatusCreated && !notFound { + // save back compatibility, but keep billing pinned to the parent's CURRENT owner: + // the caller recomputes billedIdentity on every sign, and a stale pool assignment + // would let ownership transfers multiply the per-owner children allowance + if entry.BilledIdentity != billedIdentity || entry.ParentSpaceId != parentSpaceId { + _, err = s.spaces.UpdateOne(txCtx, bson.D{{"_id", spaceId}}, bson.D{{"$set", bson.D{ + {"parentSpaceId", parentSpaceId}, + {"billedIdentity", billedIdentity}, + }}}) + } + return err + } + var inserted bool + if notFound { + if _, err = s.spaces.InsertOne(txCtx, insertNewSpaceOp{ + Identity: identity.Account(), + Status: SpaceStatusCreated, + SpaceId: spaceId, + Type: spaceType, + ParentSpaceId: parentSpaceId, + BilledIdentity: billedIdentity, + }); err != nil { + return err + } else { + inserted = true + } + } + if !inserted { + if force { + _, err = s.setStatusTx(txCtx, StatusChange{ + Identity: identity, + Status: SpaceStatusCreated, + SpaceId: spaceId, + }, entry.Status) + if err != nil { + return err + } + } else { + return coordinatorproto.ErrSpaceIsDeleted + } + } + if err = s.checkLimitTx(txCtx, identity); err != nil { + return err + } + count, err := s.spaces.CountDocuments(txCtx, bson.D{ + {"billedIdentity", billedIdentity}, + {"status", SpaceStatusCreated}, + }) + if err != nil { + return err + } + if uint32(count) > billedLimit { + return coordinatorproto.ErrSpaceLimitReached + } + return nil + }) +} + func (s *spaceStatus) MakeShareable(ctx context.Context, spaceId string, spaceType SpaceType, limit uint32) (err error) { return s.db.Tx(ctx, func(txCtx mongo.SessionContext) error { entry, err := s.Status(txCtx, spaceId) @@ -576,10 +663,20 @@ func (s *spaceStatus) checkLimitTx(txCtx mongo.SessionContext, identity crypto.P } func (s *spaceStatus) ChangeOwner(ctx context.Context, spaceId, ownerId string) (err error) { - _, err = s.spaces.UpdateOne(ctx, bson.D{{"_id", spaceId}}, bson.D{{"$set", bson.D{ - {"identity", ownerId}, - }}}) - return + return s.db.Tx(ctx, func(txCtx mongo.SessionContext) error { + if _, err := s.spaces.UpdateOne(txCtx, bson.D{{"_id", spaceId}}, bson.D{{"$set", bson.D{ + {"identity", ownerId}, + }}}); err != nil { + return err + } + // children consume the legalOwner's quota, so their billing pool follows the parent's + // owner: left behind, the old and new owners would hold two disjoint pools and every + // transfer would multiply the per-owner children allowance + _, err := s.spaces.UpdateMany(txCtx, bson.D{ + {"parentSpaceId", spaceId}, + }, bson.D{{"$set", bson.D{{"billedIdentity", ownerId}}}}) + return err + }) } func (s *spaceStatus) Init(a *app.App) (err error) { diff --git a/spacestatus/spacestatus_test.go b/spacestatus/spacestatus_test.go index f418389..6030ccd 100644 --- a/spacestatus/spacestatus_test.go +++ b/spacestatus/spacestatus_test.go @@ -851,6 +851,45 @@ func TestSpaceStatus_ChangeOwner(t *testing.T) { assert.Equal(t, newIdentity.Account(), status.Identity) } +func TestSpaceStatus_ChildBillingFollowsOwnershipTransfer(t *testing.T) { + fx := newFixture(t, 0, 0) + defer fx.Finish(t) + + _, creator, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + _, ownerA, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + _, ownerB, err := crypto.GenerateRandomEd25519KeyPair() + require.NoError(t, err) + + const parentId = "parent.id" + require.NoError(t, fx.NewStatus(ctx, parentId, ownerA, SpaceTypeRegular, false)) + + // two children fill A's pool of 2 + require.NoError(t, fx.NewChildStatus(ctx, "child.1", creator, parentId, ownerA.Account(), 2, SpaceTypeRegular, false)) + require.NoError(t, fx.NewChildStatus(ctx, "child.2", creator, parentId, ownerA.Account(), 2, SpaceTypeRegular, false)) + require.ErrorIs(t, fx.NewChildStatus(ctx, "child.3", creator, parentId, ownerA.Account(), 2, SpaceTypeRegular, false), + coordinatorproto.ErrSpaceLimitReached) + + // parent ownership A -> B migrates the children's billing pool with it: B does + // not start with an empty pool, so a transfer cannot multiply the allowance + require.NoError(t, fx.ChangeOwner(ctx, parentId, ownerB.Account())) + for _, id := range []string{"child.1", "child.2"} { + entry, err := fx.Status(ctx, id) + require.NoError(t, err) + assert.Equal(t, ownerB.Account(), entry.BilledIdentity, id) + } + require.ErrorIs(t, fx.NewChildStatus(ctx, "child.3", creator, parentId, ownerB.Account(), 2, SpaceTypeRegular, false), + coordinatorproto.ErrSpaceLimitReached) + + // a re-sign of an existing child refreshes a stale pool assignment to the + // billed identity the caller resolved from the parent's current owner + require.NoError(t, fx.NewChildStatus(ctx, "child.1", creator, parentId, ownerA.Account(), 2, SpaceTypeRegular, false)) + entry, err := fx.Status(ctx, "child.1") + require.NoError(t, err) + assert.Equal(t, ownerA.Account(), entry.BilledIdentity) +} + type fixture struct { SpaceStatus a *app.App