summaryrefslogtreecommitdiff
path: root/src/rpc2.h
blob: 52e14ec1610a2e9178d898ffc59c0241b35c3cea (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
#pragma once

#include <al/array.h>

#include "util/packet.h"

#include "packet_stream.h"

struct nn_rpc_connection;
struct nn_rpc_callback {
    u32 id;
    void (*callback)(void *, struct nn_rpc_connection *, struct nn_packet *);
    void *userdata;
};

enum {
    NNWT_SKEW_POSITIVE = 0,
    NNWT_SKEW_NEGATIVE
};

struct nn_skew {
    u64 ts;
    bool direction;
};

struct nn_rpc_connection {
    struct nn_packet_stream *stream;
    array(struct nn_rpc_callback) callbacks;
    u32 outgoing;
    bool flushing;
    struct {
        struct nn_skew s;
        s8 query_active;
        u32 repeat;
        array(struct nn_skew) results;
    } skew;
    struct nn_rpc *rpc;
};

struct nn_rpc_command {
    s8 op;
    bool (*callback)(void *, struct nn_rpc_connection *, struct nn_packet *, struct nn_packet *);
    void *userdata;
};

struct nn_rpc {
    struct nn_event_loop *loop;
    u32 increment;
    array(struct nn_rpc_command) commands;
    bool do_query_skew;
    struct nn_rpc_connection *conn; // client
    array(struct nn_rpc_connection *) connections;
    void (*connection_callback)(void *, struct nn_rpc_connection *);
    void (*ready_callback)(void *, struct nn_rpc_connection *);
    void (*connection_closed_callback)(void *, struct nn_rpc_connection *);
    void *userdata;
};

bool nn_rpc_init(struct nn_rpc *rpc, struct nn_event_loop *loop,
    void (*connection_callback)(void *, struct nn_rpc_connection *),
    void (*ready_callback)(void *, struct nn_rpc_connection *),
    void (*connection_closed_callback)(void *, struct nn_rpc_connection *), void *userdata);
void nn_rpc_query_clock_skews(struct nn_rpc *rpc, bool do_query_skew);
void nn_rpc_add_command(struct nn_rpc *rpc, struct nn_rpc_command *command);
void nn_rpc_add_stream(struct nn_rpc *rpc, struct nn_packet_stream *stream);
void nn_rpc_prepare_client(struct nn_rpc *rpc);
bool nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port);
bool nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port);
void nn_rpc_disconnect(struct nn_rpc *rpc);
struct nn_packet *nn_rpc_get_packet(struct nn_rpc *rpc, s8 op);
void nn_rpc_free(struct nn_rpc *rpc);

// The order of commands per-connection will be respected.
void nn_rpc_connection_command(struct nn_rpc_connection *conn, struct nn_packet *packet,
    void (*callback)(void *, struct nn_rpc_connection *conn, struct nn_packet *), void *userdata);
void nn_rpc_conn_flush(struct nn_rpc_connection *conn);
void nn_rpc_conn_disconnect(struct nn_rpc_connection *conn);