# How can I return a live websocket connection in HTTP.jl?

**URL:** <https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752>\
**Category:** General Usage\
**Tags:** question\
**Created:** [August 20, 2021, 9:00pm UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752 "2021-08-20T21:00:14Z")\
**Posts on this page:** 18\
**Page:** 1

<div class="post-metadata">

**Author:** ![PyDataBlog](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pydatablog/32/12327_2.png) [@PyDataBlog](https://discourse.julialang.org/u/PyDataBlog)\
**Post date:** [August 20, 2021, 9:00pm UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/1 "2021-08-20T21:00:14Z")

</div>

I have a live streaming connection opened via HTTP.jl’s Websocket with this snippet

```julia
using HTTP
using JSON3

HTTP.WebSockets.open("wss://ws.coincap.io/trades/binance") do ws
    while !eof(ws)
        println(readavailable(ws) |> JSON3.read)
    end
end

```

My goal is to make a function that returns the live streaming connection in `ws` such that this connection can be used by other processes in my codebase.

Here’s my current non-working attempt to break this operation with a function:

```julia
function open_socket()
    HTTP.WebSockets.open("wss://ws.coincap.io/trades/binance") do ws
        return ws
    end
end

con = open_socket()

while !eof(con)
    println(readavailable(con) |> JSON3.read)
end

```

Thanks in advance!

---

<div class="post-metadata">

**Author:** ![xiaodai](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/xiaodai/32/15937_2.png) [@xiaodai](https://discourse.julialang.org/u/xiaodai)\
**Post date:** [August 21, 2021, 12:15am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/2 "2021-08-21T00:15:57Z")

</div>

> [@PyDataBlog](#):
>
> ```julia
> function open_socket()
> HTTP.WebSockets.open("wss://ws.coincap.io/trades/binance") 
> end
> 
> ```

Just guessing…

---

<div class="post-metadata">

**Author:** ![PyDataBlog](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pydatablog/32/12327_2.png) [@PyDataBlog](https://discourse.julialang.org/u/PyDataBlog)\
**Post date:** [August 21, 2021, 8:48am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/3 "2021-08-21T08:48:06Z")

</div>

I tried that and got this error:

```julia
function open_socket()
    return HTTP.WebSockets.open("wss://ws.coincap.io/trades/binance")
end

con = open_socket()

```

```julia-repl
julia> con = open_socket()
ERROR: MethodError: no method matching open(::String)
You may have intended to import Base.open
Closest candidates are:
  open(::Function, ::Any; binary, verbose, headers, kw...) at /Users/../.julia/packages/HTTP/D0FSE/src/WebSockets.jl:92
Stacktrace:
 [1] open_socket()
   @ Main ./REPL[14]:2
 [2] top-level scope
   @ REPL[15]:1

```

---

<div class="post-metadata">

**Author:** ![Sukera](https://avatars.discourse-cdn.com/v4/letter/s/ce7236/32.png) [@Sukera](https://discourse.julialang.org/u/Sukera)\
**Post date:** [August 21, 2021, 9:17am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/4 "2021-08-21T09:17:22Z")

</div>

> [@PyDataBlog](#):
>
> My goal is to make a function that returns the live streaming connection in `ws` such that this connection can be used by other processes in my codebase.

Why not just pass `ws` to your function directly in the `do` block? Or put it into a `Channel`, from which you pull on another task, if you really want to move sockets around like that.

You don’t have to use `do` either - you can pass your favorite processing function to `open` directly as well:

```julia
my_ws_processing_func(ws) = ...

HTTP.open(my_ws_processing_func, "wss://ws.coincap.io/trades/binance")

```

[`do` after all is just fancy syntax](https://docs.julialang.org/en/v1/base/base/#do) for passing an anonymous function to the caller as its first argument.

---

<div class="post-metadata">

**Author:** ![PyDataBlog](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pydatablog/32/12327_2.png) [@PyDataBlog](https://discourse.julialang.org/u/PyDataBlog)\
**Post date:** [August 21, 2021, 9:50am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/5 "2021-08-21T09:50:05Z")

</div>

The primary issue is that I need that `ws` connection returned to be used in some external function. I am using [Rocket.jl](https://github.com/biaslab/Rocket.jl) to process this stream

---

<div class="post-metadata">

**Author:** ![Sukera](https://avatars.discourse-cdn.com/v4/letter/s/ce7236/32.png) [@Sukera](https://discourse.julialang.org/u/Sukera)\
**Post date:** [August 21, 2021, 9:56am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/6 "2021-08-21T09:56:43Z")

</div>

Do you have an example of what you would do if you were able to get at that directly? Some pseudocode maybe?

---

<div class="post-metadata">

**Author:** ![PyDataBlog](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pydatablog/32/12327_2.png) [@PyDataBlog](https://discourse.julialang.org/u/PyDataBlog)\
**Post date:** [August 21, 2021, 10:02am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/7 "2021-08-21T10:02:19Z")

</div>

Sure. Something like this:

```julia
using Rocket
using HTTP
using JSON3

struct WebSocketObservable <: ScheduledSubscribable{Any}
    # Any field here
    api_key::AbstractString
    tickers::AbstractString
end

struct WebsocketSubscription <: Teardown
    #.. some fields ...
    live_socket::Any
    tickers::AbstractString
end

struct MyActor <: Rocket.Actor{Any} end

Rocket.getscheduler(::WebSocketObservable) = AsyncScheduler()

function Rocket.on_subscribe!(source::WebSocketObservable, actor, scheduler)

    @async begin
        try
            HTTP.WebSockets.open("wss://socket.polygon.io/crypto") do ws
                while isopen(ws)
                    # Authenticate & Subscribe to ticker stream
                    write(ws, JSON3.read("""{"action":"auth", "params":$(source.api_key)}"""))
                    write(ws, JSON3.read("""{"action":"subscribe", "params":$(source.tickers)}"""))

                    if !eof(ws)
                        next!(actor, readavailable(ws) |> JSON3.read, scheduler)
                    else
                        complete!(actor, scheduler)
                    end
                end
            end
        catch e
            if !e.is_a(EOFError)
                error!(actor, e, scheduler)
            end
        end
    #return WebsocketSubscription(ws, source.tickers)
    end

    return WebsocketSubscription(ws, source.tickers) # can't get access to ws in the above
end

function Rocket.on_unsubscribe!(subscription::WebsocketSubscription)
    # stop listening for web socket here
    write(subscription.live_socket, JSON3.read("""{"action":"unsubscribe", "params":$(subscription.tickers)}"""))

    # usually unsubscription returns nothing
    return nothing
end

obsv = WebSocketObservable("API_KEY", "XT.*")
subscription = subscribe!(obsv, logger())
unsubscribe!(subscription)

```

As you can see, I need access to `ws` to be able to unsubscribe from the stream.

---

<div class="post-metadata">

**Author:** ![Sukera](https://avatars.discourse-cdn.com/v4/letter/s/ce7236/32.png) [@Sukera](https://discourse.julialang.org/u/Sukera)\
**Post date:** [August 21, 2021, 10:13am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/8 "2021-08-21T10:13:46Z")

</div>

[Since `HTTP.open` closes the web socket for you](https://github.com/JuliaWeb/HTTP.jl/blob/a2b467e24c9bbd45691f9d0f57b57ee7463bd15a/src/WebSockets.jl#L92-L128) when the passed function exits, I’d use a different mechanism for communicating that the socket should be closed rather than closing the socket manually. E.g. pass a channel to the function handling the websocket, `put!` a message that you want to exit in it and `return` based on that. The idea is to _not_ pass the socket around and have a bunch of places reading/writing from/to it, but only one instead. That makes your code easier to maintain, since you don’t have to hunt down erronous writes (and safer as well, since you don’t have multiple tasks potentially writing to the socket at the same time…).

You’ll also have to handle the socket closing prematurely, e.g. if the server doesn’t respond anymore and communicate that to Rocket.jl as well.

---

<div class="post-metadata">

**Author:** ![PyDataBlog](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pydatablog/32/12327_2.png) [@PyDataBlog](https://discourse.julialang.org/u/PyDataBlog)\
**Post date:** [August 21, 2021, 10:45am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/9 "2021-08-21T10:45:10Z")

</div>

Never used websockets or channels before. Is it possible to have a practical example?

---

<div class="post-metadata">

**Author:** ![Sukera](https://avatars.discourse-cdn.com/v4/letter/s/ce7236/32.png) [@Sukera](https://discourse.julialang.org/u/Sukera)\
**Post date:** [August 21, 2021, 10:58am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/10 "2021-08-21T10:58:21Z")

</div>

Something like this:

```julia
struct WebSocketObservable <: ScheduledSubscribable{Any}
    # Any field here
    api_key::AbstractString
    tickers::AbstractString
    # utility
    comm::Channel
end

struct WebsocketSubscription <: Teardown
    #.. some fields ...
    live_socket::Any
    tickers::AbstractString

    # utility
    comm::Channel
end

function handle_socket(c::Channel)
    try
        HTTP.WebSockets.open(...) do ws
            while isopen(ws) && isopen(c)
                # do regular socket processing
            end

           if !isopen(ws) # something failed
               # [...]
           end
           if !isopen(c) # we closed down regularly, have to check whether the socket is open!
               if isready(c) # we have data available
                  ret = take!(c)
                  # handle close, e.g. by doing this
                  write(ws, JSON3.read("""{"action":"unsubscribe", "params":$(subscription.tickers)}"""))

               else
                   # closed channel but no message? => Error?
               end
           end
        end
    catch
      # do catching stuff
    end
end

function Rocket.on_subscribe!(source::WebSocketObservable, actor, scheduler)
    chan = source.comm
    # untested, may have to interpolate `chan` here via `$chan`. Check `?@async` for more info.
    @async handle_socket(chan)
    return WebsocketSubscription(chan, source.tickers)
end

function Rocket.on_unsubscribe!(subscription::WebsocketSubscription)
    # stop listening for web socket here
    put!(subscription.comm, "we're done here")
    close(subscription.comm)

    # usually unsubscription returns nothing
    return nothing
end

```

You may have to interpolate `chan` into the `@async` expression via `$chan` - check `?@async` for more information.

You could also use the `Channel` directly for passing messages that should be written to the websocket (and close the socket if it’s a closing message). One way would be to have a Channel that holds `Tuple{Bool, String}`, the first being a boolean indicating whether to terminate the socket (usually false) and the second being the message to be sent. That would move the decision making logic of “what to send” out of the “how to send” logic, increasing seperation between these two concepts.

---

<div class="post-metadata">

**Author:** ![PyDataBlog](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pydatablog/32/12327_2.png) [@PyDataBlog](https://discourse.julialang.org/u/PyDataBlog)\
**Post date:** [August 22, 2021, 10:20am UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/11 "2021-08-22T10:20:13Z")

</div>

> [@Sukera](#):
>
> Something like this:

Thanks for your suggestions. I tried that and was unsuccessful. Here’s the attempt:

```julia
using Rocket
using HTTP
using JSON3

struct WebSocketObservable <: ScheduledSubscribable{Any}
    # Any field here
    api_key::AbstractString
    tickers::AbstractString
    # utility
    comm::AbstractChannel
end

struct WebsocketSubscription <: Teardown
    # utility
    comm::AbstractChannel
    #.. some fields ...
    tickers::AbstractString
end

struct MyActor <: Rocket.Actor{Any} end

Rocket.getscheduler(::WebSocketObservable) = AsyncScheduler()

function handle_socket(source::WebSocketObservable, c::Channel, actor, scheduler)
    try
        HTTP.WebSockets.open("wss://socket.polygon.io/crypto") do ws

            while isopen(ws) && isopen(c)
                # Authenticate & Subscribe to ticker stream
                write(ws, JSON3.read("""{"action":"auth", "params":$(source.api_key)}"""))
                write(ws, JSON3.read("""{"action":"subscribe", "params":$(source.tickers)}"""))

                if !eof(ws)
                    next!(actor, readavailable(ws) |> JSON3.read, scheduler)
                else
                    complete!(actor, scheduler)
                end
            end

        # if !isopen(ws) # something failed
        # # [...]
        # end

           if !isopen(c) # we closed down regularly, have to check whether the socket is open!
               if isready(c) # we have data available
                  ret = take!(c)
                  # handle close, e.g. by doing this
                  write(ws, JSON3.read("""{"action":"unsubscribe", "params":$(source.tickers)}"""))

               else
                   # closed channel but no message? => Error?
                   nothing # not sure what to do here
               end
           end

        end

    catch e
        if !e.is_a(EOFError)
            error!(actor, e, scheduler)
        end
    end
end

function Rocket.on_subscribe!(source::WebSocketObservable, actor, scheduler)
    chan = source.comm
    # untested, may have to interpolate `chan` here via `$chan`. Check `?@async` for more info.
    @async handle_socket(source.tickers, chan, actor, scheduler) # tried both `$chan` and `chan`
    return WebsocketSubscription(chan, source.tickers)
end

function Rocket.on_unsubscribe!(subscription::WebsocketSubscription)
    # stop listening for web socket here
    put!(subscription.comm, "we're done here")
    close(subscription.comm)

    # usually unsubscription returns nothing
    return nothing
end

obsv = WebSocketObservable("API_KEY", "XT.*", Channel())
subscription = subscribe!(obsv, logger()) # logger() actor just logs any available data in the subscription
unsubscribe!(subscription)

```

---

<div class="post-metadata">

**Author:** ![Sukera](https://avatars.discourse-cdn.com/v4/letter/s/ce7236/32.png) [@Sukera](https://discourse.julialang.org/u/Sukera)\
**Post date:** [August 22, 2021, 12:22pm UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/12 "2021-08-22T12:22:18Z")

</div>

> [@PyDataBlog](#):
>
> Thanks for your suggestions. I tried that and was unsuccessful. Here’s the attempt:

What do you mean by that? What did you expect to happen, what did happen? My solution was not meant as drop-in replacement for your code, I just formulated what I thought of in some form that may reasonably run, or failing that, convey what I meant with the approach.

> [@PyDataBlog](#):
>
> tried both `$chan` and `chan`

Have you interpolated the other arguments as well?

---

<div class="post-metadata">

**Author:** ![PyDataBlog](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pydatablog/32/12327_2.png) [@PyDataBlog](https://discourse.julialang.org/u/PyDataBlog)\
**Post date:** [August 22, 2021, 12:34pm UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/13 "2021-08-22T12:34:33Z")

</div>

> [@Sukera](#):
>
> Have you interpolated the other arguments as well?

Tried it with all args and it was the same no output after subscription scenario

> [@Sukera](#):
>
> What do you mean by that? What did you expect to happen, what did happen? My solution was not meant as drop-in replacement for your code, I just formulated what I thought of in some form that may reasonably run, or failing that, convey what I meant with the approach.

If a connection is successful, usually some data start to stream through. I really like your approach tho. Looks cleaner.

---

<div class="post-metadata">

**Author:** ![Sukera](https://avatars.discourse-cdn.com/v4/letter/s/ce7236/32.png) [@Sukera](https://discourse.julialang.org/u/Sukera)\
**Post date:** [August 22, 2021, 1:34pm UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/14 "2021-08-22T13:34:17Z")

</div>

Could be that you’re catching an error and just ignoring it - maybe try printing all catched errors?

---

<div class="post-metadata">

**Author:** ![PyDataBlog](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pydatablog/32/12327_2.png) [@PyDataBlog](https://discourse.julialang.org/u/PyDataBlog)\
**Post date:** [August 22, 2021, 1:54pm UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/15 "2021-08-22T13:54:46Z")

</div>

> [@PyDataBlog](#):
>
> `"wss://ws.coincap.io/trades/binance"`

Ok let’s try this idea on a public web socket connection:

```julia
using Rocket
using HTTP
using JSON3

struct WebSocketObservable <: ScheduledSubscribable{Any}
    # Any field here

    # utility
    comm::AbstractChannel
end

struct WebsocketSubscription <: Teardown
    # utility
    comm::AbstractChannel
end

Rocket.getscheduler(::WebSocketObservable) = AsyncScheduler()

function handle_socket(source::WebSocketObservable, c::Channel, actor, scheduler)
    try
        HTTP.WebSockets.open("wss://ws.coincap.io/trades/binance") do ws

            while isopen(ws) && isopen(c)

                if !eof(ws)
                    next!(actor, readavailable(ws) |> JSON3.read, scheduler)
                else
                    complete!(actor, scheduler)
                end
            end

        # if !isopen(ws) # something failed
        # # [...]
        # end

        if !isopen(c) # we closed down regularly, have to check whether the socket is open!
            if isready(c) # we have data available
                ret = take!(c)
            else
                # closed channel but no message? => Error?
                println(c)
                nothing # not sure what to do here
            end
        end

        end

    catch e
        println(e)
        if !e.is_a(EOFError)
            error!(actor, e, scheduler)
        end
    end
end

function Rocket.on_subscribe!(source::WebSocketObservable, actor, scheduler)
    chan = source.comm
    # untested, may have to interpolate `chan` here via `$chan`. Check `?@async` for more info.
    @async handle_socket(source.tickers, $chan, $actor, $scheduler) # tried both `$chan` and `chan`
    return WebsocketSubscription(chan)
end

function Rocket.on_unsubscribe!(subscription::WebsocketSubscription)
    # stop listening for web socket here
    put!(subscription.comm, "we're done here")
    close(subscription.comm)

    # usually unsubscription returns nothing
    return nothing
end

obsv = WebSocketObservable(Channel())
# logger() actor just logs any available data in the subscription
subscription = subscribe!(obsv, logger())
unsubscribe!(subscription)

```

---

<div class="post-metadata">

**Author:** ![Sukera](https://avatars.discourse-cdn.com/v4/letter/s/ce7236/32.png) [@Sukera](https://discourse.julialang.org/u/Sukera)\
**Post date:** [August 22, 2021, 2:01pm UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/16 "2021-08-22T14:01:30Z")

</div>

> [@PyDataBlog](#):
>
> ` @async handle_socket(source.tickers, $chan, $actor, $scheduler)`

Not 100% sure what you mean, as I’m not running your code and you’re not telling me what (if any) error you’re getting, but I suspect you’ll want to interpolate `source.tickers` as well. Have you checked with the docstring of `@async`, as my comment suggested?

---

<div class="post-metadata">

**Author:** ![PyDataBlog](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pydatablog/32/12327_2.png) [@PyDataBlog](https://discourse.julialang.org/u/PyDataBlog)\
**Post date:** [August 22, 2021, 2:26pm UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/17 "2021-08-22T14:26:19Z")

</div>

> [@Sukera](#):
>
> Not 100% sure what you mean, as I’m not running your code and you’re not telling me what (if any) error you’re getting, but I suspect you’ll want to interpolate `source.tickers` as well. Have you checked with the docstring of `@async` , as my comment suggested?

According to the docs, it needs to be interpolated (which I did). Tried to print the errors but nothing prints. No errors. Subscription is initialised and that’s it nothing happens.

---

<div class="post-metadata">

**Author:** ![gentstats](https://avatars.discourse-cdn.com/v4/letter/g/838e76/32.png) [@gentstats](https://discourse.julialang.org/u/gentstats)\
**Post date:** [December 20, 2023, 7:09pm UTC](https://discourse.julialang.org/t/how-can-i-return-a-live-websocket-connection-in-http-jl/66752/18 "2023-12-20T19:09:15Z")

</div>

Maybe too late for the purpose of this post, but today I’d like to read something like this:

```julia
using HTTP: WebSockets, send
using JuliaTrader: startWSStreamer, streamMessage

function open_socket(url)
    conn = false
    @async WebSockets.open(url) do ws
        conn = Ref(ws)
        while !WebSockets.isclosed(conn[])
            sleep(1) # Keep the connection open
        end
    end
    while !isa(conn[], WebSockets.WebSocket)
        sleep(0.01)
    end
    return conn
end

# A WS Server that streams messages to all connected clients
startWSStreamer("0.0.0.0", 54321)

conn = open_socket("ws://0.0.0.0:54321")
@async begin
    for msg in conn[]
        println("Client received: $msg")
    end
end

streamMessage("hello from server")
send(conn[], "hello from client")

```

```julia
julia> streamMessage("hello from server")
Client received: hello from server
julia> HTTP.send(conn[], "hello from client")
Server received: hello from client

```
