Grant a seat's traffic by caller and by kind, so an ask's asker and a channel's kind are facts the bus enforces (hq ADR 0259)
This commit is contained in:
@@ -53,9 +53,26 @@ func assertBusObjects(ctx context.Context, inv *inventory.Inventory, r broker.Ra
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// And the work queues of seats that name their caller or their kind, with each holder's worker
|
||||
// (novox/hq ADR 0259 §3): an ask queues until the router takes it, a channel's work until that kind
|
||||
// takes it.
|
||||
trafficStreams, trafficWorkers, err := seatTrafficObjects(ctx, inv)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, s := range trafficStreams {
|
||||
if err := r.EnsureStream(s); err != nil {
|
||||
return nil, fmt.Errorf("asserting the work queue %s: %w", s.Name, err)
|
||||
}
|
||||
}
|
||||
// Every one tried, and every failure named: one module's consumer the bus refuses is no reason
|
||||
// the modules after it in the list hear nothing (novox/hq issue 208, where this runs on each send).
|
||||
var failed []error
|
||||
for _, c := range trafficWorkers {
|
||||
if err := r.EnsureConsumer(c); err != nil {
|
||||
failed = append(failed, fmt.Errorf("the worker %s on %s: %w", c.Name, c.Stream, err))
|
||||
}
|
||||
}
|
||||
for _, c := range consumers {
|
||||
if err := r.EnsureConsumer(c.Consumer); err != nil {
|
||||
failed = append(failed, fmt.Errorf("how %s on %s hears what it consumes: %w", c.Module, c.Node, err))
|
||||
@@ -130,6 +147,21 @@ func moduleConsumers(ctx context.Context, inv *inventory.Inventory) ([]broker.Mo
|
||||
return broker.ConsumersOf(users), nil
|
||||
}
|
||||
|
||||
// seatTrafficObjects is the work queues and workers of seats that name their caller or their kind, from
|
||||
// the records the user list is composed from.
|
||||
func seatTrafficObjects(ctx context.Context, inv *inventory.Inventory) ([]broker.Stream, []broker.Consumer, error) {
|
||||
records, err := inv.BusRecords(ctx)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
users, err := broker.Users(records)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
streams, workers := broker.SeatTrafficObjects(users)
|
||||
return streams, workers, nil
|
||||
}
|
||||
|
||||
// moduleConsumerCount is how many modules hear what they consume, for the raise's one line.
|
||||
func moduleConsumerCount(ctx context.Context, inv *inventory.Inventory) (int, error) {
|
||||
consumers, err := moduleConsumers(ctx, inv)
|
||||
|
||||
Reference in New Issue
Block a user