Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Support
Keyboard shortcuts
?
Submit feedback
Contribute to GitLab
Sign in / Register
Toggle navigation
W
wendelin.core
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
Analytics
Analytics
Repository
Value Stream
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Create a new issue
Commits
Issue Boards
Open sidebar
Joshua
wendelin.core
Commits
f4a422b7
Commit
f4a422b7
authored
Jul 16, 2019
by
Kirill Smelkov
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
.
parent
02470041
Changes
2
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
18 additions
and
35 deletions
+18
-35
wcfs/__init__.py
wcfs/__init__.py
+18
-1
wcfs/wcfs_test.py
wcfs/wcfs_test.py
+0
-34
No files found.
wcfs/__init__.py
View file @
f4a422b7
...
@@ -105,8 +105,25 @@ class WatchLink(object):
...
@@ -105,8 +105,25 @@ class WatchLink(object):
wlink
.
_wrx
=
os
.
fdopen
(
wh
,
'rb'
)
wlink
.
_wrx
=
os
.
fdopen
(
wh
,
'rb'
)
wlink
.
_wtx
=
os
.
fdopen
(
wh2
,
'wb'
)
wlink
.
_wtx
=
os
.
fdopen
(
wh2
,
'wb'
)
# XXX vvv -> test only?
wlink
.
rx_eof
=
chan
()
# becomes ready when wcfs closes its tx side
wlink
.
fatalv
=
[]
# fatal messages received from wcfs
# XXX ...
# inv.protocol message IO
wlink
.
_acceptq
=
chan
()
# (stream, msg) server originated messages go here
wlink
.
_rxmu
=
threading
.
Lock
()
wlink
.
_rxtab
=
{}
# stream -> rxq server replies go via here
wlink
.
_accepted
=
set
()
# of stream streams we accepted but did not replied yet
wlink
.
_txmu
=
threading
.
Lock
()
# serializes writes
wlink
.
_txclosed
=
False
serveCtx
,
wlink
.
_serveCancel
=
context
.
with_cancel
(
context
.
background
())
wlink
.
_serveWG
=
sync
.
WorkGroup
(
serveCtx
)
wlink
.
_serveWG
.
go
(
wlink
.
_serveRX
)
# this tWatchLink currently watches the following files at particular state.
wlink
.
_watching
=
{}
# {} foid -> tWatch
# ---- WatchLink message IO ----
# ---- WatchLink message IO ----
...
...
wcfs/wcfs_test.py
View file @
f4a422b7
...
@@ -725,41 +725,7 @@ class tWatchLink(wcfs.WatchLink):
...
@@ -725,41 +725,7 @@ class tWatchLink(wcfs.WatchLink):
def
__init__
(
t
,
tdb
):
def
__init__
(
t
,
tdb
):
super
(
tWatchLink
,
t
).
__init__
(
tdb
.
wc
)
super
(
tWatchLink
,
t
).
__init__
(
tdb
.
wc
)
t
.
tdb
=
tdb
t
.
tdb
=
tdb
# head/watch handle.
#
# python/stdio lock file object on read/write, however we need both
# read and write to be working simultaneously.
# -> use 2 separate file objects for rx and tx.
#
# fdopen takes ownership of file descriptor and closes it when file
# object is closed -> dup fd so that each file object has its own fd.
wh
=
os
.
open
(
tdb
.
wc
.
_path
(
"head/watch"
),
os
.
O_RDWR
)
wh2
=
os
.
dup
(
wh
)
t
.
_wrx
=
os
.
fdopen
(
wh
,
'rb'
)
t
.
_wtx
=
os
.
fdopen
(
wh2
,
'wb'
)
t
.
rx_eof
=
chan
()
# becomes ready when wcfs closes its tx side
t
.
fatalv
=
[]
# fatal messages received from wcfs
# inv.protocol message IO
t
.
_acceptq
=
chan
()
# (stream, msg) server originated messages go here
t
.
_rxmu
=
threading
.
Lock
()
t
.
_rxtab
=
{}
# stream -> rxq server replies go via here
t
.
_accepted
=
set
()
# of stream streams we accepted but did not replied yet
t
.
_txmu
=
threading
.
Lock
()
# serializes writes
t
.
_txclosed
=
False
serveCtx
,
t
.
_serveCancel
=
context
.
with_cancel
(
context
.
background
())
t
.
_serveWG
=
sync
.
WorkGroup
(
serveCtx
)
t
.
_serveWG
.
go
(
t
.
_serveRX
)
# this tWatchLink currently watches the following files at particular state.
t
.
_watching
=
{}
# {} foid -> tWatch
tdb
.
_wlinks
.
add
(
t
)
tdb
.
_wlinks
.
add
(
t
)
def
_closeTX
(
t
):
def
_closeTX
(
t
):
...
...
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