Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 18 additions & 3 deletions src/core/mesh-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -373,9 +373,24 @@ export class MeshStore implements CommsStore {
for (const event of events) {
if (seen.has(JSON.stringify(event))) continue;
this.queueDelivery(agentId, event);
// A returning peer replays its own pending queue: events pushed
// while its process was down fire onDelivery now (#28).
this.fireLocalDelivery(agentId, event);
// Replay fires only for events with consumption evidence (#28):
// messages are consumed by reading (readBy), invites by acceptance
// or decline (no longer in the invited list). Transient
// notifications (member_joined, room_members, connection_request,
// delivery status) carry no consumption evidence, so replaying
// them could only ever duplicate-notify; the state they describe
// arrives via the synced room and agent records instead. They
// still merge into the queue, so drain bridges see them.
if (event.type === "room_message" || event.type === "dm") {
this.fireLocalDelivery(agentId, event);
} else if (event.type === "room_invite") {
const stillInvited = this.rooms
.get(event.room)
?.invited.includes(agentId);
if (stillInvited === true) {
this.fireLocalDelivery(agentId, event);
}
}
}
}
}
Expand Down
63 changes: 60 additions & 3 deletions src/test/downtime-replay.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -70,8 +70,15 @@ void test("events queued while the target was down replay on its first snapshot"
true,
);

// A second snapshot of the same state does not duplicate any push (the
// join notification replays too, exactly once).
// Transient notifications carry no consumption evidence, so they merge into the queue (drain bridges still see them) but never replay-fire.
assert.equal(deliveries.filter((ev) => ev.type === "room_members").length, 0);
const queue = snapshotOf(returned).deliveryQueues[targetAgent.id] ?? [];
assert.equal(
queue.some((ev) => ev.type === "room_members"),
true,
);

// A second snapshot of the same state does not duplicate the push.
returned.applyStateSync(snapshotOf(sender));
assert.equal(
deliveries.filter(
Expand All @@ -81,7 +88,6 @@ void test("events queued while the target was down replay on its first snapshot"
).length,
1,
);
assert.equal(deliveries.filter((ev) => ev.type === "room_members").length, 1);
});

void test("a queue is bounded oldest-first so downtime cannot grow it without limit", async () => {
Expand Down Expand Up @@ -123,3 +129,54 @@ void test("a queue is bounded oldest-first so downtime cannot grow it without li
true,
);
});

void test("a pending invite replays until accepted or declined", async () => {
const sender = makeStore();
const author = await sender.registerAgent({
name: "owner",
harness: "pi",
cwd: "/tmp/p",
pid: process.pid,
visibility: "visible",
tags: [],
});
const target = makeStore();
const targetAgent = await target.registerAgent({
name: "invitee",
harness: "claude-code",
cwd: "/tmp/t",
pid: process.pid,
visibility: "visible",
tags: [],
});
sender.applyStateSync(target.serialise());
await sender.createRoom({
name: "private-room",
type: "private",
owner: author.id,
description: "x",
});
await sender.inviteToRoom("private-room", targetAgent.id, author.id);

const returned = makeStore();
const deliveries: DeliveryEvent[] = [];
returned.onDelivery = (_id, ev) => {
deliveries.push(ev);
};
returned.peerId = targetAgent.id;

// Still on the invited list: the invite replays.
returned.applyStateSync(snapshotOf(sender));
assert.equal(deliveries.filter((ev) => ev.type === "room_invite").length, 1);

// Declined (no longer invited): the same snapshot no longer replays it.
await sender.declineInvite("private-room", targetAgent.id, "not now");
const declined = makeStore();
const deliveries2: DeliveryEvent[] = [];
declined.onDelivery = (_id, ev) => {
deliveries2.push(ev);
};
declined.peerId = targetAgent.id;
declined.applyStateSync(snapshotOf(sender));
assert.equal(deliveries2.filter((ev) => ev.type === "room_invite").length, 0);
});