runtime: port missing __channel_* primitives into el_seed.c
El SDK CI - dev / build-and-test (pull_request) Failing after 3m44s
El SDK CI - dev / build-and-test (pull_request) Failing after 3m44s
runtime/channel.el has always called __channel_new/__channel_send/ __channel_recv/__channel_try_recv/__channel_close, but these were only ever implemented in the pre-restructure lang/el-compiler/runtime/el_runtime.c. When the canonical runtime was consolidated onto the release copy (lang/runtime/el_runtime.c) and el_seed.c became the sole C dependency, the channel implementation was never carried forward — __mutex_new made the move, __channel_* did not. Any El program using Go-style channels currently fails to link on dev. Ported the working buffered-MPMC-channel implementation (mutex+condvar+ circular buffer, bounded and unbounded modes) from the old el_runtime.c verbatim, adapted only to el_seed.c's arena API (seed_arena_track in place of el_arena_track). Declared in el_seed.h alongside the existing mutex primitives.
This commit is contained in:
@@ -907,6 +907,219 @@ void __mutex_unlock(el_val_t m) {
|
||||
pthread_mutex_unlock(&_el_mutexes[slot]);
|
||||
}
|
||||
|
||||
/* ── Channels ─────────────────────────────────────────────────────────────── *
|
||||
* Buffered MPMC channel backed by a mutex + condvar + circular buffer.
|
||||
* Ported from the pre-restructure el_runtime.c (b2aac4b) — runtime/channel.el
|
||||
* has always called these five primitives, but they were never carried
|
||||
* forward into el_seed.c when el_runtime.c was consolidated onto the
|
||||
* canonical release copy. Native channels were silently unlinkable on dev
|
||||
* until this port.
|
||||
*
|
||||
* __channel_new(capacity) -> Int (handle)
|
||||
* __channel_send(ch, msg) — blocks if full (capacity > 0) or never (unbounded)
|
||||
* __channel_recv(ch) -> String — blocks until a message is available
|
||||
* __channel_try_recv(ch) -> String — non-blocking, returns "" if empty
|
||||
* __channel_close(ch) — signal no more sends; recv drains remaining
|
||||
*
|
||||
* Bounded channels (cap > 0): circular buffer, sender blocks when full.
|
||||
* Unbounded channels (cap == 0): dynamic array, sender never blocks.
|
||||
*/
|
||||
#define EL_CHANNEL_MAX 64
|
||||
#define EL_CHANNEL_BUF 1024
|
||||
|
||||
typedef struct {
|
||||
char** buf;
|
||||
int cap; /* 0 = unbounded (grows dynamically) */
|
||||
int head, tail, count;
|
||||
int dyn_cap; /* allocated slots for unbounded mode */
|
||||
int closed;
|
||||
pthread_mutex_t mu;
|
||||
pthread_cond_t not_empty;
|
||||
pthread_cond_t not_full;
|
||||
} ElChannel;
|
||||
|
||||
static ElChannel _channels[EL_CHANNEL_MAX];
|
||||
static int _channel_count = 0;
|
||||
static pthread_mutex_t _channel_alloc_mu = PTHREAD_MUTEX_INITIALIZER;
|
||||
|
||||
el_val_t __channel_new(el_val_t capacity_v) {
|
||||
int cap = (int)(int64_t)capacity_v;
|
||||
if (cap < 0) cap = 0;
|
||||
|
||||
pthread_mutex_lock(&_channel_alloc_mu);
|
||||
if (_channel_count >= EL_CHANNEL_MAX) {
|
||||
pthread_mutex_unlock(&_channel_alloc_mu);
|
||||
fprintf(stderr, "[__channel_new] channel table full\n");
|
||||
return EL_INT(-1);
|
||||
}
|
||||
int slot = _channel_count++;
|
||||
pthread_mutex_unlock(&_channel_alloc_mu);
|
||||
|
||||
ElChannel* ch = &_channels[slot];
|
||||
memset(ch, 0, sizeof(*ch));
|
||||
ch->cap = cap;
|
||||
ch->closed = 0;
|
||||
ch->head = 0;
|
||||
ch->tail = 0;
|
||||
ch->count = 0;
|
||||
|
||||
if (cap > 0) {
|
||||
/* Bounded: fixed circular buffer. */
|
||||
ch->buf = (char**)malloc((size_t)cap * sizeof(char*));
|
||||
ch->dyn_cap = cap;
|
||||
} else {
|
||||
/* Unbounded: start with EL_CHANNEL_BUF slots, grow as needed. */
|
||||
ch->buf = (char**)malloc(EL_CHANNEL_BUF * sizeof(char*));
|
||||
ch->dyn_cap = EL_CHANNEL_BUF;
|
||||
}
|
||||
if (!ch->buf) {
|
||||
fprintf(stderr, "[__channel_new] out of memory\n");
|
||||
return EL_INT(-1);
|
||||
}
|
||||
|
||||
pthread_mutex_init(&ch->mu, NULL);
|
||||
pthread_cond_init(&ch->not_empty, NULL);
|
||||
pthread_cond_init(&ch->not_full, NULL);
|
||||
|
||||
return EL_INT(slot);
|
||||
}
|
||||
|
||||
el_val_t __channel_send(el_val_t ch_v, el_val_t msg_v) {
|
||||
int slot = (int)(int64_t)ch_v;
|
||||
if (slot < 0 || slot >= EL_CHANNEL_MAX) return EL_STR("");
|
||||
ElChannel* ch = &_channels[slot];
|
||||
|
||||
const char* msg = EL_CSTR(msg_v);
|
||||
if (!msg) msg = "";
|
||||
char* copy = strdup(msg); /* channel owns the string */
|
||||
|
||||
pthread_mutex_lock(&ch->mu);
|
||||
|
||||
if (ch->closed) {
|
||||
/* Send on closed channel is a no-op (drop the message). */
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
free(copy);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
if (ch->cap > 0) {
|
||||
/* Bounded: block while full. */
|
||||
while (ch->count >= ch->cap && !ch->closed) {
|
||||
pthread_cond_wait(&ch->not_full, &ch->mu);
|
||||
}
|
||||
if (ch->closed) {
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
free(copy);
|
||||
return EL_STR("");
|
||||
}
|
||||
ch->buf[ch->tail] = copy;
|
||||
ch->tail = (ch->tail + 1) % ch->cap;
|
||||
ch->count++;
|
||||
} else {
|
||||
/* Unbounded: grow the buffer if needed. */
|
||||
if (ch->count >= ch->dyn_cap) {
|
||||
int new_cap = ch->dyn_cap * 2;
|
||||
char** grown = (char**)realloc(ch->buf, (size_t)new_cap * sizeof(char*));
|
||||
if (!grown) {
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
free(copy);
|
||||
fprintf(stderr, "[__channel_send] out of memory growing channel\n");
|
||||
return EL_STR("");
|
||||
}
|
||||
/* The circular buffer may have wrapped. Linearise it first.
|
||||
* In unbounded mode head is always 0 (we append at tail, drain
|
||||
* from head), so a simple memmove isn't needed — but if the
|
||||
* buffer did wrap (tail < head after growth), we need to fix up.
|
||||
* Simplest safe path: if tail wrapped, move the head..old_cap
|
||||
* segment to new_cap..new_cap+(old_cap-head). */
|
||||
if (ch->tail < ch->head) {
|
||||
/* Wrapped: [head..old_cap) is the front, [0..tail) is the back. */
|
||||
int front = ch->dyn_cap - ch->head;
|
||||
memmove(grown + ch->dyn_cap, grown + ch->head, (size_t)front * sizeof(char*));
|
||||
ch->head = ch->dyn_cap;
|
||||
}
|
||||
ch->buf = grown;
|
||||
ch->dyn_cap = new_cap;
|
||||
}
|
||||
ch->buf[ch->tail] = copy;
|
||||
ch->tail = (ch->tail + 1) % ch->dyn_cap;
|
||||
ch->count++;
|
||||
}
|
||||
|
||||
pthread_cond_signal(&ch->not_empty);
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
el_val_t __channel_recv(el_val_t ch_v) {
|
||||
int slot = (int)(int64_t)ch_v;
|
||||
if (slot < 0 || slot >= EL_CHANNEL_MAX) return EL_STR("");
|
||||
ElChannel* ch = &_channels[slot];
|
||||
|
||||
pthread_mutex_lock(&ch->mu);
|
||||
|
||||
/* Block until there is a message or the channel is closed and drained. */
|
||||
while (ch->count == 0 && !ch->closed) {
|
||||
pthread_cond_wait(&ch->not_empty, &ch->mu);
|
||||
}
|
||||
|
||||
if (ch->count == 0) {
|
||||
/* Closed and empty — signal EOF. */
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
int buf_cap = (ch->cap > 0) ? ch->cap : ch->dyn_cap;
|
||||
char* msg = ch->buf[ch->head];
|
||||
ch->head = (ch->head + 1) % buf_cap;
|
||||
ch->count--;
|
||||
|
||||
pthread_cond_signal(&ch->not_full);
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
|
||||
/* Hand the string to the arena so it is freed after the request. */
|
||||
seed_arena_track(msg);
|
||||
return EL_STR(msg);
|
||||
}
|
||||
|
||||
el_val_t __channel_try_recv(el_val_t ch_v) {
|
||||
int slot = (int)(int64_t)ch_v;
|
||||
if (slot < 0 || slot >= EL_CHANNEL_MAX) return EL_STR("");
|
||||
ElChannel* ch = &_channels[slot];
|
||||
|
||||
pthread_mutex_lock(&ch->mu);
|
||||
|
||||
if (ch->count == 0) {
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
int buf_cap = (ch->cap > 0) ? ch->cap : ch->dyn_cap;
|
||||
char* msg = ch->buf[ch->head];
|
||||
ch->head = (ch->head + 1) % buf_cap;
|
||||
ch->count--;
|
||||
|
||||
pthread_cond_signal(&ch->not_full);
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
|
||||
seed_arena_track(msg);
|
||||
return EL_STR(msg);
|
||||
}
|
||||
|
||||
el_val_t __channel_close(el_val_t ch_v) {
|
||||
int slot = (int)(int64_t)ch_v;
|
||||
if (slot < 0 || slot >= EL_CHANNEL_MAX) return EL_STR("");
|
||||
ElChannel* ch = &_channels[slot];
|
||||
|
||||
pthread_mutex_lock(&ch->mu);
|
||||
ch->closed = 1;
|
||||
/* Wake all blocked recvers and senders so they can observe the close. */
|
||||
pthread_cond_broadcast(&ch->not_empty);
|
||||
pthread_cond_broadcast(&ch->not_full);
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
/* ── Subprocess ──────────────────────────────────────────────────────────── */
|
||||
|
||||
el_val_t __exec(el_val_t cmd) {
|
||||
|
||||
@@ -139,6 +139,13 @@ el_val_t __mutex_new(void);
|
||||
void __mutex_lock(el_val_t m);
|
||||
void __mutex_unlock(el_val_t m);
|
||||
|
||||
/* Buffered MPMC channel (runtime/channel.el). capacity=0 means unbounded. */
|
||||
el_val_t __channel_new(el_val_t capacity);
|
||||
el_val_t __channel_send(el_val_t ch, el_val_t msg); /* blocks if bounded+full */
|
||||
el_val_t __channel_recv(el_val_t ch); /* blocks until available */
|
||||
el_val_t __channel_try_recv(el_val_t ch); /* non-blocking, "" if empty */
|
||||
el_val_t __channel_close(el_val_t ch);
|
||||
|
||||
/* ── Subprocess ──────────────────────────────────────────────────────────── */
|
||||
|
||||
el_val_t __exec(el_val_t cmd); /* popen, capture all stdout, return String */
|
||||
|
||||
Reference in New Issue
Block a user