• 设为首页
  • 点击收藏
  • 手机版
    手机扫一扫访问
    迪恩网络手机版
  • 关注官方公众号
    微信扫一扫关注
    公众号

Scala Unpooled类代码示例

原作者: [db:作者] 来自: [db:来源] 收藏 邀请

本文整理汇总了Scala中io.netty.buffer.Unpooled的典型用法代码示例。如果您正苦于以下问题:Scala Unpooled类的具体用法?Scala Unpooled怎么用?Scala Unpooled使用的例子?那么恭喜您, 这里精选的类代码示例或许可以为您提供帮助。



在下文中一共展示了Unpooled类的20个代码示例,这些例子默认根据受欢迎程度排序。您可以为喜欢或者感觉有用的代码点赞,您的评价将有助于我们的系统推荐出更棒的Scala代码示例。

示例1: RespondToExpectContinue

//设置package包名称以及导入依赖的类
package com.twitter.finagle.netty4.http.handler

import io.netty.buffer.Unpooled
import io.netty.channel.ChannelHandler.Sharable
import io.netty.channel.{ChannelHandlerContext, ChannelInboundHandlerAdapter}
import io.netty.handler.codec.http._



@Sharable
private[http] object RespondToExpectContinue extends ChannelInboundHandlerAdapter {

  def newContinue(): FullHttpResponse =
    new DefaultFullHttpResponse(
      HttpVersion.HTTP_1_1,
      HttpResponseStatus.CONTINUE,
      Unpooled.EMPTY_BUFFER)

  override def channelRead(ctx: ChannelHandlerContext, msg: Any): Unit = {
    msg match {
      case http: HttpMessage if HttpUtil.is100ContinueExpected(http) =>
        ctx.writeAndFlush(newContinue())
        http.headers.remove(HttpHeaderNames.EXPECT)
      case _ =>
    }
    super.channelRead(ctx, msg)
  }
} 
开发者ID:wenkeyang,项目名称:finagle,代码行数:29,代码来源:RespondToExpectContinue.scala


示例2: FixedLengthMessageAggregatorTest

//设置package包名称以及导入依赖的类
package com.twitter.finagle.netty4.http.handler

import com.twitter.conversions.storage._
import io.netty.buffer.Unpooled
import io.netty.channel.embedded.EmbeddedChannel
import io.netty.handler.codec.http._
import org.junit.runner.RunWith
import org.scalatest.FunSuite
import org.scalatest.junit.JUnitRunner

@RunWith(classOf[JUnitRunner])
class FixedLengthMessageAggregatorTest extends FunSuite {

  test("full messages pass through") {
    val agg = new FixedLengthMessageAggregator(10.megabytes)
    val channel: EmbeddedChannel = new EmbeddedChannel(new HttpRequestEncoder(), agg)
    val content = Unpooled.wrappedBuffer(new Array[Byte](11))
    val req = new DefaultFullHttpRequest(HttpVersion.HTTP_1_1, HttpMethod.POST, "/", content)
    assert(channel.writeInbound(req))

    val reqObserved = channel.readInbound[FullHttpRequest]()
    assert(reqObserved.method == HttpMethod.POST)
    assert(reqObserved.content == req.content)
  }

  test("chunked messages aren't aggregated") {
    val agg = new FixedLengthMessageAggregator(10.megabytes)
    val channel: EmbeddedChannel = new EmbeddedChannel(new HttpRequestEncoder(), agg)
    val content = Unpooled.wrappedBuffer(new Array[Byte](11))
    val head = new DefaultHttpRequest(HttpVersion.HTTP_1_1, HttpMethod.POST, "/")
    HttpUtil.setTransferEncodingChunked(head, true)

    val body = new DefaultLastHttpContent(content)
    assert(channel.writeInbound(head))
    assert(channel.writeInbound(body))

    val reqObserved = channel.readInbound[HttpRequest]()
    assert(reqObserved.method == HttpMethod.POST)

    val bodyObserved = channel.readInbound[HttpContent]()
    assert(bodyObserved.content == content)
  }

  test("fixed length messages which are chunked are aggregated") {
    val agg = new FixedLengthMessageAggregator(10.megabytes)
    val channel: EmbeddedChannel = new EmbeddedChannel(new HttpRequestEncoder(), agg)
    val content = Unpooled.wrappedBuffer(new Array[Byte](11))
    val head = new DefaultHttpRequest(HttpVersion.HTTP_1_1, HttpMethod.POST, "/")
    HttpUtil.setContentLength(head, content.readableBytes)

    val body = new DefaultLastHttpContent(content)

    assert(!channel.writeInbound(head))
    assert(channel.writeInbound(body))

    val reqObserved = channel.readInbound[FullHttpRequest]()
    assert(reqObserved.method == HttpMethod.POST)
    assert(reqObserved.content == content)
  }
} 
开发者ID:wenkeyang,项目名称:finagle,代码行数:61,代码来源:FixedLengthMessageAggregatorTest.scala


示例3: ByteBufManagerTest

//设置package包名称以及导入依赖的类
package com.twitter.finagle.netty4.http

import com.twitter.conversions.storage._
import com.twitter.finagle.ChannelBufferUsageException
import com.twitter.finagle.http.codec.ChannelBufferUsageTracker
import io.netty.buffer.Unpooled
import io.netty.channel.{ChannelPromise, ChannelHandlerContext}
import org.junit.runner.RunWith
import org.scalatest.FunSuite
import org.scalatest.junit.JUnitRunner
import org.scalatest.mock.MockitoSugar

@RunWith(classOf[JUnitRunner])
class ByteBufManagerTest extends FunSuite with MockitoSugar {

  def usageTrackerFactory() = {
    val usageTracker = new ChannelBufferUsageTracker(1000.bytes)
    assert(usageTracker.usageLimit == (1000.bytes))
    usageTracker
  }

  val ctx = mock[ChannelHandlerContext]
  val p = mock[ChannelPromise]

  test("tracks buffer usage between writes") {
    val tracker = usageTrackerFactory()
    val handler = new ByteBufManager(tracker)

    handler.channelRead(ctx, Unpooled.buffer(100))
    assert(tracker.currentUsage == 100.bytes)
    handler.channelRead(ctx, Unpooled.buffer(50))
    assert(tracker.currentUsage == 150.bytes)

    handler.write(ctx, Unpooled.EMPTY_BUFFER, p)
    assert(tracker.currentUsage == 0.bytes)

    handler.channelRead(ctx, Unpooled.buffer(123))
    assert(tracker.currentUsage == 123.bytes)
    handler.close(ctx, p)
    assert(tracker.currentUsage == 0.bytes)
  }

  test("throws if aggregate limit is exceeded") {
    val tracker = usageTrackerFactory()
    val handler = new ByteBufManager(tracker)

    handler.channelRead(ctx, Unpooled.buffer(1000))

    intercept[ChannelBufferUsageException] {
      handler.channelRead(ctx, Unpooled.buffer(1))
    }
  }
} 
开发者ID:wenkeyang,项目名称:finagle,代码行数:54,代码来源:ByteBufManagerTest.scala


示例4: ThriftClientArrayToByteBufEncoder

//设置package包名称以及导入依赖的类
package com.twitter.finagle.thrift.transport.netty4

import com.twitter.finagle.thrift.ThriftClientRequest
import com.twitter.finagle.thrift.transport.ExceptionFactory
import io.netty.buffer.Unpooled
import io.netty.channel.ChannelHandler.Sharable
import io.netty.channel.{
  ChannelHandler, ChannelHandlerContext, ChannelOutboundHandlerAdapter,
  ChannelPromise, CombinedChannelDuplexHandler}


  @Sharable
  private object ThriftClientArrayToByteBufEncoder extends ChannelOutboundHandlerAdapter {

    override def write(ctx: ChannelHandlerContext, msg: Any, promise: ChannelPromise): Unit = msg match {
      case request: ThriftClientRequest =>
        val buf = Unpooled.wrappedBuffer(request.message)
        ctx.writeAndFlush(buf, promise)

      case other =>
        val ex = ExceptionFactory.wrongClientWriteType(other)
        promise.setFailure(ex)
        throw ex
    }
  }
} 
开发者ID:wenkeyang,项目名称:finagle,代码行数:27,代码来源:ClientByteBufCodec.scala


示例5: ServerByteBufCodec

//设置package包名称以及导入依赖的类
package com.twitter.finagle.thrift.transport.netty4

import com.twitter.finagle.thrift.transport.ExceptionFactory
import io.netty.buffer.Unpooled
import io.netty.channel.{
  ChannelHandler, ChannelHandlerContext, ChannelOutboundHandlerAdapter,
  ChannelPromise, CombinedChannelDuplexHandler}
import io.netty.channel.ChannelHandler.Sharable


private[netty4] object ServerByteBufCodec {

  def apply(): ChannelHandler = {
    val encoder = ThriftServerArrayToByteBufEncoder
    val decoder = ThriftByteBufToArrayDecoder
    new CombinedChannelDuplexHandler(decoder, encoder)
  }

  @Sharable
  private object ThriftServerArrayToByteBufEncoder extends ChannelOutboundHandlerAdapter {
    override def write(ctx: ChannelHandlerContext, msg: Any, promise: ChannelPromise): Unit =
      msg match {
        case array: Array[Byte] =>
          val buf = Unpooled.wrappedBuffer(array)
          ctx.writeAndFlush(buf, promise)

        case other =>
          val ex = ExceptionFactory.wrongServerWriteType(other)
          promise.setFailure(ex)
          throw ex
      }
  }
} 
开发者ID:wenkeyang,项目名称:finagle,代码行数:34,代码来源:ServerByteBufCodec.scala


示例6: ThriftBufferedTransportDecoderTest

//设置package包名称以及导入依赖的类
package com.twitter.finagle.thrift.transport.netty4

import com.twitter.finagle.thrift.Protocols
import com.twitter.finagle.thrift.transport.AbstractBufferedTransportDecoderTest
import io.netty.buffer.{ByteBuf, Unpooled}
import io.netty.channel.embedded.EmbeddedChannel

class ThriftBufferedTransportDecoderTest extends AbstractBufferedTransportDecoderTest {

  private def getArray(buf: ByteBuf): Array[Byte] = {
    val out = new Array[Byte](buf.readableBytes())
    buf.readBytes(out)
    assert(buf.readableBytes() == 0)
    out
  }
  private def getDecoder = new ThriftBufferedTransportDecoder(Protocols.factory())

  def decode(arrays: Seq[Array[Byte]]): Vector[Array[Byte]] = {
    val data = arrays.map(Unpooled.wrappedBuffer(_))
    val channel = new EmbeddedChannel()

    channel.pipeline().addLast(getDecoder)
    channel.writeInbound(data:_*)

    var acc = Vector.empty[Array[Byte]]
    while (!channel.inboundMessages().isEmpty) {
      acc :+= getArray(channel.readInbound[ByteBuf]())
    }

    acc
  }
} 
开发者ID:wenkeyang,项目名称:finagle,代码行数:33,代码来源:ThriftBufferedTransportDecoderTest.scala


示例7: BufCodecTest

//设置package包名称以及导入依赖的类
package com.twitter.finagle.netty4.codec

import com.twitter.finagle.Failure
import com.twitter.io.{Buf, Charsets}
import io.netty.buffer.{ByteBuf, Unpooled}
import io.netty.channel.embedded.EmbeddedChannel
import org.junit.runner.RunWith
import org.scalatest.FunSuite
import org.scalatest.junit.JUnitRunner

@RunWith(classOf[JUnitRunner])
class BufCodecTest extends FunSuite {
  test("decode") {
    val ch = new EmbeddedChannel(new BufCodec)
    val q = ch.inboundMessages

    ch.writeInbound(Unpooled.wrappedBuffer("hello".getBytes(Charsets.Utf8)))
    assert(q.size == 1)
    assert(q.poll() == Buf.Utf8("hello"))
    assert(q.size == 0)

    intercept[Failure] { ch.writeInbound(new Object) }
  }

  test("encode") {
    val ch = new EmbeddedChannel(new BufCodec)
    val q = ch.outboundMessages

    ch.writeOutbound(Buf.Utf8("hello"))
    assert(q.size == 1)
    assert(q.peek().isInstanceOf[ByteBuf])
    val bb = q.poll().asInstanceOf[ByteBuf]
    assert(bb.toString(Charsets.Utf8) == "hello")
    assert(q.size == 0)

    val channelFuture = ch.write(new Object)
    assert(channelFuture.cause.isInstanceOf[Failure])
  }
} 
开发者ID:wenkeyang,项目名称:finagle,代码行数:40,代码来源:BufCodecTest.scala


示例8: UdpPipHandler

//设置package包名称以及导入依赖的类
package tk.dasb.handler

import java.net.InetSocketAddress

import io.netty.buffer.{ByteBuf, Unpooled}
import io.netty.channel.socket.DatagramPacket
import io.netty.channel.{ChannelHandlerContext, ChannelInboundHandlerAdapter}
import io.netty.util.ByteProcessor
import tk.dasb.protocol.UdpReq
import tk.dasb.util.{FutureConv, Log}


class UdpPipHandler(var clientAddress:InetSocketAddress) extends ChannelInboundHandlerAdapter with Log with FutureConv {


  def isClientAddressInValid() :Boolean =  clientAddress == null || clientAddress.getPort == 0


  def procClientToProxy(ctx: ChannelHandlerContext, msg:UdpReq) = {
    val newMessage = new DatagramPacket(Unpooled.wrappedBuffer(msg.data),msg.dstAddress)
    ctx.writeAndFlush(newMessage)
    if(isClientAddressInValid){
      clientAddress = msg.raw.sender
      log.debug("reset client address:{}",clientAddress)
    }
  }


  def procRemoteToProxy(ctx: ChannelHandlerContext, msg:DatagramPacket) = {

    if(!isClientAddressInValid ){
      val msgHead = ctx.alloc.buffer(16,32)
      msgHead.writeShort(0).writeByte(0).writeByte(1)
      msgHead.writeBytes(msg.sender.getAddress.getAddress).writeShort(msg.sender.getPort)
      val newMessage = new DatagramPacket(Unpooled.wrappedBuffer(msgHead,msg.content),clientAddress)
      log.debug("send back {}",newMessage.content.readableBytes())
      ctx.writeAndFlush(newMessage)

    }else{
      log.warn("unknow client address,drop messsage")
      msg.release()
    }
  }

  override def channelRead(ctx: ChannelHandlerContext, msg: scala.Any): Unit = msg match {
      case udp:UdpReq => procClientToProxy(ctx,udp)
      case udp:DatagramPacket => procRemoteToProxy(ctx,udp)
  }

} 
开发者ID:xqdxqd,项目名称:tiny-sock5-server,代码行数:51,代码来源:UdpPipHandler.scala


示例9: HttpBasicAuth

//设置package包名称以及导入依赖的类
package strd.net.http

import java.nio.charset.Charset

import io.netty.buffer.Unpooled
import io.netty.handler.codec.base64.Base64
import io.netty.handler.codec.http._

import scala.concurrent.Future




object HttpBasicAuth {
  def unautorized = {
    HttpResp( HttpResponseStatus.UNAUTHORIZED,
      headers = Map(HttpHeaders.Names.WWW_AUTHENTICATE -> Seq("Basic realm=\"bad user or password\"")),
      body = Some( Content("not authorized") ))
  }

  def authSync(http: HttpReq, usersPass: Seq[String], process: HttpReq => HttpResp): HttpResp = {
    http.headers.getFirst(HttpHeaders.Names.AUTHORIZATION).map( header => {
      if (header.contains("Basic ")) {
        val auth = header.replace("Basic ", "")
        val pass = Base64.decode( Unpooled.wrappedBuffer(auth.getBytes) ).toString(Charset.forName("UTF8") )
        if (usersPass.contains(pass)) {
          process(http)
        } else {
          unautorized
        }
      } else {
        unautorized
      }
    }) getOrElse unautorized
  }

  def authAsync(http: HttpReq, usersPass: Seq[String], process: HttpReq => Future[HttpResp]): Future[HttpResp] = {
    http.headers.getFirst(HttpHeaders.Names.AUTHORIZATION).map( header => {
      if (header.contains("Basic ")) {
        val auth = header.replace("Basic ", "")
        val pass = Base64.decode( Unpooled.wrappedBuffer(auth.getBytes) ).toString(Charset.forName("UTF8") )
        if (usersPass.contains(pass)) {
          process(http)
        } else {
          Future.successful(unautorized)
        }
      } else {
        Future.successful(unautorized)
      }
    }) getOrElse Future.successful(unautorized)
  }

} 
开发者ID:lembrd,项目名称:strd-net-http,代码行数:54,代码来源:HttpBasicAuth.scala


示例10: XmlHttpResponseEncoder

//设置package包名称以及导入依赖的类
package com.scxmpp.bosh

import java.util
import com.scxmpp.netty.XmlElementEncoder
import com.scxmpp.xml.XmlElement
import io.netty.buffer.Unpooled
import io.netty.channel.ChannelHandlerContext
import io.netty.handler.codec.http.{HttpUtil, HttpResponseStatus, HttpVersion, DefaultFullHttpResponse}
import io.netty.util.CharsetUtil

class XmlHttpResponseEncoder extends XmlElementEncoder {

  override def encode(ctx: ChannelHandlerContext, msg: XmlElement, out: util.List[Object]) = {
    val string = elementAsString(msg)
    val response = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK,
      Unpooled.copiedBuffer(string, CharsetUtil.UTF_8))
    HttpUtil.setContentLength(response, string.length)
    response.headers.add("Access-Control-Allow-Origin", "*")
    out.add(response)
  }
} 
开发者ID:madprogrammer,项目名称:scxmppd,代码行数:22,代码来源:XmlHttpResponseEncoder.scala


示例11: UnpoolHandler

//设置package包名称以及导入依赖的类
package com.twitter.finagle.buoyant.h2
package netty4

import io.netty.buffer.{ByteBuf, ByteBufHolder, EmptyByteBuf, Unpooled}
import io.netty.channel.ChannelHandler.Sharable
import io.netty.channel.{ChannelHandlerContext, ChannelInboundHandlerAdapter}


@Sharable
object UnpoolHandler extends ChannelInboundHandlerAdapter {
  private[this] final def copyOnHeapAndRelease(bb: ByteBuf): ByteBuf = {
    try {
      if (bb.readableBytes > 0) Unpooled.buffer(bb.readableBytes, bb.capacity).writeBytes(bb)
      else Unpooled.EMPTY_BUFFER
    } finally {
      val _ = bb.release()
    }
  }

  override def channelRead(ctx: ChannelHandlerContext, msg: Any): Unit = msg match {
    case bb: ByteBuf =>
      val _ = ctx.fireChannelRead(copyOnHeapAndRelease(bb))

    // This case is special since it helps to avoid unnecessary `replace`
    // when the underlying content is already `EmptyByteBuffer`.
    case bbh: ByteBufHolder if bbh.content.isInstanceOf[EmptyByteBuf] =>
      val _ = ctx.fireChannelRead(bbh)

    case bbh: ByteBufHolder =>
      val onHeapContent = copyOnHeapAndRelease(bbh.content)
      val _ = ctx.fireChannelRead(bbh.replace(onHeapContent))

    case _ =>
      val _ = ctx.fireChannelRead(msg)
  }
} 
开发者ID:linkerd,项目名称:linkerd,代码行数:37,代码来源:UnpoolHandler.scala


示例12: CodecLenses

//设置package包名称以及导入依赖的类
package dnsclient 

object CodecLenses {
  import scodec.bits._ 
  import scodec._
  import io.netty.buffer.Unpooled
  import io.netty.channel.socket.DatagramPacket
  import java.net.InetSocketAddress
  import scodec.interop.scalaz._
  import scalaz._
  import scalaz.syntax.profunctor._
  import scalaz.syntax.std.tuple._
  import scalaz.syntax.functor._
  import scalaz.concurrent.Task

  val datagramIso: Iso[(InetSocketAddress, Array[Byte]), DatagramPacket] = Iso(t => new DatagramPacket(Unpooled.copiedBuffer(t._2), t._1), dgp => { 
    val bytebuf = dgp.content
    val bytearr = new Array[Byte](bytebuf.readableBytes())
    bytebuf.getBytes(bytebuf.readerIndex(), bytearr)
    (dgp.recipient, bytearr)
  })
  
  def codecIso[A](implicit codec: Codec[A]): Iso[Err \/ A, Err \/ BitVector] = Iso(aa => aa.flatMap(codec.encode(_).toDisjunction), bv => bv.flatMap(b => codec.decode(b).map(_.value).toDisjunction))
  
  def bvToBa: Iso[BitVector, Array[Byte]] = Iso(b => b.toByteArray, ba => ByteVector(ba).toBitVector)
  import scalaz._
  import Scalaz._
  //hmm not so efficient
  def addressAbv[A](implicit c: Codec[A]): Iso[Err \/ (InetSocketAddress, A), Err \/ (InetSocketAddress, BitVector)] = Iso.apply[Err \/ (InetSocketAddress, A), Err \/ (InetSocketAddress, BitVector)](
      err => err.flatMap(t => codecIso[A].get(t._2.right[Err]).map(a => (t._1, a))), 
      err => err.flatMap(t => codecIso[A].rget(t._2.right[Err]).map(a => (t._1, a)))
    )

  def codecToBytes[A](implicit codec: Codec[A]): Iso[Err \/ (InetSocketAddress, A),  Err \/ DatagramPacket] = {
    val niso = addressAbv[A]
    niso compose (bvToBa.first[InetSocketAddress].choiceRight[Err] compose datagramIso.choiceRight[Err])
  }
} 
开发者ID:vmarquez,项目名称:dnsclient,代码行数:39,代码来源:CodecOptics.scala


示例13: PortClient

//设置package包名称以及导入依赖的类
package io.dac.mara.ports

import io.netty.bootstrap.Bootstrap
import io.netty.buffer.Unpooled
import io.netty.channel._
import io.netty.channel.nio.NioEventLoopGroup
import io.netty.channel.socket.SocketChannel
import io.netty.channel.socket.nio._
import io.netty.handler.codec.string.{StringDecoder, StringEncoder}
import io.netty.util.CharsetUtil


class PortClient {
  val group = new NioEventLoopGroup()

  def connect(group: EventLoopGroup): ChannelFuture = {
    val bootstrap = new Bootstrap()
      .group(group)
      .channel(classOf[NioSocketChannel])
      .option[java.lang.Boolean](ChannelOption.TCP_NODELAY, true)
      .handler(new ChannelInitializer[SocketChannel]{
        override def initChannel(ch: SocketChannel): Unit = {
          val pipeline = ch.pipeline()
          println("Adding handler to pipeline")
          pipeline.addLast(new StringEncoder(CharsetUtil.UTF_8))
          pipeline.addLast(new StringDecoder(CharsetUtil.UTF_8))
          pipeline.addLast(new PortClientHandler)
        }
      })

    bootstrap.connect("127.0.0.1", 12345).sync()
  }

  def start: Unit = {
    try {
      val future = connect(group)
      val channel = future.channel()

      val message = "D0\"hello, world\""
      val buffer = Unpooled.buffer(message.length)

      message.foreach { ch =>
        buffer.writeByte(ch.toInt)
      }

      channel.writeAndFlush(buffer).sync()

    } finally {
      Thread.sleep(10000)
      group.shutdownGracefully()
      println("Shut down")
    }
  }
}

object PortClient {
  def main(args: Array[String]): Unit = {
    val client = new PortClient
    client.start
  }
} 
开发者ID:dacjames,项目名称:mara2-lang,代码行数:62,代码来源:PortClient.scala


示例14: GrpcGatewayHandler

//设置package包名称以及导入依赖的类
package grpcgateway.handlers

import java.nio.charset.StandardCharsets

import com.trueaccord.scalapb.GeneratedMessage
import com.trueaccord.scalapb.json.JsonFormat
import io.grpc.ManagedChannel
import io.netty.buffer.Unpooled
import io.netty.channel.ChannelHandler.Sharable
import io.netty.channel.{ChannelFutureListener, ChannelHandlerContext, ChannelInboundHandlerAdapter}
import io.netty.handler.codec.http._

import scala.concurrent.{ExecutionContext, Future}
import scala.util.{Failure, Success}

@Sharable
abstract class GrpcGatewayHandler(channel: ManagedChannel)(implicit ec: ExecutionContext)
  extends ChannelInboundHandlerAdapter {

  def name: String

  def shutdown(): Unit =
    if (!channel.isShutdown) channel.shutdown()

  def unaryCall(method: HttpMethod, uri: String, body: String): Future[GeneratedMessage]

  override def channelRead(ctx: ChannelHandlerContext, msg: scala.Any): Unit = msg match {
    case req: FullHttpRequest =>
      val body = req.content().toString(StandardCharsets.UTF_8)
      unaryCall(req.method(), req.uri(), body)
        .map(JsonFormat.toJsonString)
        .map(_.getBytes(StandardCharsets.UTF_8))
        .onComplete {
          case Success(json) =>
            val res = new DefaultFullHttpResponse(
              req.protocolVersion(),
              HttpResponseStatus.OK,
              Unpooled.copiedBuffer(json)
            )
            res.headers().set(HttpHeaderNames.CONTENT_TYPE, "application/json")
            HttpUtil.setContentLength(res, json.length)
            HttpUtil.setKeepAlive(res, HttpUtil.isKeepAlive(req))
            ctx.writeAndFlush(res)

          case Failure(e) =>
            val status = e match {
              case _: UnsupportedOperationException => HttpResponseStatus.NOT_FOUND
              case _: NoSuchElementException => HttpResponseStatus.BAD_REQUEST
              case _ => HttpResponseStatus.INTERNAL_SERVER_ERROR
            }
            val res = new DefaultHttpResponse(req.protocolVersion(), status)
            ctx.writeAndFlush(res).addListener(ChannelFutureListener.CLOSE)
        }

    case _ => super.channelRead(ctx, msg)
  }
} 
开发者ID:btlines,项目名称:grpcgateway,代码行数:58,代码来源:GrpcGatewayHandler.scala


示例15: SwaggerHandler

//设置package包名称以及导入依赖的类
package grpcgateway.handlers

import javax.activation.MimetypesFileTypeMap

import io.grpc.internal.IoUtils
import io.netty.buffer.Unpooled
import io.netty.channel.ChannelHandler.Sharable
import io.netty.channel.{ChannelFutureListener, ChannelHandlerContext, ChannelInboundHandlerAdapter}
import io.netty.handler.codec.http._

object SwaggerHandler {
  val SwaggerPath = "META-INF/resources/webjars/swagger-ui/3.0.10/"
  val DocsPrefix = "/docs/"
  val SpecsPrefix = "/specs/"
}

@Sharable
class SwaggerHandler extends ChannelInboundHandlerAdapter {
  import SwaggerHandler._
  override def channelRead(ctx: ChannelHandlerContext, msg: scala.Any): Unit = msg match {
    case req: FullHttpRequest =>
      val queryString = new QueryStringDecoder(req.uri())
      val path = queryString.path()

      if (path.startsWith(DocsPrefix)) {
        val filename = SwaggerPath + path.substring(DocsPrefix.length)
        writeResource(ctx, req, filename)
      } else if (path.startsWith(SpecsPrefix)) {
        val filename = path.substring(1) // remove heading slash
        writeResource(ctx, req, filename)
      } else super.channelRead(ctx, msg)
    case _ => super.channelRead(ctx, msg)
  }

  private def writeResource(ctx: ChannelHandlerContext, req: FullHttpRequest, filename: String): Unit = {
    Thread.currentThread().getContextClassLoader.getResourceAsStream(filename) match {
      case null =>
        val res = new DefaultFullHttpResponse(req.protocolVersion(), HttpResponseStatus.NOT_FOUND)
        ctx.writeAndFlush(res).addListener(ChannelFutureListener.CLOSE)
      case resource =>
        val bytes = IoUtils.toByteArray(resource)
        val res = new DefaultFullHttpResponse(
          req.protocolVersion(),
          HttpResponseStatus.OK,
          Unpooled.copiedBuffer(bytes)
        )
        res.headers().set(HttpHeaderNames.CONTENT_TYPE, new MimetypesFileTypeMap().getContentType(filename))
        HttpUtil.setContentLength(res, bytes.length)
        HttpUtil.setKeepAlive(res, HttpUtil.isKeepAlive(req))
        ctx.writeAndFlush(res)
    }
  }
} 
开发者ID:btlines,项目名称:grpcgateway,代码行数:54,代码来源:SwaggerHandler.scala


示例16: Akka1

//设置package包名称以及导入依赖的类
package jp.kenichi.lrcon.test.akka1

import jp.kenichi.lrcon.Service
import io.netty.channel.ChannelHandlerContext
import io.netty.handler.codec.http.{HttpRequest, DefaultFullHttpResponse}
import io.netty.handler.codec.http.HttpVersion.HTTP_1_1
import io.netty.handler.codec.http.HttpResponseStatus.OK
import io.netty.buffer.Unpooled
import io.netty.util.AsciiString
import akka.actor.{ActorSystem, Actor, Props}
import akka.event.Logging
import akka.util.Timeout

class Akka1

object Akka1 extends Service {
	val system = ActorSystem("akka1")
	class MyActor extends Actor {
		val log = Logging(context.system, this)
		def receive = {
			case (ctx: ChannelHandlerContext, req: HttpRequest) =>
				log.debug(s"req.uri = ${req.uri}")
				val content = "Hello ?? from Akka1".getBytes("UTF-8")
				val resp = new DefaultFullHttpResponse(HTTP_1_1, OK, Unpooled.wrappedBuffer(content))
				resp.headers.set(CONTENT_TYPE, "text/plain")
				resp.headers.setInt(CONTENT_LENGTH, resp.content.readableBytes)
				ctx.writeAndFlush(resp) // TODO: keep-alive
		}
	}
	val myActor = system.actorOf(Props[MyActor], "myActor")

	val CONTENT_TYPE = new AsciiString("Content-Type")
	val CONTENT_LENGTH = new AsciiString("Content-Length")

	def channelRead(ctx: ChannelHandlerContext, msg: AnyRef) = myActor ! (ctx, msg)
} 
开发者ID:ken1ma,项目名称:lrcon,代码行数:37,代码来源:Akka1.scala


示例17: App2

//设置package包名称以及导入依赖的类
package jp.kenichi.lrcon.test.app2

import jp.kenichi.lrcon.Service
import io.netty.channel.ChannelHandlerContext
import io.netty.handler.codec.http.{HttpRequest, DefaultFullHttpResponse}
import io.netty.handler.codec.http.HttpVersion.HTTP_1_1
import io.netty.handler.codec.http.HttpResponseStatus.OK
import io.netty.buffer.Unpooled
import io.netty.util.AsciiString
import org.slf4j.LoggerFactory

class App2

object App2 extends Service {
	val log = LoggerFactory.getLogger(classOf[App2])

	val CONTENT_TYPE = new AsciiString("Content-Type")
	val CONTENT_LENGTH = new AsciiString("Content-Length")

	def channelRead(ctx: ChannelHandlerContext, msg: AnyRef) = msg match {
		case req: HttpRequest =>
			log.debug(s"req.uri = ${req.uri}")
			val content = "Hello ?? from App2".getBytes("UTF-8")
			val resp = new DefaultFullHttpResponse(HTTP_1_1, OK, Unpooled.wrappedBuffer(content))
			resp.headers.set(CONTENT_TYPE, "text/plain")
			resp.headers.setInt(CONTENT_LENGTH, resp.content.readableBytes)
			ctx.writeAndFlush(resp) // TODO: keep-alive
	}
} 
开发者ID:ken1ma,项目名称:lrcon,代码行数:30,代码来源:App2.scala


示例18: Akka2

//设置package包名称以及导入依赖的类
package jp.kenichi.lrcon.test.akka2

import jp.kenichi.lrcon.Service
import io.netty.channel.ChannelHandlerContext
import io.netty.handler.codec.http.{HttpRequest, DefaultFullHttpResponse}
import io.netty.handler.codec.http.HttpVersion.HTTP_1_1
import io.netty.handler.codec.http.HttpResponseStatus.OK
import io.netty.buffer.Unpooled
import io.netty.util.AsciiString
import akka.actor.{ActorSystem, Actor, Props}
import akka.event.Logging
import akka.util.Timeout

class Akka2

object Akka2 extends Service {
	val system = ActorSystem("akka2")
	class MyActor extends Actor {
		val log = Logging(context.system, this)
		def receive = {
			case (ctx: ChannelHandlerContext, req: HttpRequest) =>
				log.debug(s"req.uri = ${req.uri}")
				val content = "Hello ?? from Akka2".getBytes("UTF-8")
				val resp = new DefaultFullHttpResponse(HTTP_1_1, OK, Unpooled.wrappedBuffer(content))
				resp.headers.set(CONTENT_TYPE, "text/plain")
				resp.headers.setInt(CONTENT_LENGTH, resp.content.readableBytes)
				ctx.writeAndFlush(resp) // TODO: keep-alive
		}
	}
	val myActor = system.actorOf(Props[MyActor], "myActor")

	val CONTENT_TYPE = new AsciiString("Content-Type")
	val CONTENT_LENGTH = new AsciiString("Content-Length")

	def channelRead(ctx: ChannelHandlerContext, msg: AnyRef) = myActor ! (ctx, msg)
} 
开发者ID:ken1ma,项目名称:lrcon,代码行数:37,代码来源:Akka2.scala


示例19: App1

//设置package包名称以及导入依赖的类
package jp.kenichi.lrcon.test.app1

import jp.kenichi.lrcon.Service
import io.netty.channel.ChannelHandlerContext
import io.netty.handler.codec.http.{HttpRequest, DefaultFullHttpResponse}
import io.netty.handler.codec.http.HttpVersion.HTTP_1_1
import io.netty.handler.codec.http.HttpResponseStatus.OK
import io.netty.buffer.Unpooled
import io.netty.util.AsciiString
import org.slf4j.LoggerFactory

class App1

object App1 extends Service {
	val log = LoggerFactory.getLogger(classOf[App1])

	val CONTENT_TYPE = new AsciiString("Content-Type")
	val CONTENT_LENGTH = new AsciiString("Content-Length")

	def channelRead(ctx: ChannelHandlerContext, msg: AnyRef) = msg match {
		case req: HttpRequest =>
			log.debug(s"req.uri = ${req.uri}")
			val content = "Hello ?? from App1".getBytes("UTF-8")
			val resp = new DefaultFullHttpResponse(HTTP_1_1, OK, Unpooled.wrappedBuffer(content))
			resp.headers.set(CONTENT_TYPE, "text/plain")
			resp.headers.setInt(CONTENT_LENGTH, resp.content.readableBytes)
			ctx.writeAndFlush(resp) // TODO: keep-alive
	}
} 
开发者ID:ken1ma,项目名称:lrcon,代码行数:30,代码来源:App1.scala


示例20: NSQProtocol

//设置package包名称以及导入依赖的类
package com.github.mitallast.nsq.protocol

import com.typesafe.config.ConfigFactory
import io.netty.buffer.{ByteBuf, Unpooled}
import io.netty.channel.local.LocalAddress
import io.netty.util.CharsetUtil
import org.slf4j.LoggerFactory

object NSQProtocol {

  val log = LoggerFactory.getLogger(getClass)
  val config = ConfigFactory.load("nsq")
  val localAddr = new LocalAddress("nsq.id")
  val json ="""{"client_id":"test","hostname":"localhost","feature_negotiation":true,"heartbeat_interval":1000,"user_agent":"test"}"""
  val timestamp = System.currentTimeMillis()
  val attempts = 2
  val messageId = "WCKHEOWCMPWECHWQ"
  val topic = "scala.nsq.test"

  def buf(buffers: ByteBuf*) = Unpooled.wrappedBuffer(buffers: _*)

  def buf(value: Int) = Unpooled.buffer(4).writeInt(value)

  def buf(value: Long) = Unpooled.buffer(8).writeLong(value)

  def buf(value: String) = Unpooled.copiedBuffer(value, CharsetUtil.US_ASCII)

  def buf(value: Array[Byte]) = Unpooled.wrappedBuffer(value)

  def requestBuf(header: String, data: ByteBuf): ByteBuf = buf(buf(header), buf(data.readableBytes()), data)

  def requestBuf(header: String, data: String): ByteBuf = requestBuf(header, buf(data))

  def responseBuf(data: String) = {
    val b = buf(data)
    buf(buf(b.readableBytes() + 4), buf(0), b)
  }

  def errorBuf(data: String) = {
    val b = buf(data)
    buf(buf(b.readableBytes() + 4), buf(1), b)
  }

  def messageBuf(data: String): ByteBuf = messageBuf(buf(data))

  def messageBuf(data: ByteBuf): ByteBuf = Unpooled.buffer()
    .writeInt(4 + 8 + 2 + 16 + data.readableBytes())
    .writeInt(2)
    .writeLong(timestamp)
    .writeShort(attempts)
    .writeBytes(messageId.getBytes(CharsetUtil.US_ASCII))
    .writeBytes(data)
} 
开发者ID:mitallast,项目名称:scala-nsq,代码行数:54,代码来源:NSQProtocol.scala



注:本文中的io.netty.buffer.Unpooled类示例整理自Github/MSDocs等源码及文档管理平台,相关代码片段筛选自各路编程大神贡献的开源项目,源码版权归原作者所有,传播和使用请参考对应项目的License;未经允许,请勿转载。


鲜花

握手

雷人

路过

鸡蛋
该文章已有0人参与评论

请发表评论

全部评论

专题导读
上一篇:
Scala DefaultFrontendGlobal类代码示例发布时间:2022-05-23
下一篇:
Scala Prop类代码示例发布时间:2022-05-23
热门推荐
热门话题
阅读排行榜

扫描微信二维码

查看手机版网站

随时了解更新最新资讯

139-2527-9053

在线客服(服务时间 9:00~18:00)

在线QQ客服
地址:深圳市南山区西丽大学城创智工业园
电邮:jeky_zhao#qq.com
移动电话:139-2527-9053

Powered by 互联科技 X3.4© 2001-2213 极客世界.|Sitemap