mirror of
https://github.com/sbrow/zomq.git
synced 2026-08-26 10:53:32 -04:00
190 lines
4.4 KiB
Odin
190 lines
4.4 KiB
Odin
package zmq
|
|
|
|
import "core:c"
|
|
import "core:os"
|
|
import "core:strings"
|
|
|
|
foreign import lib "system:zmq"
|
|
|
|
Context :: distinct rawptr
|
|
Socket :: distinct rawptr
|
|
|
|
@(default_calling_convention = "c", link_prefix = "zmq_")
|
|
foreign lib {
|
|
// TODO: Try to access errno directly
|
|
errno :: proc() -> c.int ---
|
|
|
|
ctx_new :: proc() -> Context ---
|
|
@(link_name = "zmq_ctx_term")
|
|
_ctx_term :: proc(ctx: Context) -> c.int ---
|
|
@(link_name = "zmq_socket")
|
|
_socket :: proc(ctx: Context, type: Socket_Type) -> Socket ---
|
|
@(link_name = "zmq_close")
|
|
_close :: proc(s: Socket) -> c.int ---
|
|
@(link_name = "zmq_bind")
|
|
_bind :: proc(s: Socket, addr: cstring) -> c.int ---
|
|
@(link_name = "zmq_connect")
|
|
_connect :: proc(s: Socket, addr: cstring) -> c.int ---
|
|
@(link_name = "zmq_send")
|
|
_send :: proc(s: Socket, buf: rawptr, len: c.size_t, flags: Flags) -> c.int ---
|
|
|
|
@(link_name = "zmq_send_const")
|
|
_send_const :: proc(s: Socket, buf: rawptr, len: c.size_t, flags: Flags) -> c.int ---
|
|
|
|
@(link_name = "zmq_recv")
|
|
_recv :: proc(s: Socket, buf: rawptr, len: c.size_t, flags: Flags) -> c.int ---
|
|
@(link_name = "zmq_poll")
|
|
_poll :: proc(items: rawptr, nitems: c.int, timeout: c.long) -> c.int ---
|
|
@(link_name = "zmq_strerror")
|
|
_strerror :: proc(errnum: c.int) -> cstring ---
|
|
@(link_name = "zmq_version")
|
|
_version :: proc(major, minor, patch: ^c.int) ---
|
|
}
|
|
|
|
Socket_Type :: enum c.int {
|
|
PAIR = 0,
|
|
PUB = 1,
|
|
SUB = 2,
|
|
REQ = 3,
|
|
REP = 4,
|
|
DEALER = 5,
|
|
ROUTER = 6,
|
|
PULL = 7,
|
|
PUSH = 8,
|
|
XPUB = 9,
|
|
XSUB = 10,
|
|
STREAM = 11,
|
|
}
|
|
|
|
Flag :: enum c.int {
|
|
DONTWAIT,
|
|
SNDMORE,
|
|
}
|
|
|
|
Flags :: bit_set[Flag;c.int]
|
|
|
|
NoBlock :: Flags{.DONTWAIT}
|
|
|
|
// zmq_poll event flags (zmq.h ZMQ_POLLIN etc.). bit_set uses enum values as
|
|
// bit indices, so these are positions (0..3), not the ZMQ_POLL* masks (1,2,4,8):
|
|
// {.POLLIN} -> bit 0 -> mask 1 == ZMQ_POLLIN, {.POLLOUT} -> bit 1 -> mask 2, etc.
|
|
Poll_Event :: enum c.short {
|
|
POLLIN = 0,
|
|
POLLOUT = 1,
|
|
POLLERR = 2,
|
|
POLLPRI = 3,
|
|
}
|
|
Poll_Events :: bit_set[Poll_Event;c.short]
|
|
|
|
// Mirrors C zmq_pollitem_t (16 B: ptr / int / short / short).
|
|
// socket non-nil => poll the ØMQ socket; otherwise poll fd.
|
|
Poll_Item :: struct {
|
|
socket: Socket,
|
|
fd: c.int,
|
|
events: Poll_Events,
|
|
revents: Poll_Events,
|
|
}
|
|
|
|
Error :: union #shared_nil {
|
|
os.Platform_Error,
|
|
ZMQ_Error,
|
|
}
|
|
|
|
strerror :: proc(e: Error) -> cstring {
|
|
switch v in e {
|
|
case os.Platform_Error:
|
|
return _strerror(c.int(v))
|
|
case ZMQ_Error:
|
|
return _strerror(denormalize_from_zmq(c.int(v)))
|
|
case:
|
|
return _strerror(0)
|
|
}
|
|
}
|
|
|
|
version :: proc() -> (major, minor, patch: int) {
|
|
ma, mi, pa: c.int = 0, 0, 0
|
|
_version(&ma, &mi, &pa)
|
|
return int(ma), int(mi), int(pa)
|
|
}
|
|
|
|
ctx_term :: proc(ctx: Context) -> Error {
|
|
return _ctx_term(ctx) == 0 ? nil : last_error()
|
|
}
|
|
|
|
socket :: proc(ctx: Context, type: Socket_Type) -> (s: Socket, err: Error) {
|
|
s = _socket(ctx, type)
|
|
|
|
if s == nil {
|
|
err = last_error()
|
|
}
|
|
|
|
return s, err
|
|
}
|
|
|
|
close :: proc(s: Socket) -> Error {
|
|
return _close(s) == 0 ? nil : last_error()
|
|
}
|
|
|
|
bind :: proc(s: Socket, addr: cstring) -> Error {
|
|
return _bind(s, addr) == 0 ? nil : last_error()
|
|
}
|
|
|
|
connect :: proc(s: Socket, addr: cstring) -> Error {
|
|
return _connect(s, addr) == 0 ? nil : last_error()
|
|
}
|
|
|
|
send :: proc {
|
|
send_bytes,
|
|
send_const_string,
|
|
send_string,
|
|
}
|
|
|
|
send_bytes :: proc(s: Socket, buf: []u8, flags: Flags = {}) -> (int, Error) {
|
|
n := _send(s, rawptr(&buf[0]), c.size_t(len(buf)), flags)
|
|
if n < 0 {
|
|
return -1, last_error()
|
|
}
|
|
return int(n), nil
|
|
}
|
|
|
|
send_string :: proc(s: Socket, buf: string, flags: Flags = {}) -> (int, Error) {
|
|
return send_bytes(s, transmute([]u8)buf, flags)
|
|
}
|
|
|
|
send_const :: proc {
|
|
send_const_string,
|
|
send_const_bytes,
|
|
}
|
|
|
|
send_const_string :: proc(s: Socket, $str: string, flags: Flags = {}) -> (int, Error) {
|
|
return send_const_bytes(s, transmute([]byte)str, flags)
|
|
}
|
|
|
|
// buf must have @(rodata) or global lifetime.
|
|
send_const_bytes :: proc(s: Socket, buf: []byte, flags: Flags = {}) -> (int, Error) {
|
|
n := _send_const(s, rawptr(&buf[0]), c.size_t(len(buf)), flags)
|
|
if n < 0 {
|
|
return -1, last_error()
|
|
}
|
|
|
|
return int(n), nil
|
|
}
|
|
|
|
recv :: proc(s: Socket, buf: []u8, flags: Flags = Flags{}) -> (int, Error) {
|
|
n := _recv(s, rawptr(&buf[0]), c.size_t(len(buf)), flags)
|
|
if n < 0 {
|
|
return -1, last_error()
|
|
}
|
|
return int(n), nil
|
|
}
|
|
|
|
// timeout_ms -1 = block forever, 0 = non-blocking.
|
|
poll :: proc(items: []Poll_Item, timeout_ms: c.long) -> (int, Error) {
|
|
n := _poll(rawptr(&items[0]), c.int(len(items)), timeout_ms)
|
|
if n < 0 {
|
|
return -1, last_error()
|
|
}
|
|
return int(n), nil
|
|
}
|
|
|