Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
474f68b34c | ||
|
|
1a724f20fe | ||
|
|
0da0bb2157 | ||
|
|
118e333ff8 |
@@ -40,7 +40,14 @@ type Consumer struct {
|
|||||||
AckWaitSeconds int
|
AckWaitSeconds int
|
||||||
// MaxDeliver before the message is dead-lettered; zero for the mesh's default.
|
// MaxDeliver before the message is dead-lettered; zero for the mesh's default.
|
||||||
MaxDeliver int
|
MaxDeliver int
|
||||||
Why string
|
// MaxAckPending is how many deliveries the server lets stand unacknowledged at once; zero for
|
||||||
|
// the server's default, which is many. **One, for a consumer handled one at a time**
|
||||||
|
// (novox/hq issue 175): a handler that builds for minutes keeps its own message alive with a
|
||||||
|
// heartbeat, but everything handed over behind it times out unacknowledged and comes back —
|
||||||
|
// and a merge that came back rebuilt what it had just built, five times over on 2026-09-30.
|
||||||
|
// With one outstanding, the server holds the rest, and the heartbeat is keeping the message.
|
||||||
|
MaxAckPending int
|
||||||
|
Why string
|
||||||
}
|
}
|
||||||
|
|
||||||
// seatStreamName is the stream holding a seat's inbound work. Named after the seat rather than
|
// seatStreamName is the stream holding a seat's inbound work. Named after the seat rather than
|
||||||
|
|||||||
@@ -179,6 +179,7 @@ func (j *JetStream) EnsureConsumer(c Consumer) error {
|
|||||||
AckPolicy: nats.AckExplicitPolicy,
|
AckPolicy: nats.AckExplicitPolicy,
|
||||||
AckWait: time.Duration(c.AckWaitSeconds) * time.Second,
|
AckWait: time.Duration(c.AckWaitSeconds) * time.Second,
|
||||||
MaxDeliver: c.MaxDeliver,
|
MaxDeliver: c.MaxDeliver,
|
||||||
|
MaxAckPending: c.MaxAckPending,
|
||||||
DeliverGroup: c.Queue,
|
DeliverGroup: c.Queue,
|
||||||
DeliverSubject: "",
|
DeliverSubject: "",
|
||||||
Description: c.Why,
|
Description: c.Why,
|
||||||
|
|||||||
@@ -244,8 +244,13 @@ func MeshConsumers() []Consumer {
|
|||||||
Push: true,
|
Push: true,
|
||||||
AckWaitSeconds: 30,
|
AckWaitSeconds: 30,
|
||||||
MaxDeliver: 5,
|
MaxDeliver: 5,
|
||||||
Why: "the two events the mesh's own controller reacts to; after max-deliver it " +
|
// One at a time (novox/hq issue 175): acting on a merge builds for minutes, and an
|
||||||
"dead-letters, because an announcement it cannot act on will not become actionable",
|
// announcement handed over behind it must wait on the server, not time out on the
|
||||||
|
// client and come back to be acted on again.
|
||||||
|
MaxAckPending: 1,
|
||||||
|
Why: "the two events the mesh's own controller reacts to, one at a time; after " +
|
||||||
|
"max-deliver it dead-letters, because an announcement it cannot act on will not " +
|
||||||
|
"become actionable",
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -260,3 +260,14 @@ func TestNoTwoConsumersDeliverOntoTheSameSubject(t *testing.T) {
|
|||||||
seen[subject] = c.Name + " on " + c.Stream
|
seen[subject] = c.Name + " on " + c.Stream
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// The controller's events consumer is handed one announcement at a time (novox/hq issue 175): a
|
||||||
|
// merge's handler builds for minutes, and what is queued behind it must wait on the server rather
|
||||||
|
// than time out on the client and be acted on twice.
|
||||||
|
func TestTheControllerTakesOneAnnouncementAtATime(t *testing.T) {
|
||||||
|
for _, c := range MeshConsumers() {
|
||||||
|
if c.Stream == "EVENTS" && c.Name == ControllerName && c.MaxAckPending != 1 {
|
||||||
|
t.Fatalf("the events consumer may have %d outstanding; one announcement at a time", c.MaxAckPending)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+14
-1
@@ -5,7 +5,20 @@
|
|||||||
-- than for the mesh; ADR 0121 decided the rename and deferred it because a delivering seat that stops
|
-- than for the mesh; ADR 0121 decided the rename and deferred it because a delivering seat that stops
|
||||||
-- resolving mid-flight takes a provision away from every consumer. ADR 0122 removed that risk: a seat's
|
-- resolving mid-flight takes a provision away from every consumer. ADR 0122 removed that risk: a seat's
|
||||||
-- former name is an alias that resolves to it forever, a held record follows the rename by cascade, and
|
-- former name is an alias that resolves to it forever, a held record follows the rename by cascade, and
|
||||||
-- a claim written with the old name still holds. So the rename is one update and one alias.
|
-- a claim written with the old name still holds.
|
||||||
|
--
|
||||||
|
-- **Both rows may exist when this runs.** A controller whose compiled defaults already carry the new
|
||||||
|
-- name seeds it as a new seat the moment it can, and on the mesh this was written for that happened
|
||||||
|
-- before the rename: the first form of this migration renamed into a duplicate key and the control
|
||||||
|
-- node's prepare failed on every attempt (2026-09-30). So: if the new row is already there, the old
|
||||||
|
-- row's holding moves to it and the old row goes; otherwise the old row is renamed. Either way the old
|
||||||
|
-- name becomes an alias.
|
||||||
|
update seat_holding set seat = 'mesh-artifact-store'
|
||||||
|
where seat = 'the-artifact-store'
|
||||||
|
and exists (select 1 from seat where name = 'mesh-artifact-store');
|
||||||
|
delete from seat
|
||||||
|
where name = 'the-artifact-store'
|
||||||
|
and exists (select 1 from seat where name = 'mesh-artifact-store');
|
||||||
update seat set name = 'mesh-artifact-store' where name = 'the-artifact-store';
|
update seat set name = 'mesh-artifact-store' where name = 'the-artifact-store';
|
||||||
insert into seat_alias (alias, seat) values ('the-artifact-store', 'mesh-artifact-store')
|
insert into seat_alias (alias, seat) values ('the-artifact-store', 'mesh-artifact-store')
|
||||||
on conflict (alias) do update set seat = excluded.seat;
|
on conflict (alias) do update set seat = excluded.seat;
|
||||||
|
|||||||
Reference in New Issue
Block a user