mirror of
https://github.com/troglobit/finit.git
synced 2026-10-01 13:33:09 +07:00
libink: broker state belongs to the connection that has a broker
Parked calls and outbound calls awaiting a reply only ever happen on a connection talking to a broker, but the parked array sat on the server and the pending array on every connection. A server with no broker carried 4 KiB of slots it could never fill, and both were reachable from code paths that have no business in them. Move both behind one struct, allocated on the first park or call and freed with the connection. link_server_t goes from 4400 to 168 bytes; link_connection_t barely moves, its buffers dominate, but an ordinary peer no longer carries reply-tracking it never uses. Tokens are now per bus rather than per server, so link_uid_resolved() takes the connection the answer is about. Every resolver already has it: it is the first argument to both the resolver and the reply callback. Signed-off-by: Joachim Wiberg <troglobit@gmail.com>
This commit is contained in:
+24
-5
@@ -35,6 +35,7 @@ int link_connection_call(link_connection_t *conn, const char *destination,
|
||||
uint8_t hdr[LINK_CALL_HDR_MAX];
|
||||
ssize_t blen = 0;
|
||||
ssize_t hlen;
|
||||
struct link_bus *bus;
|
||||
uint32_t serial;
|
||||
int i;
|
||||
|
||||
@@ -43,8 +44,12 @@ int link_connection_call(link_connection_t *conn, const char *destination,
|
||||
return -1;
|
||||
}
|
||||
|
||||
bus = __bus_get(conn);
|
||||
if (!bus)
|
||||
return -1;
|
||||
|
||||
for (i = 0; i < LINK_PENDING_CAP; i++) {
|
||||
if (!conn->pending[i].used)
|
||||
if (!bus->pending[i].used)
|
||||
break;
|
||||
}
|
||||
if (i == LINK_PENDING_CAP) {
|
||||
@@ -80,14 +85,28 @@ int link_connection_call(link_connection_t *conn, const char *destination,
|
||||
__dbg("calling %s.%s on %s, serial %u", interface ? interface : "-",
|
||||
member, destination ? destination : "peer", serial);
|
||||
|
||||
conn->pending[i].used = 1;
|
||||
conn->pending[i].serial = serial;
|
||||
conn->pending[i].cb = cb;
|
||||
conn->pending[i].userdata = userdata;
|
||||
bus->pending[i].used = 1;
|
||||
bus->pending[i].serial = serial;
|
||||
bus->pending[i].cb = cb;
|
||||
bus->pending[i].userdata = userdata;
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
struct link_bus *__bus_get(link_connection_t *conn)
|
||||
{
|
||||
if (!conn->bus)
|
||||
conn->bus = calloc(1, sizeof(*conn->bus));
|
||||
|
||||
return conn->bus;
|
||||
}
|
||||
|
||||
void __bus_free(link_connection_t *conn)
|
||||
{
|
||||
free(conn->bus);
|
||||
conn->bus = NULL;
|
||||
}
|
||||
|
||||
void link_connection_close(link_connection_t *conn)
|
||||
{
|
||||
size_t i;
|
||||
|
||||
+39
-30
@@ -336,8 +336,14 @@ static void deliver_reply(link_connection_t *conn, const struct link_msg *m)
|
||||
void *userdata;
|
||||
int i;
|
||||
|
||||
if (!conn->bus) {
|
||||
__dbg("unsolicited reply, serial %u", m->reply_serial);
|
||||
return;
|
||||
}
|
||||
|
||||
for (i = 0; i < LINK_PENDING_CAP; i++) {
|
||||
if (conn->pending[i].used && conn->pending[i].serial == m->reply_serial)
|
||||
if (conn->bus->pending[i].used &&
|
||||
conn->bus->pending[i].serial == m->reply_serial)
|
||||
break;
|
||||
}
|
||||
if (i == LINK_PENDING_CAP) {
|
||||
@@ -345,9 +351,9 @@ static void deliver_reply(link_connection_t *conn, const struct link_msg *m)
|
||||
return;
|
||||
}
|
||||
|
||||
cb = conn->pending[i].cb;
|
||||
userdata = conn->pending[i].userdata;
|
||||
conn->pending[i].used = 0;
|
||||
cb = conn->bus->pending[i].cb;
|
||||
userdata = conn->bus->pending[i].userdata;
|
||||
conn->bus->pending[i].used = 0;
|
||||
|
||||
if (!cb)
|
||||
return;
|
||||
@@ -364,26 +370,30 @@ static void deliver_reply(link_connection_t *conn, const struct link_msg *m)
|
||||
static struct link_parked *park(link_connection_t *conn, const uint8_t *frame,
|
||||
size_t len, link_authz_t *tok)
|
||||
{
|
||||
link_server_t *srv = conn->server;
|
||||
struct link_bus *bus;
|
||||
int i;
|
||||
|
||||
if (!frame || !len || len > LINK_PARKED_MSG_MAX)
|
||||
return NULL;
|
||||
|
||||
bus = __bus_get(conn);
|
||||
if (!bus)
|
||||
return NULL;
|
||||
|
||||
for (i = 0; i < LINK_PARKED_CAP; i++) {
|
||||
if (!srv->parked[i].tok)
|
||||
if (!bus->parked[i].tok)
|
||||
break;
|
||||
}
|
||||
if (i == LINK_PARKED_CAP)
|
||||
return NULL;
|
||||
|
||||
srv->parked[i].tok = ++srv->next_tok;
|
||||
srv->parked[i].conn = conn;
|
||||
srv->parked[i].len = len;
|
||||
memcpy(srv->parked[i].buf, frame, len);
|
||||
*tok = srv->parked[i].tok;
|
||||
bus->parked[i].tok = ++bus->next_tok;
|
||||
bus->parked[i].conn = conn;
|
||||
bus->parked[i].len = len;
|
||||
memcpy(bus->parked[i].buf, frame, len);
|
||||
*tok = bus->parked[i].tok;
|
||||
|
||||
return &srv->parked[i];
|
||||
return &bus->parked[i];
|
||||
}
|
||||
|
||||
static void unpark(struct link_parked *p)
|
||||
@@ -394,46 +404,46 @@ static void unpark(struct link_parked *p)
|
||||
|
||||
void __dispatch_forget_conn(link_connection_t *conn)
|
||||
{
|
||||
link_server_t *srv = conn->server;
|
||||
struct link_bus *bus = conn->bus;
|
||||
int i;
|
||||
|
||||
if (!bus)
|
||||
return;
|
||||
|
||||
/* Drop parked calls first. A pending callback below may try to
|
||||
* resolve one, and resuming a dispatch on a connection that is
|
||||
* being torn down is no use to anyone; an invalidated slot makes
|
||||
* that resolve a no-op instead. */
|
||||
if (srv) {
|
||||
for (i = 0; i < LINK_PARKED_CAP; i++) {
|
||||
if (srv->parked[i].conn == conn)
|
||||
unpark(&srv->parked[i]);
|
||||
}
|
||||
}
|
||||
for (i = 0; i < LINK_PARKED_CAP; i++)
|
||||
unpark(&bus->parked[i]);
|
||||
|
||||
for (i = 0; i < LINK_PENDING_CAP; i++) {
|
||||
if (conn->pending[i].used && conn->pending[i].cb)
|
||||
conn->pending[i].cb(conn, NULL, conn->pending[i].userdata);
|
||||
conn->pending[i].used = 0;
|
||||
if (bus->pending[i].used && bus->pending[i].cb)
|
||||
bus->pending[i].cb(conn, NULL, bus->pending[i].userdata);
|
||||
bus->pending[i].used = 0;
|
||||
}
|
||||
|
||||
__bus_free(conn);
|
||||
}
|
||||
|
||||
static int dispatch_call(link_connection_t *conn, const struct link_msg *m,
|
||||
const uint8_t *frame, size_t framelen,
|
||||
const uid_t *known_uid);
|
||||
|
||||
void link_uid_resolved(link_server_t *server, link_authz_t tok, uid_t uid)
|
||||
void link_uid_resolved(link_connection_t *conn, link_authz_t tok, uid_t uid)
|
||||
{
|
||||
uint8_t buf[LINK_PARKED_MSG_MAX];
|
||||
struct link_parked *p = NULL;
|
||||
link_connection_t *conn;
|
||||
struct link_msg msg;
|
||||
size_t len;
|
||||
int i;
|
||||
|
||||
if (!server || !tok)
|
||||
if (!conn || !conn->bus || !tok)
|
||||
return;
|
||||
|
||||
for (i = 0; i < LINK_PARKED_CAP; i++) {
|
||||
if (server->parked[i].tok == tok) {
|
||||
p = &server->parked[i];
|
||||
if (conn->bus->parked[i].tok == tok) {
|
||||
p = &conn->bus->parked[i];
|
||||
break;
|
||||
}
|
||||
}
|
||||
@@ -442,12 +452,11 @@ void link_uid_resolved(link_server_t *server, link_authz_t tok, uid_t uid)
|
||||
|
||||
/* Copy the message out and free the slot before dispatching:
|
||||
* the handler may park a call of its own. */
|
||||
conn = p->conn;
|
||||
len = p->len;
|
||||
len = p->len;
|
||||
memcpy(buf, p->buf, len);
|
||||
unpark(p);
|
||||
|
||||
if (!conn || __msg_parse(buf, len, &msg) <= 0)
|
||||
if (__msg_parse(buf, len, &msg) <= 0)
|
||||
return;
|
||||
|
||||
__dbg("resumed %s, caller uid %d", msg.member ? msg.member : "call", (int)uid);
|
||||
|
||||
+27
-11
@@ -70,6 +70,27 @@ struct link_parked {
|
||||
uint8_t buf[LINK_PARKED_MSG_MAX];
|
||||
};
|
||||
|
||||
/* Calls in flight in either direction: inbound ones held while we ask
|
||||
* who sent them, outbound ones waiting for their reply. Both belong
|
||||
* to a conversation with a broker, so this hangs off the connection
|
||||
* and is allocated on first use. An ordinary peer, which only ever
|
||||
* calls in and is identified by SO_PEERCRED, never gets one.
|
||||
*
|
||||
* Tokens are handed out per bus, which is all link_uid_resolved()
|
||||
* needs: it is told the connection the answer belongs to. */
|
||||
struct link_bus {
|
||||
struct link_parked parked[LINK_PARKED_CAP];
|
||||
link_authz_t next_tok;
|
||||
|
||||
/* Outbound calls we made on this connection, awaiting replies. */
|
||||
struct {
|
||||
int used;
|
||||
uint32_t serial;
|
||||
link_reply_cb_t cb;
|
||||
void *userdata;
|
||||
} pending[LINK_PENDING_CAP];
|
||||
};
|
||||
|
||||
struct link_server {
|
||||
int fd;
|
||||
char path[LINK_PATH_MAX];
|
||||
@@ -83,8 +104,6 @@ struct link_server {
|
||||
/* Set by link_server_set_authorizer(); see link.h. */
|
||||
link_authorizer_t authorizer;
|
||||
void *authz_userdata;
|
||||
struct link_parked parked[LINK_PARKED_CAP];
|
||||
link_authz_t next_tok;
|
||||
};
|
||||
|
||||
/* The reply being assembled inside a method handler.
|
||||
@@ -147,15 +166,8 @@ struct link_connection {
|
||||
|
||||
uint32_t next_serial;
|
||||
|
||||
/* Outbound calls we made on this connection, awaiting replies.
|
||||
* Only a broker connection uses these today, to ask the bus
|
||||
* driver who a sender is. */
|
||||
struct {
|
||||
int used;
|
||||
uint32_t serial;
|
||||
link_reply_cb_t cb;
|
||||
void *userdata;
|
||||
} pending[LINK_PENDING_CAP];
|
||||
/* Allocated on the first park or outbound call, see above. */
|
||||
struct link_bus *bus;
|
||||
|
||||
struct link_server *server; /* back-pointer for dispatch */
|
||||
};
|
||||
@@ -174,6 +186,10 @@ int __auth_process(link_connection_t *conn);
|
||||
void __auth_generate_guid(char out[33]);
|
||||
int __auth_client(int fd, uid_t uid);
|
||||
|
||||
/* connection.c — the per-connection bus state, made on demand. */
|
||||
struct link_bus *__bus_get (link_connection_t *conn);
|
||||
void __bus_free(link_connection_t *conn);
|
||||
|
||||
/* dispatch.c */
|
||||
int __dispatch_message(link_connection_t *conn, const struct link_msg *m, size_t framelen);
|
||||
void __dispatch_forget_conn(link_connection_t *conn);
|
||||
|
||||
+3
-2
@@ -146,8 +146,9 @@ void link_server_set_authorizer(link_server_t *server, link_authorizer_t cb,
|
||||
/* Complete a deferred resolve and resume the parked call. Pass
|
||||
* (uid_t)-1 to say the caller could not be identified, which denies
|
||||
* it. Resolving a handle twice, or one whose connection has since
|
||||
* closed, does nothing. */
|
||||
void link_uid_resolved(link_server_t *server, link_authz_t tok, uid_t uid);
|
||||
* closed, does nothing. The connection is the one the resolver was
|
||||
* asked about; a reply callback is handed it as its first argument. */
|
||||
void link_uid_resolved(link_connection_t *conn, link_authz_t tok, uid_t uid);
|
||||
|
||||
/* Called with the reply to an outbound link_connection_call(). `reply`
|
||||
* is NULL if the connection dropped before one arrived. */
|
||||
|
||||
+1
-3
@@ -1361,8 +1361,6 @@ static void uid_reply_cb(link_connection_t *conn, const link_reply_t *reply, voi
|
||||
uid_t uid = (uid_t)-1;
|
||||
uint32_t val;
|
||||
|
||||
(void)conn;
|
||||
|
||||
if (!reply) {
|
||||
dbg("connection dropped before %s was identified", q->sender);
|
||||
} else if (reply->type == LINK_MSG_METHOD_RETURN &&
|
||||
@@ -1375,7 +1373,7 @@ static void uid_reply_cb(link_connection_t *conn, const link_reply_t *reply, voi
|
||||
reply->error_name ? reply->error_name : "unexpected reply");
|
||||
}
|
||||
|
||||
link_uid_resolved(server, q->tok, uid);
|
||||
link_uid_resolved(conn, q->tok, uid);
|
||||
free(q);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user