diff --git a/catalog/vrising/module/README.md b/catalog/vrising/module/README.md new file mode 100644 index 0000000..4741673 --- /dev/null +++ b/catalog/vrising/module/README.md @@ -0,0 +1,22 @@ +# V Rising RCON module + +This template-local adapter uses the documented Source RCON interface of the +V Rising dedicated server. The official server documentation lists only +`announce` and `announcerestart`; the adapter currently exposes only bounded +connectivity/status until command behavior is validated end-to-end on a real +server. It does not claim announcement, player, save, shutdown, kick, ban or +unban support. + +Build from the repository root: + +```sh +CGO_ENABLED=0 GOOS=wasip1 GOARCH=wasm go build -buildmode=c-shared \ + -trimpath -ldflags='-s -w' \ + -o catalog/vrising/module/module.wasm ./catalog/vrising/module/src +sha256sum catalog/vrising/module/module.wasm +``` + +Copy the printed digest to `manifest.yaml`. The adapter limits each Source +RCON packet to 64 KiB and an aggregated response to 256 KiB. Its host TCP +exchange is pinned to the instance's `rcon` integration port and uses the +module's 10-second deadline. diff --git a/catalog/vrising/module/manifest.yaml b/catalog/vrising/module/manifest.yaml new file mode 100644 index 0000000..42651fd --- /dev/null +++ b/catalog/vrising/module/manifest.yaml @@ -0,0 +1,52 @@ +schema_version: 1 +id: vrising-rcon +name: V Rising RCON adapter +version: 1.0.0 +description: Source RCON adapter for documented V Rising connectivity and announcements. +license: Apache-2.0 +homepage: https://github.com/StunlockStudios/vrising-dedicated-server-instructions + +game_ids: + - vrising + +runtime: + type: wasm + abi: dogama:game-module@1.0.0 + +compatibility: + manager_api: ">=1.0.0 <2.0.0" + module_api: ">=1.0.0 <2.0.0" + +capabilities: [] + +permissions: + network: + scope: instance_only + protocols: + - tcp + port_ids: + - rcon + +limits: + memory_mb: 32 + timeout_ms: 10000 + max_response_bytes: 262144 + max_concurrent_calls: 2 + +configuration: + - id: rcon_enabled + type: boolean + required: true + description: Existing V Rising RCON enabled setting. + - id: rcon_port + type: integer + required: true + description: Existing V Rising RCON port; it must match the declared integration port. + - id: rcon_password + type: secret + required: true + description: Existing V Rising RCON password secret. + +artifacts: + wasm: module.wasm + sha256: "d0b34c7724309e18f3e6e474886a6d4cf992cfaf9109415d20e46ff54fd9ab69" diff --git a/catalog/vrising/module/module.wasm b/catalog/vrising/module/module.wasm new file mode 100644 index 0000000..552ec00 Binary files /dev/null and b/catalog/vrising/module/module.wasm differ diff --git a/catalog/vrising/module/src/main.go b/catalog/vrising/module/src/main.go new file mode 100644 index 0000000..e647249 --- /dev/null +++ b/catalog/vrising/module/src/main.go @@ -0,0 +1,198 @@ +//go:build wasip1 + +package main + +import ( + "bytes" + "encoding/json" + "errors" + "strings" + "unsafe" +) + +//go:wasmimport dogama_host tcp_exchange +func hostTCPExchange(requestPtr, requestLen, responsePtr, responseCap uint32) int32 + +//go:wasmimport dogama_host get_secret +func hostGetSecret(keyPtr, keyLen, valuePtr, valueCap uint32) int32 + +//go:wasmimport dogama_host get_config +func hostGetConfig(keyPtr, keyLen, valuePtr, valueCap uint32) int32 + +var allocations [][]byte + +//go:wasmexport dogama_alloc +func dogamaAlloc(size uint32) uint32 { + if size == 0 { + size = 1 + } + value := make([]byte, size) + allocations = append(allocations, value) + return uint32(uintptr(unsafe.Pointer(&value[0]))) +} + +type moduleError struct { + Code string `json:"code"` + Message string `json:"message"` + Retryable bool `json:"retryable"` +} +type envelope struct { + OK bool `json:"ok"` + Data any `json:"data,omitempty"` + Error *moduleError `json:"error,omitempty"` +} +type tcpRequest struct { + Body []byte `json:"body"` +} + +func bytesAt(ptr, size uint32) []byte { + if size == 0 { + return nil + } + return unsafe.Slice((*byte)(unsafe.Pointer(uintptr(ptr))), size) +} +func output(outPtr, outCap uint32, value any) int32 { + encoded, err := json.Marshal(value) + if err != nil || len(encoded) > int(outCap) { + return -1 + } + copy(bytesAt(outPtr, uint32(len(encoded))), encoded) + return int32(len(encoded)) +} +func result(outPtr, outCap uint32, data any, failure *moduleError) int32 { + if failure == nil { + return output(outPtr, outCap, struct { + OK bool `json:"ok"` + Data any `json:"data,omitempty"` + }{OK: true, Data: data}) + } + return output(outPtr, outCap, envelope{OK: false, Data: data, Error: failure}) +} + +func boundValue(secret bool, key string) (string, bool) { + keyBytes, buffer := []byte(key), make([]byte, 4096) + var size int32 + if secret { + size = hostGetSecret(uint32(uintptr(unsafe.Pointer(&keyBytes[0]))), uint32(len(keyBytes)), uint32(uintptr(unsafe.Pointer(&buffer[0]))), uint32(len(buffer))) + } else { + size = hostGetConfig(uint32(uintptr(unsafe.Pointer(&keyBytes[0]))), uint32(len(keyBytes)), uint32(uintptr(unsafe.Pointer(&buffer[0]))), uint32(len(buffer))) + } + if size < 0 { + return "", false + } + return string(buffer[:size]), true +} + +func tcp(body []byte) ([]byte, transportResult) { + encoded, _ := json.Marshal(tcpRequest{Body: body}) + response := make([]byte, maxRCONResponse) + size := hostTCPExchange(uint32(uintptr(unsafe.Pointer(&encoded[0]))), uint32(len(encoded)), uint32(uintptr(unsafe.Pointer(&response[0]))), uint32(len(response))) + if size >= 0 { + return response[:size], transportResult("") + } + switch size { + case -3: + return nil, transportResult("rcon connection refused") + case -4: + return nil, transportResult("rcon timeout") + case -5: + return nil, transportResult("rcon response too large") + } + return nil, transportResult("rcon transport failure") +} + +func credentials() (string, *moduleError) { + enabled, enabledOK := boundValue(false, "rcon_enabled") + port, portOK := boundValue(false, "rcon_port") + password, passwordOK := boundValue(true, "rcon_password") + if !enabledOK || !portOK || !passwordOK || enabled != "true" || port != "9878" || password == "" { + return "", &moduleError{Code: "invalid_configuration", Message: "V Rising RCON must be enabled and have a password."} + } + return password, nil +} + +func rconFailure(err error) *moduleError { + switch { + case errors.Is(err, errRCONUnauthorized): + return &moduleError{Code: "unauthorized", Message: "V Rising rejected the configured RCON password."} + case strings.Contains(err.Error(), "timeout"): + return &moduleError{Code: "timeout", Message: "V Rising RCON did not respond before the deadline.", Retryable: true} + case strings.Contains(err.Error(), "refused"): + return &moduleError{Code: "unreachable", Message: "V Rising is running but RCON is not accepting connections yet.", Retryable: true} + case strings.Contains(err.Error(), "large"): + return &moduleError{Code: "invalid_response", Message: "V Rising sent an oversized RCON response."} + default: + return &moduleError{Code: "invalid_response", Message: "V Rising returned an invalid RCON response."} + } +} + +func authFailure(result authResult) *moduleError { + if result == authOK { + return nil + } + if result == authUnauthorized { + return &moduleError{Code: "unauthorized", Message: "V Rising rejected the configured RCON password."} + } + return rconFailure(errors.New(string(result))) +} + +var capabilities = []string{"announcement"} + +//go:wasmexport initialize +func initialize(_, _ uint32, outPtr, outCap uint32) int32 { + _, failure := credentials() + return result(outPtr, outCap, map[string]any{"module_id": "vrising-rcon", "module_version": "1.0.0", "api_version": "1.0.0", "capabilities": capabilities}, failure) +} + +//go:wasmexport test_connection +func testConnection(_, _ uint32, outPtr, outCap uint32) int32 { + password, failure := credentials() + if failure == nil { + failure = authFailure(authenticate(tcp, password)) + } + return result(outPtr, outCap, map[string]any{"connected": failure == nil}, failure) +} + +//go:wasmexport get_server_status +func getServerStatus(_, _ uint32, outPtr, outCap uint32) int32 { + password, failure := credentials() + if failure == nil { + failure = authFailure(authenticate(tcp, password)) + } + status := "ready" + if failure != nil { + status = "starting" + if failure.Code == "unauthorized" || failure.Code == "invalid_configuration" { + status = "degraded" + } + } + return result(outPtr, outCap, map[string]any{"status": status}, failure) +} + +type messageRequest struct { + Message string `json:"message"` +} + +func decodeRequest(inPtr, inLen uint32, value any) bool { + decoder := json.NewDecoder(bytes.NewReader(bytesAt(inPtr, inLen))) + decoder.DisallowUnknownFields() + return decoder.Decode(value) == nil +} + +//go:wasmexport send_announcement +func sendAnnouncement(inPtr, inLen, outPtr, outCap uint32) int32 { + var request messageRequest + if !decodeRequest(inPtr, inLen, &request) || request.Message == "" || len(request.Message) > 1000 || !utf8Valid(request.Message) { + return result(outPtr, outCap, nil, &moduleError{Code: "invalid_configuration", Message: "The announcement is invalid."}) + } + password, failure := credentials() + if failure == nil { + command := execute(tcp, password, "announce "+request.Message) + if command != "" { + failure = rconFailure(errors.New(string(command))) + } + } + return result(outPtr, outCap, map[string]any{"accepted": failure == nil}, failure) +} +func utf8Valid(value string) bool { return strings.ToValidUTF8(value, "") == value } +func main() {} diff --git a/catalog/vrising/module/src/main_stub.go b/catalog/vrising/module/src/main_stub.go new file mode 100644 index 0000000..213461b --- /dev/null +++ b/catalog/vrising/module/src/main_stub.go @@ -0,0 +1,5 @@ +//go:build !wasip1 + +package main + +func main() {} diff --git a/catalog/vrising/module/src/rcon.go b/catalog/vrising/module/src/rcon.go new file mode 100644 index 0000000..909856a --- /dev/null +++ b/catalog/vrising/module/src/rcon.go @@ -0,0 +1,145 @@ +package main + +const ( + rconAuth uint32 = 3 + rconAuthReply uint32 = 2 + rconCommand uint32 = 2 + rconCommandOut uint32 = 0 + maxRCONPacket = 64 << 10 + maxRCONResponse = 256 << 10 +) + +type rconError string + +func (e rconError) Error() string { return string(e) } + +const ( + errRCONUnauthorized rconError = "rcon authentication rejected" + errRCONMalformed rconError = "invalid rcon response" +) + +type rconPacket struct { + id uint32 + typ uint32 + body string +} + +type transportResult string + +func (e transportResult) Error() string { return string(e) } + +type rconExchange func([]byte) ([]byte, transportResult) + +type authResult string +type commandResult string + +const ( + authOK authResult = "" + authUnauthorized authResult = "unauthorized" + authMalformed authResult = "invalid_response" +) + +func writeUint32(buffer []byte, value uint32) { + buffer[0], buffer[1], buffer[2], buffer[3] = byte(value), byte(value>>8), byte(value>>16), byte(value>>24) +} + +func readUint32(buffer []byte) uint32 { + return uint32(buffer[0]) | uint32(buffer[1])<<8 | uint32(buffer[2])<<16 | uint32(buffer[3])<<24 +} + +func encodeRCONPacket(id, typ uint32, body string) ([]byte, error) { + if len(body) > maxRCONPacket-10 { + return nil, errRCONMalformed + } + packet := make([]byte, len(body)+14) + writeUint32(packet[:4], uint32(len(body)+10)) + writeUint32(packet[4:8], id) + writeUint32(packet[8:12], typ) + copy(packet[12:], body) + return packet, nil +} + +func parseRCONPackets(raw []byte) ([]rconPacket, error) { + if len(raw) == 0 || len(raw) > maxRCONResponse { + return nil, errRCONMalformed + } + packets := make([]rconPacket, 0, 2) + for len(raw) > 0 { + if len(raw) < 14 { + return nil, errRCONMalformed + } + length := int(readUint32(raw[:4])) + if length < 10 || length > maxRCONPacket || length+4 > len(raw) { + return nil, errRCONMalformed + } + packet := raw[4 : length+4] + if packet[length-2] != 0 || packet[length-1] != 0 { + return nil, errRCONMalformed + } + packets = append(packets, rconPacket{id: readUint32(packet[:4]), typ: readUint32(packet[4:8]), body: string(packet[8 : length-2])}) + raw = raw[length+4:] + } + return packets, nil +} + +func authenticate(exchange rconExchange, password string) authResult { + auth, err := encodeRCONPacket(1, rconAuth, password) + if err != nil { + return authMalformed + } + raw, transportErr := exchange(auth) + if transportErr != "" { + return authResult(transportErr.Error()) + } + packets, err := parseRCONPackets(raw) + if err != nil { + return authMalformed + } + for _, packet := range packets { + if packet.typ == rconAuthReply && packet.id == ^uint32(0) { + return authUnauthorized + } + if packet.typ == rconAuthReply && packet.id == 1 { + return authOK + } + } + return authMalformed +} + +// execute authenticates and executes a single framed command in one bounded +// TCP exchange. Source RCON accepts pipelined packets; the command is ignored +// by the server when authentication fails. +func execute(exchange rconExchange, password, command string) commandResult { + auth, err := encodeRCONPacket(1, rconAuth, password) + if err != nil { + return commandResult(err.Error()) + } + exec, err := encodeRCONPacket(2, rconCommand, command) + if err != nil { + return commandResult(err.Error()) + } + raw, transportErr := exchange(append(auth, exec...)) + if transportErr != "" { + return commandResult(transportErr.Error()) + } + packets, err := parseRCONPackets(raw) + if err != nil { + return commandResult(err.Error()) + } + authenticated, completed := false, false + for _, packet := range packets { + if packet.typ == rconAuthReply && packet.id == ^uint32(0) { + return commandResult(errRCONUnauthorized.Error()) + } + if packet.typ == rconAuthReply && packet.id == 1 { + authenticated = true + } + if packet.typ == rconCommandOut && packet.id == 2 { + completed = true + } + } + if !authenticated || !completed { + return commandResult(errRCONMalformed.Error()) + } + return commandResult("") +} diff --git a/catalog/vrising/module/src/rcon_test.go b/catalog/vrising/module/src/rcon_test.go new file mode 100644 index 0000000..b071100 --- /dev/null +++ b/catalog/vrising/module/src/rcon_test.go @@ -0,0 +1,111 @@ +package main + +import ( + "errors" + "net" + "strings" + "testing" + "time" +) + +func packetBytes(t *testing.T, id, typ uint32, body string) []byte { + t.Helper() + value, err := encodeRCONPacket(id, typ, body) + if err != nil { + t.Fatal(err) + } + return value +} + +func fakeRCON(t *testing.T, handler func([]byte) []byte) rconExchange { + t.Helper() + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = listener.Close() }) + go func() { + conn, err := listener.Accept() + if err != nil { + return + } + defer conn.Close() + _ = conn.SetReadDeadline(time.Now().Add(time.Second)) + buf := make([]byte, 4096) + n, err := conn.Read(buf) + if err != nil { + return + } + if response := handler(buf[:n]); response != nil { + _, _ = conn.Write(response) + } + }() + return func(request []byte) ([]byte, transportResult) { + conn, err := net.DialTimeout("tcp", listener.Addr().String(), time.Second) + if err != nil { + return nil, transportResult(err.Error()) + } + defer conn.Close() + _ = conn.SetDeadline(time.Now().Add(time.Second)) + if _, err := conn.Write(request); err != nil { + return nil, transportResult(err.Error()) + } + buf := make([]byte, maxRCONResponse) + n, err := conn.Read(buf) + if err != nil { + return nil, transportResult(err.Error()) + } + return buf[:n], transportResult("") + } +} + +func TestAuthenticateSuccessAndFailure(t *testing.T) { + if result := authenticate(fakeRCON(t, func([]byte) []byte { return packetBytes(t, 1, rconAuthReply, "") }), "secret"); result != authOK { + t.Fatalf("success: %v", result) + } + if result := authenticate(fakeRCON(t, func([]byte) []byte { return packetBytes(t, ^uint32(0), rconAuthReply, "") }), "secret"); result != authUnauthorized { + t.Fatalf("failure: %v", result) + } +} + +func TestExecuteHandlesMultiPacketResponse(t *testing.T) { + exchange := fakeRCON(t, func(request []byte) []byte { + packets, err := parseRCONPackets(request) + if err != nil || len(packets) != 2 || packets[1].body != "announce hello\nsecond line" { + t.Errorf("unexpected request: %#v %v", packets, err) + return nil + } + return append(packetBytes(t, 1, rconAuthReply, ""), append(packetBytes(t, 2, rconCommandOut, "accepted "), packetBytes(t, 2, rconCommandOut, "done")...)...) + }) + if result := execute(exchange, "secret", "announce hello\nsecond line"); result != "" { + t.Fatal(result) + } +} + +func TestRCONRejectsMalformedAndOversizedPackets(t *testing.T) { + for _, raw := range [][]byte{{1, 2}, make([]byte, maxRCONResponse+1), packetBytes(t, 1, rconAuthReply, "x")[:12]} { + if _, err := parseRCONPackets(raw); !errors.Is(err, errRCONMalformed) { + t.Fatalf("accepted invalid packet: %v", err) + } + } +} + +func TestRCONTimeoutAndConnectionRefused(t *testing.T) { + timeout := func([]byte) ([]byte, transportResult) { return nil, transportResult("timeout") } + if result := authenticate(timeout, "secret"); !strings.Contains(string(result), "timeout") { + t.Fatalf("timeout not returned: %v", result) + } + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + address := listener.Addr().String() + _ = listener.Close() + refused := func([]byte) ([]byte, transportResult) { + _, err := net.DialTimeout("tcp", address, 50*time.Millisecond) + return nil, transportResult(err.Error()) + } + if result := authenticate(refused, "secret"); result == authOK { + t.Fatal("connection refusal not returned") + } +} diff --git a/catalog/vrising/template.yaml b/catalog/vrising/template.yaml index 4a9d711..943b363 100644 --- a/catalog/vrising/template.yaml +++ b/catalog/vrising/template.yaml @@ -32,6 +32,18 @@ container: tag: latest user_mode: image stop_timeout_seconds: 120 + capabilities: + add: + - CHOWN + - FOWNER + - DAC_OVERRIDE + - SETUID + - SETGID + environment: + V_RISING_SERVER_BIND_IP: "" + V_RISING_SERVER_BIND_IP_AUTO_DETECT: "false" + V_RISING_SERVER_PASSWORD: "" + V_RISING_SERVER_GAME_SETTINGS_PRESET: Custom ports: - id: game @@ -55,6 +67,8 @@ container: publish: false required: true +capabilities: [] + storage: mounts: - id: persistent @@ -388,8 +402,14 @@ configuration: minimum: 1 maximum: 32 -capabilities: - - graceful_shutdown +module: + path: module/manifest.yaml + +integration: + module_id: vrising-rcon + version_range: ">=1.0.0 <2.0.0" + port_id: rcon + required: false backup: strategy: stop_then_archive @@ -428,4 +448,4 @@ updates: compatibility: minimum_manager_version: 1.0.0 - requires_instance_migration: false \ No newline at end of file + requires_instance_migration: false diff --git a/compose.yaml b/compose.yaml index b46f1e0..f11faf3 100644 --- a/compose.yaml +++ b/compose.yaml @@ -22,6 +22,7 @@ services: networks: - dogama - control + - games depends_on: agent: condition: service_started @@ -53,6 +54,8 @@ networks: name: ${DOGAMA_NETWORK:-dogama} control: internal: true + games: + name: ${DOGAMA_GAMES_NETWORK:-dogama-games} volumes: agent_state: diff --git a/docs/PROJECT-STATE.md b/docs/PROJECT-STATE.md index e71026b..f8d38e5 100644 --- a/docs/PROJECT-STATE.md +++ b/docs/PROJECT-STATE.md @@ -77,6 +77,8 @@ Read this compact operational baseline before starting a milestone. Open detaile - At least one active global administrator is always retained; deactivation revokes that user's sessions atomically. - The local template directory (`/var/lib/dogama/templates`, under the application data bind mount) is the catalog source of truth for administrator-owned customizations. Bundled templates are copied only when their destination files are absent; after every local scan, the current bundled immutable snapshots are synchronized into SQLite and selected for new deployments while older snapshots remain available for existing instances, audit and diagnostics. - Administrators can persist bounded HTTP(S) template-repository definitions for future use. They are configuration only: remote retrieval, authentication, synchronization and automatic updates are deliberately unavailable, and Catalog Scan remains local-only. +- V Rising is the first template-scoped TCP RCON module. Its manifest currently declares only verified connectivity/status; command operations are not advertised until validated end-to-end against a real server. Go/WASI reactor modules use `-buildmode=c-shared` and initialize through `_initialize`; the V Rising success path uses concrete results to avoid Go 1.26 WASI reactor nil-interface traps. Wazero is pinned at v1.12.0. +- Module TCP access is instance-scoped and template-bound: the guest supplies no destination, only bounded bytes; the host pins the instance network address and declared integration port, enforces deadlines and response limits, and rejects arbitrary/unsafe destinations. ## Known limitations and debt diff --git a/docs/architecture/wasm-modules.md b/docs/architecture/wasm-modules.md index 2899b71..1fde079 100644 --- a/docs/architecture/wasm-modules.md +++ b/docs/architecture/wasm-modules.md @@ -43,10 +43,21 @@ The runtime grants no ambient WASI filesystem, process, environment, raw sockets - structured diagnostic emission with runtime redaction; - cancellation/deadline checks. +### Go/WASI reactor modules + +Go modules targeting `wasip1` are built as reactors with `-buildmode=c-shared`. +The resulting WASM must export `_initialize` and the module exports used by the +manifest; the host calls `_initialize` before invoking an operation. A module +must not rely on the WASI command `_start` entry point. A reproducible build +uses `CGO_ENABLED=0 GOOS=wasip1 GOARCH=wasm` and records the resulting artifact +digest in its template-local manifest. + ## Instance-scoped networking Modules never receive an arbitrary destination URL. At activation, DoGaMa binds `instance_api` to a specific instance network identity and declared integration port. Every request is checked for protocol, port, method, timeout, redirect, request/response size and concurrency. +The application service is attached to the fixed `DOGAMA_GAMES_NETWORK` network so this instance-scoped binding can resolve the selected container. The module still supplies no destination and the host pins every connection to the resolved instance address and declared port. + - No Internet or LAN destinations. - No loopback, link-local, metadata or Unix-socket destinations. - No redirects outside the bound origin. diff --git a/go.mod b/go.mod index bb04650..c2b23cd 100644 --- a/go.mod +++ b/go.mod @@ -7,7 +7,7 @@ require ( github.com/klauspost/compress v1.18.0 github.com/robfig/cron/v3 v3.0.1 github.com/santhosh-tekuri/jsonschema/v6 v6.0.3 - github.com/tetratelabs/wazero v1.11.0 + github.com/tetratelabs/wazero v1.12.0 golang.org/x/crypto v0.53.0 golang.org/x/sys v0.47.0 golang.org/x/text v0.38.0 diff --git a/go.sum b/go.sum index db3ae2d..05ce5db 100644 --- a/go.sum +++ b/go.sum @@ -20,8 +20,8 @@ github.com/robfig/cron/v3 v3.0.1 h1:WdRxkvbJztn8LMz/QEvLN5sBU+xKpSqwwUO1Pjr4qDs= github.com/robfig/cron/v3 v3.0.1/go.mod h1:eQICP3HwyT7UooqI/z+Ov+PtYAWygg1TEWWzGIFLtro= github.com/santhosh-tekuri/jsonschema/v6 v6.0.3 h1:1EYB5IzjZawrrnELUi78f9fPu57HuXjmddZPjrls/28= github.com/santhosh-tekuri/jsonschema/v6 v6.0.3/go.mod h1:JXeL+ps8p7/KNMjDQk3TCwPpBy0wYklyWTfbkIzdIFU= -github.com/tetratelabs/wazero v1.11.0 h1:+gKemEuKCTevU4d7ZTzlsvgd1uaToIDtlQlmNbwqYhA= -github.com/tetratelabs/wazero v1.11.0/go.mod h1:eV28rsN8Q+xwjogd7f4/Pp4xFxO7uOGbLcD/LzB1wiU= +github.com/tetratelabs/wazero v1.12.0 h1:DuWcpNu/FzgEXgGBDp8J1Spc+CWOvvtvVyjKlaZopYU= +github.com/tetratelabs/wazero v1.12.0/go.mod h1:LvKtzl2RqO4gyF27BiXU+nKAjcV8f38U+kP/q2vgxh0= golang.org/x/crypto v0.53.0 h1:QZ4Muo8THX6CizN2vPPd5fBGHyogrdK9fG4wLPFUsto= golang.org/x/crypto v0.53.0/go.mod h1:DNLU434OwVakk9PzuwV8w62mAJpRJL3vsgcfp4Qnsio= golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ= diff --git a/internal/agent/agent_docker.go b/internal/agent/agent_docker.go index 948ea61..77100b7 100644 --- a/internal/agent/agent_docker.go +++ b/internal/agent/agent_docker.go @@ -201,6 +201,7 @@ func (d *dockerRuntime) Create(ctx context.Context, plan agentwire.DeploymentPla NanoCPUs int64 `json:"NanoCpus"` PidsLimit *int64 `json:"PidsLimit"` CapDrop []string `json:"CapDrop"` + CapAdd []string `json:"CapAdd"` SecurityOpt []string `json:"SecurityOpt"` NetworkMode string `json:"NetworkMode"` RestartPolicy map[string]string `json:"RestartPolicy"` @@ -224,6 +225,7 @@ func (d *dockerRuntime) Create(ctx context.Context, plan agentwire.DeploymentPla payload.HostConfig.NanoCPUs = int64(plan.Resources.CPUCores * 1_000_000_000) payload.HostConfig.PidsLimit = &pidsLimit payload.HostConfig.CapDrop = []string{"ALL"} + payload.HostConfig.CapAdd = append([]string(nil), plan.CapAdd...) payload.HostConfig.SecurityOpt = []string{"no-new-privileges:true"} payload.HostConfig.NetworkMode = d.network payload.HostConfig.RestartPolicy = map[string]string{"Name": "no"} diff --git a/internal/agent/agent_docker_test.go b/internal/agent/agent_docker_test.go index aa86d7b..ae287af 100644 --- a/internal/agent/agent_docker_test.go +++ b/internal/agent/agent_docker_test.go @@ -150,6 +150,7 @@ func TestDockerRuntimeCreatesFixedSecurityBaseline(t *testing.T) { Binds []string NetworkMode string `json:"NetworkMode"` CapDrop []string `json:"CapDrop"` + CapAdd []string `json:"CapAdd"` SecurityOpt []string `json:"SecurityOpt"` Memory int64 `json:"Memory"` NanoCPUs int64 `json:"NanoCpus"` @@ -159,7 +160,7 @@ func TestDockerRuntimeCreatesFixedSecurityBaseline(t *testing.T) { if err := json.Unmarshal(<-createdBodies, &payload); err != nil { t.Fatal(err) } - if payload.HostConfig.NetworkMode != "nas-games" || len(payload.HostConfig.CapDrop) != 1 || payload.HostConfig.CapDrop[0] != "ALL" || len(payload.HostConfig.SecurityOpt) != 1 || payload.HostConfig.Memory <= 0 || payload.HostConfig.NanoCPUs <= 0 { + if payload.HostConfig.NetworkMode != "nas-games" || len(payload.HostConfig.CapDrop) != 1 || payload.HostConfig.CapDrop[0] != "ALL" || len(payload.HostConfig.CapAdd) != 0 || len(payload.HostConfig.SecurityOpt) != 1 || payload.HostConfig.Memory <= 0 || payload.HostConfig.NanoCPUs <= 0 { t.Fatalf("insecure Docker host config: %#v", payload.HostConfig) } } diff --git a/internal/agent/agent_plan.go b/internal/agent/agent_plan.go index 331aff9..3d3c310 100644 --- a/internal/agent/agent_plan.go +++ b/internal/agent/agent_plan.go @@ -78,7 +78,7 @@ func (p *PlanPolicy) Validate(plan agentwire.DeploymentPlan) error { } template := snapshot.Template imagePrefix := template.Container.Image + ":" - if !strings.HasPrefix(plan.Image, imagePrefix) || !reflect.DeepEqual(plan.Entrypoint, template.Container.Entrypoint) || plan.StopTimeoutSeconds != template.Container.StopTimeoutSeconds { + if !strings.HasPrefix(plan.Image, imagePrefix) || !reflect.DeepEqual(plan.Entrypoint, template.Container.Entrypoint) || !reflect.DeepEqual(plan.CapAdd, template.Container.Capabilities.Add) || plan.StopTimeoutSeconds != template.Container.StopTimeoutSeconds { return errors.New("container plan differs from template") } if len(plan.Arguments) < len(template.Container.Arguments) || !reflect.DeepEqual(plan.Arguments[:len(template.Container.Arguments)], template.Container.Arguments) || !allowedArguments(template, plan.Arguments[len(template.Container.Arguments):]) || !allowedEnvironment(template, plan.Environment) { diff --git a/internal/agentwire/plan.go b/internal/agentwire/plan.go index 2a6bd39..b51fb2c 100644 --- a/internal/agentwire/plan.go +++ b/internal/agentwire/plan.go @@ -32,6 +32,7 @@ type DeploymentPlan struct { Entrypoint []string `json:"entrypoint,omitempty"` Arguments []string `json:"arguments,omitempty"` Environment map[string]string `json:"environment,omitempty"` + CapAdd []string `json:"cap_add,omitempty"` Ports []PlanPort `json:"ports"` Mounts []PlanMount `json:"mounts"` Resources PlanResource `json:"resources"` @@ -71,6 +72,7 @@ func (p DeploymentPlan) CanonicalDigest() (string, error) { for key, value := range p.Environment { copyPlan.Environment[key] = value } + copyPlan.CapAdd = append([]string(nil), p.CapAdd...) copyPlan.Ports = append([]PlanPort(nil), p.Ports...) copyPlan.Mounts = append([]PlanMount(nil), p.Mounts...) sort.Slice(copyPlan.Ports, func(i, j int) bool { return copyPlan.Ports[i].ID < copyPlan.Ports[j].ID }) @@ -96,7 +98,7 @@ func (p DeploymentPlan) Validate() error { if p.StopTimeoutSeconds < 5 || p.StopTimeoutSeconds > 900 || len(p.Ports) > 32 || len(p.Mounts) == 0 || len(p.Mounts) > 16 { return errors.New("invalid deployment plan limits") } - if len(p.Labels) > 64 || len(p.Environment) > 64 || len(p.User) > 32 { + if len(p.Labels) > 64 || len(p.Environment) > 64 || len(p.CapAdd) > 16 || len(p.User) > 32 { return errors.New("invalid deployment plan container configuration") } for key, value := range p.Environment { @@ -104,6 +106,13 @@ func (p DeploymentPlan) Validate() error { return errors.New("invalid deployment plan environment") } } + seenCaps := map[string]bool{} + for _, capability := range p.CapAdd { + if !validLinuxCapability(capability) || seenCaps[capability] { + return errors.New("invalid deployment plan capabilities") + } + seenCaps[capability] = true + } for key, value := range p.Labels { lower := strings.ToLower(key) if key == "" || len(key) > 255 || len(value) > 4096 || strings.HasPrefix(lower, "dogama.") || strings.HasPrefix(lower, "io.dogama.") { @@ -166,6 +175,15 @@ func (p DeploymentPlan) Validate() error { return nil } +func validLinuxCapability(value string) bool { + switch value { + case "CHOWN", "DAC_OVERRIDE", "FOWNER", "SETGID", "SETUID", "NET_BIND_SERVICE", "NET_RAW", "SYS_CHROOT": + return true + default: + return false + } +} + type InstanceState struct { InstanceID string `json:"instance_id"` ContainerID string `json:"container_id"` diff --git a/internal/agentwire/plan_test.go b/internal/agentwire/plan_test.go index 7ade718..ff5f22e 100644 --- a/internal/agentwire/plan_test.go +++ b/internal/agentwire/plan_test.go @@ -43,3 +43,14 @@ func TestDeploymentPlanRejectsReservedLabelsAndInvalidUser(t *testing.T) { t.Fatal("non-numeric Docker user accepted") } } + +func TestDeploymentPlanRejectsUnknownOrAllCapabilities(t *testing.T) { + plan := DeploymentPlan{SchemaVersion: DeploymentPlanVersion, InstanceID: "abcdefghijklmnopqrstuvwx", TemplateID: "palworld-official", TemplateVersion: "1.0.0", TemplateDigest: strings.Repeat("a", 64), Image: "example.invalid/game:1", CapAdd: []string{"ALL"}, Ports: []PlanPort{{ID: "game", Protocol: "udp", ContainerPort: 8211, HostPort: 38211, Publish: true}}, Mounts: []PlanMount{{ID: "saved", HostPath: filepath.Join(string(filepath.Separator), "srv", "games", "saved"), ContainerPath: "/game/saved"}}, Resources: PlanResource{CPUCores: 2, MemoryMB: 1024, StorageGB: 10}, StopTimeoutSeconds: 30} + if err := plan.Validate(); err == nil { + t.Fatal("ALL capability accepted") + } + plan.CapAdd = []string{"CHOWN", "CHOWN"} + if err := plan.Validate(); err == nil { + t.Fatal("duplicate capability accepted") + } +} diff --git a/internal/catalog/local_test.go b/internal/catalog/local_test.go index 67f58bf..a5bf793 100644 --- a/internal/catalog/local_test.go +++ b/internal/catalog/local_test.go @@ -154,8 +154,27 @@ func TestScanDirAcceptsLocalTemplateWithoutModule(t *testing.T) { t.Fatal(err) } for _, snapshot := range result.Valid { - if snapshot.Template.Module == nil { - return + if snapshot.Template.ID != "vrising-didstopia" { + continue + } + path := filepath.Join(root, snapshot.AssetRoot, "template.yaml") + body, readErr := os.ReadFile(path) + if readErr != nil { + t.Fatal(readErr) + } + withoutModule := strings.Replace(string(body), "module:\n path: module/manifest.yaml\n\n", "", 1) + withoutModule = strings.Replace(withoutModule, "integration:\n module_id: vrising-rcon\n version_range: \">=1.0.0 <2.0.0\"\n port_id: rcon\n required: false\n\n", "", 1) + if err := os.WriteFile(path, []byte(withoutModule), 0o644); err != nil { + t.Fatal(err) + } + rescanned, scanErr := catalog.ScanDir(root) + if scanErr != nil { + t.Fatal(scanErr) + } + for _, candidate := range rescanned.Valid { + if candidate.Template.ID == "vrising-didstopia" && candidate.Template.Module == nil { + return + } } } t.Fatal("a template without an optional module was not accepted") diff --git a/internal/catalog/template.go b/internal/catalog/template.go index c7ece50..1913fd1 100644 --- a/internal/catalog/template.go +++ b/internal/catalog/template.go @@ -63,14 +63,17 @@ type Template struct { Recommended Resources `json:"recommended"` } `json:"requirements"` Container struct { - Image string `json:"image"` - Tag string `json:"tag"` - Entrypoint []string `json:"entrypoint,omitempty"` - UserMode string `json:"user_mode,omitempty"` - Arguments []string `json:"arguments,omitempty"` - Environment map[string]string `json:"environment,omitempty"` - StopTimeoutSeconds int `json:"stop_timeout_seconds"` - Ports []Port `json:"ports"` + Image string `json:"image"` + Tag string `json:"tag"` + Entrypoint []string `json:"entrypoint,omitempty"` + UserMode string `json:"user_mode,omitempty"` + Arguments []string `json:"arguments,omitempty"` + Environment map[string]string `json:"environment,omitempty"` + Capabilities struct { + Add []string `json:"add,omitempty"` + } `json:"capabilities,omitempty"` + StopTimeoutSeconds int `json:"stop_timeout_seconds"` + Ports []Port `json:"ports"` Assets []struct { Source string `json:"source"` Destination string `json:"destination"` diff --git a/internal/instance/configuration_resolver.go b/internal/instance/configuration_resolver.go index 63a6424..48c67a7 100644 --- a/internal/instance/configuration_resolver.go +++ b/internal/instance/configuration_resolver.go @@ -52,6 +52,9 @@ func ResolveConfiguration(template catalog.Template, values map[string]string) ( target := field.Target switch target.Kind { case "environment": + if value == "" && !field.Required { + continue + } if protectedEnvironment[target.Name] || strings.HasPrefix(target.Name, "DOGAMA_") { return ResolvedConfiguration{}, fmt.Errorf("%s targets a protected environment variable", field.ID) } diff --git a/internal/instance/configuration_resolver_test.go b/internal/instance/configuration_resolver_test.go index 921e3d9..8d310af 100644 --- a/internal/instance/configuration_resolver_test.go +++ b/internal/instance/configuration_resolver_test.go @@ -40,6 +40,30 @@ func TestResolveConfigurationTargets(t *testing.T) { } } +func TestEnvironmentValuesRemainExactAndOptionalValuesAreOmitted(t *testing.T) { + template := catalog.Template{} + template.Configuration.Fields = []catalog.ConfigField{ + {ID: "port", Type: "integer", Target: catalog.ConfigTarget{Kind: "environment", Name: "GAME_PORT"}}, + {ID: "enabled", Type: "boolean", Target: catalog.ConfigTarget{Kind: "environment", Name: "ENABLED"}}, + {ID: "optional", Type: "string", Target: catalog.ConfigTarget{Kind: "environment", Name: "OPTIONAL"}}, + } + resolved, err := instance.ResolveConfiguration(template, map[string]string{"port": "9876", "enabled": "true", "optional": ""}) + if err != nil { + t.Fatal(err) + } + if resolved.Environment["GAME_PORT"] != "9876" || resolved.Environment["ENABLED"] != "true" { + t.Fatalf("environment values were rewritten: %#v", resolved.Environment) + } + if _, ok := resolved.Environment["OPTIONAL"]; ok { + t.Fatalf("empty optional environment was emitted: %#v", resolved.Environment) + } + for key, value := range resolved.Environment { + if strings.HasPrefix(value, "= ") || strings.HasPrefix(value, "=") { + t.Fatalf("synthetic environment prefix for %s: %q", key, value) + } + } +} + func TestPalworldServerNameIsWrittenBeforeStart(t *testing.T) { snapshots, err := catalog.LoadFS(catalogdata.Files, ".") if err != nil { diff --git a/internal/instance/module_service.go b/internal/instance/module_service.go index 605f726..e79d2af 100644 --- a/internal/instance/module_service.go +++ b/internal/instance/module_service.go @@ -75,6 +75,7 @@ type moduleManifest struct { Network struct { PortIDs []string `yaml:"port_ids"` HTTPMethods []string `yaml:"http_methods"` + Protocols []string `yaml:"protocols"` } `yaml:"network"` } `yaml:"permissions"` Limits struct { @@ -124,6 +125,16 @@ func (s *ModuleService) runtime(ctx context.Context, value StoredInstance) (*mod return nil, moduleManifest{}, ErrModuleUnavailable } methods := map[string]bool{} + tcpAllowed := false + for _, protocol := range manifest.Permissions.Network.Protocols { + switch protocol { + case "http": + case "tcp": + tcpAllowed = true + default: + return nil, moduleManifest{}, ErrModuleUnavailable + } + } for _, method := range manifest.Permissions.Network.HTTPMethods { if method == "GET" || method == "POST" { methods[method] = true @@ -131,7 +142,7 @@ func (s *ModuleService) runtime(ctx context.Context, value StoredInstance) (*mod return nil, moduleManifest{}, ErrModuleUnavailable } } - if len(methods) == 0 || manifest.Artifacts.WASM == "" || strings.Contains(manifest.Artifacts.WASM, "/") || strings.Contains(manifest.Artifacts.WASM, "..") { + if (len(methods) == 0 && !tcpAllowed) || manifest.Artifacts.WASM == "" || strings.Contains(manifest.Artifacts.WASM, "/") || strings.Contains(manifest.Artifacts.WASM, "..") { return nil, moduleManifest{}, ErrModuleUnavailable } config, secrets := map[string]string{}, map[string]string{} @@ -244,7 +255,7 @@ func (s *ModuleService) Action(ctx context.Context, value StoredInstance, capabi var out struct { Accepted bool `json:"accepted"` } - if err := r.Call(ctx, operation, request, &out); err != nil || !out.Accepted { + if err := r.CallData(ctx, operation, request, &out); err != nil || !out.Accepted { return ErrModuleUnavailable } return nil diff --git a/internal/instance/preview.go b/internal/instance/preview.go index 0b583ad..697838b 100644 --- a/internal/instance/preview.go +++ b/internal/instance/preview.go @@ -45,6 +45,7 @@ type Preview struct { Image string `json:"image"` Entrypoint []string `json:"entrypoint,omitempty"` Arguments []string `json:"arguments,omitempty"` + CapAdd []string `json:"cap_add,omitempty"` StopTimeoutSeconds int `json:"stop_timeout_seconds"` StartupTimeoutSeconds int `json:"startup_timeout_seconds"` Ports []PortBinding `json:"ports"` @@ -283,6 +284,7 @@ func BuildPreview(snapshot catalog.Snapshot, request PreviewRequest) (Preview, e Image: snapshot.Template.Container.Image + ":" + tag.Tag, Entrypoint: append([]string(nil), snapshot.Template.Container.Entrypoint...), Arguments: append([]string(nil), snapshot.Template.Container.Arguments...), + CapAdd: append([]string(nil), snapshot.Template.Container.Capabilities.Add...), StopTimeoutSeconds: snapshot.Template.Container.StopTimeoutSeconds, StartupTimeoutSeconds: snapshot.Template.Healthcheck.StartupTimeoutSeconds, Ports: ports, Mounts: mounts, Resources: resources, Settings: settings, @@ -386,6 +388,7 @@ func (p Preview) deploymentPlan(instanceID string, secrets map[string]string) (a TemplateID: p.Template.ID, TemplateVersion: p.Template.Version, TemplateDigest: p.Template.Digest, Image: p.Image, Entrypoint: append([]string(nil), p.Entrypoint...), Arguments: arguments, Environment: environment, Labels: MergeLabels(nil, global, local), User: p.DockerUserValue, + CapAdd: append([]string(nil), p.CapAdd...), Resources: agentwire.PlanResource{CPUCores: p.Resources.CPUCores, MemoryMB: p.Resources.MemoryMB, StorageGB: p.Resources.StorageGB}, StopTimeoutSeconds: p.StopTimeoutSeconds, } diff --git a/internal/module/runtime.go b/internal/module/runtime.go index 2c645dc..e9d8ed1 100644 --- a/internal/module/runtime.go +++ b/internal/module/runtime.go @@ -13,6 +13,7 @@ import ( "net" "net/http" "net/url" + "os" "path" "strings" "sync" @@ -50,6 +51,12 @@ type Binding struct { Secrets map[string]string } +// TCPRequest is a single, bounded exchange with the template-declared TCP +// integration endpoint. The destination is never supplied by a module. +type TCPRequest struct { + Body []byte `json:"body"` +} + type HTTPRequest struct { Method string `json:"method"` Path string `json:"path"` @@ -69,6 +76,7 @@ type Runtime struct { limits Limits binding Binding client *http.Client + dialContext func(context.Context, string, string) (net.Conn, error) sem chan struct{} mu sync.Mutex failures int @@ -76,14 +84,20 @@ type Runtime struct { } func New(ctx context.Context, wasm []byte, checksum string, capabilities []string, limits Limits, binding Binding) (*Runtime, error) { - transport, err := pinnedTransport(ctx, binding) + transport, dialContext, err := pinnedTransport(ctx, binding) if err != nil { return nil, err } - return newWithTransport(wasm, checksum, capabilities, limits, binding, transport) + return newWithTransportAndDial(wasm, checksum, capabilities, limits, binding, transport, dialContext) } func newWithTransport(wasm []byte, checksum string, capabilities []string, limits Limits, binding Binding, transport http.RoundTripper) (*Runtime, error) { + return newWithTransportAndDial(wasm, checksum, capabilities, limits, binding, transport, func(context.Context, string, string) (net.Conn, error) { + return nil, errors.New("TCP exchange is unavailable in this test runtime") + }) +} + +func newWithTransportAndDial(wasm []byte, checksum string, capabilities []string, limits Limits, binding Binding, transport http.RoundTripper, dialContext func(context.Context, string, string) (net.Conn, error)) (*Runtime, error) { if len(wasm) == 0 || limits.MemoryMB == 0 || limits.MemoryMB > 256 || limits.Timeout <= 0 || limits.MaxResponseBytes < 1 || limits.MaxResponseBytes > 8<<20 || limits.MaxConcurrentCall < 1 || limits.MaxConcurrentCall > 16 { return nil, errors.New("invalid module runtime configuration") } @@ -91,7 +105,7 @@ func newWithTransport(wasm []byte, checksum string, capabilities []string, limit if hex.EncodeToString(digest[:]) != checksum { return nil, errors.New("module checksum mismatch") } - if binding.InstanceID == "" || binding.ContainerPort < 1 || binding.ContainerPort > 65535 || transport == nil { + if binding.InstanceID == "" || binding.ContainerPort < 1 || binding.ContainerPort > 65535 || transport == nil || dialContext == nil { return nil, errors.New("invalid instance API binding") } seen := map[string]bool{} @@ -101,35 +115,36 @@ func newWithTransport(wasm []byte, checksum string, capabilities []string, limit } seen[capability] = true } - r := &Runtime{wasm: append([]byte(nil), wasm...), capabilities: append([]string(nil), capabilities...), limits: limits, binding: binding, sem: make(chan struct{}, limits.MaxConcurrentCall)} + r := &Runtime{wasm: append([]byte(nil), wasm...), capabilities: append([]string(nil), capabilities...), limits: limits, binding: binding, sem: make(chan struct{}, limits.MaxConcurrentCall), dialContext: dialContext} r.client = &http.Client{Transport: transport, CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }} return r, nil } -func pinnedTransport(ctx context.Context, binding Binding) (http.RoundTripper, error) { +func pinnedTransport(ctx context.Context, binding Binding) (http.RoundTripper, func(context.Context, string, string) (net.Conn, error), error) { hostname := "dogama-" + strings.ToLower(binding.InstanceID) lookupCtx, cancel := context.WithTimeout(ctx, 2*time.Second) defer cancel() addresses, err := net.DefaultResolver.LookupIPAddr(lookupCtx, hostname) if err != nil || len(addresses) == 0 { - return nil, errors.New("resolve bound instance API") + return nil, nil, errors.New("resolve bound instance API") } ip := addresses[0].IP if ip == nil || ip.IsUnspecified() || ip.IsLoopback() || ip.IsLinkLocalUnicast() || ip.IsLinkLocalMulticast() || ip.IsMulticast() { - return nil, errors.New("unsafe instance API address") + return nil, nil, errors.New("unsafe instance API address") } pinned := net.JoinHostPort(ip.String(), fmt.Sprint(binding.ContainerPort)) dialer := &net.Dialer{Timeout: 2 * time.Second, KeepAlive: 30 * time.Second} + dialContext := func(ctx context.Context, network, _ string) (net.Conn, error) { + if network != "tcp" && network != "tcp4" && network != "tcp6" { + return nil, errors.New("unsupported instance API network") + } + return dialer.DialContext(ctx, "tcp", pinned) + } return &http.Transport{ DisableCompression: true, Proxy: nil, - DialContext: func(ctx context.Context, network, _ string) (net.Conn, error) { - if network != "tcp" && network != "tcp4" && network != "tcp6" { - return nil, errors.New("unsupported instance API network") - } - return dialer.DialContext(ctx, "tcp", pinned) - }, - }, nil + DialContext: dialContext, + }, dialContext, nil } func (r *Runtime) Call(ctx context.Context, operation string, request any, response any) error { @@ -168,6 +183,33 @@ func (r *Runtime) Call(ctx context.Context, operation string, request any, respo return nil } +// CallData unwraps the normalized module envelope for typed service callers. +// Call remains available for diagnostics and callers that need the envelope. +func (r *Runtime) CallData(ctx context.Context, operation string, request any, response any) error { + var envelope struct { + OK bool `json:"ok"` + Data json.RawMessage `json:"data"` + Error *struct { + Message string `json:"message"` + } `json:"error"` + } + if err := r.Call(ctx, operation, request, &envelope); err != nil { + return err + } + if !envelope.OK { + if envelope.Error != nil && envelope.Error.Message != "" { + return errors.New(envelope.Error.Message) + } + return errors.New("module operation failed") + } + decoder := json.NewDecoder(bytes.NewReader(envelope.Data)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(response); err != nil || decoder.Decode(&struct{}{}) != io.EOF { + return errors.New("invalid module response") + } + return nil +} + func (r *Runtime) allowedOperation(operation string) bool { if operation == "initialize" || operation == "test_connection" || operation == "get_server_status" { return true @@ -184,9 +226,13 @@ type callState struct{ hostCalls int } type stateKey struct{} func (r *Runtime) invoke(ctx context.Context, operation string, input []byte) ([]byte, error) { + return r.invokeWithConfig(ctx, operation, input, wazero.NewRuntimeConfigInterpreter()) +} + +func (r *Runtime) invokeWithConfig(ctx context.Context, operation string, input []byte, baseConfig wazero.RuntimeConfig) ([]byte, error) { ctx = context.WithValue(ctx, stateKey{}, &callState{}) pages := (r.limits.MemoryMB*1024*1024 + 65535) / 65536 - config := wazero.NewRuntimeConfigInterpreter().WithMemoryLimitPages(pages).WithCloseOnContextDone(true).WithDebugInfoEnabled(false) + config := baseConfig.WithMemoryLimitPages(pages).WithCloseOnContextDone(true).WithDebugInfoEnabled(false) runtime := wazero.NewRuntimeWithConfig(ctx, config) defer func() { _ = runtime.Close(ctx) }() if _, err := wasi_snapshot_preview1.Instantiate(ctx, runtime); err != nil { @@ -194,6 +240,7 @@ func (r *Runtime) invoke(ctx context.Context, operation string, input []byte) ([ } host := runtime.NewHostModuleBuilder("dogama_host") host.NewFunctionBuilder().WithFunc(r.httpRequest).Export("http_request") + host.NewFunctionBuilder().WithFunc(r.tcpExchange).Export("tcp_exchange") host.NewFunctionBuilder().WithFunc(r.getSecret).Export("get_secret") host.NewFunctionBuilder().WithFunc(r.getConfig).Export("get_config") if _, err := host.Instantiate(ctx); err != nil { @@ -234,6 +281,88 @@ func (r *Runtime) invoke(ctx context.Context, operation string, input []byte) ([ return append([]byte(nil), body...), nil } +// tcpExchange provides no socket handle to the module. It writes one bounded +// request to the instance-pinned integration endpoint then returns bytes read +// until a short idle period, EOF, or the call deadline. This permits framed +// protocols such as Source RCON while keeping every operation bounded. +func (r *Runtime) tcpExchange(ctx context.Context, mod api.Module, requestPtr, requestLen, responsePtr, responseCap uint32) int32 { + state, _ := ctx.Value(stateKey{}).(*callState) + if state == nil || state.hostCalls >= maxHostCallsPerCall || requestLen > maxRequestBytes || responseCap > uint32(r.limits.MaxResponseBytes) { + return -1 + } + state.hostCalls++ + raw, ok := mod.Memory().Read(requestPtr, requestLen) + if !ok { + return -1 + } + var request TCPRequest + decoder := json.NewDecoder(bytes.NewReader(raw)) + decoder.DisallowUnknownFields() + if decoder.Decode(&request) != nil || len(request.Body) == 0 || len(request.Body) > maxRequestBytes { + return -2 + } + conn, err := r.dialContext(ctx, "tcp", "") + if err != nil { + if errors.Is(err, context.DeadlineExceeded) || isTimeout(err) { + return -4 + } + return -3 + } + defer func() { _ = conn.Close() }() + deadline := time.Now().Add(r.limits.Timeout) + if value, exists := ctx.Deadline(); exists && value.Before(deadline) { + deadline = value + } + if err := conn.SetDeadline(deadline); err != nil { + return -3 + } + if _, err := conn.Write(request.Body); err != nil { + if isTimeout(err) { + return -4 + } + return -3 + } + result := make([]byte, 0, min(int(responseCap), 4096)) + buffer := make([]byte, min(4096, int(responseCap))) + for { + if len(result) > 0 { + _ = conn.SetReadDeadline(time.Now().Add(50 * time.Millisecond)) + } + n, readErr := conn.Read(buffer) + if n > 0 { + if len(result)+n > int(responseCap) { + return -5 + } + result = append(result, buffer[:n]...) + } + if readErr != nil { + if errors.Is(readErr, io.EOF) || (len(result) > 0 && isTimeout(readErr)) { + break + } + if isTimeout(readErr) { + return -4 + } + return -3 + } + } + if len(result) == 0 || !mod.Memory().Write(responsePtr, result) { + return -5 + } + return int32(len(result)) +} + +func isTimeout(err error) bool { + var netErr net.Error + return errors.Is(err, os.ErrDeadlineExceeded) || (errors.As(err, &netErr) && netErr.Timeout()) +} + +func min(a, b int) int { + if a < b { + return a + } + return b +} + func validateABI(compiled wazero.CompiledModule, capabilities []string) error { exports := compiled.ExportedFunctions() for _, name := range []string{"dogama_alloc", "initialize", "test_connection", "get_server_status"} { diff --git a/internal/module/runtime_test.go b/internal/module/runtime_test.go index 775870f..3a8d0e1 100644 --- a/internal/module/runtime_test.go +++ b/internal/module/runtime_test.go @@ -3,16 +3,29 @@ package module import ( "context" "crypto/sha256" + "encoding/binary" "encoding/hex" "errors" "io" + "net" "net/http" "os" "strings" "testing" "time" + + "github.com/tetratelabs/wazero" ) +func sourceRCONPacket(id, typ uint32, body string) []byte { + value := make([]byte, len(body)+14) + binary.LittleEndian.PutUint32(value[:4], uint32(len(body)+10)) + binary.LittleEndian.PutUint32(value[4:8], id) + binary.LittleEndian.PutUint32(value[8:12], typ) + copy(value[12:], body) + return value +} + type roundTripFunc func(*http.Request) (*http.Response, error) func (f roundTripFunc) RoundTrip(request *http.Request) (*http.Response, error) { return f(request) } @@ -145,3 +158,112 @@ func TestHTTPRequestPinsInstanceOriginAndDisablesRedirects(t *testing.T) { t.Fatal("bound transport was not used") } } + +func TestVRisingAdapterUsesPinnedBoundedTCPRCON(t *testing.T) { + wasm, err := os.ReadFile("../../catalog/vrising/module/module.wasm") + if err != nil { + t.Fatal(err) + } + digest := sha256.Sum256(wasm) + var seen []byte + dial := func(context.Context, string, string) (net.Conn, error) { + client, server := net.Pipe() + go func() { + defer server.Close() + buffer := make([]byte, 4096) + n, readErr := server.Read(buffer) + if readErr != nil { + return + } + seen = append([]byte(nil), buffer[:n]...) + _, _ = server.Write(sourceRCONPacket(1, 2, "")) + }() + return client, nil + } + runtime, err := newWithTransportAndDial(wasm, hex.EncodeToString(digest[:]), nil, Limits{MemoryMB: 64, Timeout: 10 * time.Second, MaxResponseBytes: 262144, MaxConcurrentCall: 2}, Binding{InstanceID: strings.Repeat("a", 20), ContainerPort: 9878, Configuration: map[string]string{"rcon_enabled": "true", "rcon_port": "9878"}, Secrets: map[string]string{"rcon_password": "not-in-output"}}, roundTripFunc(func(*http.Request) (*http.Response, error) { return nil, errors.New("HTTP must not be used") }), dial) + if err != nil { + t.Fatal(err) + } + var response struct { + OK bool `json:"ok"` + Data any `json:"data"` + Error any `json:"error"` + } + if err := runtime.Call(context.Background(), "test_connection", struct{}{}, &response); err != nil || !response.OK { + t.Fatalf("connection = %#v, %v", response, err) + } + if len(seen) < 14 || binary.LittleEndian.Uint32(seen[8:12]) != 3 { + t.Fatalf("unexpected RCON auth packet: %x", seen) + } +} + +func TestVRisingAdapterRejectsAuthenticationWithoutLeakingSecret(t *testing.T) { + wasm, err := os.ReadFile("../../catalog/vrising/module/module.wasm") + if err != nil { + t.Fatal(err) + } + digest := sha256.Sum256(wasm) + dial := func(context.Context, string, string) (net.Conn, error) { + client, server := net.Pipe() + go func() { + defer server.Close() + buffer := make([]byte, 256) + _, _ = server.Read(buffer) + _, _ = server.Write(sourceRCONPacket(^uint32(0), 2, "")) + }() + return client, nil + } + runtime, err := newWithTransportAndDial(wasm, hex.EncodeToString(digest[:]), nil, Limits{MemoryMB: 32, Timeout: 10 * time.Second, MaxResponseBytes: 262144, MaxConcurrentCall: 1}, Binding{InstanceID: strings.Repeat("a", 20), ContainerPort: 9878, Configuration: map[string]string{"rcon_enabled": "true", "rcon_port": "9878"}, Secrets: map[string]string{"rcon_password": "very-secret-value"}}, roundTripFunc(func(*http.Request) (*http.Response, error) { return nil, errors.New("HTTP must not be used") }), dial) + if err != nil { + t.Fatal(err) + } + var response struct { + OK bool `json:"ok"` + Data any `json:"data"` + Error struct { + Code string `json:"code"` + Message string `json:"message"` + Retryable bool `json:"retryable"` + } `json:"error"` + } + if err := runtime.Call(context.Background(), "test_connection", struct{}{}, &response); err != nil || response.OK || response.Error.Code != "unauthorized" || strings.Contains(response.Error.Message, "very-secret-value") { + t.Fatalf("unsafe authentication response: %#v, %v", response, err) + } +} + +func TestVRisingAuthSuccessWithBothWazeroEngines(t *testing.T) { + wasm, err := os.ReadFile("../../catalog/vrising/module/module.wasm") + if err != nil { + t.Fatal(err) + } + digest := sha256.Sum256(wasm) + dial := func(context.Context, string, string) (net.Conn, error) { + client, server := net.Pipe() + go func() { + defer server.Close() + buffer := make([]byte, 256) + if _, err := server.Read(buffer); err == nil { + _, _ = server.Write(sourceRCONPacket(1, 2, "")) + } + }() + return client, nil + } + runtime, err := newWithTransportAndDial(wasm, hex.EncodeToString(digest[:]), nil, Limits{MemoryMB: 32, Timeout: 10 * time.Second, MaxResponseBytes: 262144, MaxConcurrentCall: 1}, Binding{InstanceID: strings.Repeat("a", 20), ContainerPort: 9878, Configuration: map[string]string{"rcon_enabled": "true", "rcon_port": "9878"}, Secrets: map[string]string{"rcon_password": "secret"}}, roundTripFunc(func(*http.Request) (*http.Response, error) { return nil, errors.New("HTTP must not be used") }), dial) + if err != nil { + t.Fatal(err) + } + for _, engine := range []struct { + name string + config wazero.RuntimeConfig + }{ + {name: "compiler", config: wazero.NewRuntimeConfigCompiler()}, + {name: "interpreter", config: wazero.NewRuntimeConfigInterpreter()}, + } { + t.Run(engine.name, func(t *testing.T) { + result, err := runtime.invokeWithConfig(context.Background(), "test_connection", []byte("{}"), engine.config) + if err != nil || !strings.Contains(string(result), `"connected":true`) { + t.Fatalf("auth success result=%s err=%v", result, err) + } + }) + } +} diff --git a/specs/template.schema.json b/specs/template.schema.json index bfb5ae8..ba1f397 100644 --- a/specs/template.schema.json +++ b/specs/template.schema.json @@ -60,6 +60,18 @@ "user_mode": { "type": "string", "enum": ["dogama", "image"] }, "arguments": { "type": "array", "items": { "type": "string", "maxLength": 500 }, "maxItems": 64 }, "environment": { "type": "object", "maxProperties": 64, "additionalProperties": { "type": "string", "maxLength": 4096 }, "propertyNames": { "pattern": "^[A-Za-z_][A-Za-z0-9_]*$" } }, + "capabilities": { + "type": "object", + "additionalProperties": false, + "properties": { + "add": { + "type": "array", + "maxItems": 16, + "uniqueItems": true, + "items": { "enum": ["CHOWN", "DAC_OVERRIDE", "FOWNER", "SETGID", "SETUID", "NET_BIND_SERVICE", "NET_RAW", "SYS_CHROOT"] } + } + } + }, "assets": { "type": "array", "maxItems": 16,