mirror of
https://github.com/sbrow/zomq.git
synced 2026-08-26 10:53:32 -04:00
refactor(zmq): Fixed style errors.
This commit is contained in:
@@ -2,23 +2,41 @@ package main
|
|||||||
|
|
||||||
import "core:bufio"
|
import "core:bufio"
|
||||||
import "core:fmt"
|
import "core:fmt"
|
||||||
|
import "core:log"
|
||||||
import "core:os"
|
import "core:os"
|
||||||
import zmq "zmq"
|
import "zmq"
|
||||||
|
|
||||||
main :: proc() {
|
main :: proc() {
|
||||||
|
context.logger = log.create_console_logger(.Debug, {})
|
||||||
|
|
||||||
major, minor, patch := zmq.version()
|
major, minor, patch := zmq.version()
|
||||||
fmt.printfln("zmq echo (libzmq %d.%d.%d) over inproc; type lines, Ctrl-D to quit", major, minor, patch)
|
log.infof(
|
||||||
|
"zmq echo (libzmq %d.%d.%d) over inproc; type lines, Ctrl-D to quit",
|
||||||
|
major,
|
||||||
|
minor,
|
||||||
|
patch,
|
||||||
|
)
|
||||||
|
|
||||||
ctx := zmq.new_context()
|
ctx := zmq.ctx_new()
|
||||||
defer zmq.term(ctx)
|
defer zmq.ctx_term(ctx)
|
||||||
|
|
||||||
rep := must_socket(ctx, .REP, "rep socket")
|
rep, serr := zmq.socket(ctx, .REP)
|
||||||
|
if serr != nil {
|
||||||
|
log.panic(zmq.strerror(serr))
|
||||||
|
}
|
||||||
defer zmq.close(rep)
|
defer zmq.close(rep)
|
||||||
must(zmq.bind(rep, "inproc://echo"), "bind")
|
if err := zmq.bind(rep, "inproc://echo"); err != nil {
|
||||||
|
log.panic(zmq.strerror(err))
|
||||||
|
}
|
||||||
|
|
||||||
req := must_socket(ctx, .REQ, "req socket")
|
req, cerr := zmq.socket(ctx, .REQ)
|
||||||
|
if cerr != nil {
|
||||||
|
log.panic(zmq.strerror(cerr))
|
||||||
|
}
|
||||||
defer zmq.close(req)
|
defer zmq.close(req)
|
||||||
must(zmq.connect(req, "inproc://echo"), "connect")
|
if err := zmq.connect(req, "inproc://echo"); err != nil {
|
||||||
|
log.panic(zmq.strerror(err))
|
||||||
|
}
|
||||||
|
|
||||||
echo(req, rep)
|
echo(req, rep)
|
||||||
}
|
}
|
||||||
@@ -28,47 +46,31 @@ echo :: proc(req, rep: zmq.Socket) {
|
|||||||
bufio.scanner_init(&scanner, os.to_reader(os.stdin))
|
bufio.scanner_init(&scanner, os.to_reader(os.stdin))
|
||||||
defer bufio.scanner_destroy(&scanner)
|
defer bufio.scanner_destroy(&scanner)
|
||||||
|
|
||||||
buf: [1024]u8
|
buf: [255]u8
|
||||||
for bufio.scan(&scanner) {
|
for bufio.scan(&scanner) {
|
||||||
line := bufio.scanner_text(&scanner)
|
line := bufio.scanner_text(&scanner)
|
||||||
if len(line) == 0 {
|
if len(line) == 0 {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if _, err := zmq.send(req, transmute([]u8)line); !zmq.is_ok(err) {
|
if _, err := zmq.send(req, transmute([]u8)line); err != nil {
|
||||||
fmt.eprintfln("send: %s", zmq.error_string(err))
|
fmt.eprintfln("send: %s", zmq.strerror(err))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
n, rerr := zmq.recv(rep, buf[:])
|
n, rerr := zmq.recv(rep, buf[:])
|
||||||
if !zmq.is_ok(rerr) {
|
if rerr != nil {
|
||||||
fmt.eprintfln("rep recv: %s", zmq.error_string(rerr))
|
fmt.eprintfln("rep recv: %s", zmq.strerror(rerr))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if _, werr := zmq.send(rep, buf[:n]); !zmq.is_ok(werr) {
|
if _, werr := zmq.send(rep, buf[:n]); werr != nil {
|
||||||
fmt.eprintfln("rep send: %s", zmq.error_string(werr))
|
fmt.eprintfln("rep send: %s", zmq.strerror(werr))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
m, qerr := zmq.recv(req, buf[:])
|
m, qerr := zmq.recv(req, buf[:])
|
||||||
if !zmq.is_ok(qerr) {
|
if qerr != nil {
|
||||||
fmt.eprintfln("req recv: %s", zmq.error_string(qerr))
|
fmt.eprintfln("req recv: %s", zmq.strerror(qerr))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
fmt.println(string(buf[:m]))
|
fmt.println(string(buf[:m]))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
must :: proc(err: zmq.Error, at: string) {
|
|
||||||
if !zmq.is_ok(err) {
|
|
||||||
fatal(at, err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
must_socket :: proc(ctx: zmq.Context, type: zmq.Socket_Type, at: string) -> zmq.Socket {
|
|
||||||
s, err := zmq.socket(ctx, type)
|
|
||||||
must(err, at)
|
|
||||||
return s
|
|
||||||
}
|
|
||||||
|
|
||||||
fatal :: proc(at: string, err: zmq.Error) {
|
|
||||||
fmt.eprintfln("%s: %s", at, zmq.error_string(err))
|
|
||||||
os.exit(1)
|
|
||||||
}
|
|
||||||
|
|||||||
+43
-29
@@ -6,21 +6,32 @@ foreign import lib "system:zmq"
|
|||||||
|
|
||||||
@(default_calling_convention = "c")
|
@(default_calling_convention = "c")
|
||||||
foreign lib {
|
foreign lib {
|
||||||
@(link_name = "zmq_ctx_new") _ctx_new :: proc() -> rawptr ---
|
@(link_name = "zmq_ctx_new")
|
||||||
@(link_name = "zmq_ctx_term") _ctx_term :: proc(ctx: rawptr) -> c.int ---
|
_ctx_new :: proc() -> rawptr ---
|
||||||
@(link_name = "zmq_socket") _socket :: proc(ctx: rawptr, type: Socket_Type) -> rawptr ---
|
@(link_name = "zmq_ctx_term")
|
||||||
@(link_name = "zmq_close") _close :: proc(s: rawptr) -> c.int ---
|
_ctx_term :: proc(ctx: rawptr) -> c.int ---
|
||||||
@(link_name = "zmq_bind") _bind :: proc(s: rawptr, addr: cstring) -> c.int ---
|
@(link_name = "zmq_socket")
|
||||||
@(link_name = "zmq_connect") _connect :: proc(s: rawptr, addr: cstring) -> c.int ---
|
_socket :: proc(ctx: rawptr, type: Socket_Type) -> rawptr ---
|
||||||
@(link_name = "zmq_send") _send :: proc(s: rawptr, buf: rawptr, len: c.size_t, flags: c.int) -> c.int ---
|
@(link_name = "zmq_close")
|
||||||
@(link_name = "zmq_recv") _recv :: proc(s: rawptr, buf: rawptr, len: c.size_t, flags: c.int) -> c.int ---
|
_close :: proc(s: rawptr) -> c.int ---
|
||||||
@(link_name = "zmq_errno") _errno :: proc() -> c.int ---
|
@(link_name = "zmq_bind")
|
||||||
@(link_name = "zmq_strerror") _strerror :: proc(errnum: c.int) -> cstring ---
|
_bind :: proc(s: rawptr, addr: cstring) -> c.int ---
|
||||||
@(link_name = "zmq_version") _version :: proc(major, minor, patch: ^c.int) ---
|
@(link_name = "zmq_connect")
|
||||||
|
_connect :: proc(s: rawptr, addr: cstring) -> c.int ---
|
||||||
|
@(link_name = "zmq_send")
|
||||||
|
_send :: proc(s: rawptr, buf: rawptr, len: c.size_t, flags: c.int) -> c.int ---
|
||||||
|
@(link_name = "zmq_recv")
|
||||||
|
_recv :: proc(s: rawptr, buf: rawptr, len: c.size_t, flags: c.int) -> c.int ---
|
||||||
|
@(link_name = "zmq_errno")
|
||||||
|
_errno :: proc() -> c.int ---
|
||||||
|
@(link_name = "zmq_strerror")
|
||||||
|
_strerror :: proc(errnum: c.int) -> cstring ---
|
||||||
|
@(link_name = "zmq_version")
|
||||||
|
_version :: proc(major, minor, patch: ^c.int) ---
|
||||||
}
|
}
|
||||||
|
|
||||||
Context :: distinct rawptr
|
Context :: distinct rawptr
|
||||||
Socket :: distinct rawptr
|
Socket :: distinct rawptr
|
||||||
|
|
||||||
Socket_Type :: enum c.int {
|
Socket_Type :: enum c.int {
|
||||||
PAIR = 0,
|
PAIR = 0,
|
||||||
@@ -42,21 +53,19 @@ Flag :: enum c.int {
|
|||||||
SNDMORE = 2,
|
SNDMORE = 2,
|
||||||
}
|
}
|
||||||
|
|
||||||
Flags :: bit_set[Flag; c.int]
|
Flags :: bit_set[Flag;c.int]
|
||||||
|
|
||||||
NoBlock :: Flags{.DONTWAIT}
|
NoBlock :: Flags{.DONTWAIT}
|
||||||
|
|
||||||
Error :: distinct c.int
|
Error :: enum c.int {
|
||||||
|
Ok,
|
||||||
|
}
|
||||||
|
|
||||||
OK :: Error(0)
|
strerror :: proc(e: Error) -> string {
|
||||||
|
|
||||||
is_ok :: proc(e: Error) -> bool { return c.int(e) == 0 }
|
|
||||||
|
|
||||||
error_string :: proc(e: Error) -> string {
|
|
||||||
return string(_strerror(c.int(e)))
|
return string(_strerror(c.int(e)))
|
||||||
}
|
}
|
||||||
|
|
||||||
last_error :: proc() -> Error { return Error(_errno()) }
|
last_error :: proc() -> Error {return Error(_errno())}
|
||||||
|
|
||||||
version :: proc() -> (major, minor, patch: int) {
|
version :: proc() -> (major, minor, patch: int) {
|
||||||
ma, mi, pa: c.int = 0, 0, 0
|
ma, mi, pa: c.int = 0, 0, 0
|
||||||
@@ -64,12 +73,12 @@ version :: proc() -> (major, minor, patch: int) {
|
|||||||
return int(ma), int(mi), int(pa)
|
return int(ma), int(mi), int(pa)
|
||||||
}
|
}
|
||||||
|
|
||||||
new_context :: proc() -> Context {
|
ctx_new :: proc() -> Context {
|
||||||
return Context(_ctx_new())
|
return Context(_ctx_new())
|
||||||
}
|
}
|
||||||
|
|
||||||
term :: proc(ctx: Context) -> Error {
|
ctx_term :: proc(ctx: Context) -> Error {
|
||||||
return _ctx_term(rawptr(ctx)) == 0 ? OK : last_error()
|
return _ctx_term(rawptr(ctx)) == 0 ? .Ok : last_error()
|
||||||
}
|
}
|
||||||
|
|
||||||
socket :: proc(ctx: Context, type: Socket_Type) -> (Socket, Error) {
|
socket :: proc(ctx: Context, type: Socket_Type) -> (Socket, Error) {
|
||||||
@@ -77,11 +86,11 @@ socket :: proc(ctx: Context, type: Socket_Type) -> (Socket, Error) {
|
|||||||
if s == nil {
|
if s == nil {
|
||||||
return nil, last_error()
|
return nil, last_error()
|
||||||
}
|
}
|
||||||
return Socket(s), OK
|
return Socket(s), .Ok
|
||||||
}
|
}
|
||||||
|
|
||||||
close :: proc(s: Socket) -> Error {
|
close :: proc(s: Socket) -> Error {
|
||||||
return _close(rawptr(s)) == 0 ? OK : last_error()
|
return _close(rawptr(s)) == 0 ? .Ok : last_error()
|
||||||
}
|
}
|
||||||
|
|
||||||
bind :: proc(s: Socket, addr: string) -> Error {
|
bind :: proc(s: Socket, addr: string) -> Error {
|
||||||
@@ -97,7 +106,7 @@ send :: proc(s: Socket, buf: []u8, flags: Flags = Flags{}) -> (int, Error) {
|
|||||||
if n < 0 {
|
if n < 0 {
|
||||||
return -1, last_error()
|
return -1, last_error()
|
||||||
}
|
}
|
||||||
return int(n), OK
|
return int(n), .Ok
|
||||||
}
|
}
|
||||||
|
|
||||||
recv :: proc(s: Socket, buf: []u8, flags: Flags = Flags{}) -> (int, Error) {
|
recv :: proc(s: Socket, buf: []u8, flags: Flags = Flags{}) -> (int, Error) {
|
||||||
@@ -105,15 +114,20 @@ recv :: proc(s: Socket, buf: []u8, flags: Flags = Flags{}) -> (int, Error) {
|
|||||||
if n < 0 {
|
if n < 0 {
|
||||||
return -1, last_error()
|
return -1, last_error()
|
||||||
}
|
}
|
||||||
return int(n), OK
|
return int(n), .Ok
|
||||||
}
|
}
|
||||||
|
|
||||||
endpoint :: proc(s: rawptr, addr: string, f: proc "c" (s: rawptr, addr: cstring) -> c.int) -> Error {
|
endpoint :: proc(
|
||||||
|
s: rawptr,
|
||||||
|
addr: string,
|
||||||
|
f: proc "c" (s: rawptr, addr: cstring) -> c.int,
|
||||||
|
) -> Error {
|
||||||
buf: [256]u8
|
buf: [256]u8
|
||||||
n := copy(buf[:], addr)
|
n := copy(buf[:], addr)
|
||||||
if n >= len(buf) {
|
if n >= len(buf) {
|
||||||
return Error(7)
|
return Error(7)
|
||||||
}
|
}
|
||||||
buf[n] = 0
|
buf[n] = 0
|
||||||
return f(s, cstring(rawptr(&buf[0]))) == 0 ? OK : last_error()
|
return f(s, cstring(rawptr(&buf[0]))) == 0 ? .Ok : last_error()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user