|
@@ -11,11 +11,12 @@ type Bus struct {
|
|
}
|
|
}
|
|
|
|
|
|
func Connect(url string) (*Bus, error) {
|
|
func Connect(url string) (*Bus, error) {
|
|
|
|
+ nats.RegisterEncoder("v1", &MessageEncoderV1{})
|
|
conn, err := nats.Connect(url)
|
|
conn, err := nats.Connect(url)
|
|
if err != nil {
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to connect to the broker: %w", err)
|
|
return nil, fmt.Errorf("failed to connect to the broker: %w", err)
|
|
}
|
|
}
|
|
- encodedConn, err := nats.NewEncodedConn(conn, nats.JSON_ENCODER) // TODO
|
|
|
|
|
|
+ encodedConn, err := nats.NewEncodedConn(conn, "v1")
|
|
if err != nil {
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to configure the bus encoding: %w", err)
|
|
return nil, fmt.Errorf("failed to configure the bus encoding: %w", err)
|
|
}
|
|
}
|
|
@@ -23,41 +24,36 @@ func Connect(url string) (*Bus, error) {
|
|
}
|
|
}
|
|
|
|
|
|
func (bus *Bus) SendEmail(request *SendEmail) error {
|
|
func (bus *Bus) SendEmail(request *SendEmail) error {
|
|
- envelope := request // TODO
|
|
|
|
- if err := bus.NATS.Publish("email.outbound", envelope); err != nil {
|
|
|
|
- return fmt.Errorf("failed to enqueue outbound email: %w", err)
|
|
|
|
- }
|
|
|
|
- return nil
|
|
|
|
|
|
+ return bus.Publish("email.outbound", request)
|
|
}
|
|
}
|
|
|
|
|
|
func (bus *Bus) RegisterID(request *RegisterID) error {
|
|
func (bus *Bus) RegisterID(request *RegisterID) error {
|
|
- envelope := request // TODO
|
|
|
|
- if err := bus.NATS.Publish("id.register", envelope); err != nil {
|
|
|
|
- return fmt.Errorf("failed to enqueue ID registration: %w", err)
|
|
|
|
- }
|
|
|
|
- return nil
|
|
|
|
|
|
+ return bus.Publish("id.register", request)
|
|
}
|
|
}
|
|
|
|
|
|
func (bus *Bus) VerifyEmail(request *VerifyEmail) error {
|
|
func (bus *Bus) VerifyEmail(request *VerifyEmail) error {
|
|
- envelope := request // TODO
|
|
|
|
- if err := bus.NATS.Publish("id.verify", envelope); err != nil {
|
|
|
|
- return fmt.Errorf("failed to enqueue ID verification: %w", err)
|
|
|
|
- }
|
|
|
|
- return nil
|
|
|
|
|
|
+ return bus.Publish("id.verify", request)
|
|
}
|
|
}
|
|
|
|
|
|
func (bus *Bus) GrantAccess(request *GrantAccess) error {
|
|
func (bus *Bus) GrantAccess(request *GrantAccess) error {
|
|
- envelope := request // TODO
|
|
|
|
- if err := bus.NATS.Publish("access.grant", envelope); err != nil {
|
|
|
|
- return fmt.Errorf("failed to enqueue access grant: %w", err)
|
|
|
|
- }
|
|
|
|
- return nil
|
|
|
|
|
|
+ return bus.Publish("access.grant", request)
|
|
}
|
|
}
|
|
|
|
|
|
func (bus *Bus) RevokeAccess(request *RevokeAccess) error {
|
|
func (bus *Bus) RevokeAccess(request *RevokeAccess) error {
|
|
- envelope := request // TODO
|
|
|
|
- if err := bus.NATS.Publish("access.revoke", envelope); err != nil {
|
|
|
|
- return fmt.Errorf("failed to enqueue access revocation: %w", err)
|
|
|
|
|
|
+ return bus.Publish("access.revoke", request)
|
|
|
|
+}
|
|
|
|
+
|
|
|
|
+func (bus *Bus) Publish(stream string, request interface{}) error {
|
|
|
|
+ rawEvent, err := NewRawEvent(0, request) // TODO
|
|
|
|
+ if err != nil {
|
|
|
|
+ return err
|
|
|
|
+ }
|
|
|
|
+ return bus.PublishRawEvent(stream, rawEvent)
|
|
|
|
+}
|
|
|
|
+
|
|
|
|
+func (bus *Bus) PublishRawEvent(stream string, rawEvent *RawEvent) error {
|
|
|
|
+ if err := bus.NATS.Publish(stream, rawEvent); err != nil {
|
|
|
|
+ return fmt.Errorf("failed to publish message: %w", err)
|
|
}
|
|
}
|
|
return nil
|
|
return nil
|
|
}
|
|
}
|