I'm putting together a core.async flow that takes bytes read from a socket as the flow's source process, parses them into clojure maps in an intermediate process, and those maps are used to update program state in an atom in the flow's sink process. [Socket Source] -> [Bytes->Map Parser] -> [Map->Atom State Sink] Based on the example flow, it seems that I need to write the socket's bytes to a port (async/chan) registered to the Socket Source, as a param. Sounds straightforward. But this socket could be recreated, for instance, if it disconnects. Because the port is registered with flow/in-ports during :init, it seems like I might not want to just directly wire the socket's output to the source process's input port because I can't cause :init to happen again to re-register. If I need to reconnect the socket, I could: - recreate the whole flow - bind the port to separate (async/chan) which gets its data from the actual socket, giving it a layer of indirection. I'd also have to flow/inject into a control channel to cause any stale bytes to be cleared from processes. It seems to me that recreating the flow instead of doing maintenance on it is simpler for my use case. But flow is new to me. Is recreating a flow OK, if it is done infrequently? Are there other gotyas?
There shouldn't be any problems with recreating flows. Especially if you want to throw away intermediate state, then creating a new flow can make sense.
Great! Makes my life simpler
One gotcha is that I'm not sure there's any clear way to know if a flow has fully stopped, https://clojurians.slack.com/archives/C05423W6H/p1770232006065359 So if you do create a new flow that shares state with the old flow, you may have to do an awkward dance to make sure the old flow has stopped using the shared resource before the new flow begins.
It's also possible to have a reconnecting source proc that is mostly written using classic core.async, but can be managed via flow. However, if you can get away with just recreating the flow, that can be a straightforward option.
I meant classic* core.async, not "class core.async"
You're saying I could move most of the socket code into the source process? That's an interesting idea! I had been keeping the source process as "clean" as possible, but perhaps I could pass in something like a protocol
I'd probably still have to send messages down the flow after a disconnect, to clear out the stale data
Since you can replace a "socket source" proc with any process that produces bytes, I think it's ok if the implementation isn't super clean
True enough. And that solve my puzzling over what actually goes into the source proc
in general, it's nice to have multiple input procs that can feed into your pipeline (test noise proc, file source proc, socket source proc, etc).
Good point
If you really are clearing stale data, then creating a new flow might still be the easiest option.
I was thinking alternative inputs would be fed in via the (async/chan) port
you can do either. it's nice to have input sources managed by flow so you can start/pause/stop everything
I see!
> I'd probably still have to send messages down the flow after a disconnect, to clear out the stale data
I've mostly been encoding messages as maps with an :op key so you can send control messages like :eof down through the flow.
That's an idea. I was also thinking of making a control channel, and/or using a UUID associated with the socket
Not sure if it's considered hygienic for a process to reach out to (via inject) other processes though 😕
(outside of a channel, that is)
I've previously thought about having control ports between procs, but I don't think there are guarantees about the order of reads/writes across chans. If you have: :new-stream, :data, :data, :data, :eof, :new-stream :data :data vs data and control: data chan -- :data :data :data :data :data control chan -- :new-stream, :eof, :new-stream, :eof I'm not sure if there are any guarantees that data that comes from a new stream will arrive after the corresponding control message.
Good point