Skip to content
Merged
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,9 @@ export class GroupMembershipItemKind implements ItemKind<GroupMembershipGrant> {
readonly kind = "endpointGroupMembership";
readonly priority = PRIORITY_BANDS.membership;

// capacity is per-group (groupTable) but the key is per-endpoint; the group slot is gated by groupKeyMap.
readonly excludeFromAdmission = true;

#commands(node: ClientNode, localEndpoint: number) {
return node.endpoints.for(localEndpoint).commandsOf(GroupsClient);
}
Expand Down
7 changes: 6 additions & 1 deletion packages/node-manager/src/task/Task.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
*/

import { ImplementationError } from "@matter/general";
import { ChangeEntry, TaskPhase, TaskState, TaskStatus } from "./types.js";
import { ChangeEntry, PlannedChange, TaskPhase, TaskState, TaskStatus } from "./types.js";

export interface TaskPersistence {
type: string;
Expand Down Expand Up @@ -55,6 +55,11 @@ export abstract class Task<P = unknown> {
};
}

/** Intents this task will create, derived from params, for pre-flight capacity admission. Removals omit. */
plannedChanges(): PlannedChange[] {
return new Array<PlannedChange>();
}

/** Deterministic internal id from type + params. Subclasses override with their own key. */
static idFor(_params: unknown): string {
throw new ImplementationError("idFor must be implemented by the Task subclass");
Expand Down
52 changes: 49 additions & 3 deletions packages/node-manager/src/task/TaskManagerBehavior.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,14 +7,15 @@
import { ReconcilerBehavior } from "#ReconcilerBehavior.js";
import { Logger, Mutex, Observable } from "@matter/general";
import { DatatypeModel, FieldElement } from "@matter/model";
import { Agent, Behavior, ClientNode, Node, ServerNode } from "@matter/node";
import { TaskCancelledSignal, TaskSuspendedSignal } from "./errors.js";
import { Agent, Behavior, ClientNode, DesiredStateBehavior, itemMapKey, Node, ServerNode } from "@matter/node";
import { TaskCancelledSignal, TaskCapacityExceededError, TaskSuspendedSignal } from "./errors.js";
import { ADD_NODE_TO_GROUP_TYPE, AddNodeToGroup } from "./groups/AddNodeToGroup.js";
import { REMOVE_NODE_FROM_GROUP_TYPE, RemoveNodeFromGroup } from "./groups/RemoveNodeFromGroup.js";
import { Revert, REVERT_TYPE } from "./Revert.js";
import { GateControl, RunningTaskContext } from "./RunningTaskContext.js";
import { Task, TaskPersistence } from "./Task.js";
import { TaskCtor, TaskRegistry } from "./TaskRegistry.js";
import { TaskState, TaskStatus } from "./types.js";
import { PlannedChange, TaskState, TaskStatus } from "./types.js";

const TERMINAL_STATES: ReadonlySet<TaskState> = new Set<TaskState>(["completed", "failed", "cancelled"]);

Expand Down Expand Up @@ -68,6 +69,7 @@ export class TaskManagerBehavior extends Behavior {
/** Built-in task types registered before the resume pass. */
protected registerBuiltins(): void {
this.internal.registry.register(ADD_NODE_TO_GROUP_TYPE, AddNodeToGroup);
this.internal.registry.register(REMOVE_NODE_FROM_GROUP_TYPE, RemoveNodeFromGroup);
this.internal.registry.register(REVERT_TYPE, Revert);
}

Expand Down Expand Up @@ -213,8 +215,52 @@ export class TaskManagerBehavior extends Behavior {
gate.wake.emit();
}

/**
* Reject a task before any node mutation if its planned changes would overflow a target's device capacity.
* Runs before the first persist/phase; the thrown error ends the task `failed` with an empty changeSet.
*/
async #admit(task: Task): Promise<void> {
const planned = task.plannedChanges();
if (planned.length === 0) {
return;
}
const byNodeKind = new Map<string, PlannedChange[]>();
for (const pc of planned) {
const k = `${pc.peerId}\0${pc.kind}`;
let group = byNodeKind.get(k);
if (group === undefined) {
group = new Array<PlannedChange>();
byNodeKind.set(k, group);
}
group.push(pc);
}
for (const group of byNodeKind.values()) {
const { peerId, kind } = group[0];
const peer = this.resolvePeerNode(peerId);
if (peer === undefined) {
continue; // unresolvable peer: the phase gate will park; capacity is re-checked on device write
}
const itemKind = await this.endpoint.act(agent => this.taskReconciler(agent).itemKind(kind));
if (itemKind?.excludeFromAdmission) {
continue; // capacity counts a coarser resource another kind already gates (e.g. membership vs group)
}
const capacity = await itemKind?.capacity?.(peer);
if (capacity === undefined) {
continue; // kind reports no capacity limit (e.g. groupKey) — the device write is the gate
}
const items = peer.stateOf(DesiredStateBehavior).items;
const added = group.filter(pc => items[itemMapKey(pc.kind, pc.key)] === undefined).length;
if (capacity.used + added > capacity.limit) {
throw new TaskCapacityExceededError(
`Task ${task.id}: ${kind} on ${peerId} exceeds capacity — needs ${added} slot(s) but only ${capacity.limit - capacity.used} free`,
);
}
}
}

async #drive(task: Task): Promise<void> {
try {
await this.#admit(task); // fail-fast before any node is touched
// Persist before first phase so a crash-resume sees the task.
await this.#persist(task);
while (task.progress.phaseIndex < task.phases.length && task.progress.state === "running") {
Expand Down
3 changes: 3 additions & 0 deletions packages/node-manager/src/task/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@ export class TaskPeerUnavailableError extends TaskError {}
/** A task's forward work failed terminally; the manager spawns a revert task to roll back its changeSet. */
export class TaskFailedError extends TaskError {}

/** A task's planned changes would exceed a node's device capacity for some item kind. */
export class TaskCapacityExceededError extends TaskError {}

/** Internal signal a running gate throws when cancel is requested, so #drive stops cleanly (not "failed"). */
export class TaskCancelledSignal extends TaskError {}

Expand Down
41 changes: 32 additions & 9 deletions packages/node-manager/src/task/groups/AddNodeToGroup.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,8 @@
import { GroupId } from "@matter/types";
import { GroupKeyManagement } from "@matter/types/clusters/group-key-management";
import { Task } from "../Task.js";
import { TaskContext, TaskPhase } from "../types.js";
import { PlannedChange, TaskContext, TaskPhase } from "../types.js";
import { membershipKey } from "./keys.js";

export const ADD_NODE_TO_GROUP_TYPE = "addNodeToGroup";

Expand All @@ -31,19 +32,35 @@ export class AddNodeToGroup extends Task<AddNodeToGroupParams> {
readonly type = ADD_NODE_TO_GROUP_TYPE;

static override idFor(params: AddNodeToGroupParams): string {
return `${ADD_NODE_TO_GROUP_TYPE}:${params.peerId}:${params.groupId}`;
return `${ADD_NODE_TO_GROUP_TYPE}:${params.peerId}:${params.groupId}:${params.endpoint}`;
}

get phases(): TaskPhase[] {
return [{ name: "provision", run: ctx => this.#provision(ctx) }];
}

async #provision(ctx: TaskContext): Promise<void> {
override plannedChanges(): PlannedChange[] {
const p = this.params;
const peer = ctx.resolvePeer(p.peerId);
const groupId = GroupId(p.groupId);
return [
{ peerId: p.peerId, kind: "groupKey", key: String(p.groupKeySetId), intent: this.#keySet() },
{
peerId: p.peerId,
kind: "groupKeyMap",
key: String(p.groupId),
intent: { groupId: GroupId(p.groupId), groupKeySetId: p.groupKeySetId },
},
{
peerId: p.peerId,
kind: "endpointGroupMembership",
key: membershipKey(p.groupId, p.endpoint),
intent: { localEndpoint: p.endpoint, groupId: GroupId(p.groupId), groupName: p.groupName },
},
];
}

const keySet = {
#keySet() {
const p = this.params;
return {
groupKeySetId: p.groupKeySetId,
groupKeySecurityPolicy: p.groupKeySecurityPolicy,
epochKey0: p.epochKey0,
Expand All @@ -53,8 +70,14 @@ export class AddNodeToGroup extends Task<AddNodeToGroupParams> {
epochKey2: null,
epochStartTime2: null,
};
}

await ctx.setIntent(peer, "groupKey", String(p.groupKeySetId), keySet, "converge");
async #provision(ctx: TaskContext): Promise<void> {
const p = this.params;
const peer = ctx.resolvePeer(p.peerId);
const groupId = GroupId(p.groupId);

await ctx.setIntent(peer, "groupKey", String(p.groupKeySetId), this.#keySet(), "converge");
await ctx.setIntent(
peer,
"groupKeyMap",
Expand All @@ -65,15 +88,15 @@ export class AddNodeToGroup extends Task<AddNodeToGroupParams> {
await ctx.setIntent(
peer,
"endpointGroupMembership",
String(p.groupId),
membershipKey(p.groupId, p.endpoint),
{ localEndpoint: p.endpoint, groupId, groupName: p.groupName },
"converge",
);

await ctx.awaitCommitted([
{ peer, kind: "groupKey", key: String(p.groupKeySetId) },
{ peer, kind: "groupKeyMap", key: String(p.groupId) },
{ peer, kind: "endpointGroupMembership", key: String(p.groupId) },
{ peer, kind: "endpointGroupMembership", key: membershipKey(p.groupId, p.endpoint) },
]);
}
}
68 changes: 68 additions & 0 deletions packages/node-manager/src/task/groups/RemoveNodeFromGroup.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
/**
* @license
* Copyright 2022-2026 Matter.js Authors
* SPDX-License-Identifier: Apache-2.0
*/

import { ClientNode, DesiredStateBehavior, itemMapKey } from "@matter/node";
import { Task } from "../Task.js";
import { TaskContext, TaskPhase } from "../types.js";
import { membershipKey } from "./keys.js";

export const REMOVE_NODE_FROM_GROUP_TYPE = "removeNodeFromGroup";

export interface RemoveNodeFromGroupParams {
peerId: string;
endpoint: number;
groupId: number;
}

/**
* Removes a peer endpoint from a group: drops the membership, then the group-to-key-set map and the key set
* itself — but only while no other group still references them ({@link TaskContext.removeIntentIfUnreferenced}).
* Dependents-first (membership, then map, then key set) so each reference check sees the prior removal.
*/
export class RemoveNodeFromGroup extends Task<RemoveNodeFromGroupParams> {
readonly type = REMOVE_NODE_FROM_GROUP_TYPE;

static override idFor(p: RemoveNodeFromGroupParams): string {
return `${REMOVE_NODE_FROM_GROUP_TYPE}:${p.peerId}:${p.groupId}:${p.endpoint}`;
}

get phases(): TaskPhase[] {
return [{ name: "remove", run: ctx => this.#remove(ctx) }];
}

async #remove(ctx: TaskContext): Promise<void> {
const p = this.params;
const peer = ctx.tryResolvePeer(p.peerId);
if (peer === undefined) {
return; // decommissioned: intent is GC'd with the node
}

// The keyset id is unreadable once the map intent is gone, so capture it before removal.
const keySetId = this.#mappedKeySetId(peer, p.groupId);

const removed = new Array<{ kind: string; key: string }>();
if (
await ctx.removeIntentIfUnreferenced(peer, "endpointGroupMembership", membershipKey(p.groupId, p.endpoint))
) {
removed.push({ kind: "endpointGroupMembership", key: membershipKey(p.groupId, p.endpoint) });
}
if (await ctx.removeIntentIfUnreferenced(peer, "groupKeyMap", String(p.groupId))) {
removed.push({ kind: "groupKeyMap", key: String(p.groupId) });
}
if (keySetId !== undefined && (await ctx.removeIntentIfUnreferenced(peer, "groupKey", String(keySetId)))) {
removed.push({ kind: "groupKey", key: String(keySetId) });
}

if (removed.length > 0) {
await ctx.awaitGate([peer], () => removed.every(r => ctx.itemAbsent(peer, r.kind, r.key)));
}
}

#mappedKeySetId(peer: ClientNode, groupId: number): number | undefined {
const item = peer.stateOf(DesiredStateBehavior).items[itemMapKey("groupKeyMap", String(groupId))];
return (item?.intent as { groupKeySetId?: number } | undefined)?.groupKeySetId;
}
}
10 changes: 10 additions & 0 deletions packages/node-manager/src/task/groups/keys.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
/**
* @license
* Copyright 2022-2026 Matter.js Authors
* SPDX-License-Identifier: Apache-2.0
*/

/** Membership intent key: per (group, endpoint) so one peer can join a group on several endpoints. */
export function membershipKey(groupId: number, endpoint: number): string {
return `${groupId}:${endpoint}`;
}
8 changes: 8 additions & 0 deletions packages/node-manager/src/task/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,14 @@ export interface ChangeEntry {
prior?: { intent: unknown; mode: ItemMode };
}

/** An intent a task will create, derived from its params, for pre-flight capacity admission. */
export interface PlannedChange {
peerId: string;
kind: string;
key: string;
intent: unknown;
}

export interface TaskPhase {
name: string;
run(ctx: TaskContext): Promise<void>;
Expand Down
Loading