History

I started my adventures with the Actor Model Paradigm back in 2021. I was working for Revemarketing at the time and we were working on something that required state maintenance and calculations based on said state. Our Architect (Vijay) designed the service using OpenFaas and Akka. And so my adventure began. The product had to be scrapped before we started making any revenue from it, but I’d learnt of a new architecture pattern.

The Problem

Fast forward to next year, around March of 2022, I was working with UKG on UKG Talk, and the push notifications service needed to be refactored to be able to scale. The service was written in Node.js. It handled push notifications for Posts, chats, group chats etc. and how it worked was something like this:

  • Read event details from PubSub
  • Figure out who to send the notifications for the event in question
  • Loop through the list of users
    • Figure out if the user has the notification setting turned on to get notified about the event in question
    • Send Notification if the setting is on

Since the fan out was happening after reading the event from PubSub, the for loop was the bottleneck.

The Research

I figured, if we could parallelize after the fanout, it would speed things up, by a lot, depending on how much scale we have.

Though there were other ways of solving that problem, perhaps even better ones than the one I was thinking, but in my mind, the first pattern that fit here, and perfectly, was the Actor model paradigm.

My first instinct was to rewrite the thing in Java using Akka, but we had one other Java resource in the team at the time, and since the team was all Node.js, we figured, let’s see if we have any other Node.js based frameworks.

And so I did. I found Comedy, I found NACT and I found akka.js. But unfortunately, none would work for us…

Comedy

Comedy is a solid framework, has Javascript support. But lacks these features (at the time I was researching):

  • Stateful Actors
  • Kubernetes API Support
  • Singleton Actors

NACT

NACT is also awesome. But it lacked these features (at the time I was researching):

  • Clustering
  • Kubernetes API Support
  • Singleton Actors

akka.js

akka.js is a port of Akka from the Scala world. The port works on top of scala.js. But, as evidenced in the akka.js-examples, the application logic needs to be written in Scala, and not Javascript. It will run on Node, but the code will not be native Javascript. So this was a dead end too.

The solution

So I sat down on a (boring) Friday evening, with a couple of beers, and started writing. By 3 am of Monday morning, I was done with a minimum viable framework. I wrote a ping pong example too, to test if it works. I named it JsAM (Javascript Actor Model Framework). The comms were all a beautiful spaghetti of callbacks and callback caching and Node’s event loop.

I pitched it to my boss and our architect on that Monday, and we saw how it fit our use case. Showed it our VP Praveen the next day and got his blessing too. In the next week, I wrote a few more examples, specifically testing out how much speed I could get it to for a million actors. Faced a few roadblocks, fixed those roadblocks, got a framework.

In a few weeks, I had an Actor Model framework with all the features we need.

The features

  • Clustering - through Leader Elections
  • Stateful Actors
  • Singleton Actors
  • Queued Actors (we only have queued actors)
  • Management & Monitoring endpoints - Prometheus support was later added to track these
  • Actor Communications through gRPC
  • Actor Respawning - In case nodes go down
  • Cluster Rebalancing - In case nodes are added or removed
  • Kubernetes API support
  • Actor State Compression - This is optional, and since it’s all node.js, caused more grief than good
  • Actor System Node Caching - For when transferring 1 Meg of data 10 times (the number of nodes) is more efficient than transferring it a million times (the number of actors).

I made a cool logo and everything. 😊

I spent almost every (personal) weekend that year, testing and perfecting the framework, ironing out the creases, bugfix after bugfix after bugfix. It was grueling, trying to debug a distributed system, but it was worth every swear I hurled at the computer. šŸ˜…

It was a labor of love…

Roadblocks

Some of the roadblocks were thus:

  • Comms were too slow. I was using got, and all actor to actor communications were Rest calls, and turns out it wasn’t as fast (efficient?) - Replaced it with grpc-node.
  • Actor branching was too slow - Even with gRPC, one actor reaching out to a million in one of the examples was taking a while. But that was more a design problem than a problem with the framework. The solution I came up with was to put another layer between the main asking actor and the million asked actors. So, ultimately the asker becomes a grandparent, which would ask a thousand parent actors, which would then ask a thousand of their child actors, thereby reducing the parallelization need from a million to a thousand. And I used async.parallelLimit for it.
  • For Kubernetes deployments, if a pod goes down and comes back up, it gets assigned the same IP. Since we were tracking nodes by IP and port (and port won’t change for k8s deployments), the cluster couldn’t tell the difference between the dead instance and the new instance. That was an easy enough fix once identified. I just added a UUID run identifier, and included that in the pings and elections.
  • Messages being passed around were too big - This one came up during the implementation of the notifications service. Each message was about a Megabyte (we had fields in there nobody knew why). Now considering an event where we have to notify a million users, that would mean a Gigabyte of bandwidth wasted just transferring that one event. So instead, I implemented Actor System Caching, where we’d transfer that message to all nodes, and they’ll cache it. When asking the million actors, we’d just send the cache key, with which they could look it up from their Actor System cache (in node), and process further.
  • Actor Respawning failed sometimes - This has been a persistent thorn in my side. And as of a few weeks ago, the reason for the final nail in the coffin of JsAM. This was a problem with both Rebalancing and Respawning. When an actor is transferred, and that actor is being asked, even with the retries, it would still go to the old address instead trying to find the new one. What I could figure out, was that ping wasn’t detecting the failure fast enough, or if it did, the actor didn’t. I never could figure out if it was a problem with the code or if it was because of the load, but in all the examples I ran (with the latest iteration after bugfixes), it never failed once. Yet in actual use, it did fail and stop processing altogether. We put workarounds in place, so we couldn’t have rolling updates and the underlying k8s nodes can’t go down, or autoscaling.

The Future

Since it’s going to not be used in a few months, and all my efforts to get my boss to consider open sourcing it have failed, it will die in purgatory, never having been seen, or cloned, checked out or used. šŸ˜”

This was a very personal journey, to find meaning, to contribute something to this world. And so, the loss feels that much deeper. Yes, I realize that I’m compensating for a lack of said meaning in my life, but still… It was fucking glorious.

But, after having mourned, I will write a new one. This time in Typescript, because, y’know… Types. And hopefully that will help catch some problems. Others, I’m hoping the open source community will help with.