mirror of
https://github.com/openglow-org/forgefirm.git
synced 2026-09-27 16:51:12 -07:00
forgectrl: streams preempt - the newest camera request wins the mux
A stream request for the other camera kicks current stream clients via a generation counter (their streams end cleanly; viewers freeze on the last frame) and switches. The index page retry consults /cam/status first so a preempted view does not steal the camera back.
This commit is contained in:
@@ -91,6 +91,8 @@ static struct {
|
||||
cam_id_t home_cam; /* camera the engine serves for streaming -
|
||||
* what arbitration must compare against */
|
||||
int clients;
|
||||
uint64_t kick_gen; /* bumped to preempt all current stream
|
||||
* clients (their streams end cleanly) */
|
||||
struct timespec last_activity;
|
||||
|
||||
/* published stream frame (half-res JPEG) */
|
||||
@@ -727,9 +729,15 @@ static int ensure_engine(cam_id_t cam, char *err, size_t errlen)
|
||||
|
||||
if (running && cur != cam) {
|
||||
if (clients > 0) {
|
||||
/* Grace: a client that just disconnected releases its pin
|
||||
* only when the MHD send fails on the next frame - absorb
|
||||
* that (page navigations) instead of failing instantly. */
|
||||
/* Last request wins: preempt the current stream clients
|
||||
* (single-operator machine - the newest ask is the
|
||||
* operator). Kicked clients wake, end their streams
|
||||
* cleanly (viewers freeze on their last frame), and
|
||||
* release their pins; wait for that to drain. */
|
||||
pthread_mutex_lock(&eng.lock);
|
||||
eng.kick_gen++;
|
||||
pthread_cond_broadcast(&eng.frame_cv);
|
||||
pthread_mutex_unlock(&eng.lock);
|
||||
struct timespec t0, t;
|
||||
now_ts(&t0);
|
||||
do {
|
||||
@@ -741,8 +749,8 @@ static int ensure_engine(cam_id_t cam, char *err, size_t errlen)
|
||||
} while (clients > 0 && ts_diff(&t, &t0) < SWITCH_GRACE_S);
|
||||
if (clients > 0) {
|
||||
snprintf(err, errlen,
|
||||
"camera busy: %d client(s) streaming %s",
|
||||
clients, camdefs[cur].name);
|
||||
"camera switch timed out: %d client(s) still "
|
||||
"attached to %s", clients, camdefs[cur].name);
|
||||
return -1;
|
||||
}
|
||||
}
|
||||
@@ -878,6 +886,7 @@ int cam_snapshot(cam_id_t cam, int full, int quality,
|
||||
|
||||
struct cam_client {
|
||||
uint64_t last_seq;
|
||||
uint64_t gen; /* kick generation at open; a bump ends the stream */
|
||||
uint8_t *buf;
|
||||
size_t cap;
|
||||
};
|
||||
@@ -897,6 +906,7 @@ cam_client_t *cam_client_open(cam_id_t cam, char *err, size_t errlen)
|
||||
}
|
||||
pthread_mutex_lock(&eng.lock);
|
||||
eng.clients++;
|
||||
c->gen = eng.kick_gen;
|
||||
now_ts(&eng.last_activity);
|
||||
pthread_mutex_unlock(&eng.lock);
|
||||
pthread_mutex_unlock(&eng.ctl);
|
||||
@@ -907,7 +917,7 @@ long cam_client_next(cam_client_t *c, const uint8_t **jpeg)
|
||||
{
|
||||
pthread_mutex_lock(&eng.lock);
|
||||
int timeouts = 0;
|
||||
while (eng.running && eng.seq <= c->last_seq) {
|
||||
while (eng.running && c->gen == eng.kick_gen && eng.seq <= c->last_seq) {
|
||||
struct timespec deadline;
|
||||
clock_gettime(CLOCK_REALTIME, &deadline);
|
||||
deadline.tv_sec += CLIENT_WAIT_S;
|
||||
@@ -915,7 +925,7 @@ long cam_client_next(cam_client_t *c, const uint8_t **jpeg)
|
||||
== ETIMEDOUT && ++timeouts >= 2)
|
||||
break;
|
||||
}
|
||||
if (!eng.running || eng.seq <= c->last_seq) {
|
||||
if (!eng.running || c->gen != eng.kick_gen || eng.seq <= c->last_seq) {
|
||||
pthread_mutex_unlock(&eng.lock);
|
||||
return -1;
|
||||
}
|
||||
|
||||
@@ -36,9 +36,11 @@ void cam_engine_shutdown(void);
|
||||
int cam_snapshot(cam_id_t cam, int full, int quality,
|
||||
uint8_t **jpeg, size_t *len, char *err, size_t errlen);
|
||||
|
||||
/* Stream client: open pins the engine to a camera (starting or switching it
|
||||
* if needed), next blocks for a frame newer than the last one returned and
|
||||
* copies it into a client-owned buffer, close releases the pin. */
|
||||
/* Stream client: open makes the engine serve `cam` (starting it, or
|
||||
* preempting the current clients and switching - last request wins; the
|
||||
* preempted clients' next() returns -1 so their streams end cleanly).
|
||||
* next blocks for a frame newer than the last one returned and copies it
|
||||
* into a client-owned buffer, close releases the pin. */
|
||||
typedef struct cam_client cam_client_t;
|
||||
|
||||
cam_client_t *cam_client_open(cam_id_t cam, char *err, size_t errlen);
|
||||
|
||||
@@ -12,11 +12,11 @@
|
||||
* GET /cam/snapshot?cam=&res=full|half&q= single JPEG (default full res)
|
||||
* GET /cam/status JSON engine status
|
||||
*
|
||||
* The two cameras share the hardware mux, so streaming clients pin the
|
||||
* selection: a STREAM request for the other camera waits a short grace
|
||||
* for the pin to drain, then returns 409. A SNAPSHOT of the other camera
|
||||
* never fails busy - the engine borrows the mux for one frame and the
|
||||
* stream freezes for a few seconds instead.
|
||||
* The two cameras share the hardware mux; the newest request wins it. A
|
||||
* STREAM request for the other camera preempts the current stream
|
||||
* clients (their streams end cleanly - viewers freeze on the last frame)
|
||||
* and switches. A SNAPSHOT of the other camera does not switch: the
|
||||
* engine borrows the mux for one frame and the stream freezes briefly.
|
||||
* Environment: FORGECTRL_PORT (8080), FORGECTRL_STREAM_Q (75),
|
||||
* FORGECTRL_LAMP (132).
|
||||
*
|
||||
@@ -246,11 +246,17 @@ static const char index_html[] =
|
||||
"msg=document.getElementById('msg');"
|
||||
"function setCam(c){cam=c;retries=0;msg.textContent='';"
|
||||
"v.src='/cam/stream?cam='+c+'&t='+Date.now();}"
|
||||
"v.onerror=function(){if(retries++<5){"
|
||||
"msg.textContent='stream retrying...';"
|
||||
"setTimeout(function(){v.src='/cam/stream?cam='+cam+'&t='+Date.now();},"
|
||||
"700);}else{msg.textContent="
|
||||
"'stream unavailable (camera busy from another viewer?)';}};"
|
||||
"function reload(){v.src='/cam/stream?cam='+cam+'&t='+Date.now();}"
|
||||
/* Retry only when the engine still serves our camera - if another
|
||||
* viewer preempted it, retrying would steal it right back. */
|
||||
"v.onerror=function(){fetch('/cam/status').then(function(r){"
|
||||
"return r.json();}).then(function(s){"
|
||||
"if(s.cam===cam&&retries++<5){msg.textContent='stream retrying...';"
|
||||
"setTimeout(reload,700);}else if(s.cam!==cam){msg.textContent="
|
||||
"'stream taken by another viewer ('+s.cam+') - press a button to resume';}"
|
||||
"else{msg.textContent='stream error - press a stream button to retry';}"
|
||||
"}).catch(function(){msg.textContent="
|
||||
"'service unreachable';});};"
|
||||
"function peek(){"
|
||||
"msg.textContent='head peek (stream pauses a few seconds)...';"
|
||||
"p.style.display='inline';p.onload=function(){msg.textContent='';};"
|
||||
|
||||
Reference in New Issue
Block a user