Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Support
Keyboard shortcuts
?
Submit feedback
Contribute to GitLab
Sign in / Register
Toggle navigation
N
neoppod
Project overview
Project overview
Details
Activity
Releases
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Issues
0
Issues
0
List
Boards
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Analytics
Analytics
CI / CD
Repository
Value Stream
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
Levin Zimmermann
neoppod
Commits
d252963a
Commit
d252963a
authored
Jul 09, 2020
by
Kirill Smelkov
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
.
parent
18b22971
Changes
4
Hide whitespace changes
Inline
Side-by-side
Showing
4 changed files
with
75 additions
and
58 deletions
+75
-58
go/zodb/storage/zeo/marshal.go
go/zodb/storage/zeo/marshal.go
+24
-22
go/zodb/storage/zeo/zeo.go
go/zodb/storage/zeo/zeo.go
+35
-20
go/zodb/storage/zeo/zeo_test.go
go/zodb/storage/zeo/zeo_test.go
+7
-7
go/zodb/storage/zeo/zrpc.go
go/zodb/storage/zeo/zrpc.go
+9
-9
No files found.
go/zodb/storage/zeo/marshal.go
View file @
d252963a
...
...
@@ -35,11 +35,13 @@ import (
"lab.nexedi.com/kirr/neo/go/zodb/internal/pickletools"
)
type
encoding
byte
// Z - pickles, M - msgpack
// ---- message encode/decode ----
// pktEncode encodes message into raw packet.
func
(
zl
*
zLink
)
pktEncode
(
m
msg
)
*
pktBuf
{
switch
zl
.
encoding
{
func
(
e
encoding
)
pktEncode
(
m
msg
)
*
pktBuf
{
switch
e
{
case
'Z'
:
return
pktEncodeZ
(
m
)
case
'M'
:
return
pktEncodeM
(
m
)
default
:
panic
(
"bug"
)
...
...
@@ -47,8 +49,8 @@ func (zl *zLink) pktEncode(m msg) *pktBuf {
}
// pktDecode decodes raw packet into message.
func
(
zl
*
zLink
)
pktDecode
(
pkb
*
pktBuf
)
(
msg
,
error
)
{
switch
zl
.
encoding
{
func
(
e
encoding
)
pktDecode
(
pkb
*
pktBuf
)
(
msg
,
error
)
{
switch
e
{
case
'Z'
:
return
pktDecodeZ
(
pkb
)
case
'M'
:
return
pktDecodeM
(
pkb
)
default
:
panic
(
"bug"
)
...
...
@@ -240,8 +242,8 @@ func derrf(format string, argv ...interface{}) error {
// ---- encode/decode for data types ----
// xuint64Unpack tries to decode packed 8-byte string as bigendian uint64
func
(
zl
*
zLink
)
xuint64Unpack
(
xv
interface
{})
(
uint64
,
bool
)
{
switch
zl
.
encoding
{
func
(
e
encoding
)
xuint64Unpack
(
xv
interface
{})
(
uint64
,
bool
)
{
switch
e
{
default
:
panic
(
"bug"
)
...
...
@@ -270,11 +272,11 @@ func (zl *zLink) xuint64Unpack(xv interface{}) (uint64, bool) {
}
// xuint64Pack packs v into big-endian 8-byte string
func
(
zl
*
zLink
)
xuint64Pack
(
v
uint64
)
interface
{}
{
func
(
e
encoding
)
xuint64Pack
(
v
uint64
)
interface
{}
{
var
b
[
8
]
byte
binary
.
BigEndian
.
PutUint64
(
b
[
:
],
v
)
switch
zl
.
encoding
{
switch
e
{
default
:
panic
(
"bug"
)
...
...
@@ -288,28 +290,28 @@ func (zl *zLink) xuint64Pack(v uint64) interface{} {
}
}
func
(
zl
*
zLink
)
tidPack
(
tid
zodb
.
Tid
)
interface
{}
{
return
zl
.
xuint64Pack
(
uint64
(
tid
))
func
(
e
encoding
)
tidPack
(
tid
zodb
.
Tid
)
interface
{}
{
return
e
.
xuint64Pack
(
uint64
(
tid
))
}
func
(
zl
*
zLink
)
oidPack
(
oid
zodb
.
Oid
)
interface
{}
{
return
zl
.
xuint64Pack
(
uint64
(
oid
))
func
(
e
encoding
)
oidPack
(
oid
zodb
.
Oid
)
interface
{}
{
return
e
.
xuint64Pack
(
uint64
(
oid
))
}
func
(
zl
*
zLink
)
tidUnpack
(
xv
interface
{})
(
zodb
.
Tid
,
bool
)
{
v
,
ok
:=
zl
.
xuint64Unpack
(
xv
)
func
(
e
encoding
)
asTid
(
xv
interface
{})
(
zodb
.
Tid
,
bool
)
{
v
,
ok
:=
e
.
xuint64Unpack
(
xv
)
return
zodb
.
Tid
(
v
),
ok
}
func
(
zl
*
zLink
)
oidUnpack
(
xv
interface
{})
(
zodb
.
Oid
,
bool
)
{
v
,
ok
:=
zl
.
xuint64Unpack
(
xv
)
func
(
e
encoding
)
asOid
(
xv
interface
{})
(
zodb
.
Oid
,
bool
)
{
v
,
ok
:=
e
.
xuint64Unpack
(
xv
)
return
zodb
.
Oid
(
v
),
ok
}
// asTuple tries to decode object as tuple. XXX
func
(
zl
*
zLink
)
asTuple
(
xt
interface
{})
(
tuple
,
bool
)
{
switch
zl
.
encoding
{
func
(
e
encoding
)
asTuple
(
xt
interface
{})
(
tuple
,
bool
)
{
switch
e
{
default
:
panic
(
"bug"
)
...
...
@@ -326,8 +328,8 @@ func (zl *zLink) asTuple(xt interface{}) (tuple, bool) {
}
// asBytes tries to decode object as raw bytes.
func
(
zl
*
zLink
)
asBytes
(
xb
interface
{})
([]
byte
,
bool
)
{
switch
zl
.
encoding
{
func
(
e
encoding
)
asBytes
(
xb
interface
{})
([]
byte
,
bool
)
{
switch
e
{
default
:
panic
(
"bug"
)
...
...
@@ -347,8 +349,8 @@ func (zl *zLink) asBytes(xb interface{}) ([]byte, bool) {
}
// asString tries to decode object as string.
func
(
zl
*
zLink
)
asString
(
xs
interface
{})
(
string
,
bool
)
{
switch
zl
.
encoding
{
func
(
e
encoding
)
asString
(
xs
interface
{})
(
string
,
bool
)
{
switch
e
{
default
:
panic
(
"bug"
)
...
...
go/zodb/storage/zeo/zeo.go
View file @
d252963a
...
...
@@ -36,7 +36,7 @@ import (
)
type
zeo
struct
{
srv
*
zLink
// XXX rename -> link?
link
*
zLink
// driver client <- watcher: database commits | errors.
watchq
chan
<-
zodb
.
Event
...
...
@@ -63,7 +63,7 @@ func (z *zeo) Sync(ctx context.Context) (head zodb.Tid, err error) {
return
zodb
.
InvalidTid
,
err
}
head
,
ok
:=
z
.
srv
.
tidUnpack
(
xhead
)
head
,
ok
:=
z
.
link
.
enc
.
asTid
(
xhead
)
if
!
ok
{
return
zodb
.
InvalidTid
,
rpc
.
ereplyf
(
"got %v; expect tid"
,
xhead
)
}
...
...
@@ -83,19 +83,20 @@ func (z *zeo) Load(ctx context.Context, xid zodb.Xid) (*mem.Buf, zodb.Tid, error
func
(
z
*
zeo
)
_Load
(
ctx
context
.
Context
,
xid
zodb
.
Xid
)
(
*
mem
.
Buf
,
zodb
.
Tid
,
error
)
{
rpc
:=
z
.
rpc
(
"loadBefore"
)
xres
,
err
:=
rpc
.
call
(
ctx
,
z
.
srv
.
oidPack
(
xid
.
Oid
),
z
.
srv
.
tidPack
(
xid
.
At
+
1
))
// XXX at2Before
enc
:=
z
.
link
.
enc
xres
,
err
:=
rpc
.
call
(
ctx
,
enc
.
oidPack
(
xid
.
Oid
),
enc
.
tidPack
(
xid
.
At
+
1
))
// XXX at2Before
if
err
!=
nil
{
return
nil
,
0
,
err
}
// (data, serial, next_serial | None)
res
,
ok
:=
z
.
srv
.
asTuple
(
xres
)
res
,
ok
:=
enc
.
asTuple
(
xres
)
if
!
ok
||
len
(
res
)
!=
3
{
return
nil
,
0
,
rpc
.
ereplyf
(
"got %#v; expect 3-tuple"
,
xres
)
}
data
,
ok1
:=
z
.
srv
.
asBytes
(
res
[
0
])
serial
,
ok2
:=
z
.
srv
.
tidUnpack
(
res
[
1
])
data
,
ok1
:=
enc
.
asBytes
(
res
[
0
])
serial
,
ok2
:=
enc
.
asTid
(
res
[
1
])
// next_serial (res[2]) - just ignore
if
!
(
ok1
&&
ok2
)
{
...
...
@@ -115,20 +116,21 @@ func (z *zeo) Iterate(ctx context.Context, tidMin, tidMax zodb.Tid) zodb.ITxnIte
func
(
z
*
zeo
)
invalidateTransaction
(
arg
interface
{})
(
err
error
)
{
defer
xerr
.
Context
(
&
err
,
"invalidateTransaction"
)
t
,
ok
:=
z
.
srv
.
asTuple
(
arg
)
enc
:=
z
.
link
.
enc
t
,
ok
:=
enc
.
asTuple
(
arg
)
if
!
ok
||
len
(
t
)
!=
2
{
return
fmt
.
Errorf
(
"got %#v; expect 2-tuple"
,
arg
)
}
// (tid, oidv)
tid
,
ok1
:=
z
.
srv
.
tidUnpack
(
t
[
0
])
xoidt
,
ok2
:=
z
.
srv
.
asTuple
(
t
[
1
])
tid
,
ok1
:=
enc
.
asTid
(
t
[
0
])
xoidt
,
ok2
:=
enc
.
asTuple
(
t
[
1
])
if
!
(
ok1
&&
ok2
)
{
return
fmt
.
Errorf
(
"got (%T, %T); expect (tid, []oid)"
,
t
...
)
}
oidv
:=
[]
zodb
.
Oid
{}
for
_
,
xoid
:=
range
xoidt
{
oid
,
ok
:=
z
.
srv
.
oidUnpack
(
xoid
)
oid
,
ok
:=
enc
.
asOid
(
xoid
)
if
!
ok
{
return
fmt
.
Errorf
(
"non-oid %#v in oidv"
,
xoid
)
}
...
...
@@ -187,7 +189,7 @@ func ereplyf(addr, method, format string, argv ...interface{}) *errorUnexpectedR
// rpc returns rpc object handy to make calls/create errors
func
(
z
*
zeo
)
rpc
(
method
string
)
rpc
{
return
rpc
{
zl
:
z
.
srv
,
method
:
method
}
return
rpc
{
zl
:
z
.
link
,
method
:
method
}
}
type
rpc
struct
{
...
...
@@ -242,7 +244,7 @@ func (r rpc) excError(exc string, argv tuple) error {
return
r
.
ereplyf
(
"poskeyerror: got %#v; expect 1-tuple"
,
argv
...
)
}
oid
,
ok
:=
r
.
zl
.
oidUnpack
(
argv
[
0
])
oid
,
ok
:=
r
.
zl
.
enc
.
asOid
(
argv
[
0
])
if
!
ok
{
return
r
.
ereplyf
(
"poskeyerror: got (%v); expect (oid)"
,
argv
[
0
])
}
...
...
@@ -259,14 +261,15 @@ func (r rpc) excError(exc string, argv tuple) error {
// zeo5Error decodes arg of reply with msgExcept flag set and returns
// corresponding error.
func
(
r
rpc
)
zeo5Error
(
arg
interface
{})
error
{
enc
:=
r
.
zl
.
enc
// ('type', (arg1, arg2, arg3, ...))
texc
,
ok
:=
r
.
zl
.
asTuple
(
arg
)
texc
,
ok
:=
enc
.
asTuple
(
arg
)
if
!
ok
||
len
(
texc
)
!=
2
{
return
r
.
ereplyf
(
"except5: got %#v; expect 2-tuple"
,
arg
)
}
exc
,
ok1
:=
r
.
zl
.
asString
(
texc
[
0
])
argv
,
ok2
:=
r
.
zl
.
asTuple
(
texc
[
1
])
exc
,
ok1
:=
enc
.
asString
(
texc
[
0
])
argv
,
ok2
:=
enc
.
asTuple
(
texc
[
1
])
if
!
(
ok1
&&
ok2
)
{
return
r
.
ereplyf
(
"except5: got (%T, %T); expect (str, tuple)"
,
texc
...
)
}
...
...
@@ -382,7 +385,19 @@ func openByURL(ctx context.Context, u *url.URL, opt *zodb.DriverOptions) (_ zodb
return
nil
,
zodb
.
InvalidTid
,
fmt
.
Errorf
(
"TODO write mode not implemented"
)
}
zl
,
err
:=
dialZLink
(
ctx
,
net
,
addr
)
// XXX + methodTable {invalidateTransaction tid, oidv} -> ...
z
:=
&
zeo
{
watchq
:
opt
.
Watchq
,
url
:
url
}
//zl, err := dialZLink(ctx, net, addr) // XXX + methodTable {invalidateTransaction tid, oidv} -> ...
zl
,
err
:=
dialZLink
(
ctx
,
net
,
addr
,
/*
// notifyTab
map[string]func(interface{})error {
"invalidateTransaction": z.invalidateTransaction,
},
// serveTab
nil,
*/
)
if
err
!=
nil
{
return
nil
,
zodb
.
InvalidTid
,
err
}
...
...
@@ -394,7 +409,7 @@ func openByURL(ctx context.Context, u *url.URL, opt *zodb.DriverOptions) (_ zodb
}()
z
:=
&
zeo
{
srv
:
zl
,
watchq
:
opt
.
Watchq
,
url
:
url
}
z
.
link
=
zl
rpc
:=
z
.
rpc
(
"register"
)
xlastTid
,
err
:=
rpc
.
call
(
ctx
,
storageID
,
opt
.
ReadOnly
)
...
...
@@ -404,7 +419,7 @@ func openByURL(ctx context.Context, u *url.URL, opt *zodb.DriverOptions) (_ zodb
// register returns last_tid in ZEO5 but nothing earlier.
// if so we have to retrieve last_tid in another RPC.
if
z
.
srv
.
ver
<
"5"
{
if
z
.
link
.
ver
<
"5"
{
rpc
=
z
.
rpc
(
"lastTransaction"
)
xlastTid
,
err
=
rpc
.
call
(
ctx
)
if
err
!=
nil
{
...
...
@@ -412,7 +427,7 @@ func openByURL(ctx context.Context, u *url.URL, opt *zodb.DriverOptions) (_ zodb
}
}
lastTid
,
ok
:=
zl
.
tidUnpack
(
xlastTid
)
lastTid
,
ok
:=
zl
.
enc
.
asTid
(
xlastTid
)
if
!
ok
{
return
nil
,
zodb
.
InvalidTid
,
rpc
.
ereplyf
(
"got %v; expect tid"
,
xlastTid
)
}
...
...
@@ -458,7 +473,7 @@ func openByURL(ctx context.Context, u *url.URL, opt *zodb.DriverOptions) (_ zodb
}
func
(
z
*
zeo
)
Close
()
error
{
err
:=
z
.
srv
.
Close
()
err
:=
z
.
link
.
Close
()
if
z
.
watchq
!=
nil
{
close
(
z
.
watchq
)
}
...
...
go/zodb/storage/zeo/zeo_test.go
View file @
d252963a
...
...
@@ -43,7 +43,7 @@ type ZEOSrv interface {
Addr
()
string
// unix-socket address of the server
Close
()
error
Encoding
()
byte
// encoding used on the wire - 'M' or 'Z'
Encoding
()
encoding
// encoding used on the wire - 'M' or 'Z'
}
// ZEOPySrv represents running ZEO/py server.
...
...
@@ -134,10 +134,10 @@ func (z *ZEOPySrv) Close() (err error) {
return
err
}
func
(
z
*
ZEOPySrv
)
Encoding
()
byte
{
enc
oding
:=
byte
(
'Z'
)
if
z
.
opt
.
msgpack
{
enc
oding
=
byte
(
'M'
)
}
return
enc
oding
func
(
z
*
ZEOPySrv
)
Encoding
()
encoding
{
enc
:=
encoding
(
'Z'
)
if
z
.
opt
.
msgpack
{
enc
=
encoding
(
'M'
)
}
return
enc
}
...
...
@@ -227,8 +227,8 @@ func TestHandshake(t *testing.T) {
}()
ewant
:=
zsrv
.
Encoding
()
if
zlink
.
enc
oding
!=
ewant
{
t
.
Fatalf
(
"handshake: encoding=%c ; want %c"
,
zlink
.
enc
oding
,
ewant
)
if
zlink
.
enc
!=
ewant
{
t
.
Fatalf
(
"handshake: encoding=%c ; want %c"
,
zlink
.
enc
,
ewant
)
}
})
}
...
...
go/zodb/storage/zeo/zrpc.go
View file @
d252963a
...
...
@@ -73,8 +73,8 @@ type zLink struct {
down1
sync
.
Once
errClose
error
// error got from .link.Close()
ver
string
// protocol version in use (without "Z" or "M" prefix)
enc
oding
byte
// protocol encoding in use ('Z' or 'M')
ver
string
// protocol version in use (without "Z" or "M" prefix)
enc
encoding
// protocol encoding in use ('Z' or 'M')
}
// (called after handshake)
...
...
@@ -120,7 +120,7 @@ func (zl *zLink) Close() error {
// serveRecv handles receives from underlying link and dispatches them to calls
// waiting
results
.
// waiting
for results, to notify and serve handlers
.
func
(
zl
*
zLink
)
serveRecv
()
{
defer
zl
.
serveWg
.
Done
()
for
{
...
...
@@ -143,7 +143,7 @@ func (zl *zLink) serveRecv() {
// serveRecv1 handles 1 incoming packet.
func
(
zl
*
zLink
)
serveRecv1
(
pkb
*
pktBuf
)
error
{
// decode packet
m
,
err
:=
zl
.
pktDecode
(
pkb
)
m
,
err
:=
zl
.
enc
.
pktDecode
(
pkb
)
if
err
!=
nil
{
return
err
}
...
...
@@ -245,7 +245,7 @@ func (zl *zLink) Call(ctx context.Context, method string, argv ...interface{}) (
zl
.
callMu
.
Unlock
()
// (msgid, async, method, argv)
pkb
:=
zl
.
pktEncode
(
msg
{
pkb
:=
zl
.
enc
.
pktEncode
(
msg
{
msgid
:
callID
,
flags
:
0
,
method
:
method
,
...
...
@@ -279,7 +279,7 @@ func (zl *zLink) reply(msgid int64, res interface{}) (err error) {
}
}()
pkb
:=
zl
.
pktEncode
(
msg
{
pkb
:=
zl
.
enc
.
pktEncode
(
msg
{
msgid
:
msgid
,
flags
:
msgAsync
,
method
:
".reply"
,
...
...
@@ -480,7 +480,7 @@ func handshake(ctx context.Context, conn net.Conn) (_ *zLink, err error) {
}
// use wire encoding preferred by server
enc
oding
:=
proto
[
0
]
enc
:=
encoding
(
proto
[
0
])
// extract peer version from protocol string and choose actual
// version to use as min(peer, mybest)
...
...
@@ -505,14 +505,14 @@ func handshake(ctx context.Context, conn net.Conn) (_ *zLink, err error) {
// version selected - now send it back to server as
// corresponding handshake reply.
pkb
=
allocPkb
()
pkb
.
WriteString
(
fmt
.
Sprintf
(
"%c%s"
,
enc
oding
,
ver
))
pkb
.
WriteString
(
fmt
.
Sprintf
(
"%c%s"
,
enc
,
ver
))
err
=
zl
.
sendPkt
(
pkb
)
if
err
!=
nil
{
return
fmt
.
Errorf
(
"tx: %s"
,
err
)
}
zl
.
ver
=
ver
zl
.
enc
oding
=
encoding
zl
.
enc
=
enc
close
(
hok
)
return
nil
})
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment