-
Notifications
You must be signed in to change notification settings - Fork 298
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge remote-tracking branch 'eugeneia/mcpring-group-freelist' into i…
…pfix-rss
- Loading branch information
Showing
14 changed files
with
612 additions
and
21 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,49 @@ | ||
-- Use of this source code is governed by the Apache 2.0 license; see COPYING. | ||
|
||
module(...,package.seeall) | ||
|
||
local shm = require("core.shm") | ||
local interlink = require("lib.interlink") | ||
|
||
local Receiver = {name="apps.interlink.Receiver"} | ||
|
||
function Receiver:new (_, name) | ||
local self = {} | ||
self.shm_name = "group/interlink/"..name..".interlink" | ||
self.backlink = "interlink/receiver/"..name..".interlink" | ||
self.interlink = interlink.attach_receiver(self.shm_name) | ||
shm.alias(self.backlink, self.shm_name) | ||
return setmetatable(self, {__index=Receiver}) | ||
end | ||
|
||
function Receiver:pull () | ||
local o, r, n = self.output.output, self.interlink, 0 | ||
if not o then return end -- don’t forward packets until connected | ||
while not interlink.empty(r) and n < engine.pull_npackets do | ||
link.transmit(o, interlink.extract(r)) | ||
n = n + 1 | ||
end | ||
interlink.pull(r) | ||
end | ||
|
||
function Receiver:stop () | ||
interlink.detach_receiver(self.interlink, self.shm_name) | ||
shm.unlink(self.backlink) | ||
end | ||
|
||
-- Detach receivers to prevent leaking interlinks opened by pid. | ||
-- | ||
-- This is an internal API function provided for cleanup during | ||
-- process termination. | ||
function Receiver.shutdown (pid) | ||
for _, name in ipairs(shm.children("/"..pid.."/interlink/receiver")) do | ||
local backlink = "/"..pid.."/interlink/receiver/"..name..".interlink" | ||
local shm_name = "/"..pid.."/group/interlink/"..name..".interlink" | ||
-- Call protected in case /<pid>/group is already unlinked. | ||
local ok, r = pcall(interlink.open, shm_name) | ||
if ok then interlink.detach_receiver(r, shm_name) end | ||
shm.unlink(backlink) | ||
end | ||
end | ||
|
||
return Receiver |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,33 @@ | ||
#!snabb snsh | ||
|
||
-- Use of this source code is governed by the Apache 2.0 license; see COPYING. | ||
|
||
local worker = require("core.worker") | ||
local interlink = require("lib.interlink") | ||
local Receiver = require("apps.interlink.receiver") | ||
local Sink = require("apps.basic.basic_apps").Sink | ||
|
||
-- Synopsis: selftest.snabb [duration] | ||
local DURATION = tonumber(main.parameters[1]) or 10 | ||
|
||
worker.start("source", [[require("apps.interlink.test_source").start("test")]]) | ||
|
||
local c = config.new() | ||
|
||
config.app(c, "test", Receiver) | ||
config.app(c, "sink", Sink) | ||
config.link(c, "test.output->sink.input") | ||
|
||
engine.configure(c) | ||
engine.main({duration=DURATION, report={showlinks=true}}) | ||
|
||
for w, s in pairs(worker.status()) do | ||
print(("worker %s: pid=%s alive=%s status=%s"):format( | ||
w, s.pid, s.alive, s.status)) | ||
end | ||
local stats = link.stats(engine.app_table["sink"].input.input) | ||
print(stats.txpackets / 1e6 / DURATION .. " Mpps") | ||
|
||
-- test teardown | ||
engine.configure(config.new()) | ||
engine.main({duration=0.1}) |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,15 @@ | ||
-- Use of this source code is governed by the Apache 2.0 license; see COPYING. | ||
|
||
module(...,package.seeall) | ||
|
||
local Transmitter = require("apps.interlink.transmitter") | ||
local Source = require("apps.basic.basic_apps").Source | ||
|
||
function start (name) | ||
local c = config.new() | ||
config.app(c, name, Transmitter) | ||
config.app(c, "source", Source) | ||
config.link(c, "source.output -> "..name..".input") | ||
engine.configure(c) | ||
engine.main() | ||
end |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,49 @@ | ||
-- Use of this source code is governed by the Apache 2.0 license; see COPYING. | ||
|
||
module(...,package.seeall) | ||
|
||
local shm = require("core.shm") | ||
local interlink = require("lib.interlink") | ||
|
||
local Transmitter = {name="apps.interlink.Transmitter"} | ||
|
||
function Transmitter:new (_, name) | ||
local self = {} | ||
self.shm_name = "group/interlink/"..name..".interlink" | ||
self.backlink = "interlink/transmitter/"..name..".interlink" | ||
self.interlink = interlink.attach_transmitter(self.shm_name) | ||
shm.alias(self.backlink, self.shm_name) | ||
return setmetatable(self, {__index=Transmitter}) | ||
end | ||
|
||
function Transmitter:push () | ||
local i, r = self.input.input, self.interlink | ||
while not (interlink.full(r) or link.empty(i)) do | ||
local p = link.receive(i) | ||
packet.account_free(p) -- stimulate breathing | ||
interlink.insert(r, p) | ||
end | ||
interlink.push(r) | ||
end | ||
|
||
function Transmitter:stop () | ||
interlink.detach_transmitter(self.interlink, self.shm_name) | ||
shm.unlink(self.backlink) | ||
end | ||
|
||
-- Detach transmitters to prevent leaking interlinks opened by pid. | ||
-- | ||
-- This is an internal API function provided for cleanup during | ||
-- process termination. | ||
function Transmitter.shutdown (pid) | ||
for _, name in ipairs(shm.children("/"..pid.."/interlink/transmitter")) do | ||
local backlink = "/"..pid.."/interlink/transmitter/"..name..".interlink" | ||
local shm_name = "/"..pid.."/group/interlink/"..name..".interlink" | ||
-- Call protected in case /<pid>/group is already unlinked. | ||
local ok, r = pcall(interlink.open, shm_name) | ||
if ok then interlink.detach_transmitter(r, shm_name) end | ||
shm.unlink(backlink) | ||
end | ||
end | ||
|
||
return Transmitter |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.