RunNativeProcess runs the native process mode
(binaryName string)
| 76 | |
| 77 | // RunNativeProcess runs the native process mode |
| 78 | func RunNativeProcess(binaryName string) { |
| 79 | appCtx, appCtxCancel := context.WithCancel(context.Background()) |
| 80 | defer appCtxCancel() |
| 81 | |
| 82 | logger := nativeLogger.With().Int("pid", os.Getpid()).Logger() |
| 83 | setProcTitle("starting") |
| 84 | |
| 85 | // Parse native options |
| 86 | var proxyOptions nativeProxyOptions |
| 87 | if err := env.Parse(&proxyOptions); err != nil { |
| 88 | logger.Fatal().Err(err).Msg("failed to parse native proxy options") |
| 89 | } |
| 90 | |
| 91 | // Connect to video stream socket |
| 92 | conn, err := net.Dial("unix", proxyOptions.VideoStreamUnixSocket) |
| 93 | if err != nil { |
| 94 | logger.Fatal().Err(err).Msg("failed to connect to video stream socket") |
| 95 | } |
| 96 | logger.Info().Str("videoStreamSocketPath", proxyOptions.VideoStreamUnixSocket).Msg("connected to video stream socket") |
| 97 | |
| 98 | nativeOptions := proxyOptions.toNativeOptions() |
| 99 | nativeOptions.OnVideoFrameReceived = func(frame []byte, duration time.Duration) { |
| 100 | // Write 4-byte frame length prefix, then frame data |
| 101 | var frameSizeBuffer [4]byte |
| 102 | binary.LittleEndian.PutUint32(frameSizeBuffer[:], uint32(len(frame))) |
| 103 | |
| 104 | if _, err := conn.Write(frameSizeBuffer[:]); err != nil { |
| 105 | logger.Fatal().Err(err).Msg("failed to write frame size to video stream socket") |
| 106 | } |
| 107 | if _, err := conn.Write(frame); err != nil { |
| 108 | logger.Fatal().Err(err).Msg("failed to write frame to video stream socket") |
| 109 | } |
| 110 | } |
| 111 | nativeOptions.OnVideoStateChange = func(state VideoState) { |
| 112 | updateProcessTitle(&state) |
| 113 | } |
| 114 | |
| 115 | // Create native instance |
| 116 | nativeInstance := NewNative(*nativeOptions) |
| 117 | gspt.SetProcTitle("jetkvm: [native] initializing") |
| 118 | |
| 119 | // Start native instance |
| 120 | if err := nativeInstance.Start(); err != nil { |
| 121 | logger.Fatal().Err(err).Msg("failed to start native instance") |
| 122 | } |
| 123 | |
| 124 | grpcLogger := logger.With().Str("socketPath", fmt.Sprintf("@%v", proxyOptions.CtrlUnixSocket)).Logger() |
| 125 | setProcTitle("starting gRPC server") |
| 126 | // Create gRPC server |
| 127 | grpcServer := NewGRPCServer(nativeInstance, &grpcLogger) |
| 128 | |
| 129 | logger.Info().Msg("starting gRPC server") |
| 130 | // Start gRPC server |
| 131 | server, lis, err := StartGRPCServer(grpcServer, fmt.Sprintf("@%v", proxyOptions.CtrlUnixSocket), &logger) |
| 132 | if err != nil { |
| 133 | logger.Fatal().Err(err).Msg("failed to start gRPC server") |
| 134 | } |
| 135 | setProcTitle("ready") |
no test coverage detected