package protocol import ( "errors" "git.netflux.io/rob/octoplex/internal/domain" "git.netflux.io/rob/octoplex/internal/event" pb "git.netflux.io/rob/octoplex/internal/generated/grpc" "google.golang.org/protobuf/types/known/timestamppb" ) // EventToProto converts an event to a protobuf message. func EventToProto(ev event.Event) *pb.Event { switch evt := ev.(type) { case event.AppStateChangedEvent: return buildAppStateChangeEvent(evt) case event.DestinationAddedEvent: return buildDestinationAddedEvent(evt) case event.AddDestinationFailedEvent: return buildAddDestinationFailedEvent(evt) case event.DestinationStreamExitedEvent: return buildDestinationStreamExitedEvent(evt) case event.StartDestinationFailedEvent: return buildStartDestinationFailedEvent(evt) case event.DestinationRemovedEvent: return buildDestinationRemovedEvent(evt) case event.RemoveDestinationFailedEvent: return buildRemoveDestinationFailedEvent(evt) case event.FatalErrorOccurredEvent: return buildFatalErrorOccurredEvent(evt) case event.OtherInstanceDetectedEvent: return buildOtherInstanceDetectedEvent(evt) case event.MediaServerStartedEvent: return buildMediaServerStartedEvent(evt) default: panic("unknown event type") } } func buildAppStateChangeEvent(evt event.AppStateChangedEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_AppStateChanged{ AppStateChanged: &pb.AppStateChangedEvent{ AppState: &pb.AppState{ Source: &pb.Source{ Container: containerToProto(evt.State.Source.Container), Live: evt.State.Source.Live, LiveChangedAt: timestamppb.New(evt.State.Source.LiveChangedAt), Tracks: evt.State.Source.Tracks, ExitReason: evt.State.Source.ExitReason, }, Destinations: destinationsToProto(evt.State.Destinations), BuildInfo: &pb.BuildInfo{ GoVersion: evt.State.BuildInfo.GoVersion, Version: evt.State.BuildInfo.Version, Commit: evt.State.BuildInfo.Commit, Date: evt.State.BuildInfo.Date, }, }, }, }, } } func buildDestinationAddedEvent(evt event.DestinationAddedEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_DestinationAdded{ DestinationAdded: &pb.DestinationAddedEvent{Url: evt.URL}, }, } } func buildAddDestinationFailedEvent(evt event.AddDestinationFailedEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_AddDestinationFailed{ AddDestinationFailed: &pb.AddDestinationFailedEvent{Url: evt.URL, Error: evt.Err.Error()}, }, } } func buildDestinationStreamExitedEvent(evt event.DestinationStreamExitedEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_DestinationStreamExited{ DestinationStreamExited: &pb.DestinationStreamExitedEvent{Name: evt.Name, Error: evt.Err.Error()}, }, } } func buildStartDestinationFailedEvent(evt event.StartDestinationFailedEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_StartDestinationFailed{ StartDestinationFailed: &pb.StartDestinationFailedEvent{Url: evt.URL, Message: evt.Message}, }, } } func buildDestinationRemovedEvent(evt event.DestinationRemovedEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_DestinationRemoved{ DestinationRemoved: &pb.DestinationRemovedEvent{Url: evt.URL}, }, } } func buildRemoveDestinationFailedEvent(evt event.RemoveDestinationFailedEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_RemoveDestinationFailed{ RemoveDestinationFailed: &pb.RemoveDestinationFailedEvent{Url: evt.URL, Error: evt.Err.Error()}, }, } } func buildFatalErrorOccurredEvent(evt event.FatalErrorOccurredEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_FatalError{ FatalError: &pb.FatalErrorEvent{Message: evt.Message}, }, } } func buildOtherInstanceDetectedEvent(_ event.OtherInstanceDetectedEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_OtherInstanceDetected{ OtherInstanceDetected: &pb.OtherInstanceDetectedEvent{}, }, } } func buildMediaServerStartedEvent(evt event.MediaServerStartedEvent) *pb.Event { return &pb.Event{ EventType: &pb.Event_MediaServerStarted{ MediaServerStarted: &pb.MediaServerStartedEvent{RtmpUrl: evt.RTMPURL, RtmpsUrl: evt.RTMPSURL}, }, } } // EventFromProto converts a protobuf message to an event. func EventFromProto(pbEv *pb.Event) event.Event { if pbEv == nil || pbEv.EventType == nil { panic("invalid or nil pb.Event") } switch evt := pbEv.EventType.(type) { case *pb.Event_AppStateChanged: return parseAppStateChangedEvent(evt.AppStateChanged) case *pb.Event_DestinationAdded: return parseDestinationAddedEvent(evt.DestinationAdded) case *pb.Event_AddDestinationFailed: return parseAddDestinationFailedEvent(evt.AddDestinationFailed) case *pb.Event_DestinationStreamExited: return parseDestinationStreamExitedEvent(evt.DestinationStreamExited) case *pb.Event_StartDestinationFailed: return parseStartDestinationFailedEvent(evt.StartDestinationFailed) case *pb.Event_DestinationRemoved: return parseDestinationRemovedEvent(evt.DestinationRemoved) case *pb.Event_RemoveDestinationFailed: return parseRemoveDestinationFailedEvent(evt.RemoveDestinationFailed) case *pb.Event_FatalError: return parseFatalErrorOccurredEvent(evt.FatalError) case *pb.Event_OtherInstanceDetected: return parseOtherInstanceDetectedEvent(evt.OtherInstanceDetected) case *pb.Event_MediaServerStarted: return parseMediaServerStartedEvent(evt.MediaServerStarted) default: panic("unknown pb.Event type") } } func parseAppStateChangedEvent(evt *pb.AppStateChangedEvent) event.Event { if evt == nil || evt.AppState == nil || evt.AppState.Source == nil { panic("invalid AppStateChangedEvent") } return event.AppStateChangedEvent{ State: domain.AppState{ Source: domain.Source{ Container: protoToContainer(evt.AppState.Source.Container), Live: evt.AppState.Source.Live, LiveChangedAt: evt.AppState.Source.LiveChangedAt.AsTime(), Tracks: evt.AppState.Source.Tracks, ExitReason: evt.AppState.Source.ExitReason, }, Destinations: protoToDestinations(evt.AppState.Destinations), BuildInfo: domain.BuildInfo{ GoVersion: evt.AppState.BuildInfo.GoVersion, Version: evt.AppState.BuildInfo.Version, Commit: evt.AppState.BuildInfo.Commit, Date: evt.AppState.BuildInfo.Date, }, }, } } func parseDestinationAddedEvent(evt *pb.DestinationAddedEvent) event.Event { if evt == nil { panic("nil DestinationAddedEvent") } return event.DestinationAddedEvent{URL: evt.Url} } func parseAddDestinationFailedEvent(evt *pb.AddDestinationFailedEvent) event.Event { if evt == nil { panic("nil AddDestinationFailedEvent") } return event.AddDestinationFailedEvent{URL: evt.Url, Err: errors.New(evt.Error)} } func parseDestinationStreamExitedEvent(evt *pb.DestinationStreamExitedEvent) event.Event { if evt == nil { panic("nil DestinationStreamExitedEvent") } return event.DestinationStreamExitedEvent{Name: evt.Name, Err: errors.New(evt.Error)} } func parseStartDestinationFailedEvent(evt *pb.StartDestinationFailedEvent) event.Event { if evt == nil { panic("nil StartDestinationFailedEvent") } return event.StartDestinationFailedEvent{URL: evt.Url, Message: evt.Message} } func parseDestinationRemovedEvent(evt *pb.DestinationRemovedEvent) event.Event { if evt == nil { panic("nil DestinationRemovedEvent") } return event.DestinationRemovedEvent{URL: evt.Url} } func parseRemoveDestinationFailedEvent(evt *pb.RemoveDestinationFailedEvent) event.Event { if evt == nil { panic("nil RemoveDestinationFailedEvent") } return event.RemoveDestinationFailedEvent{URL: evt.Url, Err: errors.New(evt.Error)} } func parseFatalErrorOccurredEvent(evt *pb.FatalErrorEvent) event.Event { if evt == nil { panic("nil FatalErrorEvent") } return event.FatalErrorOccurredEvent{Message: evt.Message} } func parseOtherInstanceDetectedEvent(_ *pb.OtherInstanceDetectedEvent) event.Event { return event.OtherInstanceDetectedEvent{} } func parseMediaServerStartedEvent(evt *pb.MediaServerStartedEvent) event.Event { if evt == nil { panic("nil MediaServerStartedEvent") } return event.MediaServerStartedEvent{RTMPURL: evt.RtmpUrl, RTMPSURL: evt.RtmpsUrl} }