feat(vrising): support template runtime requirements #44
@@ -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.
|
||||
@@ -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"
|
||||
Binary file not shown.
@@ -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() {}
|
||||
@@ -0,0 +1,5 @@
|
||||
//go:build !wasip1
|
||||
|
||||
package main
|
||||
|
||||
func main() {}
|
||||
@@ -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("")
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
requires_instance_migration: false
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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=
|
||||
|
||||
@@ -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"}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
|
||||
+144
-15
@@ -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"} {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user