Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
330f95a375 | ||
|
|
e744e62d03 |
@@ -10,6 +10,7 @@ import (
|
|||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
"net"
|
||||||
"os"
|
"os"
|
||||||
"os/exec"
|
"os/exec"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
@@ -96,18 +97,16 @@ func runServe(args []string) error {
|
|||||||
defer func() { _ = closeUsers() }()
|
defer func() { _ = closeUsers() }()
|
||||||
|
|
||||||
// 5. Local worker agents before the server, so they can claim immediately.
|
// 5. Local worker agents before the server, so they can claim immediately.
|
||||||
coordinatorURL := "http://" + *addr
|
// They always dial the loopback address: an --addr of 0.0.0.0 is not a
|
||||||
agents, err := spawnAgents(ctx, log, *dataDir, *workers, coordinatorURL, workerToken, venvPython)
|
// connectable target from the same host.
|
||||||
|
agentURL, resolvedPublic := serveURLs(*addr, *publicURL)
|
||||||
|
agents, err := spawnAgents(ctx, log, *dataDir, *workers, agentURL, workerToken, venvPython)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
defer stopAgents(agents)
|
defer stopAgents(agents)
|
||||||
|
|
||||||
// 6. The coordinator server itself.
|
// 6. The coordinator server itself.
|
||||||
coordinatorPublicURL := *publicURL
|
|
||||||
if coordinatorPublicURL == "" {
|
|
||||||
coordinatorPublicURL = "http://" + *addr
|
|
||||||
}
|
|
||||||
cfg := infra.Config{
|
cfg := infra.Config{
|
||||||
Addr: *addr,
|
Addr: *addr,
|
||||||
DatabaseEngine: "sqlite",
|
DatabaseEngine: "sqlite",
|
||||||
@@ -115,8 +114,10 @@ func runServe(args []string) error {
|
|||||||
Token: workerToken,
|
Token: workerToken,
|
||||||
JWTSecret: jwtSecret,
|
JWTSecret: jwtSecret,
|
||||||
UserserviceURL: "http://" + usersAddr,
|
UserserviceURL: "http://" + usersAddr,
|
||||||
PublicCoordinatorURL: coordinatorPublicURL,
|
PublicCoordinatorURL: resolvedPublic,
|
||||||
PublicUserserviceURL: "http://" + usersAddr,
|
// The exchange is fronted by the coordinator's own proxy, so the UI
|
||||||
|
// falls back to the coordinator origin for USERSERVICE_URL.
|
||||||
|
PublicUserserviceURL: "",
|
||||||
LogLevel: "info",
|
LogLevel: "info",
|
||||||
StorageDir: filepath.Join(*dataDir, "artifacts"),
|
StorageDir: filepath.Join(*dataDir, "artifacts"),
|
||||||
DocsDir: *docsDir,
|
DocsDir: *docsDir,
|
||||||
@@ -132,12 +133,13 @@ func runServe(args []string) error {
|
|||||||
WorkerOfflineAfter: 1 * time.Minute,
|
WorkerOfflineAfter: 1 * time.Minute,
|
||||||
AutoMigrate: true,
|
AutoMigrate: true,
|
||||||
}
|
}
|
||||||
|
browserURL, _ := serveURLs(*addr, *publicURL)
|
||||||
if *open {
|
if *open {
|
||||||
openBrowser("http://" + *addr + "/ui/admin")
|
openBrowser(browserURL + "/ui/admin")
|
||||||
}
|
}
|
||||||
|
|
||||||
// Print the login once the server is about to start.
|
// Print the login once the server is about to start.
|
||||||
fmt.Printf("\nSciMesh is starting at http://%s/ui\n", *addr)
|
fmt.Printf("\nSciMesh is starting at %s/ui\n", browserURL)
|
||||||
fmt.Printf(" admin login: %s / %s\n", *email, *password)
|
fmt.Printf(" admin login: %s / %s\n", *email, *password)
|
||||||
if runtimeStatus(venvPython) {
|
if runtimeStatus(venvPython) {
|
||||||
fmt.Printf(" scientific runtime: ready (%s)\n", venvPython)
|
fmt.Printf(" scientific runtime: ready (%s)\n", venvPython)
|
||||||
@@ -337,3 +339,33 @@ func openBrowser(target string) {
|
|||||||
// #nosec G204 -- target is the local UI URL the operator asked to open.
|
// #nosec G204 -- target is the local UI URL the operator asked to open.
|
||||||
_ = exec.CommandContext(context.Background(), command, target).Start()
|
_ = exec.CommandContext(context.Background(), command, target).Start()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// serveURLs derives the two addresses of a serve instance from the listen
|
||||||
|
// address and the optional --public-url flag:
|
||||||
|
//
|
||||||
|
// - the agent URL is always the loopback form of the port, because spawned
|
||||||
|
// local workers share the host and 0.0.0.0 is not connectable from it;
|
||||||
|
// - the public URL is what browsers and remote workers are told. An explicit
|
||||||
|
// --public-url wins; a listen host that is a real address is used as-is;
|
||||||
|
// a wildcard host (0.0.0.0, ::, or empty) yields an empty public URL, so
|
||||||
|
// the UI falls back to the browser's own origin (the coordinator's LAN
|
||||||
|
// address as the browser sees it).
|
||||||
|
func serveURLs(addr, publicURL string) (agentURL, resolvedPublic string) {
|
||||||
|
host, port, err := net.SplitHostPort(addr)
|
||||||
|
if err != nil {
|
||||||
|
// No port in the listen address: assume the default and treat the
|
||||||
|
// whole string as a host (e.g. a bare wildcard).
|
||||||
|
host, port = addr, "8080"
|
||||||
|
}
|
||||||
|
agentURL = "http://127.0.0.1:" + port
|
||||||
|
if publicURL != "" {
|
||||||
|
return agentURL, publicURL
|
||||||
|
}
|
||||||
|
host = strings.Trim(host, "[]")
|
||||||
|
switch host {
|
||||||
|
case "", "0.0.0.0", "::":
|
||||||
|
return agentURL, ""
|
||||||
|
default:
|
||||||
|
return agentURL, "http://" + addr
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -0,0 +1,22 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import "testing"
|
||||||
|
|
||||||
|
func TestServeURLs(t *testing.T) {
|
||||||
|
cases := []struct {
|
||||||
|
addr, public, agent, resolved string
|
||||||
|
}{
|
||||||
|
{"127.0.0.1:8080", "", "http://127.0.0.1:8080", "http://127.0.0.1:8080"},
|
||||||
|
{"0.0.0.0:8080", "", "http://127.0.0.1:8080", ""},
|
||||||
|
{":8080", "", "http://127.0.0.1:8080", ""},
|
||||||
|
{"::", "", "http://127.0.0.1:8080", ""},
|
||||||
|
{"192.168.1.10:8080", "", "http://127.0.0.1:8080", "http://192.168.1.10:8080"},
|
||||||
|
{"0.0.0.0:8080", "http://cluster.example:8080", "http://127.0.0.1:8080", "http://cluster.example:8080"},
|
||||||
|
}
|
||||||
|
for _, c := range cases {
|
||||||
|
agent, resolved := serveURLs(c.addr, c.public)
|
||||||
|
if agent != c.agent || resolved != c.resolved {
|
||||||
|
t.Errorf("serveURLs(%q, %q) = (%q, %q), want (%q, %q)", c.addr, c.public, agent, resolved, c.agent, c.resolved)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -48,6 +48,12 @@ type Client struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func NewClient(baseURL string, tokens TokenProvider, timeout time.Duration) *Client {
|
func NewClient(baseURL string, tokens TokenProvider, timeout time.Duration) *Client {
|
||||||
|
// Payload transfers get a more generous budget than control calls: a large
|
||||||
|
// shard over a slow link easily outlives the API timeout.
|
||||||
|
transferTimeout := timeout * 4
|
||||||
|
if transferTimeout < 2*time.Minute {
|
||||||
|
transferTimeout = 2 * time.Minute
|
||||||
|
}
|
||||||
return &Client{
|
return &Client{
|
||||||
baseURL: strings.TrimRight(baseURL, "/"),
|
baseURL: strings.TrimRight(baseURL, "/"),
|
||||||
tokens: tokens,
|
tokens: tokens,
|
||||||
@@ -57,7 +63,7 @@ func NewClient(baseURL string, tokens TokenProvider, timeout time.Duration) *Cli
|
|||||||
CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse },
|
CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse },
|
||||||
},
|
},
|
||||||
dlClient: &http.Client{
|
dlClient: &http.Client{
|
||||||
Timeout: timeout,
|
Timeout: transferTimeout,
|
||||||
CheckRedirect: func(req *http.Request, via []*http.Request) error {
|
CheckRedirect: func(req *http.Request, via []*http.Request) error {
|
||||||
if len(via) >= 10 {
|
if len(via) >= 10 {
|
||||||
return fmt.Errorf("too many redirects")
|
return fmt.Errorf("too many redirects")
|
||||||
|
|||||||
@@ -223,3 +223,17 @@ func sha256Of(t *testing.T, value string) string {
|
|||||||
digest := sha256.Sum256([]byte(value))
|
digest := sha256.Sum256([]byte(value))
|
||||||
return fmt.Sprintf("%x", digest)
|
return fmt.Sprintf("%x", digest)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestNewClientTransferTimeoutExceedsAPITimeout(t *testing.T) {
|
||||||
|
c := NewClient("http://coord:8080", &StaticToken{token: "t"}, 30*time.Second)
|
||||||
|
if c.apiClient.Timeout != 30*time.Second {
|
||||||
|
t.Errorf("api timeout = %v, want 30s", c.apiClient.Timeout)
|
||||||
|
}
|
||||||
|
if c.dlClient.Timeout < 2*time.Minute {
|
||||||
|
t.Errorf("transfer timeout = %v, want at least 2m", c.dlClient.Timeout)
|
||||||
|
}
|
||||||
|
short := NewClient("http://coord:8080", &StaticToken{token: "t"}, 3*time.Minute)
|
||||||
|
if short.dlClient.Timeout != 12*time.Minute {
|
||||||
|
t.Errorf("transfer timeout = %v, want 4x the api timeout", short.dlClient.Timeout)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user