Current section

53 Versions

Jump to

Compare versions

6 files changed
+415 additions
-64 deletions
  @@ -0,0 +1,343 @@
1 + ## 1. Overall System Architecture - High Performance Overview
2 +
3 + ```mermaid
4 + graph LR
5 + subgraph "Elixir/OTP Layer"
6 + Client["Client API"]
7 + Pool["Pool Manager<br/>⚡ Non-blocking async"]
8 + TaskSup["Task Supervisor<br/>⚡ Isolated execution"]
9 +
10 + subgraph "Worker Management"
11 + WorkerSup["Worker Supervisor<br/>⚡ Dynamic workers"]
12 + Starter1["Worker Starter 1<br/>🔄 Auto-restart"]
13 + Starter2["Worker Starter 2<br/>🔄 Auto-restart"]
14 + StarterN["Worker Starter N<br/>🔄 Auto-restart"]
15 + Worker1["Worker 1<br/>GenServer"]
16 + Worker2["Worker 2<br/>GenServer"]
17 + WorkerN["Worker N<br/>GenServer"]
18 + end
19 +
20 + subgraph "High-Performance Registries"
21 + Registry["Worker Registry<br/>⚡ O(1) lookups"]
22 + ProcReg["Process Registry<br/>⚡ PID tracking"]
23 + StarterReg["Starter Registry<br/>⚡ Supervisor tracking"]
24 + end
25 +
26 + subgraph "Session Store (ETS)"
27 + SessionStore["Session Store<br/>⚡ Concurrent R/W<br/>📊 Decentralized counters"]
28 + GlobalPrograms["Global Programs<br/>⚡ Public table access"]
29 + end
30 + end
31 +
32 + subgraph "External Processes"
33 + Python1["Python Process 1<br/>🐍 Port communication"]
34 + Python2["Python Process 2<br/>🐍 Port communication"]
35 + PythonN["Python Process N<br/>🐍 Port communication"]
36 + end
37 +
38 + Client -->|"⚡ Async call"| Pool
39 + Pool -->|"⚡ O(1) checkout"| Registry
40 + Pool -->|"🎯 Session affinity"| SessionStore
41 + Pool -->|"⚡ Task.async_nolink"| TaskSup
42 + TaskSup -->|Execute| Worker1
43 + TaskSup -->|Execute| Worker2
44 + TaskSup -->|Execute| WorkerN
45 +
46 + WorkerSup -->|Supervise| Starter1
47 + WorkerSup -->|Supervise| Starter2
48 + WorkerSup -->|Supervise| StarterN
49 +
50 + Starter1 -->|Auto-restart| Worker1
51 + Starter2 -->|Auto-restart| Worker2
52 + StarterN -->|Auto-restart| WorkerN
53 +
54 + Worker1 -->|Binary protocol<br/>4-byte frames| Python1
55 + Worker2 -->|Binary protocol<br/>4-byte frames| Python2
56 + WorkerN -->|Binary protocol<br/>4-byte frames| PythonN
57 +
58 + Worker1 -->|Register| Registry
59 + Worker2 -->|Register| Registry
60 + WorkerN -->|Register| Registry
61 +
62 + Worker1 -->|Track PID| ProcReg
63 + Worker2 -->|Track PID| ProcReg
64 + WorkerN -->|Track PID| ProcReg
65 +
66 + style Pool fill:#f9f,stroke:#333,stroke-width:4px,color:#000
67 + style SessionStore fill:#bbf,stroke:#333,stroke-width:4px,color:#000
68 + style Registry fill:#bfb,stroke:#333,stroke-width:4px,color:#000
69 + style TaskSup fill:#fbf,stroke:#333,stroke-width:4px,color:#000
70 + ```
71 +
72 + ## 2. Request Flow - Performance Critical Path
73 +
74 + ```mermaid
75 + sequenceDiagram
76 + participant C as Client
77 + participant P as Pool<br/>⚡ Non-blocking
78 + participant TS as TaskSupervisor<br/>⚡ Isolated
79 + participant R as Registry<br/>⚡ O(1)
80 + participant S as SessionStore<br/>⚡ ETS
81 + participant W as Worker
82 + participant E as External Process
83 +
84 + C->>P: execute(command, args)
85 +
86 + alt Session-based request
87 + P->>S: get_preferred_worker<br/>⚡ O(1) ETS lookup
88 + S-->>P: worker_id or nil
89 + end
90 +
91 + P->>R: checkout_worker<br/>⚡ O(1) via Registry
92 + R-->>P: worker_id
93 +
94 + P->>TS: Task.async_nolink<br/>⚡ Non-blocking
95 + Note over P: Pool returns immediately<br/>to handle next request
96 +
97 + TS->>W: GenServer.call
98 + W->>E: Port.command<br/>⚡ Binary protocol
99 + E-->>W: Response<br/>⚡ 4-byte framed
100 + W-->>TS: Result
101 + TS-->>C: GenServer.reply<br/>⚡ Direct to client
102 +
103 + TS->>P: checkin_worker<br/>⚡ Cast (async)
104 +
105 + alt Queued requests exist
106 + P->>P: Process next<br/>from queue
107 + else No queued requests
108 + P->>R: Mark available<br/>⚡ O(1) update
109 + end
110 + ```
111 +
112 + ## 3. ETS Tables Architecture - High Performance Storage
113 +
114 + ```mermaid
115 + graph LR
116 + subgraph "Session Store ETS Tables"
117 + subgraph "Sessions Table"
118 + ST[":snakepit_sessions<br/>⚡ read_concurrency: true<br/>⚡ write_concurrency: true<br/>⚡ decentralized_counters: true"]
119 +
120 + S1["Key: session_1<br/>Value: {last_accessed, ttl, session_data}"]
121 + S2["Key: session_2<br/>Value: {last_accessed, ttl, session_data}"]
122 + SN["Key: session_N<br/>Value: {last_accessed, ttl, session_data}"]
123 + end
124 +
125 + subgraph "Global Programs Table"
126 + GP[":snakepit_sessions_global_programs<br/>⚡ Same optimizations"]
127 +
128 + P1["Key: program_1<br/>Value: {data, timestamp}"]
129 + P2["Key: program_2<br/>Value: {data, timestamp}"]
130 + PN["Key: program_N<br/>Value: {data, timestamp}"]
131 + end
132 + end
133 +
134 + subgraph "Optimized Operations"
135 + Read["⚡ Concurrent reads<br/>No locking"]
136 + Write["⚡ Concurrent writes<br/>Decentralized counters"]
137 + Cleanup["⚡ select_delete<br/>Atomic batch cleanup"]
138 + end
139 +
140 + ST --> S1
141 + ST --> S2
142 + ST --> SN
143 +
144 + GP --> P1
145 + GP --> P2
146 + GP --> PN
147 +
148 + Read --> ST
149 + Read --> GP
150 + Write --> ST
151 + Write --> GP
152 + Cleanup --> ST
153 + Cleanup --> GP
154 +
155 + style ST fill:#bbf,stroke:#333,stroke-width:4px,color:#000
156 + style GP fill:#bbf,stroke:#333,stroke-width:4px,color:#000
157 + style Read fill:#bfb,stroke:#333,stroke-width:2px,color:#000
158 + style Write fill:#bfb,stroke:#333,stroke-width:2px,color:#000
159 + style Cleanup fill:#fbb,stroke:#333,stroke-width:2px,color:#000
160 + ```
161 +
162 + ## 4. Worker Lifecycle - Performance & Reliability
163 +
164 + ```mermaid
165 + stateDiagram-v2
166 + [*] --> Starting: Pool requests worker
167 +
168 + Starting --> Initializing: Port opened<br/>⚡ Parallel startup
169 +
170 + Initializing --> Ready: Init ping OK<br/>📊 Telemetry emitted
171 + Initializing --> Failed: Timeout/Error
172 +
173 + Ready --> Busy: Request received<br/>⚡ O(1) checkout
174 + Busy --> Ready: Response sent<br/>⚡ O(1) checkin
175 +
176 + Ready --> HealthCheck: Periodic check<br/>⏱️ Every 30s
177 + HealthCheck --> Ready: Healthy
178 + HealthCheck --> Unhealthy: Failed
179 +
180 + Unhealthy --> Restarting: Supervisor detects
181 + Failed --> Restarting: Auto-restart
182 +
183 + Restarting --> Starting: ♻️ Via Starter
184 +
185 + Ready --> Terminating: Shutdown signal
186 + Busy --> Terminating: Graceful shutdown
187 +
188 + Terminating --> [*]: Process cleaned up
189 +
190 + note right of Ready
191 + ⚡ Worker pool maintains
192 + hot workers ready for
193 + immediate use
194 + end note
195 +
196 + note right of Busy
197 + ⚡ Non-blocking async
198 + execution via Task
199 + Supervisor
200 + end note
201 + ```
202 +
203 + ## 5. Concurrent Initialization Performance
204 +
205 + ```mermaid
206 + graph TD
207 + subgraph "Sequential Startup (Traditional)"
208 + T0["Start"] --> W1S["Worker 1<br/>2s"]
209 + W1S --> W2S["Worker 2<br/>2s"]
210 + W2S --> W3S["Worker 3<br/>2s"]
211 + W3S --> W4S["Worker 4<br/>2s"]
212 + W4S --> DoneS["Ready<br/>Total: 8s"]
213 + end
214 +
215 + subgraph "Concurrent Startup (Snakepit)"
216 + T0C["Start"] --> Init["Task.async_stream"]
217 + Init --> W1C["Worker 1<br/>2s"]
218 + Init --> W2C["Worker 2<br/>2s"]
219 + Init --> W3C["Worker 3<br/>2s"]
220 + Init --> W4C["Worker 4<br/>2s"]
221 +
222 + W1C --> Collect
223 + W2C --> Collect
224 + W3C --> Collect
225 + W4C --> Collect
226 +
227 + Collect --> DoneC["Ready<br/>Total: ~2s"]
228 + end
229 +
230 + style Init fill:#f9f,stroke:#333,stroke-width:4px,color:#000
231 + style DoneC fill:#bfb,stroke:#333,stroke-width:4px,color:#000
232 + style DoneS fill:#fbb,stroke:#333,stroke-width:2px,color:#000
233 + ```
234 +
235 + ## 6. Request Queueing & Load Distribution
236 +
237 + ```mermaid
238 + graph LR
239 + subgraph "High-Performance Request Handling"
240 + subgraph "Request Queue"
241 + Q[":queue (Erlang)<br/>⚡ FIFO<br/>⚡ O(1) operations"]
242 + R1["Request 1"]
243 + R2["Request 2"]
244 + R3["Request 3"]
245 + RN["Request N"]
246 + end
247 +
248 + subgraph "Worker Pool State"
249 + Available["MapSet<br/>⚡ O(1) member check<br/>⚡ O(1) add/remove"]
250 + Busy["Map<br/>⚡ O(1) lookup"]
251 +
252 + AW1["Worker 1"]
253 + AW2["Worker 2"]
254 + BW3["Worker 3 🔴"]
255 + BW4["Worker 4 🔴"]
256 + end
257 +
258 + subgraph "Load Distribution"
259 + Check{"Worker<br/>Available?"}
260 + Assign["Assign to worker<br/>⚡ O(1)"]
261 + Queue["Queue request<br/>⚡ O(1)"]
262 + Dequeue["Process from queue<br/>⚡ O(1)"]
263 + end
264 + end
265 +
266 + R1 --> Check
267 + R2 --> Check
268 + R3 --> Check
269 + RN --> Check
270 +
271 + Check -->|Yes| Assign
272 + Check -->|No| Queue
273 +
274 + Queue --> Q
275 + Q --> Dequeue
276 +
277 + Assign --> Available
278 + Available --> AW1
279 + Available --> AW2
280 +
281 + Busy --> BW3
282 + Busy --> BW4
283 +
284 + Dequeue -->|Worker freed| Assign
285 +
286 + style Q fill:#bbf,stroke:#333,stroke-width:4px,color:#000
287 + style Available fill:#bfb,stroke:#333,stroke-width:4px,color:#000
288 + style Check fill:#f9f,stroke:#333,stroke-width:4px,color:#000
289 + ```
290 +
291 + ## 7. Process Registry - O(1) Performance
292 +
293 + ```mermaid
294 + graph LR
295 + subgraph "Registry Architecture"
296 + subgraph "Worker Registry"
297 + WR["Elixir Registry<br/>⚡ :unique keys<br/>⚡ O(1) operations"]
298 + WK1["worker_1 → PID1"]
299 + WK2["worker_2 → PID2"]
300 + WKN["worker_N → PIDN"]
301 + end
302 +
303 + subgraph "Process Registry (ETS)"
304 + PR["Process Registry<br/>⚡ :protected table<br/>⚡ read_concurrency"]
305 + PK1["worker_1 → {pid, os_pid, fingerprint}"]
306 + PK2["worker_2 → {pid, os_pid, fingerprint}"]
307 + PKN["worker_N → {pid, os_pid, fingerprint}"]
308 + end
309 +
310 + subgraph "Starter Registry"
311 + SR["Starter Registry<br/>⚡ Supervisor tracking"]
312 + SK1["worker_1 → Starter PID1"]
313 + SK2["worker_2 → Starter PID2"]
314 + SKN["worker_N → Starter PIDN"]
315 + end
316 + end
317 +
318 + subgraph "O(1) Operations"
319 + Op1["via_tuple lookup<br/>⚡ Direct to worker"]
320 + Op2["Reverse lookup<br/>⚡ PID to worker_id"]
321 + Op3["OS PID tracking<br/>⚡ Cleanup guarantee"]
322 + end
323 +
324 + WR --> WK1
325 + WR --> WK2
326 + WR --> WKN
327 +
328 + PR --> PK1
329 + PR --> PK2
330 + PR --> PKN
331 +
332 + SR --> SK1
333 + SR --> SK2
334 + SR --> SKN
335 +
336 + Op1 --> WR
337 + Op2 --> WR
338 + Op3 --> PR
339 +
340 + style WR fill:#bfb,stroke:#333,stroke-width:4px,color:#000
341 + style PR fill:#bbf,stroke:#333,stroke-width:4px,color:#000
342 + style SR fill:#fbf,stroke:#333,stroke-width:4px,color:#000
343 + ```
\ No newline at end of file
  @@ -1,7 +1,7 @@
1 1 # Snakepit 🐍
2 2
3 3 <div align="center">
4 - <img src="img/snakepit-logo.svg" alt="Snakepit Logo" width="200" height="200">
4 + <img src="assets/snakepit-logo.svg" alt="Snakepit Logo" width="200" height="200">
5 5 </div>
6 6
7 7 > A high-performance, generalized process pooler and session manager for external language integrations in Elixir
  @@ -579,7 +579,7 @@ stats = Snakepit.get_stats()
579 579
580 580
581 581 ```
582 - ┌─────────────────────────────────────────────────────────┐
582 + ┌───────────────────────────────────────────────────────┐
583 583 │ Snakepit Application │
584 584 ├───────────────────────────────────────────────────────┤
585 585 │ │
  @@ -0,0 +1,20 @@
1 + <svg width="256px" height="256px" viewBox="0 0 100 100" xmlns="http://www.w3.org/2000/svg" version="1.1">
2 + <polygon points="50,5 93,27.5 93,72.5 50,95 7,72.5 7,27.5" fill="#4e2a8e" />
3 +
4 + <g transform="translate(1, 2)">
5 + <path d="M 62,64
6 + C 72,64 72,52 62,52
7 + C 52,52 52,40 62,40
8 + C 70,40 70,30 62,30
9 + L 50,30"
10 + fill="none"
11 + stroke="#3776ab"
12 + stroke-width="7"
13 + stroke-linecap="round"
14 + stroke-linejoin="round"/>
15 +
16 + <path d="M 49,30 L 42,26 L 42,34 Z" fill="#3776ab" />
17 +
18 + <circle cx="44" cy="30" r="1.2" fill="#ffd43b"/>
19 + </g>
20 + </svg>
  @@ -1,6 +1,6 @@
1 1 {<<"links">>,[{<<"GitHub">>,<<"https://github.com/nshkrdotcom/snakepit">>}]}.
2 2 {<<"name">>,<<"snakepit">>}.
3 - {<<"version">>,<<"0.1.0">>}.
3 + {<<"version">>,<<"0.1.1">>}.
4 4 {<<"description">>,
5 5 <<"High-performance pooler and session manager for external language integrations">>}.
6 6 {<<"elixir">>,<<"~> 1.18">>}.
  @@ -8,8 +8,7 @@
8 8 {<<"licenses">>,[<<"MIT">>]}.
9 9 {<<"files">>,
10 10 [<<"lib">>,<<"lib/snakepit.ex">>,<<"lib/snakepit">>,
11 - <<"lib/snakepit/application.ex.bak">>,<<"lib/snakepit/adapter.ex">>,
12 - <<"lib/snakepit/python">>,<<"lib/snakepit/utils.ex">>,
11 + <<"lib/snakepit/adapter.ex">>,<<"lib/snakepit/utils.ex">>,
13 12 <<"lib/snakepit/bridge">>,<<"lib/snakepit/bridge/protocol.ex">>,
14 13 <<"lib/snakepit/bridge/session.ex">>,
15 14 <<"lib/snakepit/bridge/session_store.ex">>,
  @@ -26,8 +25,9 @@
26 25 <<"lib/snakepit/adapters/generic_python.ex">>,<<"priv">>,<<"priv/python">>,
27 26 <<"priv/python/generic_bridge.py">>,
28 27 <<"priv/python/example_custom_bridge.py">>,<<"priv/javascript">>,
29 - <<"priv/javascript/generic_bridge.js">>,<<".formatter.exs">>,<<"mix.exs">>,
30 - <<"README.md">>,<<"LICENSE">>,<<"CHANGELOG.md">>]}.
28 + <<"priv/javascript/generic_bridge.js">>,<<"assets">>,
29 + <<"assets/snakepit-logo.svg">>,<<".formatter.exs">>,<<"mix.exs">>,
30 + <<"README.md">>,<<"LICENSE">>,<<"CHANGELOG.md">>,<<"DIAGS.md">>]}.
31 31 {<<"requirements">>,
32 32 [[{<<"name">>,<<"jason">>},
33 33 {<<"app">>,<<"jason">>},
  @@ -1,52 +0,0 @@
1 - defmodule Snakepit.Application do
2 - @moduledoc """
3 - Application supervisor for Snakepit pooler.
4 -
5 - Starts the core infrastructure:
6 - - Registry for worker process registration
7 - - ProcessRegistry for Python PID tracking
8 - - SessionStore for session management
9 - - WorkerSupervisor for managing worker processes
10 - - Pool manager for request distribution
11 - """
12 -
13 - use Application
14 - require Logger
15 -
16 - @impl true
17 - def start(_type, _args) do
18 - # Check if pooling is enabled (default: false to prevent auto-start issues)
19 - pooling_enabled = Application.get_env(:snakepit, :pooling_enabled, false)
20 -
21 - children =
22 - if pooling_enabled do
23 - pool_config = Application.get_env(:snakepit, :pool_config, %{})
24 - pool_size = Map.get(pool_config, :pool_size, System.schedulers_online() * 2)
25 -
26 - Logger.info("🚀 Starting Snakepit with pooling enabled (size: #{pool_size})")
27 -
28 - [
29 - # Registry for worker process registration
30 - Snakepit.Pool.Registry,
31 -
32 - # Process registry for Python PID tracking
33 - Snakepit.Pool.ProcessRegistry,
34 -
35 - # Session store for session management
36 - Snakepit.Bridge.SessionStore,
37 -
38 - # Worker supervisor for managing worker processes
39 - Snakepit.Pool.WorkerSupervisor,
40 -
41 - # Main pool manager
42 - {Snakepit.Pool, [size: pool_size]}
43 - ]
44 - else
45 - Logger.info("🔧 Starting Snakepit with pooling disabled")
46 - []
47 - end
48 -
49 - opts = [strategy: :one_for_one, name: Snakepit.Supervisor]
50 - Supervisor.start_link(children, opts)
51 - end
52 - end
Loading more files…