Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 0 additions & 3 deletions api/openapi.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,6 @@ components:
- port
- protocol
- service_name
- service_port
properties:
name:
type: string
Expand All @@ -160,8 +159,6 @@ components:
type: string
service_name:
type: string
service_port:
type: integer
selector:
type: object
additionalProperties:
Expand Down
71 changes: 71 additions & 0 deletions apps/druid/core/services/runtime_ports.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
package services

import (
"fmt"

"github.com/highcard-dev/daemon/internal/core/domain"
)

func resolveRuntimePorts(ports []domain.Port, routing []domain.RuntimeRouteAssignment, requireAssignments bool) ([]domain.Port, error) {
resolved := append([]domain.Port(nil), ports...)
for index := range resolved {
port := &resolved[index]
if port.Port != 0 {
continue
}

assignedPort := 0
for _, assignment := range routing {
portName := assignment.PortName
if portName == "" {
portName = assignment.Name
}
if portName != port.Name {
continue
}
if assignment.PublicPort < 1 || assignment.PublicPort > 65535 {
return nil, fmt.Errorf("dynamic port %q has invalid public port %d", port.Name, assignment.PublicPort)
}
if assignedPort != 0 && assignedPort != assignment.PublicPort {
return nil, fmt.Errorf("dynamic port %q has conflicting public ports %d and %d", port.Name, assignedPort, assignment.PublicPort)
}
assignedPort = assignment.PublicPort
}

if assignedPort == 0 {
if requireAssignments {
return nil, fmt.Errorf("dynamic port %q has no public routing assignment", port.Name)
}
continue
}
port.Port = assignedPort
}
return resolved, nil
}

func validateDynamicPortsUnchanged(
ports []domain.Port,
currentRouting []domain.RuntimeRouteAssignment,
nextRouting []domain.RuntimeRouteAssignment,
) error {
currentPorts, err := resolveRuntimePorts(ports, currentRouting, false)
if err != nil {
return err
}
nextPorts, err := resolveRuntimePorts(ports, nextRouting, false)
if err != nil {
return err
}
for index, port := range ports {
if port.Port != 0 || currentPorts[index].Port == nextPorts[index].Port {
continue
}
return fmt.Errorf(
"cannot change dynamic port %q from %d to %d while runtime is running; stop the runtime before applying routing",
port.Name,
currentPorts[index].Port,
nextPorts[index].Port,
)
}
return nil
}
11 changes: 9 additions & 2 deletions apps/druid/core/services/runtime_session_execution.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,14 @@ func (s *RuntimeSession) runCommand(cmd string) error {
root = s.scrollService.GetCwd()
}
file := s.scrollService.GetFile()
procedureEnv, err := coreservices.BuildRuntimeProcedureEnv(file, cmd, command, coreservices.RuntimeEnvContext{
runtimePorts, err := resolveRuntimePorts(file.Ports, routing, true)
if err != nil {
s.setCommandProcedureStatus(cmd, command, domain.ScrollLockStatusError, nil)
return err
}
runtimeFile := *file
runtimeFile.Ports = runtimePorts
procedureEnv, err := coreservices.BuildRuntimeProcedureEnv(&runtimeFile, cmd, command, coreservices.RuntimeEnvContext{
ScrollID: scrollID,
ScrollName: scrollName,
Backend: s.runtimeBackend.Name(),
Expand All @@ -49,7 +56,7 @@ func (s *RuntimeSession) runCommand(cmd string) error {
ScrollID: scrollID,
Command: command,
Root: root,
GlobalPorts: file.Ports,
GlobalPorts: runtimePorts,
Routing: routing,
ProcedureEnv: procedureEnv,
ProcedureStatusObserver: func(procedure string, status domain.ScrollLockStatus, exitCode *int) {
Expand Down
82 changes: 82 additions & 0 deletions apps/druid/core/services/runtime_session_execution_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package services

import (
"errors"
"strings"
"testing"

"github.com/highcard-dev/daemon/internal/core/domain"
Expand Down Expand Up @@ -62,6 +63,9 @@ func TestRuntimeSessionRunCommandPassesRoutingAndScrollIdentity(t *testing.T) {
if len(seen.Routing) != 1 || seen.Routing[0].PublicPort != 443 {
t.Fatalf("Routing = %#v", seen.Routing)
}
if len(seen.GlobalPorts) != 1 || seen.GlobalPorts[0].Port != 8080 {
t.Fatalf("GlobalPorts = %#v, fixed port must remain 8080", seen.GlobalPorts)
}
env := seen.ProcedureEnv["web"]
if env["DRUID_SCROLL_ID"] != "scroll-a" || env["DRUID_SCROLL_NAME"] != "scroll-name" {
t.Fatalf("scroll env = %#v", env)
Expand All @@ -80,6 +84,64 @@ func TestRuntimeSessionRunCommandPassesRoutingAndScrollIdentity(t *testing.T) {
}
}

func TestRuntimeSessionRunCommandResolvesDynamicPortsFromRouting(t *testing.T) {
var seen ports.RuntimeCommand
session := newRuntimeSessionExecutionTest(t, dynamicExecutionScrollYAML(), &fakeWorkerBackend{
runCommand: func(command ports.RuntimeCommand) (*int, error) {
seen = command
return nil, nil
},
})
session.runtimeScroll.Routing = []domain.RuntimeRouteAssignment{{
Name: "main",
PortName: "main",
ExternalIP: "127.0.0.1",
PublicPort: 11000,
Protocol: "udp",
}}

if err := session.runCommand("serve"); err != nil {
t.Fatal(err)
}

if len(seen.GlobalPorts) != 1 || seen.GlobalPorts[0].Port != 11000 {
t.Fatalf("GlobalPorts = %#v, want dynamic port 11000", seen.GlobalPorts)
}
env := seen.ProcedureEnv["server"]
if env["DRUID_PORT_MAIN"] != "11000" || env["DRUID_PORT_MAIN_1"] != "11000" {
t.Fatalf("runtime env = %#v, want internal port 11000", env)
}
if env["DRUID_PORT_MAIN_PUBLIC"] != "11000" {
t.Fatalf("runtime env = %#v, want public port 11000", env)
}
}

func TestRuntimeSessionRunCommandRejectsUnassignedDynamicPort(t *testing.T) {
called := false
session := newRuntimeSessionExecutionTest(t, dynamicExecutionScrollYAML(), &fakeWorkerBackend{
runCommand: func(command ports.RuntimeCommand) (*int, error) {
called = true
return nil, nil
},
})

err := session.runCommand("serve")
if err == nil || !strings.Contains(err.Error(), `dynamic port "main" has no public routing assignment`) {
t.Fatalf("runCommand error = %v", err)
}
if called {
t.Fatal("runtime backend was called without a dynamic port assignment")
}

updated, storeErr := session.store.GetScroll(session.runtimeScroll.ID)
if storeErr != nil {
t.Fatal(storeErr)
}
if updated.Procedures["serve"]["server"].Status != domain.ScrollLockStatusError {
t.Fatalf("procedure status = %#v, want error", updated.Procedures["serve"]["server"])
}
}

func TestRuntimeSessionRunCommandPersistsProcedureStatusCallbacks(t *testing.T) {
exitCode := 0
session := newRuntimeSessionExecutionTest(t, executionScrollYAML(), &fakeWorkerBackend{
Expand Down Expand Up @@ -205,3 +267,23 @@ commands:
image: alpine:3.20
`
}

func dynamicExecutionScrollYAML() string {
return `name: scroll-name
desc: Dynamic runtime session execution test
version: 0.1.0
app_version: "1.0"
ports:
- name: main
protocol: udp
serve: serve
commands:
serve:
run: persistent
procedures:
- id: server
image: alpine:3.20
expectedPorts:
- name: main
`
}
25 changes: 24 additions & 1 deletion apps/druid/core/services/runtime_session_runtime.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,14 @@ import (
func (s *RuntimeSession) Ports() ([]domain.RuntimePortStatus, error) {
s.mu.Lock()
runtimeScroll := *s.runtimeScroll
routing := append([]domain.RuntimeRouteAssignment(nil), s.runtimeScroll.Routing...)
s.mu.Unlock()
return s.runtimeBackend.ExpectedPorts(runtimeScroll.Root, s.scrollService.GetFile().Commands, s.scrollService.GetFile().Ports)
file := s.scrollService.GetFile()
runtimePorts, err := resolveRuntimePorts(file.Ports, routing, false)
if err != nil {
return nil, err
}
return s.runtimeBackend.ExpectedPorts(runtimeScroll.Root, file.Commands, runtimePorts)
}

func (s *RuntimeSession) RoutingTargets() ([]domain.RuntimeRoutingTarget, error) {
Expand All @@ -31,6 +37,12 @@ func (s *RuntimeSession) Queue() domain.ProcedureStatusMap {

func (s *RuntimeSession) ApplyRouting(assignments []domain.RuntimeRouteAssignment) (*domain.RuntimeScroll, error) {
s.mu.Lock()
if s.runtimeScroll.Status == domain.RuntimeScrollStatusRunning || hasRunningProcedure(s.runtimeScroll.Procedures) {
if err := validateDynamicPortsUnchanged(s.scrollService.GetFile().Ports, s.runtimeScroll.Routing, assignments); err != nil {
s.mu.Unlock()
return nil, err
}
}
s.runtimeScroll.Routing = assignments
s.runtimeScroll.LastError = ""
err := s.store.UpdateScroll(s.runtimeScroll)
Expand All @@ -42,6 +54,17 @@ func (s *RuntimeSession) ApplyRouting(assignments []domain.RuntimeRouteAssignmen
return s.store.GetScroll(id)
}

func hasRunningProcedure(procedures domain.ProcedureStatusMap) bool {
for _, commandProcedures := range procedures {
for _, procedure := range commandProcedures {
if procedure.Status == domain.ScrollLockStatusRunning {
return true
}
}
}
return false
}

func (s *RuntimeSession) StopRuntime() error {
s.mu.Lock()
root := s.runtimeScroll.Root
Expand Down
42 changes: 42 additions & 0 deletions apps/druid/core/services/runtime_supervisor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -958,6 +958,48 @@ func TestRuntimeSessionApplyRoutingPersistsAssignments(t *testing.T) {
}
}

func TestRuntimeSessionApplyRoutingRejectsRunningDynamicPortRemap(t *testing.T) {
session := newRuntimeSessionForTest(t, map[string]domain.LockStatus{}, dynamicExecutionScrollYAML())
session.runtimeScroll.Status = domain.RuntimeScrollStatusRunning
session.runtimeScroll.Procedures = domain.ProcedureStatusMap{
"serve": {"server": {Status: domain.ScrollLockStatusRunning}},
}
session.runtimeScroll.Routing = []domain.RuntimeRouteAssignment{{
Name: "main",
PortName: "main",
PublicPort: 11000,
Protocol: "udp",
}}

_, err := session.ApplyRouting([]domain.RuntimeRouteAssignment{{
Name: "main",
PortName: "main",
PublicPort: 11001,
Protocol: "udp",
}})
if err == nil || !strings.Contains(err.Error(), `cannot change dynamic port "main" from 11000 to 11001 while runtime is running`) {
t.Fatalf("ApplyRouting error = %v", err)
}
if session.runtimeScroll.Routing[0].PublicPort != 11000 {
t.Fatalf("routing changed after rejection: %#v", session.runtimeScroll.Routing)
}

session.runtimeScroll.Status = domain.RuntimeScrollStatusStopped
session.runtimeScroll.Procedures = nil
updated, err := session.ApplyRouting([]domain.RuntimeRouteAssignment{{
Name: "main",
PortName: "main",
PublicPort: 11001,
Protocol: "udp",
}})
if err != nil {
t.Fatal(err)
}
if updated.Routing[0].PublicPort != 11001 {
t.Fatalf("routing = %#v, want stopped runtime remapped to 11001", updated.Routing)
}
}

func TestRuntimeSessionQueueReturnsProcedureStatuses(t *testing.T) {
session := newRuntimeSessionForTest(t, map[string]domain.LockStatus{}, cachedScrollYAML(""))
session.RememberDoneItem("start")
Expand Down
57 changes: 28 additions & 29 deletions internal/api/generated.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 0 additions & 1 deletion internal/core/domain/runtime_scroll.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,6 @@ type RuntimeRoutingTarget struct {
Protocol string `json:"protocol"`
Namespace string `json:"namespace,omitempty"`
ServiceName string `json:"service_name"`
ServicePort int `json:"service_port"`
Selector map[string]string `json:"selector,omitempty"`
}

Expand Down
Loading
Loading