-
Notifications
You must be signed in to change notification settings - Fork 160
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
runproc sample #21
Open
mcqueenorama
wants to merge
16
commits into
gocircuit:master
Choose a base branch
from
mcqueenorama:master
base: master
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
runproc sample #21
Changes from all commits
Commits
Show all changes
16 commits
Select commit
Hold shift + click to select a range
4377a7a
runproc sample
a9bb995
added a flag to prefix output lines with the anchor so it can be grep…
9c8903b
gofmt
53ec5f4
no longer require the user to specify a name for the process
cb0ea06
add a Name field to the input json so the user can specify the name o…
ca48604
make the runproc io go async so its much faster
0eb107b
change "anchors" flag to "tag" and remove scrub flag since its not ne…
a006468
get runproc ready to get the anchors itself
b9b4e09
get the standalone runproc check going
d8bd8ef
reset import paths for pull request
049ebd7
make the --all commands go async for speed
b5c5b7a
remove new Name field according to discussion and put all in runproc/…
6e1fa05
missed an instance of Name
ba9f674
fix up the prefixed writing and remove the Scrub since they are autos…
7338ad0
redo the waiting use p.Wait()
af9af76
get stderr of the process too
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
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
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
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -8,13 +8,16 @@ | |
package main | ||
|
||
import ( | ||
"bufio" | ||
"encoding/json" | ||
"fmt" | ||
"io" | ||
"io/ioutil" | ||
"os" | ||
|
||
"github.com/gocircuit/circuit/client" | ||
"github.com/gocircuit/circuit/client/docker" | ||
"github.com/gocircuit/circuit/kit/iomisc" | ||
|
||
"github.com/gocircuit/circuit/github.com/codegangsta/cli" | ||
) | ||
|
@@ -53,6 +56,112 @@ func mkproc(x *cli.Context) { | |
} | ||
} | ||
|
||
func doRun(x *cli.Context, c *client.Client, cmd client.Cmd, path string, done chan bool) { | ||
|
||
w2, _ := parseGlob(path) | ||
a2 := c.Walk(w2) | ||
_runproc(x, c, a2, cmd, done) | ||
|
||
} | ||
|
||
func runproc(x *cli.Context) { | ||
defer func() { | ||
if r := recover(); r != nil { | ||
fatalf("error, likely due to missing server or misspelled anchor: %v", r) | ||
} | ||
}() | ||
c := dial(x) | ||
args := x.Args() | ||
|
||
if len(args) != 1 && !x.Bool("all") { | ||
fatalf("runproc needs an anchor argument or use the --all flag to to execute on every host in the circuit") | ||
} | ||
buf, _ := ioutil.ReadAll(os.Stdin) | ||
var cmd client.Cmd | ||
if err := json.Unmarshal(buf, &cmd); err != nil { | ||
fatalf("command json not parsing: %v", err) | ||
} | ||
cmd.Scrub = true | ||
|
||
el := "/runproc/" + keygen(x) | ||
|
||
done := make(chan bool, 10) | ||
if x.Bool("all") { | ||
|
||
w, _ := parseGlob("/") | ||
|
||
anchor := c.Walk(w) | ||
|
||
procs := 0 | ||
|
||
for _, a := range anchor.View() { | ||
|
||
procs++ | ||
|
||
go func(x *cli.Context, cmd client.Cmd, a string, done chan bool) { | ||
|
||
doRun(x, c, cmd, a, done) | ||
|
||
}(x, cmd, a.Path()+el, done) | ||
|
||
} | ||
|
||
for ; procs > 0 ; procs-- { | ||
|
||
select { | ||
case <-done: | ||
continue | ||
} | ||
|
||
} | ||
|
||
} else { | ||
|
||
doRun(x, c, cmd, args[0]+el, done) | ||
<-done | ||
|
||
} | ||
|
||
} | ||
|
||
func _runproc(x *cli.Context, c *client.Client, a client.Anchor, cmd client.Cmd, done chan bool) { | ||
|
||
p, err := a.MakeProc(cmd) | ||
if err != nil { | ||
fatalf("mkproc error: %s", err) | ||
} | ||
|
||
stdin := p.Stdin() | ||
if err := stdin.Close(); err != nil { | ||
fatalf("error closing stdin: %v", err) | ||
} | ||
|
||
if x.Bool("tag") { | ||
|
||
stdout := iomisc.PrefixReader(a.Addr() + " ", p.Stdout()) | ||
stderr := iomisc.PrefixReader(a.Addr() + " ", p.Stderr()) | ||
|
||
stdoutScanner := bufio.NewScanner(stdout) | ||
for stdoutScanner.Scan() { | ||
fmt.Println(stdoutScanner.Text()) | ||
} | ||
|
||
stderrScanner := bufio.NewScanner(stderr) | ||
for stderrScanner.Scan() { | ||
fmt.Println(stderrScanner.Text()) | ||
} | ||
|
||
} else { | ||
|
||
io.Copy(os.Stdout, p.Stdout()) | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. What about stderr? |
||
io.Copy(os.Stderr, p.Stderr()) | ||
|
||
} | ||
p.Wait() | ||
done <- true | ||
|
||
} | ||
|
||
func mkdkr(x *cli.Context) { | ||
defer func() { | ||
if r := recover(); r != nil { | ||
|
@@ -92,7 +201,9 @@ func sgnl(x *cli.Context) { | |
fatalf("signal needs an anchor and a signal name arguments") | ||
} | ||
w, _ := parseGlob(args[1]) | ||
u, ok := c.Walk(w).Get().(interface{Signal(string) error}) | ||
u, ok := c.Walk(w).Get().(interface { | ||
Signal(string) error | ||
}) | ||
if !ok { | ||
fatalf("anchor is not a process or a docker container") | ||
} | ||
|
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Use PrefixReader or PrefixWriter from package kit/iomisc instead.