1038 lines
40 KiB
Go
1038 lines
40 KiB
Go
//go:build windows
|
|
|
|
package windows
|
|
|
|
// The Windows adapter deliberately keeps all Win32 handles in this file. The
|
|
// selector in selection.go is pure policy; this layer obtains one verified
|
|
// primary token, creates a suspended child with an explicit handle list, and
|
|
// puts it in a kill-on-close Job before releasing it.
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"os/exec"
|
|
"strings"
|
|
"sync"
|
|
"syscall"
|
|
"time"
|
|
"unsafe"
|
|
|
|
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
|
|
"github.com/rvbox/rvbox/internal/client/supervisor"
|
|
winapi "golang.org/x/sys/windows"
|
|
)
|
|
|
|
const (
|
|
logon32LogonService = 5
|
|
logon32ProviderDefault = 0
|
|
securitySystemRID = "S-1-5-18"
|
|
securityLocalServiceRID = "S-1-5-19"
|
|
disableMaxPrivilege = 0x1
|
|
mediumIntegrityRID = 0x2000
|
|
highIntegrityRID = 0x3000
|
|
)
|
|
|
|
var (
|
|
advapi32 = syscall.NewLazyDLL("advapi32.dll")
|
|
procLogonUserW = advapi32.NewProc("LogonUserW")
|
|
procCreateRestrictedToken = advapi32.NewProc("CreateRestrictedToken")
|
|
kernel32 = syscall.NewLazyDLL("kernel32.dll")
|
|
procAttachConsole = kernel32.NewProc("AttachConsole")
|
|
procFreeConsole = kernel32.NewProc("FreeConsole")
|
|
procGenerateCtrlEvent = kernel32.NewProc("GenerateConsoleCtrlEvent")
|
|
procSetCtrlHandler = kernel32.NewProc("SetConsoleCtrlHandler")
|
|
)
|
|
|
|
func stdinLineEnding() []byte { return []byte{'\r', '\n'} }
|
|
|
|
type nativeHandles struct {
|
|
process winapi.Handle
|
|
job winapi.Handle
|
|
pid uint32
|
|
close sync.Once
|
|
}
|
|
|
|
// NewSupervisor constructs the machine-wide Windows implementation. The
|
|
// service process is expected to run as LocalSystem; token selection verifies
|
|
// that assumption when a command is started and records the selected context.
|
|
func NewSupervisor(options NativeOptions) (supervisor.Supervisor, error) {
|
|
options = options.withDefaults()
|
|
return newExecSupervisor(options), nil
|
|
}
|
|
|
|
func (manager *execSupervisor) Start(ctx context.Context, spec supervisor.StartSpec) (supervisor.Process, error) {
|
|
if err := validateExecutionSource(&spec); err != nil {
|
|
return nil, err
|
|
}
|
|
if spec.Execution.GetShellType() != rvboxv1.ShellType_SHELL_CMD && spec.Execution.GetShellType() != rvboxv1.ShellType_SHELL_POWERSHELL {
|
|
return nil, ErrUnsupportedShell
|
|
}
|
|
manager.mu.Lock()
|
|
if _, exists := manager.active[spec.IssueUUID]; exists {
|
|
manager.mu.Unlock()
|
|
return nil, ErrProcessAlreadyRunning
|
|
}
|
|
manager.mu.Unlock()
|
|
|
|
plan, err := ResolveShell(spec.Execution.GetShellType(), manager.options.Shells)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if err := verifyExecutable(plan.ApplicationName); err != nil {
|
|
return nil, err
|
|
}
|
|
wrapper, err := BuildWrapper(spec.Execution, spec.ScriptBody, manager.options.MaxWrapperBytes)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var cleanup func()
|
|
fail := func(cause error) (supervisor.Process, error) {
|
|
if cleanup != nil {
|
|
cleanup()
|
|
}
|
|
return nil, cause
|
|
}
|
|
|
|
token, identity, err := manager.selectToken(spec.Execution.GetElevated())
|
|
if err != nil {
|
|
return fail(err)
|
|
}
|
|
defer token.Close()
|
|
wrapperPath, cleanup, err := materializeWrapper(spec.WorkingDirectory, wrapper, manager.options.Now(), func(path string) error {
|
|
return secureWrapperFile(path, identity.UserSID)
|
|
})
|
|
if err != nil {
|
|
return fail(err)
|
|
}
|
|
baseEnvironment, err := token.Environ(false)
|
|
if err != nil {
|
|
return fail(fmt.Errorf("build token environment: %w", err))
|
|
}
|
|
environment, err := BuildEnvironmentBlock(parseEnvironment(baseEnvironment), spec.Environment)
|
|
if err != nil {
|
|
return fail(err)
|
|
}
|
|
launch, err := plan.BuildLaunchPlan(wrapperPath, spec.WorkingDirectory, environment)
|
|
if err != nil {
|
|
return fail(err)
|
|
}
|
|
|
|
stdinRead, stdinWrite, stdoutRead, stdoutWrite, stderrRead, stderrWrite, err := createStandardPipes()
|
|
if err != nil {
|
|
return fail(err)
|
|
}
|
|
closeFiles := func() {
|
|
for _, file := range []*os.File{stdinRead, stdinWrite, stdoutRead, stdoutWrite, stderrRead, stderrWrite} {
|
|
if file != nil {
|
|
_ = file.Close()
|
|
}
|
|
}
|
|
}
|
|
pipesTransferred := false
|
|
defer func() {
|
|
// Parent-side handles are retained only after successful process
|
|
// creation. Any error path closes both ends here.
|
|
if !pipesTransferred {
|
|
closeFiles()
|
|
}
|
|
}()
|
|
|
|
job, err := createKillOnCloseJob()
|
|
if err != nil {
|
|
closeFiles()
|
|
return fail(fmt.Errorf("create command Job: %w", err))
|
|
}
|
|
if err := applyJobProfiles(job, manager.options.JobProfiles, spec.ExecutionProfiles); err != nil {
|
|
_ = winapi.CloseHandle(job)
|
|
return fail(err)
|
|
}
|
|
cleanupJob := true
|
|
defer func() {
|
|
if cleanupJob {
|
|
_ = winapi.CloseHandle(job)
|
|
}
|
|
}()
|
|
|
|
application, err := winapi.UTF16PtrFromString(launch.ApplicationName)
|
|
if err != nil {
|
|
return fail(err)
|
|
}
|
|
commandLine, err := winapi.UTF16FromString(launch.CommandLine)
|
|
if err != nil {
|
|
return fail(err)
|
|
}
|
|
workingDirectory, err := winapi.UTF16PtrFromString(launch.WorkingDirectory)
|
|
if err != nil {
|
|
return fail(err)
|
|
}
|
|
attributeList, err := winapi.NewProcThreadAttributeList(1)
|
|
if err != nil {
|
|
return fail(err)
|
|
}
|
|
defer attributeList.Delete()
|
|
childHandles := []winapi.Handle{winapi.Handle(stdinRead.Fd()), winapi.Handle(stdoutWrite.Fd()), winapi.Handle(stderrWrite.Fd())}
|
|
if err := attributeList.Update(winapi.PROC_THREAD_ATTRIBUTE_HANDLE_LIST, unsafe.Pointer(&childHandles[0]), uintptr(len(childHandles))*unsafe.Sizeof(childHandles[0])); err != nil {
|
|
return fail(err)
|
|
}
|
|
startup := winapi.StartupInfoEx{}
|
|
startup.Cb = uint32(unsafe.Sizeof(startup))
|
|
startup.Flags = winapi.STARTF_USESTDHANDLES | winapi.STARTF_USESHOWWINDOW
|
|
startup.ShowWindow = winapi.SW_HIDE
|
|
startup.StdInput = childHandles[0]
|
|
startup.StdOutput = childHandles[1]
|
|
startup.StdErr = childHandles[2]
|
|
startup.ProcThreadAttributeList = attributeList.List()
|
|
var processInfo winapi.ProcessInformation
|
|
flags := uint32(winapi.CREATE_NEW_CONSOLE | winapi.CREATE_SUSPENDED | winapi.CREATE_UNICODE_ENVIRONMENT | winapi.EXTENDED_STARTUPINFO_PRESENT)
|
|
var environmentPointer *uint16
|
|
if len(environment) > 0 {
|
|
environmentPointer = &environment[0]
|
|
}
|
|
if err := winapi.CreateProcessAsUser(token, application, &commandLine[0], nil, nil, true, flags, environmentPointer, workingDirectory, &startup.StartupInfo, &processInfo); err != nil {
|
|
return fail(fmt.Errorf("create suspended command process: %w", err))
|
|
}
|
|
// The child owns these handles after CreateProcessAsUser returns. Keep only
|
|
// the three parent ends and the process/job handles in the daemon. The
|
|
// primary thread remains suspended until the executor has durably recorded
|
|
// launch authorization and calls Process.Release.
|
|
_ = stdinRead.Close()
|
|
_ = stdoutWrite.Close()
|
|
_ = stderrWrite.Close()
|
|
if err := winapi.AssignProcessToJobObject(job, processInfo.Process); err != nil {
|
|
_ = winapi.TerminateProcess(processInfo.Process, 1)
|
|
_ = winapi.CloseHandle(processInfo.Process)
|
|
_ = winapi.CloseHandle(processInfo.Thread)
|
|
return fail(fmt.Errorf("assign command to Job: %w", err))
|
|
}
|
|
pipesTransferred = true
|
|
var threadClosed sync.Once
|
|
closeThread := func() {
|
|
threadClosed.Do(func() { _ = winapi.CloseHandle(processInfo.Thread) })
|
|
}
|
|
releaseFn := func() error {
|
|
if _, err := winapi.ResumeThread(processInfo.Thread); err != nil {
|
|
closeThread()
|
|
return fmt.Errorf("release suspended command: %w", err)
|
|
}
|
|
closeThread()
|
|
return nil
|
|
}
|
|
started := manager.options.Now()
|
|
handles := &nativeHandles{process: processInfo.Process, job: job, pid: processInfo.ProcessId}
|
|
cleanupJob = false
|
|
command := &exec.Cmd{Process: osProcess(processInfo.ProcessId)}
|
|
waitFn := func() (int32, bool, error) {
|
|
_, waitErr := winapi.WaitForSingleObject(processInfo.Process, winapi.INFINITE)
|
|
var code uint32
|
|
if err := winapi.GetExitCodeProcess(processInfo.Process, &code); err != nil && waitErr == nil {
|
|
waitErr = err
|
|
}
|
|
closeThread()
|
|
handles.close.Do(func() {
|
|
_ = winapi.CloseHandle(processInfo.Process)
|
|
_ = winapi.CloseHandle(job)
|
|
})
|
|
return int32(code), false, waitErr
|
|
}
|
|
killFn := func(code uint32) error {
|
|
err := winapi.TerminateJobObject(job, code)
|
|
closeThread()
|
|
return err
|
|
}
|
|
process := manager.registerProcess(spec.IssueUUID, identity, command, stdinWrite, stdoutRead, stderrRead, started, waitFn, killFn, releaseFn, cleanup)
|
|
process.snapshotFn = func() (supervisor.ResourceSnapshot, error) {
|
|
return queryJobSnapshot(job, manager.options.Now())
|
|
}
|
|
return process, nil
|
|
}
|
|
|
|
func osProcess(pid uint32) *os.Process {
|
|
process, err := os.FindProcess(int(pid))
|
|
if err != nil {
|
|
return &os.Process{}
|
|
}
|
|
return process
|
|
}
|
|
|
|
func verifyExecutable(path string) error {
|
|
info, err := os.Stat(path)
|
|
if err != nil {
|
|
return fmt.Errorf("stat configured shell %q: %w", path, err)
|
|
}
|
|
if !info.Mode().IsRegular() {
|
|
return fmt.Errorf("configured shell %q is not a regular file", path)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// secureWrapperFile replaces the inherited directory ACL with a protected
|
|
// DACL. The service writes the wrapper before this call; afterward only
|
|
// LocalSystem and the selected effective token SID can read it. This is done
|
|
// after token selection so active-user commands do not depend on inherited
|
|
// broad ProgramData permissions.
|
|
func secureWrapperFile(path, effectiveSID string) error {
|
|
systemSID, err := winapi.StringToSid(securitySystemRID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
effective := systemSID
|
|
if effectiveSID != "" && effectiveSID != securitySystemRID {
|
|
effective, err = winapi.StringToSid(effectiveSID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
entries := []winapi.EXPLICIT_ACCESS{{
|
|
AccessPermissions: winapi.GENERIC_ALL,
|
|
AccessMode: winapi.SET_ACCESS,
|
|
Trustee: winapi.TRUSTEE{
|
|
TrusteeForm: winapi.TRUSTEE_IS_SID,
|
|
TrusteeType: winapi.TRUSTEE_IS_WELL_KNOWN_GROUP,
|
|
TrusteeValue: winapi.TrusteeValueFromSID(systemSID),
|
|
},
|
|
}}
|
|
if effective != systemSID {
|
|
entries = append(entries, winapi.EXPLICIT_ACCESS{
|
|
AccessPermissions: winapi.GENERIC_READ,
|
|
AccessMode: winapi.SET_ACCESS,
|
|
Trustee: winapi.TRUSTEE{
|
|
TrusteeForm: winapi.TRUSTEE_IS_SID,
|
|
TrusteeType: winapi.TRUSTEE_IS_USER,
|
|
TrusteeValue: winapi.TrusteeValueFromSID(effective),
|
|
},
|
|
})
|
|
}
|
|
acl, err := winapi.ACLFromEntries(entries, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return winapi.SetNamedSecurityInfo(path, winapi.SE_FILE_OBJECT, winapi.OWNER_SECURITY_INFORMATION|winapi.DACL_SECURITY_INFORMATION|winapi.PROTECTED_DACL_SECURITY_INFORMATION, systemSID, nil, acl, nil)
|
|
}
|
|
|
|
func createStandardPipes() (*os.File, *os.File, *os.File, *os.File, *os.File, *os.File, error) {
|
|
security := &winapi.SecurityAttributes{Length: uint32(unsafe.Sizeof(winapi.SecurityAttributes{})), InheritHandle: 1}
|
|
var stdinReadHandle, stdinWriteHandle winapi.Handle
|
|
var stdoutReadHandle, stdoutWriteHandle winapi.Handle
|
|
var stderrReadHandle, stderrWriteHandle winapi.Handle
|
|
if err := winapi.CreatePipe(&stdinReadHandle, &stdinWriteHandle, security, 0); err != nil {
|
|
return nil, nil, nil, nil, nil, nil, err
|
|
}
|
|
if err := winapi.CreatePipe(&stdoutReadHandle, &stdoutWriteHandle, security, 0); err != nil {
|
|
_ = winapi.CloseHandle(stdinReadHandle)
|
|
_ = winapi.CloseHandle(stdinWriteHandle)
|
|
return nil, nil, nil, nil, nil, nil, err
|
|
}
|
|
if err := winapi.CreatePipe(&stderrReadHandle, &stderrWriteHandle, security, 0); err != nil {
|
|
for _, handle := range []winapi.Handle{stdinReadHandle, stdinWriteHandle, stdoutReadHandle, stdoutWriteHandle} {
|
|
_ = winapi.CloseHandle(handle)
|
|
}
|
|
return nil, nil, nil, nil, nil, nil, err
|
|
}
|
|
for _, handle := range []winapi.Handle{stdinWriteHandle, stdoutReadHandle, stderrReadHandle} {
|
|
if err := winapi.SetHandleInformation(handle, winapi.HANDLE_FLAG_INHERIT, 0); err != nil {
|
|
for _, closeHandle := range []winapi.Handle{stdinReadHandle, stdinWriteHandle, stdoutReadHandle, stdoutWriteHandle, stderrReadHandle, stderrWriteHandle} {
|
|
_ = winapi.CloseHandle(closeHandle)
|
|
}
|
|
return nil, nil, nil, nil, nil, nil, err
|
|
}
|
|
}
|
|
return os.NewFile(uintptr(stdinReadHandle), "rvbox-stdin-read"), os.NewFile(uintptr(stdinWriteHandle), "rvbox-stdin-write"), os.NewFile(uintptr(stdoutReadHandle), "rvbox-stdout-read"), os.NewFile(uintptr(stdoutWriteHandle), "rvbox-stdout-write"), os.NewFile(uintptr(stderrReadHandle), "rvbox-stderr-read"), os.NewFile(uintptr(stderrWriteHandle), "rvbox-stderr-write"), nil
|
|
}
|
|
|
|
func createKillOnCloseJob() (winapi.Handle, error) {
|
|
job, err := winapi.CreateJobObject(nil, nil)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
info := winapi.JOBOBJECT_EXTENDED_LIMIT_INFORMATION{}
|
|
info.BasicLimitInformation.LimitFlags = winapi.JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE
|
|
if _, err := winapi.SetInformationJobObject(job, winapi.JobObjectExtendedLimitInformation, uintptr(unsafe.Pointer(&info)), uint32(unsafe.Sizeof(info))); err != nil {
|
|
_ = winapi.CloseHandle(job)
|
|
return 0, err
|
|
}
|
|
return job, nil
|
|
}
|
|
|
|
// applyJobProfiles combines the requested dimensions and applies them before
|
|
// process creation. The profile names originate from the validated protobuf
|
|
// ExecutionSpec; unknown names and any required control without a native
|
|
// implementation are permanent pre-launch failures.
|
|
func applyJobProfiles(job winapi.Handle, configured map[string]JobProfile, requested []string) error {
|
|
if len(requested) == 0 {
|
|
return nil
|
|
}
|
|
var combined JobProfile
|
|
for _, name := range requested {
|
|
profile, ok := configured[name]
|
|
if !ok {
|
|
return fmt.Errorf("%w: execution profile %q is not configured", supervisor.ErrUnsupported, name)
|
|
}
|
|
for _, required := range profile.RequiredControls {
|
|
switch required {
|
|
case "cpu", "memory", "pids":
|
|
case "io":
|
|
return fmt.Errorf("%w: Windows Job I/O rate control is not available in this build", supervisor.ErrUnsupported)
|
|
default:
|
|
return fmt.Errorf("%w: unknown required Job control %q", supervisor.ErrUnsupported, required)
|
|
}
|
|
}
|
|
if profile.CPUPercent > combined.CPUPercent {
|
|
combined.CPUPercent = profile.CPUPercent
|
|
}
|
|
if profile.MemoryMaxBytes > 0 && (combined.MemoryMaxBytes == 0 || profile.MemoryMaxBytes < combined.MemoryMaxBytes) {
|
|
combined.MemoryMaxBytes = profile.MemoryMaxBytes
|
|
}
|
|
if profile.PIDsMax > 0 && (combined.PIDsMax == 0 || profile.PIDsMax < combined.PIDsMax) {
|
|
combined.PIDsMax = profile.PIDsMax
|
|
}
|
|
if profile.IOReadBPS > 0 || profile.IOWriteBPS > 0 {
|
|
return fmt.Errorf("%w: Windows Job I/O rate control is not available in this build", supervisor.ErrUnsupported)
|
|
}
|
|
}
|
|
limits := winapi.JOBOBJECT_EXTENDED_LIMIT_INFORMATION{}
|
|
limits.BasicLimitInformation.LimitFlags = winapi.JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE
|
|
if combined.MemoryMaxBytes > 0 {
|
|
if uint64(uintptr(combined.MemoryMaxBytes)) != combined.MemoryMaxBytes {
|
|
return fmt.Errorf("%w: memory profile exceeds native pointer size", supervisor.ErrUnsupported)
|
|
}
|
|
limits.ProcessMemoryLimit = uintptr(combined.MemoryMaxBytes)
|
|
limits.BasicLimitInformation.LimitFlags |= winapi.JOB_OBJECT_LIMIT_PROCESS_MEMORY
|
|
}
|
|
if combined.PIDsMax > 0 {
|
|
if combined.PIDsMax > uint64(^uint32(0)) {
|
|
return fmt.Errorf("%w: process-count profile exceeds Windows limit", supervisor.ErrUnsupported)
|
|
}
|
|
limits.BasicLimitInformation.ActiveProcessLimit = uint32(combined.PIDsMax)
|
|
limits.BasicLimitInformation.LimitFlags |= winapi.JOB_OBJECT_LIMIT_ACTIVE_PROCESS
|
|
}
|
|
if _, err := winapi.SetInformationJobObject(job, winapi.JobObjectExtendedLimitInformation, uintptr(unsafe.Pointer(&limits)), uint32(unsafe.Sizeof(limits))); err != nil {
|
|
return fmt.Errorf("apply Windows Job limits: %w", err)
|
|
}
|
|
if combined.CPUPercent > 0 {
|
|
// The config contract expresses CPU allowance as a percentage of one
|
|
// logical CPU (for example 200 means two logical CPUs). Windows Job
|
|
// CpuRate is hundredths of a percentage of the whole machine, so scale
|
|
// by the active processor count before applying the hard cap. A profile
|
|
// larger than this host is intentionally capped at the host capacity,
|
|
// which means it imposes no additional CPU restriction but remains a
|
|
// valid, atomically verified profile.
|
|
processors := uint64(winapi.GetActiveProcessorCount(winapi.ALL_PROCESSOR_GROUPS))
|
|
if processors == 0 {
|
|
return fmt.Errorf("%w: Windows did not report an active processor count", supervisor.ErrUnsupported)
|
|
}
|
|
if combined.CPUPercent > (^uint64(0)-processors+1)/100 {
|
|
return fmt.Errorf("%w: CPU profile %d%% overflows Windows Job rate conversion", supervisor.ErrUnsupported, combined.CPUPercent)
|
|
}
|
|
cpuRate := (combined.CPUPercent*100 + processors - 1) / processors
|
|
if cpuRate > 10000 {
|
|
cpuRate = 10000
|
|
}
|
|
cpu := struct {
|
|
ControlFlags uint32
|
|
CPUrate uint32
|
|
Weight uint32
|
|
}{ControlFlags: 0x1 | 0x4 /* ENABLE | HARD_CAP */, CPUrate: uint32(cpuRate)}
|
|
if _, err := winapi.SetInformationJobObject(job, winapi.JobObjectCpuRateControlInformation, uintptr(unsafe.Pointer(&cpu)), uint32(unsafe.Sizeof(cpu))); err != nil {
|
|
return fmt.Errorf("apply Windows Job CPU limit: %w", err)
|
|
}
|
|
var cpuReadback struct {
|
|
ControlFlags uint32
|
|
CPUrate uint32
|
|
Weight uint32
|
|
}
|
|
var cpuReturned uint32
|
|
if err := winapi.QueryInformationJobObject(job, int32(winapi.JobObjectCpuRateControlInformation), uintptr(unsafe.Pointer(&cpuReadback)), uint32(unsafe.Sizeof(cpuReadback)), &cpuReturned); err != nil {
|
|
return fmt.Errorf("verify Windows Job CPU limit: %w", err)
|
|
}
|
|
if cpuReadback.CPUrate != cpu.CPUrate || cpuReadback.ControlFlags&0x5 != 0x5 {
|
|
return errors.New("Windows Job CPU limit did not read back as requested")
|
|
}
|
|
}
|
|
// Read back every requested limit before authorization. This catches
|
|
// policy restrictions and unsupported Job implementations early.
|
|
var readback winapi.JOBOBJECT_EXTENDED_LIMIT_INFORMATION
|
|
var returned uint32
|
|
if err := winapi.QueryInformationJobObject(job, int32(winapi.JobObjectExtendedLimitInformation), uintptr(unsafe.Pointer(&readback)), uint32(unsafe.Sizeof(readback)), &returned); err != nil {
|
|
return fmt.Errorf("verify Windows Job limits: %w", err)
|
|
}
|
|
if readback.BasicLimitInformation.LimitFlags&winapi.JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE == 0 {
|
|
return errors.New("Windows Job lost kill-on-close protection")
|
|
}
|
|
if combined.MemoryMaxBytes > 0 && uint64(readback.ProcessMemoryLimit) != combined.MemoryMaxBytes {
|
|
return errors.New("Windows Job memory limit did not read back as requested")
|
|
}
|
|
if combined.PIDsMax > 0 && uint64(readback.BasicLimitInformation.ActiveProcessLimit) != combined.PIDsMax {
|
|
return errors.New("Windows Job process limit did not read back as requested")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (manager *execSupervisor) selectToken(elevated bool) (winapi.Token, supervisor.EffectiveIdentity, error) {
|
|
candidates, err := DiscoverActiveSessions()
|
|
if err != nil {
|
|
// A failed WTS enumeration is treated as no usable interactive
|
|
// session; LocalSystem still provides a deterministic service path.
|
|
candidates = nil
|
|
}
|
|
selection := Select(SelectionInput{Elevated: elevated, ActiveSessions: candidates, ActiveSystemAvailable: true, LocalServiceAvailable: true, LocalSystemAvailable: true})
|
|
attempted := make([]string, 0, len(selection.Attempts)+1)
|
|
details := make([]string, 0, len(selection.Attempts)+1)
|
|
addAttempt := func(contextName ExecutionContext, detail string) {
|
|
for _, existing := range attempted {
|
|
if existing == string(contextName) {
|
|
if detail != "" {
|
|
details = append(details, string(contextName)+": "+detail)
|
|
}
|
|
return
|
|
}
|
|
}
|
|
attempted = append(attempted, string(contextName))
|
|
if detail != "" {
|
|
details = append(details, string(contextName)+": "+detail)
|
|
}
|
|
}
|
|
for _, attempt := range selection.Attempts {
|
|
addAttempt(attempt.Context, string(attempt.Reason))
|
|
}
|
|
withEvidence := func(identity supervisor.EffectiveIdentity) supervisor.EffectiveIdentity {
|
|
identity.AttemptedContexts = append([]string(nil), attempted...)
|
|
identity.SelectionDetail = boundSelectionDetail(strings.Join(details, "; "))
|
|
return identity
|
|
}
|
|
rejection := func(cause error) error {
|
|
identity := withEvidence(supervisor.EffectiveIdentity{})
|
|
return &supervisor.StartError{Cause: cause, WindowsIdentity: identity.WindowsIdentity()}
|
|
}
|
|
if selection.Effective == nil {
|
|
if selection.Error != nil {
|
|
return 0, supervisor.EffectiveIdentity{}, rejection(selection.Error)
|
|
}
|
|
return 0, supervisor.EffectiveIdentity{}, rejection(errors.New("Windows execution context selection failed"))
|
|
}
|
|
var selected *SessionCandidate
|
|
if selection.Effective.SessionID != nil {
|
|
for index := range candidates {
|
|
if candidates[index].SessionID == *selection.Effective.SessionID {
|
|
selected = &candidates[index]
|
|
break
|
|
}
|
|
}
|
|
}
|
|
for _, attempt := range selection.Attempts {
|
|
token, identity, err := openTokenForAttempt(attempt.Context, selected)
|
|
if err == nil {
|
|
return token, withEvidence(identity), nil
|
|
}
|
|
addAttempt(attempt.Context, "native preparation failed: "+err.Error())
|
|
if !elevated {
|
|
return 0, supervisor.EffectiveIdentity{}, rejection(err)
|
|
}
|
|
}
|
|
// The pure selector stops as soon as ACTIVE_SYSTEM is available. A native
|
|
// privilege/session operation can still fail (for example, SeTcb was
|
|
// removed), so the final LOCAL_SYSTEM fallback is attempted here before
|
|
// launch preparation, never by retrying a created process.
|
|
if elevated {
|
|
addAttempt(ContextLocalSystem, "fallback")
|
|
if token, identity, err := openTokenForAttempt(ContextLocalSystem, nil); err == nil {
|
|
return token, withEvidence(identity), nil
|
|
} else {
|
|
addAttempt(ContextLocalSystem, "native preparation failed: "+err.Error())
|
|
return 0, supervisor.EffectiveIdentity{}, rejection(err)
|
|
}
|
|
}
|
|
return 0, supervisor.EffectiveIdentity{}, rejection(errors.New("all Windows execution contexts failed before launch preparation"))
|
|
}
|
|
|
|
func boundSelectionDetail(detail string) string {
|
|
if len(detail) <= 4096 {
|
|
return detail
|
|
}
|
|
return detail[:4096]
|
|
}
|
|
|
|
func openTokenForAttempt(contextName ExecutionContext, candidate *SessionCandidate) (winapi.Token, supervisor.EffectiveIdentity, error) {
|
|
switch contextName {
|
|
case ContextActiveUser, ContextActiveUserElevated, ContextActiveSystem:
|
|
if candidate == nil {
|
|
return 0, supervisor.EffectiveIdentity{}, errors.New("active execution context has no selected session")
|
|
}
|
|
var token winapi.Token
|
|
if err := winapi.WTSQueryUserToken(candidate.SessionID, &token); err != nil {
|
|
return 0, supervisor.EffectiveIdentity{}, err
|
|
}
|
|
if contextName == ContextActiveUserElevated {
|
|
if !token.IsElevated() {
|
|
linked, err := token.GetLinkedToken()
|
|
_ = token.Close()
|
|
if err != nil {
|
|
return 0, supervisor.EffectiveIdentity{}, err
|
|
}
|
|
token = linked
|
|
}
|
|
} else if contextName == ContextActiveSystem {
|
|
_ = token.Close()
|
|
serviceToken, identity, err := duplicateServiceTokenForSession(candidate.SessionID)
|
|
if err != nil {
|
|
return 0, supervisor.EffectiveIdentity{}, err
|
|
}
|
|
identity.Context = string(ContextActiveSystem)
|
|
identity.SessionUserSID = candidate.UserSID
|
|
identity.LogonSID = candidate.LogonSID
|
|
return serviceToken, identity, nil
|
|
} else if token.IsElevated() {
|
|
// A full administrator token can be returned when UAC is disabled or
|
|
// policy supplies an already-unfiltered token. Normal commands must
|
|
// still run as that user without administrator authority; create a
|
|
// restricted medium token instead of silently falling back to a
|
|
// service identity.
|
|
restricted, err := createRestrictedMediumToken(token)
|
|
_ = token.Close()
|
|
if err != nil {
|
|
return 0, supervisor.EffectiveIdentity{}, err
|
|
}
|
|
token = restricted
|
|
}
|
|
if err := verifyUserToken(token, candidate, contextName == ContextActiveUser); err != nil {
|
|
_ = token.Close()
|
|
return 0, supervisor.EffectiveIdentity{}, err
|
|
}
|
|
identity := supervisor.EffectiveIdentity{Context: string(contextName), SessionID: candidate.SessionID, SessionUserSID: candidate.UserSID, UserSID: candidate.UserSID, LogonSID: candidate.LogonSID, Elevated: contextName != ContextActiveUser, Integrity: map[bool]string{true: "high", false: "medium"}[contextName != ContextActiveUser]}
|
|
return token, identity, nil
|
|
case ContextLocalService:
|
|
token, err := logonLocalService()
|
|
if err != nil {
|
|
return 0, supervisor.EffectiveIdentity{}, err
|
|
}
|
|
return token, supervisor.EffectiveIdentity{Context: string(ContextLocalService), UserSID: securityLocalServiceRID, Elevated: false, Integrity: "medium"}, nil
|
|
case ContextLocalSystem:
|
|
return duplicateServiceToken()
|
|
default:
|
|
return 0, supervisor.EffectiveIdentity{}, errors.New("unknown Windows execution context")
|
|
}
|
|
}
|
|
|
|
// createRestrictedMediumToken turns a full administrator token into the
|
|
// normal active-user token required for elevated=false. It disables all
|
|
// privileges, disables the built-in Administrators SID, and verifies medium
|
|
// integrity before the token is returned to the launch path.
|
|
func createRestrictedMediumToken(source winapi.Token) (winapi.Token, error) {
|
|
// A full administrator token may carry several built-in groups that grant
|
|
// privileged access even after the administrator SID is hidden. Disable the
|
|
// complete known privilege-bearing built-in set before forcing medium
|
|
// integrity; membership itself is retained in the token evidence, but these
|
|
// groups cannot authorize the normal command.
|
|
disabledTypes := []winapi.WELL_KNOWN_SID_TYPE{
|
|
winapi.WinBuiltinAdministratorsSid,
|
|
winapi.WinBuiltinPowerUsersSid,
|
|
winapi.WinBuiltinAccountOperatorsSid,
|
|
winapi.WinBuiltinSystemOperatorsSid,
|
|
winapi.WinBuiltinPrintOperatorsSid,
|
|
winapi.WinBuiltinBackupOperatorsSid,
|
|
winapi.WinBuiltinRemoteDesktopUsersSid,
|
|
winapi.WinBuiltinNetworkConfigurationOperatorsSid,
|
|
winapi.WinBuiltinPerfMonitoringUsersSid,
|
|
winapi.WinBuiltinPerfLoggingUsersSid,
|
|
winapi.WinBuiltinDCOMUsersSid,
|
|
winapi.WinBuiltinCryptoOperatorsSid,
|
|
winapi.WinBuiltinHyperVAdminsSid,
|
|
winapi.WinBuiltinRemoteManagementUsersSid,
|
|
}
|
|
disabled := make([]winapi.SIDAndAttributes, len(disabledTypes))
|
|
for index, sidType := range disabledTypes {
|
|
sid, err := winapi.CreateWellKnownSid(sidType)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
disabled[index] = winapi.SIDAndAttributes{Sid: sid}
|
|
}
|
|
var restricted winapi.Token
|
|
result, _, callErr := procCreateRestrictedToken.Call(
|
|
uintptr(source), disableMaxPrivilege,
|
|
uintptr(len(disabled)), uintptr(unsafe.Pointer(&disabled[0])),
|
|
0, 0,
|
|
0, 0,
|
|
uintptr(unsafe.Pointer(&restricted)),
|
|
)
|
|
if result == 0 {
|
|
if callErr != syscall.Errno(0) {
|
|
return 0, callErr
|
|
}
|
|
return 0, syscall.GetLastError()
|
|
}
|
|
if err := setMediumIntegrity(restricted); err != nil {
|
|
_ = restricted.Close()
|
|
return 0, err
|
|
}
|
|
if restricted.IsElevated() {
|
|
_ = restricted.Close()
|
|
return 0, errors.New("restricted active-user token remained elevated")
|
|
}
|
|
if err := verifyTokenIntegrity(restricted, mediumIntegrityRID); err != nil {
|
|
_ = restricted.Close()
|
|
return 0, err
|
|
}
|
|
return restricted, nil
|
|
}
|
|
|
|
func setMediumIntegrity(token winapi.Token) error {
|
|
mediumSID, err := winapi.StringToSid("S-1-16-8192")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
sidLength := winapi.GetLengthSid(mediumSID)
|
|
headerSize := uint32(unsafe.Sizeof(winapi.Tokenmandatorylabel{}))
|
|
buffer := make([]byte, headerSize+sidLength)
|
|
label := (*winapi.Tokenmandatorylabel)(unsafe.Pointer(&buffer[0]))
|
|
label.Label.Sid = (*winapi.SID)(unsafe.Pointer(&buffer[headerSize]))
|
|
label.Label.Attributes = winapi.SE_GROUP_INTEGRITY | winapi.SE_GROUP_INTEGRITY_ENABLED
|
|
copy(buffer[headerSize:], unsafe.Slice((*byte)(unsafe.Pointer(mediumSID)), sidLength))
|
|
return winapi.SetTokenInformation(token, winapi.TokenIntegrityLevel, &buffer[0], uint32(len(buffer)))
|
|
}
|
|
|
|
func verifyUserToken(token winapi.Token, candidate *SessionCandidate, normal bool) error {
|
|
if err := verifyTokenIdentity(token, candidate.UserSID, candidate.SessionID); err != nil {
|
|
return fmt.Errorf("verify active token identity: %w", err)
|
|
}
|
|
if normal && token.IsElevated() {
|
|
return errors.New("normal active-user token is elevated")
|
|
}
|
|
if normal {
|
|
if err := verifyTokenIntegrity(token, mediumIntegrityRID); err != nil {
|
|
return fmt.Errorf("normal active-user token is not medium integrity: %w", err)
|
|
}
|
|
} else {
|
|
if !token.IsElevated() {
|
|
return errors.New("elevated active-user token is not elevated")
|
|
}
|
|
if err := verifyTokenIntegrity(token, highIntegrityRID); err != nil {
|
|
return fmt.Errorf("elevated active-user token is not high integrity: %w", err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func duplicateServiceToken() (winapi.Token, supervisor.EffectiveIdentity, error) {
|
|
return duplicateServiceTokenForSession(0)
|
|
}
|
|
|
|
func duplicateServiceTokenForSession(sessionID uint32) (winapi.Token, supervisor.EffectiveIdentity, error) {
|
|
var source winapi.Token
|
|
if err := winapi.OpenProcessToken(winapi.CurrentProcess(), winapi.TOKEN_ALL_ACCESS, &source); err != nil {
|
|
return 0, supervisor.EffectiveIdentity{}, err
|
|
}
|
|
defer source.Close()
|
|
var restorePrivilege func()
|
|
if sessionID != 0 {
|
|
var err error
|
|
restorePrivilege, err = enableTokenPrivilege(source, "SeTcbPrivilege")
|
|
if err != nil {
|
|
return 0, supervisor.EffectiveIdentity{}, fmt.Errorf("enable SeTcbPrivilege for active session token: %w", err)
|
|
}
|
|
defer restorePrivilege()
|
|
}
|
|
var target winapi.Token
|
|
if err := winapi.DuplicateTokenEx(source, winapi.TOKEN_ALL_ACCESS, nil, winapi.SecurityImpersonation, winapi.TokenPrimary, &target); err != nil {
|
|
return 0, supervisor.EffectiveIdentity{}, err
|
|
}
|
|
if sessionID != 0 {
|
|
if err := winapi.SetTokenInformation(target, winapi.TokenSessionId, (*byte)(unsafe.Pointer(&sessionID)), uint32(unsafe.Sizeof(sessionID))); err != nil {
|
|
_ = target.Close()
|
|
return 0, supervisor.EffectiveIdentity{}, err
|
|
}
|
|
var actualSession uint32
|
|
var returned uint32
|
|
if err := winapi.GetTokenInformation(target, winapi.TokenSessionId, (*byte)(unsafe.Pointer(&actualSession)), uint32(unsafe.Sizeof(actualSession)), &returned); err != nil {
|
|
_ = target.Close()
|
|
return 0, supervisor.EffectiveIdentity{}, fmt.Errorf("verify active SYSTEM token session: %w", err)
|
|
}
|
|
if returned != uint32(unsafe.Sizeof(actualSession)) || actualSession != sessionID {
|
|
_ = target.Close()
|
|
return 0, supervisor.EffectiveIdentity{}, errors.New("active SYSTEM token session did not read back as requested")
|
|
}
|
|
}
|
|
if err := verifyTokenIdentity(target, securitySystemRID, sessionID); err != nil {
|
|
_ = target.Close()
|
|
return 0, supervisor.EffectiveIdentity{}, fmt.Errorf("duplicated service token is invalid: %w", err)
|
|
}
|
|
return target, supervisor.EffectiveIdentity{Context: string(ContextLocalSystem), SessionID: sessionID, UserSID: securitySystemRID, Elevated: true, Integrity: "system"}, nil
|
|
}
|
|
|
|
// enableTokenPrivilege enables one privilege only for the short operation
|
|
// that needs it and returns a best-effort restoration closure. Windows may
|
|
// report ERROR_NOT_ALL_ASSIGNED even when AdjustTokenPrivileges itself
|
|
// succeeds; treat that as a hard capability failure rather than silently
|
|
// creating a Session-0 token for an active-session request.
|
|
func enableTokenPrivilege(token winapi.Token, name string) (func(), error) {
|
|
privilegeName, err := winapi.UTF16PtrFromString(name)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var luid winapi.LUID
|
|
if err := winapi.LookupPrivilegeValue(nil, privilegeName, &luid); err != nil {
|
|
return nil, err
|
|
}
|
|
state := winapi.Tokenprivileges{PrivilegeCount: 1}
|
|
state.Privileges[0] = winapi.LUIDAndAttributes{Luid: luid, Attributes: winapi.SE_PRIVILEGE_ENABLED}
|
|
var previous winapi.Tokenprivileges
|
|
var returned uint32
|
|
if err := winapi.AdjustTokenPrivileges(token, false, &state, uint32(unsafe.Sizeof(state)), &previous, &returned); err != nil {
|
|
return nil, err
|
|
}
|
|
if lastErr := winapi.GetLastError(); lastErr != nil && !errors.Is(lastErr, syscall.Errno(0)) {
|
|
return nil, lastErr
|
|
}
|
|
restore := func() {
|
|
if returned == 0 {
|
|
return
|
|
}
|
|
_ = winapi.AdjustTokenPrivileges(token, false, &previous, uint32(unsafe.Sizeof(previous)), nil, nil)
|
|
}
|
|
return restore, nil
|
|
}
|
|
|
|
func logonLocalService() (winapi.Token, error) {
|
|
account, _ := syscall.UTF16PtrFromString("LocalService")
|
|
domainName, _ := syscall.UTF16PtrFromString("NT AUTHORITY")
|
|
var token winapi.Token
|
|
r, _, callErr := procLogonUserW.Call(uintptr(unsafe.Pointer(account)), uintptr(unsafe.Pointer(domainName)), 0, logon32LogonService, logon32ProviderDefault, uintptr(unsafe.Pointer(&token)))
|
|
if r == 0 {
|
|
if callErr != syscall.Errno(0) {
|
|
return 0, callErr
|
|
}
|
|
return 0, syscall.GetLastError()
|
|
}
|
|
if err := verifyTokenIdentity(token, securityLocalServiceRID, 0); err != nil {
|
|
_ = token.Close()
|
|
return 0, fmt.Errorf("LocalService token failed identity verification: %w", err)
|
|
}
|
|
return token, nil
|
|
}
|
|
|
|
func verifyTokenIdentity(token winapi.Token, expectedSID string, expectedSession uint32) error {
|
|
user, err := token.GetTokenUser()
|
|
if err != nil || user.User.Sid == nil {
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return errors.New("token has no user SID")
|
|
}
|
|
if user.User.Sid.String() != expectedSID {
|
|
return fmt.Errorf("token user SID %q does not match %q", user.User.Sid.String(), expectedSID)
|
|
}
|
|
actualSession, err := tokenInformationUint32(token, winapi.TokenSessionId)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if actualSession != expectedSession {
|
|
return fmt.Errorf("token session %d does not match %d", actualSession, expectedSession)
|
|
}
|
|
tokenType, err := tokenInformationUint32(token, winapi.TokenType)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tokenType != winapi.TokenPrimary {
|
|
return fmt.Errorf("token type %d is not primary", tokenType)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func tokenInformationUint32(token winapi.Token, class uint32) (uint32, error) {
|
|
buffer, err := tokenInformation(token, class)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
if len(buffer) != int(unsafe.Sizeof(uint32(0))) {
|
|
return 0, errors.New("token information has unexpected size")
|
|
}
|
|
return *(*uint32)(unsafe.Pointer(&buffer[0])), nil
|
|
}
|
|
|
|
func tokenInformation(token winapi.Token, class uint32) ([]byte, error) {
|
|
var returned uint32
|
|
err := winapi.GetTokenInformation(token, class, nil, 0, &returned)
|
|
if returned == 0 && err != nil && !errors.Is(err, winapi.ERROR_INSUFFICIENT_BUFFER) {
|
|
return nil, err
|
|
}
|
|
if returned == 0 {
|
|
return nil, errors.New("token information returned an empty buffer")
|
|
}
|
|
buffer := make([]byte, returned)
|
|
if err := winapi.GetTokenInformation(token, class, &buffer[0], returned, &returned); err != nil {
|
|
return nil, err
|
|
}
|
|
if returned > uint32(len(buffer)) {
|
|
return nil, errors.New("token information length changed during query")
|
|
}
|
|
return buffer[:returned], nil
|
|
}
|
|
|
|
func verifyTokenIntegrity(token winapi.Token, minimumRID uint32) error {
|
|
buffer, err := tokenInformation(token, winapi.TokenIntegrityLevel)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(buffer) < int(unsafe.Sizeof(winapi.Tokenmandatorylabel{})) {
|
|
return errors.New("token integrity information is truncated")
|
|
}
|
|
label := (*winapi.Tokenmandatorylabel)(unsafe.Pointer(&buffer[0]))
|
|
if label.Label.Sid == nil || label.Label.Sid.SubAuthorityCount() == 0 {
|
|
return errors.New("token integrity SID is missing")
|
|
}
|
|
level := label.Label.Sid.SubAuthority(uint32(label.Label.Sid.SubAuthorityCount()) - 1)
|
|
if level < minimumRID {
|
|
return fmt.Errorf("token integrity level 0x%x is below 0x%x", level, minimumRID)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (manager *execSupervisor) Signal(ctx context.Context, process supervisor.Process, signal supervisor.SignalKind) (supervisor.SignalOutcome, error) {
|
|
if process == nil {
|
|
return supervisor.SignalOutcome{}, ErrProcessNotFound
|
|
}
|
|
native, ok := process.(*execProcess)
|
|
if !ok || native.killFn == nil {
|
|
return supervisor.SignalOutcome{}, ErrProcessNotFound
|
|
}
|
|
if signal != supervisor.SignalTerm && signal != supervisor.SignalKill {
|
|
return supervisor.SignalOutcome{}, errors.New("unsupported signal")
|
|
}
|
|
if signal == supervisor.SignalTerm {
|
|
// Each command has its own hidden console. The helper path is kept in
|
|
// this short-lived call and is deliberately best-effort: a session that
|
|
// has already exited or a policy that denies AttachConsole is recorded,
|
|
// then the bounded grace period ends in an explicit Job kill.
|
|
breakDelivered, breakErr := sendControlBreak(native.pid)
|
|
if breakErr != nil && ctx.Err() != nil {
|
|
return supervisor.SignalOutcome{}, ctx.Err()
|
|
}
|
|
if breakDelivered {
|
|
select {
|
|
case <-native.done:
|
|
return supervisor.SignalOutcome{Delivered: true, Detail: "CTRL_BREAK delivered", ObservedAt: manager.options.Now()}, nil
|
|
default:
|
|
}
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return supervisor.SignalOutcome{}, ctx.Err()
|
|
case <-native.done:
|
|
return supervisor.SignalOutcome{Delivered: breakDelivered, Detail: "command exited after TERM", ObservedAt: manager.options.Now()}, nil
|
|
case <-time.After(manager.options.WindowsTermGrace):
|
|
}
|
|
}
|
|
if err := native.killFn(1); err != nil {
|
|
return supervisor.SignalOutcome{}, err
|
|
}
|
|
detail := "Windows Job terminated"
|
|
if signal == supervisor.SignalTerm {
|
|
detail = "CTRL_BREAK grace expired; Windows Job terminated"
|
|
}
|
|
return supervisor.SignalOutcome{Delivered: true, Escalated: signal == supervisor.SignalTerm, Detail: detail, ObservedAt: manager.options.Now()}, nil
|
|
}
|
|
|
|
// sendControlBreak is the native equivalent of the signal-helper mode. The
|
|
// production helper is normally a separate short-lived rvbox.exe invocation;
|
|
// this direct implementation keeps the same verified PID/console boundary
|
|
// for the first service build and never addresses a process by a caller-
|
|
// supplied PID. The PID comes only from execProcess metadata.
|
|
func sendControlBreak(pid uint32) (bool, error) {
|
|
if pid == 0 {
|
|
return false, errors.New("command has no verified console PID")
|
|
}
|
|
if result, _, err := procAttachConsole.Call(uintptr(pid)); result == 0 {
|
|
if err == syscall.Errno(0) {
|
|
err = syscall.GetLastError()
|
|
}
|
|
return false, err
|
|
}
|
|
defer procFreeConsole.Call()
|
|
// Prevent the service/helper itself from acting on the generated event.
|
|
procSetCtrlHandler.Call(0, 1)
|
|
defer procSetCtrlHandler.Call(0, 0)
|
|
result, _, err := procGenerateCtrlEvent.Call(1 /* CTRL_BREAK_EVENT */, 0)
|
|
if result == 0 {
|
|
if err == syscall.Errno(0) {
|
|
err = syscall.GetLastError()
|
|
}
|
|
return false, err
|
|
}
|
|
return true, nil
|
|
}
|
|
|
|
func (manager *execSupervisor) Snapshot(ctx context.Context, process supervisor.Process) (supervisor.ResourceSnapshot, error) {
|
|
if process == nil {
|
|
return supervisor.ResourceSnapshot{}, ErrProcessNotFound
|
|
}
|
|
native, ok := process.(*execProcess)
|
|
if !ok || native.cmd == nil || native.cmd.Process == nil {
|
|
return supervisor.ResourceSnapshot{}, ErrProcessNotFound
|
|
}
|
|
if native.snapshotFn != nil {
|
|
select {
|
|
case <-ctx.Done():
|
|
return supervisor.ResourceSnapshot{}, ctx.Err()
|
|
default:
|
|
}
|
|
return native.snapshotFn()
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return supervisor.ResourceSnapshot{}, ctx.Err()
|
|
default:
|
|
}
|
|
return supervisor.ResourceSnapshot{ProcessCount: 1, ObservedAt: manager.options.Now(), Complete: false, Detail: "Windows Job accounting is available after native completion integration"}, nil
|
|
}
|
|
|
|
type jobBasicAndIOAccounting struct {
|
|
TotalUserTime int64
|
|
TotalKernelTime int64
|
|
ThisPeriodTotalUserTime int64
|
|
ThisPeriodTotalKernelTime int64
|
|
TotalPageFaultCount uint32
|
|
TotalProcesses uint32
|
|
ActiveProcesses uint32
|
|
TotalTerminatedProcesses uint32
|
|
IO winapi.IO_COUNTERS
|
|
}
|
|
|
|
func queryJobSnapshot(job winapi.Handle, now time.Time) (supervisor.ResourceSnapshot, error) {
|
|
if job == 0 || job == winapi.InvalidHandle {
|
|
return supervisor.ResourceSnapshot{}, ErrProcessNotFound
|
|
}
|
|
var accounting jobBasicAndIOAccounting
|
|
var returned uint32
|
|
if err := winapi.QueryInformationJobObject(job, int32(winapi.JobObjectBasicAndIoAccountingInformation), uintptr(unsafe.Pointer(&accounting)), uint32(unsafe.Sizeof(accounting)), &returned); err != nil {
|
|
return supervisor.ResourceSnapshot{}, err
|
|
}
|
|
var limits winapi.JOBOBJECT_EXTENDED_LIMIT_INFORMATION
|
|
if err := winapi.QueryInformationJobObject(job, int32(winapi.JobObjectExtendedLimitInformation), uintptr(unsafe.Pointer(&limits)), uint32(unsafe.Sizeof(limits)), &returned); err != nil {
|
|
return supervisor.ResourceSnapshot{}, err
|
|
}
|
|
userKernel := accounting.TotalUserTime + accounting.TotalKernelTime
|
|
var cpu time.Duration
|
|
if userKernel > 0 && userKernel <= int64(^uint64(0)>>1)/100 {
|
|
cpu = time.Duration(userKernel) * 100 * time.Nanosecond
|
|
}
|
|
return supervisor.ResourceSnapshot{CPUTime: cpu, ResidentBytes: uint64(limits.PeakJobMemoryUsed), IOReadBytes: accounting.IO.ReadTransferCount, IOWriteBytes: accounting.IO.WriteTransferCount, ProcessCount: uint64(accounting.ActiveProcesses), ObservedAt: now, Complete: accounting.ActiveProcesses == 0, Detail: "Windows Job accounting"}, nil
|
|
}
|
|
|
|
func (manager *execSupervisor) StopAll(_ context.Context) error {
|
|
manager.mu.Lock()
|
|
processes := make([]*execProcess, 0, len(manager.active))
|
|
for _, process := range manager.active {
|
|
processes = append(processes, process)
|
|
}
|
|
manager.mu.Unlock()
|
|
for _, process := range processes {
|
|
if process.killFn != nil {
|
|
if err := process.killFn(1); err != nil && !errors.Is(err, winapi.ERROR_INVALID_HANDLE) {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}
|