Skip to content

Commit 0c083cd

Browse files
omer9564claude
andcommitted
feat(natsstore): serve all tenants from a single muxed K/V bucket
Replace the bucket-per-tenant model with one muxed bucket (config.bucket), keyed "<tenant>.<key...>". The Rego interface is unchanged: data still lands at data.nats.kv.<tenant>.<...> and the builtins (nats.kv.watch_bucket / nats.kv.get_data) keep their names and arity — their first argument is now the tenant token (the leading key segment) instead of a bucket name. Core changes: - config: add required `bucket`; rename `root_bucket` -> `root_tenant` (a tenant whose subtree mounts at the OPA data root). - NATSKeyToOPAPath: derive the tenant from the key's first token (drop the out-of-band bucket arg); root mode strips the tenant token. - nats_client: open the single bucket once (cached); add tenantKeys(), a prefix-filtered lister so loads are O(tenant), never O(whole bucket). - watcher: watch "<tenant>.>" (single filter -> scopeable consumer) instead of ">"; identity is the tenant. - builtins: loadTenantAsGJSON reads only the tenant slice and strips the "<tenant>." prefix via the new pure buildTenantJSON helper. Tests (TDD for the pure logic): - reshaped path-mapping/config/JSON-builder unit tests; - new muxed_integration_test.go (real NATS): tenantKeys isolation, data placement, root-tenant stripping; - migrated the example (one DATA bucket, <tenant>.* keys) + configs; the docker-compose integration test passes end to end. No backward compatibility retained (intentional). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent 24867e8 commit 0c083cd

16 files changed

Lines changed: 424 additions & 195 deletions

.gitignore

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,3 @@
1+
2+
# build artifact (compiled binary)
3+
/opa-nats

README.md

Lines changed: 24 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -41,12 +41,14 @@ The plugin can be configured through OPA's configuration system. Here's an examp
4141
plugins:
4242
nats:
4343
server_url: "nats://localhost:4222"
44+
bucket: "DATA" # REQUIRED: the single muxed K/V bucket holding every tenant
45+
# as keys "<tenant>.<key...>"
4446
ttl: "10m"
4547
refresh_interval: "30s"
4648
max_reconnect_attempts: 10
4749
reconnect_wait: "2s"
48-
max_bucket_watchers: 10
49-
root_bucket: "" # Optional - leave empty for multi-bucket mode
50+
max_bucket_watchers: 10 # LRU cap on concurrent per-tenant watchers
51+
root_tenant: "" # Optional - a tenant whose subtree mounts at the OPA data root
5052

5153
# Authentication (choose one)
5254
credentials: "/path/to/nats.creds"
@@ -89,28 +91,35 @@ services:
8991

9092
## Built-in Functions
9193

92-
The plugin provides custom Rego built-in functions for interacting with NATS K/V:
94+
The plugin provides custom Rego built-in functions for interacting with NATS K/V.
9395

94-
### `nats.kv.watch_bucket(bucket_name)`
96+
> **Single muxed bucket model.** All tenants live in one K/V bucket (`bucket` in
97+
> config), keyed as `<tenant>.<key...>`. The first argument to both builtins is
98+
> the **tenant id** — the plugin reads only that tenant's slice (a prefix-filtered
99+
> watch), never the whole bucket. The data lands at `data.nats.kv.<tenant>.<...>`,
100+
> the same place as before, so policies are unchanged.
95101

96-
Watches a NATS K/V bucket and returns all its data. This function automatically manages bucket watchers with LRU caching.
102+
### `nats.kv.watch_bucket(tenant_id)`
103+
104+
Watches a single tenant's slice of the muxed bucket and returns its data. Watchers
105+
are managed with LRU caching (one ordered consumer per watched tenant).
97106

98107
```rego
99-
# Watch a specific bucket
108+
# Watch one tenant
100109
group_data := nats.kv.watch_bucket("550e8400-e29b-41d4-a716-446655440000")
101110
102111
# Access nested data
103112
members := group_data.members
104113
permissions := group_data.permissions
105114
```
106115

107-
### `nats.kv.get_data(bucket_name, key)`
116+
### `nats.kv.get_data(tenant_id, key)`
108117

109-
Retrieves a specific key from a NATS K/V bucket.
118+
Retrieves a specific key from a tenant's slice (key is tenant-relative).
110119

111120
```rego
112-
# Get specific data from a bucket
113-
user_data := nats.kv.get_data("users", "123e4567-e89b-12d3-a456-426614174000")
121+
# Get specific data for a tenant
122+
user_data := nats.kv.get_data("550e8400-e29b-41d4-a716-446655440000", "members")
114123
```
115124

116125
## Examples
@@ -276,8 +285,11 @@ accounts: {
276285
jetstream: enabled
277286
users: [
278287
{user: "opa", pass: "secret", permissions: {
279-
subscribe: ["$JS.API.>", "$KV.>"]
280-
publish: ["$JS.API.>", "$KV.>"]
288+
# Scope to the single muxed bucket (DATA). The plugin is a multi-tenant
289+
# cloud reader living in the shared account, so it sees all tenants;
290+
# narrow further per-deployment if desired.
291+
subscribe: ["$JS.API.>", "$KV.DATA.>"]
292+
publish: ["$JS.API.>", "$KV.DATA.>"]
281293
}}
282294
]
283295
}

examples/opa-nats/config-compose.yaml

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,13 @@
44
plugins:
55
nats:
66
server_url: "nats://nats:4222"
7+
bucket: "DATA"
78
ttl: "10m"
89
refresh_interval: "30s"
910
max_reconnect_attempts: 10
1011
reconnect_wait: "2s"
1112
max_bucket_watchers: 1
12-
root_bucket: "permit-groups-default"
13+
root_tenant: "permit-groups-default"
1314

1415
services:
1516
authz:

examples/opa-nats/config.yaml

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,12 +6,13 @@ plugins:
66
# OPA runs inside the compose network and needs to reach NATS via the
77
# docker service name `nats`.
88
server_url: "nats://localhost:4222"
9+
bucket: "DATA"
910
ttl: "10m"
1011
refresh_interval: "30s"
1112
max_reconnect_attempts: 10
1213
reconnect_wait: "2s"
1314
max_bucket_watchers: 1
14-
root_bucket: "permit-groups-default"
15+
root_tenant: "permit-groups-default"
1516

1617
services:
1718
authz:

examples/opa-nats/docker-compose.yaml

Lines changed: 36 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -59,67 +59,50 @@ services:
5959
- NATS_URL=nats://nats:4222
6060
command: >
6161
/bin/sh -c "
62-
echo 'Setting up NATS K/V buckets for UUID-based groups (each group UUID = separate bucket)...';
63-
64-
# Create individual buckets for each group UUID
65-
echo 'Creating group UUID buckets...';
66-
nats kv add 550e8400-e29b-41d4-a716-446655440000 --replicas=1 --history=10 --ttl=12h;
67-
nats kv add 660e8400-e29b-41d4-a716-446655440001 --replicas=1 --history=10 --ttl=12h;
68-
nats kv add 770e8400-e29b-41d4-a716-446655440002 --replicas=1 --history=10 --ttl=12h;
69-
70-
# Create default bucket for non-group data (users)
71-
nats kv add permit-groups-default --replicas=1 --history=10 --ttl=12h;
72-
73-
echo 'All buckets created successfully';
74-
75-
echo 'Populating Administrators group (bucket: 550e8400-e29b-41d4-a716-446655440000)...';
76-
# In this bucket, store data without the groups.{uuid} prefix since the bucket IS the group
77-
nats kv put 550e8400-e29b-41d4-a716-446655440000 metadata '{\"id\":\"550e8400-e29b-41d4-a716-446655440000\",\"name\":\"Administrators\",\"description\":\"System administrators group\",\"created_at\":\"2024-01-15T10:00:00Z\"}';
78-
nats kv put 550e8400-e29b-41d4-a716-446655440000 members.123e4567-e89b-12d3-a456-426614174000 '{\"user_id\":\"123e4567-e89b-12d3-a456-426614174000\",\"joined_at\":\"2024-01-15T10:05:00Z\",\"role\":\"admin\"}';
79-
nats kv put 550e8400-e29b-41d4-a716-446655440000 members.234e5678-e89b-12d3-a456-426614174001 '{\"user_id\":\"234e5678-e89b-12d3-a456-426614174001\",\"joined_at\":\"2024-01-15T10:10:00Z\",\"role\":\"member\"}';
80-
nats kv put 550e8400-e29b-41d4-a716-446655440000 permissions.perm-admin-api '{\"id\":\"perm-admin-api\",\"resource\":\"api\",\"action\":\"admin\",\"effect\":\"allow\"}';
81-
nats kv put 550e8400-e29b-41d4-a716-446655440000 permissions.perm-admin-users '{\"id\":\"perm-admin-users\",\"resource\":\"users\",\"action\":\"manage\",\"effect\":\"allow\"}';
82-
83-
echo 'Populating Developers group (bucket: 660e8400-e29b-41d4-a716-446655440001)...';
84-
nats kv put 660e8400-e29b-41d4-a716-446655440001 metadata '{\"id\":\"660e8400-e29b-41d4-a716-446655440001\",\"name\":\"Developers\",\"description\":\"Development team group\",\"created_at\":\"2024-01-15T10:30:00Z\"}';
85-
nats kv put 660e8400-e29b-41d4-a716-446655440001 members.345e6789-e89b-12d3-a456-426614174002 '{\"user_id\":\"345e6789-e89b-12d3-a456-426614174002\",\"joined_at\":\"2024-01-15T10:35:00Z\",\"role\":\"lead\"}';
86-
nats kv put 660e8400-e29b-41d4-a716-446655440001 members.234e5678-e89b-12d3-a456-426614174001 '{\"user_id\":\"234e5678-e89b-12d3-a456-426614174001\",\"joined_at\":\"2024-01-15T10:40:00Z\",\"role\":\"member\"}';
87-
nats kv put 660e8400-e29b-41d4-a716-446655440001 permissions.perm-dev-code '{\"id\":\"perm-dev-code\",\"resource\":\"code\",\"action\":\"write\",\"effect\":\"allow\"}';
88-
nats kv put 660e8400-e29b-41d4-a716-446655440001 permissions.perm-dev-deploy '{\"id\":\"perm-dev-deploy\",\"resource\":\"deployment\",\"action\":\"deploy\",\"effect\":\"allow\"}';
89-
90-
echo 'Populating ReadOnly Users group (bucket: 770e8400-e29b-41d4-a716-446655440002)...';
91-
nats kv put 770e8400-e29b-41d4-a716-446655440002 metadata '{\"id\":\"770e8400-e29b-41d4-a716-446655440002\",\"name\":\"ReadOnly Users\",\"description\":\"Read-only access group\",\"created_at\":\"2024-01-15T11:00:00Z\"}';
92-
nats kv put 770e8400-e29b-41d4-a716-446655440002 members.456e7890-e89b-12d3-a456-426614174003 '{\"user_id\":\"456e7890-e89b-12d3-a456-426614174003\",\"joined_at\":\"2024-01-15T11:05:00Z\",\"role\":\"member\"}';
93-
nats kv put 770e8400-e29b-41d4-a716-446655440002 permissions.perm-readonly-data '{\"id\":\"perm-readonly-data\",\"resource\":\"data\",\"action\":\"read\",\"effect\":\"allow\"}';
94-
95-
echo 'Populating default bucket with user group assignments...';
96-
# User-group assignments go to the default bucket since they are not group-specific
97-
nats kv put permit-groups-default metadata '{\"id\":\"permit-groups-default\",\"name\":\"Default Group\",\"description\":\"Default group for non-grouped data\",\"created_at\":\"2024-01-15T10:00:00Z\"}';
98-
99-
nats kv put permit-groups-default users '[\"550e8400-e29b-41d4-a716-446655440000\",\"660e8400-e29b-41d4-a716-446655440001\"]';
100-
# nats kv put permit-groups-default users.234e5678-e89b-12d3-a456-426614174001.groups '[\"150e8400-e29b-41d4-a716-446655440000\",\"660e8400-e29b-41d4-a716-446655440001\"]';
101-
# nats kv put permit-groups-default users.345e6789-e89b-12d3-a456-426614174002.groups '[\"160e8400-e29b-41d4-a716-446655440001\"]';
102-
# nats kv put permit-groups-default users.456e7890-e89b-12d3-a456-426614174003.groups '[\"170e8400-e29b-41d4-a716-446655440002\"]';
103-
# nats kv put permit-groups-default users.123e4567-e89b-12d3-a456-426614174000.groups '[\"150e8400-e29b-41d4-a716-446655440000\"]';
62+
echo 'Setting up the single muxed NATS K/V bucket DATA (keys are <tenant>.<key>)...';
63+
64+
# ONE bucket holds every tenant; the group UUID is the leading key token.
65+
nats kv add DATA --replicas=1 --history=10 --ttl=12h;
66+
echo 'Bucket DATA created successfully';
67+
68+
echo 'Populating Administrators group (tenant 550e8400-e29b-41d4-a716-446655440000)...';
69+
nats kv put DATA 550e8400-e29b-41d4-a716-446655440000.metadata '{\"id\":\"550e8400-e29b-41d4-a716-446655440000\",\"name\":\"Administrators\",\"description\":\"System administrators group\",\"created_at\":\"2024-01-15T10:00:00Z\"}';
70+
nats kv put DATA 550e8400-e29b-41d4-a716-446655440000.members.123e4567-e89b-12d3-a456-426614174000 '{\"user_id\":\"123e4567-e89b-12d3-a456-426614174000\",\"joined_at\":\"2024-01-15T10:05:00Z\",\"role\":\"admin\"}';
71+
nats kv put DATA 550e8400-e29b-41d4-a716-446655440000.members.234e5678-e89b-12d3-a456-426614174001 '{\"user_id\":\"234e5678-e89b-12d3-a456-426614174001\",\"joined_at\":\"2024-01-15T10:10:00Z\",\"role\":\"member\"}';
72+
nats kv put DATA 550e8400-e29b-41d4-a716-446655440000.permissions.perm-admin-api '{\"id\":\"perm-admin-api\",\"resource\":\"api\",\"action\":\"admin\",\"effect\":\"allow\"}';
73+
nats kv put DATA 550e8400-e29b-41d4-a716-446655440000.permissions.perm-admin-users '{\"id\":\"perm-admin-users\",\"resource\":\"users\",\"action\":\"manage\",\"effect\":\"allow\"}';
74+
75+
echo 'Populating Developers group (tenant 660e8400-e29b-41d4-a716-446655440001)...';
76+
nats kv put DATA 660e8400-e29b-41d4-a716-446655440001.metadata '{\"id\":\"660e8400-e29b-41d4-a716-446655440001\",\"name\":\"Developers\",\"description\":\"Development team group\",\"created_at\":\"2024-01-15T10:30:00Z\"}';
77+
nats kv put DATA 660e8400-e29b-41d4-a716-446655440001.members.345e6789-e89b-12d3-a456-426614174002 '{\"user_id\":\"345e6789-e89b-12d3-a456-426614174002\",\"joined_at\":\"2024-01-15T10:35:00Z\",\"role\":\"lead\"}';
78+
nats kv put DATA 660e8400-e29b-41d4-a716-446655440001.members.234e5678-e89b-12d3-a456-426614174001 '{\"user_id\":\"234e5678-e89b-12d3-a456-426614174001\",\"joined_at\":\"2024-01-15T10:40:00Z\",\"role\":\"member\"}';
79+
nats kv put DATA 660e8400-e29b-41d4-a716-446655440001.permissions.perm-dev-code '{\"id\":\"perm-dev-code\",\"resource\":\"code\",\"action\":\"write\",\"effect\":\"allow\"}';
80+
nats kv put DATA 660e8400-e29b-41d4-a716-446655440001.permissions.perm-dev-deploy '{\"id\":\"perm-dev-deploy\",\"resource\":\"deployment\",\"action\":\"deploy\",\"effect\":\"allow\"}';
81+
82+
echo 'Populating ReadOnly Users group (tenant 770e8400-e29b-41d4-a716-446655440002)...';
83+
nats kv put DATA 770e8400-e29b-41d4-a716-446655440002.metadata '{\"id\":\"770e8400-e29b-41d4-a716-446655440002\",\"name\":\"ReadOnly Users\",\"description\":\"Read-only access group\",\"created_at\":\"2024-01-15T11:00:00Z\"}';
84+
nats kv put DATA 770e8400-e29b-41d4-a716-446655440002.members.456e7890-e89b-12d3-a456-426614174003 '{\"user_id\":\"456e7890-e89b-12d3-a456-426614174003\",\"joined_at\":\"2024-01-15T11:05:00Z\",\"role\":\"member\"}';
85+
nats kv put DATA 770e8400-e29b-41d4-a716-446655440002.permissions.perm-readonly-data '{\"id\":\"perm-readonly-data\",\"resource\":\"data\",\"action\":\"read\",\"effect\":\"allow\"}';
86+
87+
echo 'Populating the root tenant (permit-groups-default mounts at the OPA data root)...';
88+
nats kv put DATA permit-groups-default.metadata '{\"id\":\"permit-groups-default\",\"name\":\"Default Group\",\"description\":\"Default group for non-grouped data\",\"created_at\":\"2024-01-15T10:00:00Z\"}';
89+
nats kv put DATA permit-groups-default.users '[\"550e8400-e29b-41d4-a716-446655440000\",\"660e8400-e29b-41d4-a716-446655440001\"]';
10490
10591
echo 'Test data added successfully';
10692
echo '';
107-
echo 'Created buckets:';
93+
echo 'Bucket:';
10894
nats kv list;
10995
echo '';
11096
echo 'Sample verification:';
111-
echo 'Administrators group metadata (bucket 550e8400-e29b-41d4-a716-446655440000):';
112-
nats kv get 550e8400-e29b-41d4-a716-446655440000 metadata;
113-
echo '';
114-
echo 'Developers group members (bucket 660e8400-e29b-41d4-a716-446655440001):';
115-
nats kv ls 660e8400-e29b-41d4-a716-446655440001 | grep members;
97+
echo 'Administrators group metadata (tenant 550e8400-...):';
98+
nats kv get DATA 550e8400-e29b-41d4-a716-446655440000.metadata;
11699
echo '';
117-
echo 'User group assignments (default bucket):';
118-
nats kv get permit-groups-default users.234e5678-e89b-12d3-a456-426614174001.groups;
100+
echo 'Developers group keys (tenant 660e8400-...):';
101+
nats kv ls DATA | grep 660e8400-e29b-41d4-a716-446655440001;
119102
echo '';
120-
echo 'Setup complete - ready for bucket-per-group UUID testing';
121-
echo 'Architecture: Each group UUID is its own NATS bucket containing that groups metadata, members, and permissions';
122-
echo 'User group assignments stored in default bucket since they span multiple groups';
103+
echo 'Setup complete - single muxed bucket DATA, keys <tenant>.<key>';
104+
echo 'Architecture: one bucket; each group UUID is the leading key token (the tenant);';
105+
echo 'the plugin reads only a tenant slice via a prefix-filtered watch.';
123106
"
124107
networks:
125108
- nats-net

examples/opa-nats/rego/test.rego

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,10 @@ package test
22

33
import rego.v1
44

5+
# input.bucket_id identifies the tenant whose slice of the single muxed bucket
6+
# we want. The Rego interface is unchanged from the bucket-per-tenant model: the
7+
# value is now the tenant token, and its data is still injected at
8+
# data.nats.kv.<bucket_id> (same place as before).
59
default bucket_watched := false
610

711
bucket_watched := nats.kv.watch_bucket(input.bucket_id)

pkg/natsstore/additional_test.go

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -164,6 +164,7 @@ func TestPluginFactory_Validate_EdgeCases(t *testing.T) {
164164
name: "valid minimal config",
165165
configData: map[string]interface{}{
166166
"server_url": "nats://localhost:4222",
167+
"bucket": "DATA",
167168
},
168169
expectError: false,
169170
},
@@ -189,7 +190,8 @@ func TestPluginFactory_Validate_EdgeCases(t *testing.T) {
189190
"max_reconnect_attempts": 5,
190191
"reconnect_wait": "1s",
191192
"max_bucket_watchers": 20,
192-
"root_bucket": "test-root",
193+
"bucket": "DATA",
194+
"root_tenant": "test-root",
193195
"username": "testuser",
194196
"password": "testpass",
195197
},

pkg/natsstore/bucket_watcher_manager.go

Lines changed: 10 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -66,20 +66,22 @@ func (gw *BucketWatcher) Start() error {
6666
return nil
6767
}
6868

69-
// Get or create the bucket for this bucket
70-
kv, err := gw.natsClient.getBucket(gw.bucketName)
69+
// Open the single muxed bucket (handle is cached).
70+
kv, err := gw.natsClient.getBucket()
7171
if err != nil {
72-
return fmt.Errorf("failed to get bucket %s: %w", gw.bucketName, err)
72+
return fmt.Errorf("failed to get bucket: %w", err)
7373
}
7474
if err := gw.dataTransformer.LoadBucketDataBulk(gw.ctx, gw.bucketName, gw.natsClient, gw.opaStore, gw.isRoot); err != nil {
75-
return fmt.Errorf("failed to load bucket data for %s: %w", gw.bucketName, err)
75+
return fmt.Errorf("failed to load data for tenant %s: %w", gw.bucketName, err)
7676
}
7777

78-
// Create watcher for all keys in this bucket
79-
watchPattern := ">" // Watch all keys in this bucket
78+
// Watch ONLY this tenant's slice. gw.bucketName is the tenant token; the
79+
// single-filter "<tenant>.>" maps to the scopeable extended consumer form
80+
// (one ordered consumer per watched tenant), never the whole bucket.
81+
watchPattern := gw.bucketName + ".>"
8082
watcher, err := kv.Watch(watchPattern, nats.Context(gw.ctx))
8183
if err != nil {
82-
return fmt.Errorf("failed to create watcher for bucket %s: %w", gw.bucketName, err)
84+
return fmt.Errorf("failed to create watcher for tenant %s: %w", gw.bucketName, err)
8385
}
8486

8587
gw.watcher = watcher
@@ -252,7 +254,7 @@ func NewBucketWatcherManager(natsClient *NATSClient, maxWatchers int, logger log
252254
dataTransformer: dataTransformer,
253255
logger: logger,
254256
maxWatchers: maxWatchers,
255-
rootBucket: config.RootBucket,
257+
rootBucket: config.RootTenant,
256258
}
257259
manager.newWatcher = func(bucketName string, opaStore storage.Store) (*BucketWatcher, error) {
258260
w, err := NewBucketWatcher(bucketName, manager.natsClient, manager.logger, manager.dataTransformer, opaStore, false)

0 commit comments

Comments
 (0)