diff --git a/internal/hub/hub.go b/internal/hub/hub.go index 90eb893..9a424c8 100644 --- a/internal/hub/hub.go +++ b/internal/hub/hub.go @@ -45,11 +45,11 @@ type Event struct { // nil/empty means "no filter on this dimension" (match everything). // An empty non-nil slice means "match nothing on this dimension". type Scope struct { - IATAs []string - RegionIATAs []string // pre-expanded from regionId by the WS handler - PayloadTypes []uint8 + IATAs []string + RegionIATAs []string // pre-expanded from regionId by the WS handler + PayloadTypes []uint8 ChannelHashes []string - Events []EventType + Events []EventType } // Client represents a connected WebSocket consumer. @@ -100,9 +100,9 @@ type Hub struct { } type subscribeMsg struct { - client *Client - scope Scope - hasScope bool // true when this is an AddScope call, false for NewClient registration + client *Client + scope Scope + hasScope bool // true when this is an AddScope call, false for NewClient registration } type unsubscribeMsg struct { diff --git a/internal/ingest/ingest.go b/internal/ingest/ingest.go index f87f70b..0af971e 100644 --- a/internal/ingest/ingest.go +++ b/internal/ingest/ingest.go @@ -193,6 +193,16 @@ type statusEvent struct { LastStatusAt int64 `json:"lastStatusAt"` // epoch ms } +// channelMessageEvent is the JSON payload for a channelMessage WS event. +type channelMessageEvent struct { + ChannelID int `json:"channelId"` + ChannelHash string `json:"channelHash"` // hex-encoded single byte + PacketHash string `json:"packetHash"` // hex-encoded + SenderName string `json:"senderName"` + Content string `json:"content"` + SentAt int64 `json:"sentAt"` // epoch ms +} + // packetObservationEvent is the JSON payload for a packetObservation WS event. // Shape matches the design doc § Server → Client events. type packetObservationEvent struct { @@ -666,10 +676,35 @@ func (w *Worker) handlePayloadTypeSideEffects(ctx context.Context, packet *meshc if err != nil { log.Printf("ingest[%s]: db: insert channel message failed: %v", w.cfg.BrokerName, err) } + + evt := channelMessageEvent{ + ChannelID: channelID, + ChannelHash: hex.EncodeToString(channelHashBytes), + PacketHash: hex.EncodeToString(packetHash), + SenderName: payload.Sender, + Content: payload.Text, + SentAt: time.Unix(int64(payload.Timestamp), 0).UnixMilli(), + } + w.broadcast(hub.EventChannelMessage, iata, 0, fmt.Sprintf("%02x", grpTxt.ChannelHash), evt) return } } +func (w *Worker) broadcast(eventType hub.EventType, iata string, payloadType uint8, channelHash string, payload any) { + b, err := json.Marshal(payload) + if err != nil { + log.Printf("ingest[%s]: failed to marshal %s event: %v", w.cfg.BrokerName, eventType, err) + return + } + w.hub.Broadcast(hub.Event{ + Type: eventType, + Payload: b, + IATA: iata, + PayloadType: payloadType, + ChannelHash: channelHash, + }) +} + // fanOut builds and broadcasts the packetObservation event to connected WS clients. func (w *Worker) fanOut(packetHash []byte, p *meshcore.Packet, iata string, isFirst bool, observerName string, observerID uuid.UUID, heardAt time.Time, rssi int16, snr float32, sourceBroker string) { evt := packetObservationEvent{} @@ -686,17 +721,7 @@ func (w *Worker) fanOut(packetHash []byte, p *meshcore.Packet, iata string, isFi evt.Observation.SNR = snr evt.Observation.SourceBroker = sourceBroker - payload, err := json.Marshal(evt) - if err != nil { - log.Printf("ingest[%s]: failed to marshal packetObservation event: %v", w.cfg.BrokerName, err) - return - } - w.hub.Broadcast(hub.Event{ - Type: hub.EventPacketObservation, - Payload: payload, - IATA: iata, - PayloadType: p.PayloadType(), - }) + w.broadcast(hub.EventPacketObservation, iata, p.PayloadType(), "", evt) } // parseNumber handles RSSI and SNR fields that different observer types send as