1147 lines
41 KiB
C
1147 lines
41 KiB
C
// context.c -- actor threading: one thread + one MPSC queue per context, plus an
|
|
// implicit "host context" (id 0) driven by the host thread via calogPump.
|
|
//
|
|
// Model (design.md sec 6):
|
|
// - Scripts run fire-and-forget on their own context threads (calogContextEval
|
|
// enqueues and returns).
|
|
// - A registered native runs on the HOST thread: a script calling one parks and
|
|
// posts a CALL onto the host queue, which calogPump services on the host thread.
|
|
// So native C code is serialized on one thread and needs no locking. An inline
|
|
// native (calogRegisterInline) instead runs on the calling thread.
|
|
// - A script function value (CalogFnT) runs on its owning engine's thread (sec 10).
|
|
// - A script error is posted to the host queue and delivered to the error handler
|
|
// during calogPump.
|
|
//
|
|
// All per-runtime actor state (the context registry, the host context, the routing
|
|
// hooks, the error sink) lives in CalogT, so independent runtimes coexist in one
|
|
// process. The dispatch reaches its runtime through the object it already holds -- the
|
|
// serving context's ->broker, a callable's ->runtime, or an explicit calog argument --
|
|
// never a process global. currentContext (a thread-local) names the calling thread's
|
|
// context; calogPump sets it to the pumped runtime's host for the drain, so ONE thread
|
|
// can host several runtimes by pumping each in turn. Because context ids number from 1
|
|
// per runtime, inline/reply decisions match the runtime too, not the id alone. A value
|
|
// or callable is still not shared across runtimes (a cross-runtime reply cannot route).
|
|
//
|
|
// Everything a thread waits on flows through its own queue; a caller that is itself
|
|
// a context pumps its own queue (servicing re-entrant calls) rather than sleeping.
|
|
// Ownership across the queue: a CALL deep-copies its args; the target frees them and
|
|
// MOVES the result back. No heap pointer is shared between threads at rest.
|
|
|
|
#define _POSIX_C_SOURCE 200809L
|
|
|
|
#include "calogInternal.h"
|
|
|
|
#include <pthread.h>
|
|
#include <stdio.h>
|
|
#include <stdlib.h>
|
|
#include <string.h>
|
|
|
|
#define CONTEXT_INITIAL_CONTEXTS 8
|
|
#define CONTEXT_INITIAL_ENGINES 4
|
|
|
|
// The host context id (script contexts get 1..N). currentContext == calog->hostContext
|
|
// means "on this runtime's host thread."
|
|
// CALOG_HOST_ID is defined in calogInternal.h (shared with value.c's finalize path).
|
|
|
|
// A script context id packs a 1-based slot index in the low CONTEXT_INDEX_BITS and a
|
|
// generation counter above it. When a slot is recycled its generation is bumped, so a
|
|
// handle minted against the old occupant no longer resolves (design.md sec 9). The id
|
|
// is 64-bit -- a 32-bit slot index + 32-bit generation -- so neither the live-context
|
|
// count nor open/close churn hits a preset ceiling in practice.
|
|
#define CONTEXT_INDEX_BITS 32
|
|
#define CONTEXT_INDEX_MASK 0xFFFFFFFFu
|
|
|
|
typedef enum MessageKindE {
|
|
messageCallE = 0,
|
|
messageReplyE = 1,
|
|
messageShutdownE = 2,
|
|
messageEvalE = 3,
|
|
messageReleaseE = 4,
|
|
messageErrorE = 5
|
|
} MessageKindE;
|
|
|
|
typedef struct ReplyBoxT {
|
|
pthread_mutex_t mutex;
|
|
pthread_cond_t cond;
|
|
bool done;
|
|
int32_t status;
|
|
CalogValueT result;
|
|
} ReplyBoxT;
|
|
|
|
typedef struct MessageT {
|
|
MessageKindE kind;
|
|
CalogNativeFnT fn; // messageCallE: the native to run
|
|
void *userData; // messageCallE: its userData
|
|
CalogValueT *args; // messageCallE: deep-copied arguments
|
|
int32_t argCount;
|
|
char *source; // messageEvalE: script source; messageErrorE: error text
|
|
CalogFnT *callable; // messageReleaseE: the callable to finalize
|
|
uint64_t replyToId; // messageErrorE: the failing context id
|
|
ReplyBoxT *replyBox;
|
|
uint64_t token;
|
|
int32_t status;
|
|
CalogValueT result;
|
|
struct MessageT *next;
|
|
} MessageT;
|
|
|
|
struct CalogContextT {
|
|
uint64_t id;
|
|
CalogT *broker;
|
|
const CalogEngineT *engine;
|
|
void *interp;
|
|
pthread_t thread;
|
|
pthread_mutex_t queueMutex;
|
|
pthread_cond_t queueCond;
|
|
MessageT *head;
|
|
MessageT *tail;
|
|
MessageT *stash;
|
|
uint64_t nextToken;
|
|
bool shuttingDown;
|
|
bool started;
|
|
bool interpDead; // set (under ctxMutex) once the interpreter is torn down
|
|
};
|
|
|
|
// currentContext is the ONLY process-global actor state -- and it is thread-local: it
|
|
// names the context (host or script) the calling thread belongs to. A foreign thread
|
|
// (never a calog thread) leaves it NULL and reaches its runtime through a passed
|
|
// handle instead.
|
|
static _Thread_local CalogContextT *currentContext = NULL;
|
|
|
|
static int32_t actorInvokeCallable(CalogFnT *callable, CalogValueT *args, int32_t argCount, CalogValueT *result);
|
|
static void actorReleaseCallable(CalogFnT *callable);
|
|
static int32_t actorRoute(CalogT *calog, CalogEntryT *entry, CalogValueT *args, int32_t argCount, CalogValueT *result);
|
|
static int32_t contextDispatch(CalogT *calog, uint64_t targetId, CalogNativeFnT fn, void *userData, CalogValueT *args, int32_t argCount, CalogValueT *result);
|
|
static void contextDispatchCall(CalogContextT *context, MessageT *message);
|
|
static void contextDispatchError(CalogT *calog, MessageT *message);
|
|
static void contextDispatchEval(CalogContextT *context, MessageT *message);
|
|
static void contextDispatchRelease(MessageT *message);
|
|
static void contextDrainQueue(CalogContextT *context);
|
|
static int32_t contextEnqueue(CalogT *calog, uint64_t targetId, MessageT *message);
|
|
static int32_t contextPostRelease(CalogT *calog, uint64_t targetId, CalogFnT *callable);
|
|
static void contextReply(CalogT *calog, MessageT *message, int32_t status, CalogValueT *result);
|
|
static int32_t contextSendBlocking(CalogT *calog, uint64_t targetId, MessageT *message, CalogValueT *result);
|
|
static void enqueueRaw(CalogContextT *context, MessageT *message);
|
|
static void hostDispatch(CalogT *calog, MessageT *message);
|
|
static uint64_t idCompose(int64_t index, uint32_t generation);
|
|
static uint32_t idGeneration(uint64_t id);
|
|
static int64_t idIndex(uint64_t id);
|
|
static MessageT *messageDequeue(CalogContextT *context);
|
|
static void messageFree(MessageT *message);
|
|
static bool onOwnerThread(CalogT *runtime, uint64_t owner);
|
|
static void postError(CalogT *calog, uint64_t contextId, const CalogValueT *result);
|
|
static int32_t pumpUntil(CalogContextT *context, uint64_t token, int32_t *outStatus, CalogValueT *result);
|
|
static int32_t registryFreePush(CalogT *calog, int64_t index);
|
|
static CalogContextT *registryResolveLocked(CalogT *calog, uint64_t id);
|
|
static void serveLoop(CalogContextT *context);
|
|
static MessageT *tryDequeue(CalogContextT *context);
|
|
static void *threadMain(void *arg);
|
|
|
|
|
|
int32_t calogActorInit(CalogT *calog) {
|
|
// Build this runtime's host context (id 0, no thread) and claim the calling thread
|
|
// as its host. All actor state hangs off CalogT, so multiple runtimes coexist.
|
|
calog->hostContext = (CalogContextT *)calloc(1, sizeof(*calog->hostContext));
|
|
if (calog->hostContext == NULL) {
|
|
return calogErrOomE;
|
|
}
|
|
calog->hostContext->id = CALOG_HOST_ID;
|
|
calog->hostContext->broker = calog;
|
|
pthread_mutex_init(&calog->hostContext->queueMutex, NULL);
|
|
pthread_cond_init(&calog->hostContext->queueCond, NULL);
|
|
currentContext = calog->hostContext;
|
|
|
|
pthread_mutex_init(&calog->ctxMutex, NULL);
|
|
calog->routeHook = actorRoute;
|
|
calog->invokeHook = actorInvokeCallable;
|
|
calog->releaseHook = actorReleaseCallable;
|
|
return calogOkE;
|
|
}
|
|
|
|
|
|
// True when the calling thread is the owner context's own thread. Ids alone are
|
|
// ambiguous across runtimes (each numbers contexts from 1), so the runtime must match
|
|
// too; a foreign thread (currentContext NULL) is never the owner.
|
|
static bool onOwnerThread(CalogT *runtime, uint64_t owner) {
|
|
return currentContext != NULL && currentContext->broker == runtime && currentContext->id == owner;
|
|
}
|
|
|
|
|
|
static int32_t actorInvokeCallable(CalogFnT *callable, CalogValueT *args, int32_t argCount, CalogValueT *result) {
|
|
CalogT *runtime;
|
|
uint64_t owner;
|
|
|
|
// Inline when already on the owner's thread; otherwise marshal to the owner (in the
|
|
// callable's own runtime). A callable's fn + userData are exactly a native call, so
|
|
// the CALL machinery carries it unchanged.
|
|
runtime = calogFnRuntime(callable);
|
|
owner = calogFnOwner(callable);
|
|
if (onOwnerThread(runtime, owner)) {
|
|
CalogNativeFnT fn;
|
|
fn = calogFnNative(callable);
|
|
return fn(args, argCount, result, calogFnUserData(callable));
|
|
}
|
|
return contextDispatch(runtime, owner, calogFnNative(callable), calogFnUserData(callable), args, argCount, result);
|
|
}
|
|
|
|
|
|
static void actorReleaseCallable(CalogFnT *callable) {
|
|
CalogT *runtime;
|
|
uint64_t owner;
|
|
|
|
// Finalize inline when already on the owner's thread; otherwise marshal the
|
|
// finalize to the owner so the engine release (luaL_unref / sq_release) runs there.
|
|
runtime = calogFnRuntime(callable);
|
|
owner = calogFnOwner(callable);
|
|
if (onOwnerThread(runtime, owner)) {
|
|
calogFnFinalize(callable);
|
|
return;
|
|
}
|
|
if (contextPostRelease(runtime, owner, callable) != calogOkE) {
|
|
// Owner unreachable -- a quiescence violation (design.md sec 9, deferred).
|
|
// Best-effort finalize inline.
|
|
calogFnFinalize(callable);
|
|
}
|
|
}
|
|
|
|
|
|
static int32_t actorRoute(CalogT *calog, CalogEntryT *entry, CalogValueT *args, int32_t argCount, CalogValueT *result) {
|
|
// An inline native, or any call already on this runtime's host thread, runs inline;
|
|
// otherwise the native is marshalled to the host thread (serviced by calogPump).
|
|
if (entry->runInline || currentContext == calog->hostContext) {
|
|
return entry->fn(args, argCount, result, entry->userData);
|
|
}
|
|
return contextDispatch(calog, CALOG_HOST_ID, entry->fn, entry->userData, args, argCount, result);
|
|
}
|
|
|
|
|
|
void calogActorShutdown(CalogT *calog) {
|
|
int64_t index;
|
|
MessageT *message;
|
|
|
|
for (index = 0; index < calog->ctxCount; index++) {
|
|
CalogContextT *context;
|
|
context = calog->ctxSlots[index].context;
|
|
if (context == NULL || !context->started) {
|
|
continue;
|
|
}
|
|
message = (MessageT *)calloc(1, sizeof(*message));
|
|
if (message != NULL) {
|
|
message->kind = messageShutdownE;
|
|
enqueueRaw(context, message);
|
|
}
|
|
}
|
|
// Join EVERY thread before freeing ANY context: a context tearing down may release a
|
|
// callable owned by a sibling (marshalled across contexts), which resolves/posts to that
|
|
// sibling -- so no sibling may be freed while another's thread still runs.
|
|
for (index = 0; index < calog->ctxCount; index++) {
|
|
CalogContextT *context;
|
|
context = calog->ctxSlots[index].context;
|
|
if (context != NULL && context->started) {
|
|
pthread_join(context->thread, NULL);
|
|
}
|
|
}
|
|
// All context threads are stopped. Unregister each under ctxMutex (so any in-flight
|
|
// registryResolveLocked from the host sees NULL rather than a freed slot), drain any
|
|
// orphaned release, then free.
|
|
for (index = 0; index < calog->ctxCount; index++) {
|
|
CalogContextT *context;
|
|
context = calog->ctxSlots[index].context;
|
|
if (context == NULL) {
|
|
continue;
|
|
}
|
|
pthread_mutex_lock(&calog->ctxMutex);
|
|
calog->ctxSlots[index].context = NULL;
|
|
pthread_mutex_unlock(&calog->ctxMutex);
|
|
contextDrainQueue(context);
|
|
pthread_mutex_destroy(&context->queueMutex);
|
|
pthread_cond_destroy(&context->queueCond);
|
|
free(context);
|
|
}
|
|
free(calog->ctxSlots);
|
|
free(calog->ctxFree);
|
|
calog->ctxSlots = NULL;
|
|
calog->ctxFree = NULL;
|
|
calog->ctxCount = 0;
|
|
calog->ctxCap = 0;
|
|
calog->ctxFreeCount = 0;
|
|
calog->ctxFreeCap = 0;
|
|
pthread_mutex_destroy(&calog->ctxMutex);
|
|
|
|
// Drain and free the host context. Only clear the calling thread's currentContext
|
|
// if it named THIS runtime's host -- the thread may still host other runtimes.
|
|
if (calog->hostContext != NULL) {
|
|
bool wasHost;
|
|
wasHost = (currentContext == calog->hostContext);
|
|
while ((message = tryDequeue(calog->hostContext)) != NULL) {
|
|
messageFree(message);
|
|
}
|
|
pthread_mutex_destroy(&calog->hostContext->queueMutex);
|
|
pthread_cond_destroy(&calog->hostContext->queueCond);
|
|
free(calog->hostContext);
|
|
calog->hostContext = NULL;
|
|
if (wasHost) {
|
|
currentContext = NULL;
|
|
}
|
|
}
|
|
calog->routeHook = NULL;
|
|
calog->invokeHook = NULL;
|
|
calog->releaseHook = NULL;
|
|
calog->errorHandler = NULL;
|
|
calog->errorUserData = NULL;
|
|
}
|
|
|
|
|
|
CalogT *calogContextBroker(const CalogContextT *context) {
|
|
return context->broker;
|
|
}
|
|
|
|
|
|
// The public runtime lifecycle: calogCreate composes the registry with the actor
|
|
// layer (which builds the host context and installs the routing hooks); calogDestroy
|
|
// tears the actor layer down (joining context threads) before freeing the registry.
|
|
CalogT *calogCreate(void) {
|
|
CalogT *calog;
|
|
|
|
calog = calogBrokerCreate();
|
|
if (calog == NULL) {
|
|
return NULL;
|
|
}
|
|
if (calogActorInit(calog) != calogOkE) {
|
|
calogBrokerDestroy(calog);
|
|
return NULL;
|
|
}
|
|
return calog;
|
|
}
|
|
|
|
|
|
void calogDestroy(CalogT *calog) {
|
|
if (calog == NULL) {
|
|
return;
|
|
}
|
|
calogActorShutdown(calog);
|
|
calogBrokerDestroy(calog);
|
|
}
|
|
|
|
|
|
// Create + start a context in one step (design: no owned-native registration happens
|
|
// between the two, so there is nothing to do in between). Returns NULL on failure.
|
|
CalogContextT *calogContextOpen(CalogT *broker, const CalogEngineT *engine) {
|
|
CalogContextT *context;
|
|
int64_t index;
|
|
uint32_t generation;
|
|
|
|
context = (CalogContextT *)calloc(1, sizeof(*context));
|
|
if (context == NULL) {
|
|
return NULL;
|
|
}
|
|
context->broker = broker;
|
|
context->engine = engine;
|
|
pthread_mutex_init(&context->queueMutex, NULL);
|
|
pthread_cond_init(&context->queueCond, NULL);
|
|
|
|
pthread_mutex_lock(&broker->ctxMutex);
|
|
if (broker->ctxFreeCount > 0) {
|
|
broker->ctxFreeCount--;
|
|
index = broker->ctxFree[broker->ctxFreeCount];
|
|
generation = broker->ctxSlots[index].generation + 1u;
|
|
} else {
|
|
if (broker->ctxCount + 1 > (int64_t)CONTEXT_INDEX_MASK) {
|
|
pthread_mutex_unlock(&broker->ctxMutex);
|
|
pthread_mutex_destroy(&context->queueMutex);
|
|
pthread_cond_destroy(&context->queueCond);
|
|
free(context);
|
|
return NULL;
|
|
}
|
|
if (broker->ctxCount == broker->ctxCap) {
|
|
int64_t newCap;
|
|
CalogRegistrySlotT *grown;
|
|
newCap = (broker->ctxCap == 0) ? CONTEXT_INITIAL_CONTEXTS : broker->ctxCap * CALOG_GROWTH_FACTOR;
|
|
grown = (CalogRegistrySlotT *)realloc(broker->ctxSlots, (size_t)newCap * sizeof(CalogRegistrySlotT));
|
|
if (grown == NULL) {
|
|
pthread_mutex_unlock(&broker->ctxMutex);
|
|
pthread_mutex_destroy(&context->queueMutex);
|
|
pthread_cond_destroy(&context->queueCond);
|
|
free(context);
|
|
return NULL;
|
|
}
|
|
broker->ctxSlots = grown;
|
|
broker->ctxCap = newCap;
|
|
}
|
|
index = broker->ctxCount;
|
|
generation = 0;
|
|
broker->ctxCount++;
|
|
}
|
|
context->id = idCompose(index, generation);
|
|
broker->ctxSlots[index].context = context;
|
|
broker->ctxSlots[index].generation = generation;
|
|
pthread_mutex_unlock(&broker->ctxMutex);
|
|
|
|
if (pthread_create(&context->thread, NULL, threadMain, context) != 0) {
|
|
pthread_mutex_lock(&broker->ctxMutex);
|
|
broker->ctxSlots[index].context = NULL;
|
|
registryFreePush(broker, index);
|
|
pthread_mutex_unlock(&broker->ctxMutex);
|
|
pthread_mutex_destroy(&context->queueMutex);
|
|
pthread_cond_destroy(&context->queueCond);
|
|
free(context);
|
|
return NULL;
|
|
}
|
|
context->started = true;
|
|
return context;
|
|
}
|
|
|
|
|
|
// Search the registered engines for a file "<baseFileName>.<ext>". The first engine
|
|
// (in registration order), then its first extension, that names a readable file wins:
|
|
// its contents are loaded fire-and-forget into a fresh context on that engine. Reads
|
|
// the file on the calling thread. Returns NULL if nothing matched or the load failed.
|
|
CalogContextT *calogContextLoad(CalogT *calog, const char *baseFileName) {
|
|
int64_t engineIndex;
|
|
|
|
for (engineIndex = 0; engineIndex < calog->engineCount; engineIndex++) {
|
|
const CalogEngineT *engine;
|
|
int32_t extIndex;
|
|
|
|
engine = calog->engines[engineIndex];
|
|
if (engine->extensions == NULL) {
|
|
continue;
|
|
}
|
|
for (extIndex = 0; engine->extensions[extIndex] != NULL; extIndex++) {
|
|
CalogContextT *context;
|
|
const char *ext;
|
|
char *path;
|
|
char *source;
|
|
FILE *file;
|
|
size_t pathSize;
|
|
long fileSize;
|
|
size_t readCount;
|
|
|
|
ext = engine->extensions[extIndex];
|
|
pathSize = strlen(baseFileName) + strlen(ext) + 2; // '.' + '\0'
|
|
path = (char *)malloc(pathSize);
|
|
if (path == NULL) {
|
|
return NULL;
|
|
}
|
|
snprintf(path, pathSize, "%s.%s", baseFileName, ext);
|
|
file = fopen(path, "rb");
|
|
free(path);
|
|
if (file == NULL) {
|
|
continue; // not this extension -- try the next
|
|
}
|
|
// First existing file wins. Read it whole, open a context, load it.
|
|
if (fseek(file, 0, SEEK_END) != 0 || (fileSize = ftell(file)) < 0) {
|
|
fclose(file);
|
|
return NULL;
|
|
}
|
|
rewind(file);
|
|
source = (char *)malloc((size_t)fileSize + 1);
|
|
if (source == NULL) {
|
|
fclose(file);
|
|
return NULL;
|
|
}
|
|
readCount = fread(source, 1, (size_t)fileSize, file);
|
|
fclose(file);
|
|
source[readCount] = '\0';
|
|
context = calogContextOpen(calog, engine);
|
|
if (context == NULL) {
|
|
free(source);
|
|
return NULL;
|
|
}
|
|
if (calogContextEval(context, source) != calogOkE) {
|
|
free(source);
|
|
calogContextClose(context);
|
|
return NULL;
|
|
}
|
|
free(source);
|
|
return context;
|
|
}
|
|
}
|
|
return NULL;
|
|
}
|
|
|
|
|
|
CalogT *calogCurrent(void) {
|
|
return currentContext != NULL ? currentContext->broker : NULL;
|
|
}
|
|
|
|
|
|
uint64_t calogCurrentId(void) {
|
|
return currentContext != NULL ? currentContext->id : CALOG_HOST_ID;
|
|
}
|
|
|
|
|
|
// Tear down a single context (quiescence assumed). The thread is stopped and joined,
|
|
// then the slot is unlinked under the registry lock so no foreign enqueue can reach
|
|
// the freed queue mutex; the freed index returns to the freelist.
|
|
void calogContextClose(CalogContextT *context) {
|
|
CalogT *broker;
|
|
int64_t index;
|
|
|
|
if (context == NULL) {
|
|
return;
|
|
}
|
|
broker = context->broker;
|
|
if (context->started) {
|
|
MessageT *message;
|
|
message = (MessageT *)calloc(1, sizeof(*message));
|
|
if (message != NULL) {
|
|
message->kind = messageShutdownE;
|
|
enqueueRaw(context, message);
|
|
}
|
|
pthread_join(context->thread, NULL);
|
|
}
|
|
pthread_mutex_lock(&broker->ctxMutex);
|
|
index = idIndex(context->id);
|
|
if (index >= 0 && index < broker->ctxCount && broker->ctxSlots[index].context == context) {
|
|
broker->ctxSlots[index].context = NULL;
|
|
registryFreePush(broker, index);
|
|
}
|
|
pthread_mutex_unlock(&broker->ctxMutex);
|
|
contextDrainQueue(context);
|
|
pthread_mutex_destroy(&context->queueMutex);
|
|
pthread_cond_destroy(&context->queueCond);
|
|
free(context);
|
|
}
|
|
|
|
|
|
static int32_t contextDispatch(CalogT *calog, uint64_t targetId, CalogNativeFnT fn, void *userData, CalogValueT *args, int32_t argCount, CalogValueT *result) {
|
|
MessageT *call;
|
|
int32_t status;
|
|
int32_t index;
|
|
|
|
calogValueNil(result);
|
|
call = (MessageT *)calloc(1, sizeof(*call));
|
|
if (call == NULL) {
|
|
return calogFail(result, calogErrOomE, "out of memory creating call message");
|
|
}
|
|
call->kind = messageCallE;
|
|
call->fn = fn;
|
|
call->userData = userData;
|
|
call->argCount = argCount;
|
|
if (argCount > 0) {
|
|
call->args = (CalogValueT *)calloc((size_t)argCount, sizeof(CalogValueT));
|
|
if (call->args == NULL) {
|
|
free(call);
|
|
return calogFail(result, calogErrOomE, "out of memory copying call arguments");
|
|
}
|
|
for (index = 0; index < argCount; index++) {
|
|
status = calogValueCopy(&call->args[index], &args[index]);
|
|
if (status != calogOkE) {
|
|
int32_t cleanup;
|
|
for (cleanup = 0; cleanup < index; cleanup++) {
|
|
calogValueFree(&call->args[cleanup]);
|
|
}
|
|
free(call->args);
|
|
free(call);
|
|
return calogFail(result, status, "failed to copy call arguments");
|
|
}
|
|
}
|
|
}
|
|
return contextSendBlocking(calog, targetId, call, result);
|
|
}
|
|
|
|
|
|
static void contextDispatchCall(CalogContextT *context, MessageT *message) {
|
|
CalogValueT result;
|
|
int32_t status;
|
|
int32_t index;
|
|
|
|
calogValueNil(&result);
|
|
status = message->fn(message->args, message->argCount, &result, message->userData);
|
|
|
|
if (message->args != NULL) {
|
|
for (index = 0; index < message->argCount; index++) {
|
|
calogValueFree(&message->args[index]);
|
|
}
|
|
free(message->args);
|
|
message->args = NULL;
|
|
}
|
|
contextReply(context->broker, message, status, &result);
|
|
}
|
|
|
|
|
|
static void contextDispatchError(CalogT *calog, MessageT *message) {
|
|
// Runs on the host thread (calogPump / a nested host pump): deliver a fire-and-
|
|
// forget script error to the handler, or log it if none is set.
|
|
if (calog->errorHandler != NULL) {
|
|
calog->errorHandler(message->replyToId, message->source != NULL ? message->source : "", calog->errorUserData);
|
|
} else {
|
|
fprintf(stderr, "calog: context %llu script error: %s\n", (unsigned long long)message->replyToId, message->source != NULL ? message->source : "");
|
|
}
|
|
free(message->source);
|
|
free(message);
|
|
}
|
|
|
|
|
|
static void contextDispatchEval(CalogContextT *context, MessageT *message) {
|
|
CalogValueT result;
|
|
int32_t status;
|
|
|
|
// Fire-and-forget: run the script, free the message, and on failure post the
|
|
// error to the host thread. No reply.
|
|
calogValueNil(&result);
|
|
if (context->engine == NULL || context->engine->runSource == NULL) {
|
|
status = calogFail(&result, calogErrUnsupportedE, "context has no runnable engine");
|
|
} else if (context->interp == NULL) {
|
|
// createInterpreter failed (threadMain ignores its status); reject cleanly.
|
|
status = calogFail(&result, calogErrUnsupportedE, "engine interpreter unavailable");
|
|
} else {
|
|
status = context->engine->runSource(context->interp, message->source, &result);
|
|
}
|
|
free(message->source);
|
|
free(message);
|
|
if (status != calogOkE) {
|
|
postError(context->broker, context->id, &result);
|
|
}
|
|
calogValueFree(&result);
|
|
}
|
|
|
|
|
|
static void contextDispatchRelease(MessageT *message) {
|
|
// Runs on the callable's owner thread: perform the engine's closure release via
|
|
// calogFnFinalize, then free the message shell.
|
|
calogFnFinalize(message->callable);
|
|
free(message);
|
|
}
|
|
|
|
|
|
static void contextDrainQueue(CalogContextT *context) {
|
|
MessageT *message;
|
|
|
|
// Called at teardown once the owning thread has stopped and the context is unregistered,
|
|
// so no new messages can arrive. A release can be orphaned here if a last drop of one of
|
|
// this context's callables raced its shutdown (enqueued after serveLoop stopped serving);
|
|
// finalize those so the callable is not leaked, and free every other pending message.
|
|
while ((message = tryDequeue(context)) != NULL) {
|
|
if (message->kind == messageReleaseE) {
|
|
calogFnFinalize(message->callable);
|
|
free(message);
|
|
} else {
|
|
messageFree(message);
|
|
}
|
|
}
|
|
}
|
|
|
|
|
|
static int32_t contextEnqueue(CalogT *calog, uint64_t targetId, MessageT *message) {
|
|
CalogContextT *target;
|
|
int64_t index;
|
|
bool inRange;
|
|
|
|
pthread_mutex_lock(&calog->ctxMutex);
|
|
target = registryResolveLocked(calog, targetId);
|
|
if (target != NULL) {
|
|
enqueueRaw(target, message);
|
|
pthread_mutex_unlock(&calog->ctxMutex);
|
|
return calogOkE;
|
|
}
|
|
// Resolve failed: an in-range index means the slot existed and has since been
|
|
// freed or recycled (a dead context); out of range means the id named nothing.
|
|
index = idIndex(targetId);
|
|
inRange = index >= 0 && index < calog->ctxCount;
|
|
pthread_mutex_unlock(&calog->ctxMutex);
|
|
if (inRange) {
|
|
return calogErrDeadE;
|
|
}
|
|
return calogErrNotFoundE;
|
|
}
|
|
|
|
|
|
// Fire-and-forget: post a callable finalize to its owner thread. On enqueue failure
|
|
// the message is freed and actorReleaseCallable finalizes the dead callable itself.
|
|
static int32_t contextPostRelease(CalogT *calog, uint64_t targetId, CalogFnT *callable) {
|
|
MessageT *message;
|
|
int32_t status;
|
|
|
|
message = (MessageT *)calloc(1, sizeof(*message));
|
|
if (message == NULL) {
|
|
return calogErrOomE;
|
|
}
|
|
message->kind = messageReleaseE;
|
|
message->callable = callable;
|
|
status = contextEnqueue(calog, targetId, message);
|
|
if (status != calogOkE) {
|
|
free(message);
|
|
}
|
|
return status;
|
|
}
|
|
|
|
|
|
int32_t calogContextEval(CalogContextT *context, const char *source) {
|
|
MessageT *eval;
|
|
char *sourceCopy;
|
|
int32_t status;
|
|
|
|
sourceCopy = strdup(source);
|
|
if (sourceCopy == NULL) {
|
|
return calogErrOomE;
|
|
}
|
|
eval = (MessageT *)calloc(1, sizeof(*eval));
|
|
if (eval == NULL) {
|
|
free(sourceCopy);
|
|
return calogErrOomE;
|
|
}
|
|
eval->kind = messageEvalE;
|
|
eval->source = sourceCopy;
|
|
// Fire-and-forget: enqueue onto the context's thread and return. The script runs
|
|
// asynchronously; errors surface via the error handler.
|
|
status = contextEnqueue(context->broker, context->id, eval);
|
|
if (status != calogOkE) {
|
|
messageFree(eval);
|
|
}
|
|
return status;
|
|
}
|
|
|
|
|
|
uint64_t calogContextId(const CalogContextT *context) {
|
|
return context->id;
|
|
}
|
|
|
|
|
|
void *calogContextInterp(CalogContextT *context) {
|
|
return context->interp;
|
|
}
|
|
|
|
|
|
bool calogContextRegistered(CalogT *runtime, uint64_t ctxId) {
|
|
bool registered;
|
|
|
|
if (runtime == NULL) {
|
|
return false;
|
|
}
|
|
pthread_mutex_lock(&runtime->ctxMutex);
|
|
registered = (registryResolveLocked(runtime, ctxId) != NULL);
|
|
pthread_mutex_unlock(&runtime->ctxMutex);
|
|
return registered;
|
|
}
|
|
|
|
|
|
void calogPump(CalogT *calog) {
|
|
CalogContextT *previous;
|
|
MessageT *message;
|
|
|
|
if (calog->hostContext == NULL) {
|
|
return;
|
|
}
|
|
// Present the calling thread as THIS runtime's host for the drain, so a native --
|
|
// and anything it calls -- resolves to this runtime. Restoring afterward lets one
|
|
// thread pump several runtimes in turn (each drain acts as the right host).
|
|
previous = currentContext;
|
|
currentContext = calog->hostContext;
|
|
// Non-blocking: run every pending host-thread message and return.
|
|
while ((message = tryDequeue(calog->hostContext)) != NULL) {
|
|
hostDispatch(calog, message);
|
|
}
|
|
currentContext = previous;
|
|
}
|
|
|
|
|
|
// Append an engine to the runtime's calogContextLoad search list (setup-time, host
|
|
// thread). On OOM the engine is simply not registered -- load just won't find it.
|
|
void calogRegisterEngine(CalogT *calog, const CalogEngineT *engine) {
|
|
if (calog->engineCount == calog->engineCap) {
|
|
int64_t newCap;
|
|
const CalogEngineT **grown;
|
|
newCap = (calog->engineCap == 0) ? CONTEXT_INITIAL_ENGINES : calog->engineCap * CALOG_GROWTH_FACTOR;
|
|
grown = (const CalogEngineT **)realloc(calog->engines, (size_t)newCap * sizeof(*grown));
|
|
if (grown == NULL) {
|
|
return;
|
|
}
|
|
calog->engines = grown;
|
|
calog->engineCap = newCap;
|
|
}
|
|
calog->engines[calog->engineCount] = engine;
|
|
calog->engineCount++;
|
|
}
|
|
|
|
|
|
void calogSetErrorHandler(CalogT *calog, CalogErrorFnT fn, void *userData) {
|
|
calog->errorHandler = fn;
|
|
calog->errorUserData = userData;
|
|
}
|
|
|
|
|
|
// The reply tail shared by CALL dispatch: hand (status, result) back to whoever is
|
|
// blocked on this message -- an external caller's reply box, or a context caller via
|
|
// a REPLY enqueued onto its queue (in the serving context's runtime). Consumes message.
|
|
static void contextReply(CalogT *calog, MessageT *message, int32_t status, CalogValueT *result) {
|
|
if (message->replyBox != NULL) {
|
|
ReplyBoxT *box;
|
|
box = message->replyBox;
|
|
pthread_mutex_lock(&box->mutex);
|
|
box->status = status;
|
|
calogValueMove(&box->result, result);
|
|
box->done = true;
|
|
pthread_cond_signal(&box->cond);
|
|
pthread_mutex_unlock(&box->mutex);
|
|
free(message);
|
|
} else {
|
|
// Reuse the request message as the REPLY: it already carries the caller's
|
|
// token and id, so the wakeup cannot be lost to an allocation failure here.
|
|
uint64_t replyToId;
|
|
replyToId = message->replyToId;
|
|
message->kind = messageReplyE;
|
|
message->status = status;
|
|
calogValueMove(&message->result, result);
|
|
if (contextEnqueue(calog, replyToId, message) != calogOkE) {
|
|
messageFree(message);
|
|
}
|
|
}
|
|
}
|
|
|
|
|
|
// Enqueue an already-built CALL to targetId and wait for its reply. A context caller
|
|
// pumps its own queue (servicing re-entrant work) instead of sleeping; a foreign
|
|
// caller blocks on a private reply box. The message is consumed either way.
|
|
static int32_t contextSendBlocking(CalogT *calog, uint64_t targetId, MessageT *message, CalogValueT *result) {
|
|
CalogContextT *caller;
|
|
int32_t status;
|
|
|
|
// The token/pump path needs the caller's queue in THIS runtime (its reply routes
|
|
// by the caller's id, resolved in calog). A thread that is foreign to calog -- a
|
|
// non-calog thread, or one currently hosting a different runtime -- takes the reply
|
|
// box instead, so a cross-runtime call cannot misroute its reply.
|
|
caller = (currentContext != NULL && currentContext->broker == calog) ? currentContext : NULL;
|
|
if (caller != NULL) {
|
|
uint64_t token;
|
|
int32_t replyStatus;
|
|
token = ++caller->nextToken;
|
|
message->replyToId = caller->id;
|
|
message->token = token;
|
|
status = contextEnqueue(calog, targetId, message);
|
|
if (status != calogOkE) {
|
|
messageFree(message);
|
|
return calogFail(result, status, "target context unavailable");
|
|
}
|
|
status = pumpUntil(caller, token, &replyStatus, result);
|
|
if (status != calogOkE) {
|
|
return calogFail(result, status, "context is shutting down");
|
|
}
|
|
return replyStatus;
|
|
}
|
|
ReplyBoxT box;
|
|
pthread_mutex_init(&box.mutex, NULL);
|
|
pthread_cond_init(&box.cond, NULL);
|
|
box.done = false;
|
|
box.status = calogOkE;
|
|
calogValueNil(&box.result);
|
|
message->replyBox = &box;
|
|
status = contextEnqueue(calog, targetId, message);
|
|
if (status != calogOkE) {
|
|
messageFree(message);
|
|
pthread_mutex_destroy(&box.mutex);
|
|
pthread_cond_destroy(&box.cond);
|
|
return calogFail(result, status, "target context unavailable");
|
|
}
|
|
pthread_mutex_lock(&box.mutex);
|
|
while (!box.done) {
|
|
pthread_cond_wait(&box.cond, &box.mutex);
|
|
}
|
|
pthread_mutex_unlock(&box.mutex);
|
|
calogValueMove(result, &box.result);
|
|
status = box.status;
|
|
pthread_mutex_destroy(&box.mutex);
|
|
pthread_cond_destroy(&box.cond);
|
|
return status;
|
|
}
|
|
|
|
|
|
static void enqueueRaw(CalogContextT *context, MessageT *message) {
|
|
pthread_mutex_lock(&context->queueMutex);
|
|
message->next = NULL;
|
|
if (context->tail != NULL) {
|
|
context->tail->next = message;
|
|
} else {
|
|
context->head = message;
|
|
}
|
|
context->tail = message;
|
|
pthread_cond_signal(&context->queueCond);
|
|
pthread_mutex_unlock(&context->queueMutex);
|
|
}
|
|
|
|
|
|
static void hostDispatch(CalogT *calog, MessageT *message) {
|
|
switch (message->kind) {
|
|
case messageCallE:
|
|
contextDispatchCall(calog->hostContext, message);
|
|
break;
|
|
case messageReleaseE:
|
|
contextDispatchRelease(message);
|
|
break;
|
|
case messageErrorE:
|
|
contextDispatchError(calog, message);
|
|
break;
|
|
default:
|
|
messageFree(message);
|
|
break;
|
|
}
|
|
}
|
|
|
|
|
|
static uint64_t idCompose(int64_t index, uint32_t generation) {
|
|
return ((uint64_t)generation << CONTEXT_INDEX_BITS) | (uint64_t)(index + 1);
|
|
}
|
|
|
|
|
|
static uint32_t idGeneration(uint64_t id) {
|
|
return (uint32_t)(id >> CONTEXT_INDEX_BITS);
|
|
}
|
|
|
|
|
|
static int64_t idIndex(uint64_t id) {
|
|
uint64_t low;
|
|
|
|
low = id & CONTEXT_INDEX_MASK;
|
|
if (low == 0) {
|
|
return -1;
|
|
}
|
|
return (int64_t)low - 1;
|
|
}
|
|
|
|
|
|
static MessageT *messageDequeue(CalogContextT *context) {
|
|
MessageT *message;
|
|
|
|
pthread_mutex_lock(&context->queueMutex);
|
|
while (context->head == NULL && !context->shuttingDown) {
|
|
pthread_cond_wait(&context->queueCond, &context->queueMutex);
|
|
}
|
|
if (context->head == NULL) {
|
|
pthread_mutex_unlock(&context->queueMutex);
|
|
return NULL;
|
|
}
|
|
message = context->head;
|
|
context->head = message->next;
|
|
if (context->head == NULL) {
|
|
context->tail = NULL;
|
|
}
|
|
pthread_mutex_unlock(&context->queueMutex);
|
|
message->next = NULL;
|
|
return message;
|
|
}
|
|
|
|
|
|
static void messageFree(MessageT *message) {
|
|
int32_t index;
|
|
|
|
if (message == NULL) {
|
|
return;
|
|
}
|
|
if (message->kind == messageCallE && message->args != NULL) {
|
|
for (index = 0; index < message->argCount; index++) {
|
|
calogValueFree(&message->args[index]);
|
|
}
|
|
free(message->args);
|
|
}
|
|
if (message->kind == messageEvalE || message->kind == messageErrorE) {
|
|
free(message->source);
|
|
}
|
|
if (message->kind == messageReplyE) {
|
|
calogValueFree(&message->result);
|
|
}
|
|
free(message);
|
|
}
|
|
|
|
|
|
// Post a fire-and-forget script error to the host thread's error handler.
|
|
static void postError(CalogT *calog, uint64_t contextId, const CalogValueT *result) {
|
|
MessageT *message;
|
|
const char *text;
|
|
|
|
text = (result->type == calogStringE) ? result->as.s.bytes : "script error";
|
|
message = (MessageT *)calloc(1, sizeof(*message));
|
|
if (message == NULL) {
|
|
return;
|
|
}
|
|
message->kind = messageErrorE;
|
|
message->replyToId = contextId;
|
|
message->source = strdup(text);
|
|
if (message->source == NULL) {
|
|
free(message);
|
|
return;
|
|
}
|
|
if (contextEnqueue(calog, CALOG_HOST_ID, message) != calogOkE) {
|
|
free(message->source);
|
|
free(message);
|
|
}
|
|
}
|
|
|
|
|
|
static int32_t pumpUntil(CalogContextT *context, uint64_t token, int32_t *outStatus, CalogValueT *result) {
|
|
for (;;) {
|
|
MessageT *message;
|
|
MessageT *prev;
|
|
MessageT *node;
|
|
|
|
// Re-check the stash: a nested pump may have stashed this token's reply.
|
|
prev = NULL;
|
|
node = context->stash;
|
|
while (node != NULL) {
|
|
if (node->token == token) {
|
|
if (prev != NULL) {
|
|
prev->next = node->next;
|
|
} else {
|
|
context->stash = node->next;
|
|
}
|
|
*outStatus = node->status;
|
|
calogValueMove(result, &node->result);
|
|
free(node);
|
|
return calogOkE;
|
|
}
|
|
prev = node;
|
|
node = node->next;
|
|
}
|
|
|
|
message = messageDequeue(context);
|
|
if (message == NULL) {
|
|
calogValueNil(result);
|
|
return calogErrNotFoundE;
|
|
}
|
|
if (message->kind == messageReplyE) {
|
|
if (message->token == token) {
|
|
*outStatus = message->status;
|
|
calogValueMove(result, &message->result);
|
|
free(message);
|
|
return calogOkE;
|
|
}
|
|
message->next = context->stash;
|
|
context->stash = message;
|
|
continue;
|
|
}
|
|
if (message->kind == messageCallE) {
|
|
contextDispatchCall(context, message);
|
|
continue;
|
|
}
|
|
if (message->kind == messageEvalE) {
|
|
contextDispatchEval(context, message);
|
|
continue;
|
|
}
|
|
if (message->kind == messageReleaseE) {
|
|
contextDispatchRelease(message);
|
|
continue;
|
|
}
|
|
if (message->kind == messageErrorE) {
|
|
contextDispatchError(context->broker, message);
|
|
continue;
|
|
}
|
|
context->shuttingDown = true;
|
|
free(message);
|
|
}
|
|
}
|
|
|
|
|
|
// Push a freed slot index onto the recycle freelist. Called under calog->ctxMutex. On
|
|
// OOM the index is simply not recycled (a leaked slot, never a correctness fault).
|
|
static int32_t registryFreePush(CalogT *calog, int64_t index) {
|
|
if (calog->ctxFreeCount == calog->ctxFreeCap) {
|
|
int64_t newCap;
|
|
int64_t *grown;
|
|
newCap = (calog->ctxFreeCap == 0) ? CONTEXT_INITIAL_CONTEXTS : calog->ctxFreeCap * CALOG_GROWTH_FACTOR;
|
|
grown = (int64_t *)realloc(calog->ctxFree, (size_t)newCap * sizeof(int64_t));
|
|
if (grown == NULL) {
|
|
return calogErrOomE;
|
|
}
|
|
calog->ctxFree = grown;
|
|
calog->ctxFreeCap = newCap;
|
|
}
|
|
calog->ctxFree[calog->ctxFreeCount] = index;
|
|
calog->ctxFreeCount++;
|
|
return calogOkE;
|
|
}
|
|
|
|
|
|
static CalogContextT *registryResolveLocked(CalogT *calog, uint64_t id) {
|
|
int64_t index;
|
|
|
|
if (id == CALOG_HOST_ID) {
|
|
return calog->hostContext;
|
|
}
|
|
index = idIndex(id);
|
|
if (index < 0 || index >= calog->ctxCount) {
|
|
return NULL;
|
|
}
|
|
if (calog->ctxSlots[index].context == NULL) {
|
|
return NULL;
|
|
}
|
|
if (calog->ctxSlots[index].generation != idGeneration(id)) {
|
|
return NULL;
|
|
}
|
|
// Interpreter already torn down: treat as gone so late releases finalize without touching
|
|
// it (the slot itself stays populated until calogContextClose/calogActorShutdown frees it).
|
|
if (calog->ctxSlots[index].context->interpDead) {
|
|
return NULL;
|
|
}
|
|
return calog->ctxSlots[index].context;
|
|
}
|
|
|
|
|
|
static void serveLoop(CalogContextT *context) {
|
|
for (;;) {
|
|
MessageT *message;
|
|
message = messageDequeue(context);
|
|
if (message == NULL) {
|
|
return;
|
|
}
|
|
if (message->kind == messageShutdownE) {
|
|
context->shuttingDown = true;
|
|
free(message);
|
|
continue;
|
|
}
|
|
if (message->kind == messageCallE) {
|
|
contextDispatchCall(context, message);
|
|
continue;
|
|
}
|
|
if (message->kind == messageEvalE) {
|
|
contextDispatchEval(context, message);
|
|
continue;
|
|
}
|
|
if (message->kind == messageReleaseE) {
|
|
contextDispatchRelease(message);
|
|
continue;
|
|
}
|
|
messageFree(message);
|
|
}
|
|
}
|
|
|
|
|
|
static MessageT *tryDequeue(CalogContextT *context) {
|
|
MessageT *message;
|
|
|
|
pthread_mutex_lock(&context->queueMutex);
|
|
message = context->head;
|
|
if (message != NULL) {
|
|
context->head = message->next;
|
|
if (context->head == NULL) {
|
|
context->tail = NULL;
|
|
}
|
|
}
|
|
pthread_mutex_unlock(&context->queueMutex);
|
|
if (message != NULL) {
|
|
message->next = NULL;
|
|
}
|
|
return message;
|
|
}
|
|
|
|
|
|
static void *threadMain(void *arg) {
|
|
CalogContextT *context;
|
|
|
|
context = (CalogContextT *)arg;
|
|
currentContext = context;
|
|
if (context->engine != NULL && context->engine->createInterpreter != NULL) {
|
|
context->engine->createInterpreter(context, &context->interp);
|
|
}
|
|
serveLoop(context);
|
|
// Mark the interpreter dead BEFORE destroying it, so registryResolveLocked stops handing
|
|
// this context out: a callable this context owns that is released from now on (e.g. one
|
|
// published via the export library and dropped after the context unloads) finalizes by
|
|
// freeing its struct directly instead of running the engine release on a dead interpreter.
|
|
pthread_mutex_lock(&context->broker->ctxMutex);
|
|
context->interpDead = true;
|
|
pthread_mutex_unlock(&context->broker->ctxMutex);
|
|
if (context->engine != NULL && context->engine->destroyInterpreter != NULL) {
|
|
context->engine->destroyInterpreter(context->interp);
|
|
}
|
|
currentContext = NULL;
|
|
return NULL;
|
|
}
|