ProduceKeyChanges creates new change events for each key
(keys []api.DeviceMessage)
| 34 | |
| 35 | // ProduceKeyChanges creates new change events for each key |
| 36 | func (p *KeyChange) ProduceKeyChanges(keys []api.DeviceMessage) error { |
| 37 | userToDeviceCount := make(map[string]int) |
| 38 | for _, key := range keys { |
| 39 | id, err := p.DB.StoreKeyChange(context.Background(), key.UserID) |
| 40 | if err != nil { |
| 41 | return err |
| 42 | } |
| 43 | key.DeviceChangeID = id |
| 44 | value, err := json.Marshal(key) |
| 45 | if err != nil { |
| 46 | return err |
| 47 | } |
| 48 | |
| 49 | m := &nats.Msg{ |
| 50 | Subject: p.Topic, |
| 51 | Header: nats.Header{}, |
| 52 | } |
| 53 | m.Header.Set(jetstream.UserID, key.UserID) |
| 54 | m.Data = value |
| 55 | |
| 56 | _, err = p.JetStream.PublishMsg(m) |
| 57 | if err != nil { |
| 58 | return err |
| 59 | } |
| 60 | |
| 61 | userToDeviceCount[key.UserID]++ |
| 62 | } |
| 63 | for userID, count := range userToDeviceCount { |
| 64 | logrus.WithFields(logrus.Fields{ |
| 65 | "user_id": userID, |
| 66 | "num_key_changes": count, |
| 67 | }).Tracef("Produced to key change topic '%s'", p.Topic) |
| 68 | } |
| 69 | return nil |
| 70 | } |
| 71 | |
| 72 | func (p *KeyChange) ProduceSigningKeyUpdate(key api.CrossSigningKeyUpdate) error { |
| 73 | output := &api.DeviceMessage{ |
nothing calls this directly
no test coverage detected