Repository navigation
Expand file tree
/
Copy pathhypergo.go
More file actions
329 lines (297 loc) · 9.56 KB
/
Copy pathhypergo.go
File metadata and controls
329 lines (297 loc) · 9.56 KB
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
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
// Package hypergo is a client for HyperDex.
package hypergo
/*
#cgo LDFLAGS: -lhyperdex-client -lhyperdex-admin
#include <netinet/in.h>
#include "hyperdex/client.h"
#include "hyperdex/admin.h"
#include "hyperdex/datastructures.h"
*/
import "C"
import (
"fmt"
"io"
"log"
"runtime"
)
// CHANNEL_BUFFER_SIZE is the size of all the returned channels' buffer.
// You can set it to 0 if you want unbuffered channels.
var CHANNEL_BUFFER_SIZE = 1
// Timeout in miliseconds.
// Negative timeout means no timeout.
var TIMEOUT = -1
// Predicates
const (
FAIL = C.HYPERPREDICATE_FAIL
EQUALS = C.HYPERPREDICATE_EQUALS
LESS_EQUAL = C.HYPERPREDICATE_LESS_EQUAL
GREATER_EQUAL = C.HYPERPREDICATE_GREATER_EQUAL
CONTAINS_LESS_THAN = C.HYPERPREDICATE_CONTAINS_LESS_THAN // alias of LENGTH_LESS_EQUAL
REGEX = C.HYPERPREDICATE_REGEX
LENGTH_EQUALS = C.HYPERPREDICATE_LENGTH_EQUALS
LENGTH_LESS_EQUAL = C.HYPERPREDICATE_LENGTH_LESS_EQUAL
LENGTH_GREATER_EQUAL = C.HYPERPREDICATE_LENGTH_GREATER_EQUAL
CONTAINS = C.HYPERPREDICATE_CONTAINS
)
// Client is the hyperdex client used to make requests to hyperdex.
type Client struct {
ptr *C.struct_hyperdex_client
arena *C.struct_hyperdex_ds_arena
requests []request
closeChan chan bool
}
// Admin is a priviledged client used to do meta operations to hyperdex
type Admin struct {
ptr *C.struct_hyperdex_admin
requests []adminRequest
closeChan chan bool
}
// Attributes represents a map of key-value attribute pairs.
//
// The value must be either a string or an int64-compatible integer
// (int, int8, int16, int32, int64, uint8, uint16, uint32).
// An incompatible type will NOT result in a panic but in a regular error return.
//
// Please note that there is no support for uint64 since its negative might be incorrectly evaluated.
// Support for uint has been dropped because it is unspecified whether it is 32 or 64 bits.
// Correspond to an array of hyperdex_client_attribute
type Attributes map[string]interface{}
type Set []interface{}
type SetStr []string
type SetI64 []int64
type SetF64 []float64
type List []interface{}
type ListStr []string
type ListI64 []int64
type ListF64 []float64
type Map map[interface{}]interface{}
type MapStrStr map[string]string
type MapStrI64 map[string]int64
type MapStrF64 map[string]float64
type MapI64Str map[int64]string
type MapI64I64 map[int64]int64
type MapI64F64 map[int64]float64
type MapF64Str map[float64]string
type MapF64I64 map[float64]int64
type MapF64F64 map[float64]float64
// A hyperdex object.
// Err contains any error that happened when trying to retrieve
// this object.
type Object struct {
Err error
Key string
Attrs Attributes
}
// Read-only channel of objects
type ObjectChannel <-chan Object
// Read-only channel of errors
type ErrorChannel <-chan error
// Correspond to a hyperdex_client_attribute_check
type Condition struct {
Attr string
Value interface{}
Predicate int
}
// An outstanding request that contains three callback functions:
// success, failure, and complete. success() is called when
// hyperdex returns HYPERDEX_CLIENT_SUCCESS, and failure is called
// in all other cases.
//
// The boolean flag isIterator signifies whether the request is
// that of an iterator, namely hyperdex_client_search and its
// variants.
//
// complete() is called after success() if the request is not
// an iterator, or after receiving HYPREDEX_SEARCH_DONE if the
// request is an iterator.
type request struct {
id int64
isIterator bool
success func()
failure func(C.enum_hyperdex_client_returncode, string)
complete func()
status *C.enum_hyperdex_client_returncode
}
type adminRequest struct {
id int64
success func()
failure func(C.enum_hyperdex_admin_returncode, string)
status *C.enum_hyperdex_admin_returncode
}
// A custom error type that allows for examining HyperDex error code.
type HyperError struct {
returnCode C.enum_hyperdex_client_returncode
msg string
}
func (e HyperError) Error() string {
return fmt.Sprintf("Error %d: %s", e.returnCode, e.msg)
}
// Set output of log. Simply a wrapper around log.SetOutput.
func SetLogOutput(w io.Writer) {
log.SetOutput(w)
}
// NewClient initializes a hyperdex client ready to use.
//
// For every call to NewClient, there must be a call to Destroy.
//
// Panics when the internal looping goroutine receives an error from hyperdex.
//
// Example:
// client, err := hyperdex_client.NewClient("127.0.0.1", 1234)
// if err != nil {
// //handle error
// }
// defer client.Destroy()
// // use client
func NewClient(ip string, port int) (*Client, error) {
C_client := C.hyperdex_client_create(C.CString(ip), C.uint16_t(port))
C_arena := C.hyperdex_ds_arena_create()
if C_client == nil {
return nil, fmt.Errorf("Could not create hyperdex_client (ip=%s, port=%d)", ip, port)
}
client := &Client{
C_client,
C_arena,
make([]request, 0, 8), // No reallocation within 8 concurrent requests to hyperdex_client_loop
make(chan bool, 1),
}
go func() {
for {
select {
// quit goroutine when client is destroyed
case <-client.closeChan:
return
default:
// check if there are pending requests
// and only if there are, call hyperdex_client_loop
if len(client.requests) > 0 {
var status C.enum_hyperdex_client_returncode
ret := int64(C.hyperdex_client_loop(client.ptr, C.int(TIMEOUT), &status))
//log.Printf("hyperdex_client_loop(%X, %d, %X) -> %d\n", unsafe.Pointer(client.ptr), hyperdex_client_loop_timeout, unsafe.Pointer(&status), ret)
if ret < 0 {
panic(newInternalError(status,
C.GoString(C.hyperdex_client_error_message(client.ptr))).Error())
}
// find processed request among pending requests
for i, req := range client.requests {
if req.id == ret {
if status == C.HYPERDEX_CLIENT_SUCCESS {
switch *req.status {
case C.HYPERDEX_CLIENT_SUCCESS:
if req.success != nil {
req.success()
}
if req.isIterator {
// We want to break out at here so that the
// request won't get removed
goto SKIP_DELETING_REQUEST
} else if req.complete != nil {
// We want to break out at here so that the
// request won't get removed
req.complete()
}
case C.HYPERDEX_CLIENT_SEARCHDONE:
if req.complete != nil {
req.complete()
}
case C.HYPERDEX_CLIENT_CMPFAIL:
if req.failure != nil {
req.failure(*req.status,
C.GoString(C.hyperdex_client_error_message(client.ptr)))
}
default:
if req.failure != nil {
req.failure(*req.status,
C.GoString(C.hyperdex_client_error_message(client.ptr)))
}
}
} else if req.failure != nil {
req.failure(status,
C.GoString(C.hyperdex_client_error_message(client.ptr)))
}
client.requests = append(client.requests[:i], client.requests[i+1:]...)
SKIP_DELETING_REQUEST:
break
}
}
}
// prevent other goroutines from starving
runtime.Gosched()
}
}
panic("Should not be reached: end of infinite loop")
}()
return client, nil
}
// Destroy closes the connection between the Client and hyperdex. It has to be used on a client that is not used anymore.
//
// For every call to NewClient, there must be a call to Destroy.
func (client *Client) Destroy() {
close(client.closeChan)
C.hyperdex_client_destroy(client.ptr)
//log.Printf("hyperdex_client_destroy(%X)\n", unsafe.Pointer(client.ptr))
}
func NewAdmin(ip string, port int) (*Admin, error) {
C_admin := C.hyperdex_admin_create(C.CString(ip), C.uint16_t(port))
if C_admin == nil {
return nil, fmt.Errorf("Could not create hyperdex_admin (ip=%s, port=%d)", ip, port)
}
admin := &Admin{
C_admin,
make([]adminRequest, 0, 8),
make(chan bool, 1),
}
go func() {
for {
select {
// quit goroutine when client is destroyed
case <-admin.closeChan:
return
default:
// check if there are pending requests
// and only if there are, call hyperdex_client_loop
if len(admin.requests) > 0 {
var status C.enum_hyperdex_admin_returncode
ret := int64(C.hyperdex_admin_loop(admin.ptr, C.int(TIMEOUT), &status))
//log.Printf("hyperdex_client_loop(%X, %d, %X) -> %d\n", unsafe.Pointer(client.ptr), hyperdex_client_loop_timeout, unsafe.Pointer(&status), ret)
if ret < 0 {
panic(newInternalError(status, "Admin error"))
}
// find processed request among pending requests
for i, req := range admin.requests {
if req.id == ret {
if status == C.HYPERDEX_ADMIN_SUCCESS {
switch *req.status {
case C.HYPERDEX_ADMIN_SUCCESS:
if req.success != nil {
req.success()
}
default:
if req.failure != nil {
req.failure(*req.status,
C.GoString(C.hyperdex_admin_error_message(admin.ptr)))
}
}
} else if req.failure != nil {
req.failure(status,
C.GoString(C.hyperdex_admin_error_message(admin.ptr)))
}
admin.requests = append(admin.requests[:i], admin.requests[i+1:]...)
}
}
}
// prevent other goroutines from starving
runtime.Gosched()
}
}
panic("Should not be reached: end of infinite loop")
}()
return admin, nil
}
// Destroy closes the connection between the Admin and hyperdex. It has to be used on a admin that is not used anymore.
//
// For every call to NewAdmin, there must be a call to Destroy.
func (admin *Admin) Destroy() {
close(admin.closeChan)
C.hyperdex_admin_destroy(admin.ptr)
}