core-async 2026-09-04

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.

👍 1

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"

which is just using core.async without core.async flow

👍 1

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).

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'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.