Read mutable Windows JSON without locking atomic writers and salvage Defender evidence

This commit is contained in:
omar
2026-09-25 22:11:52 +03:00
parent 95d11b54d4
commit 9bb6ed29ac
9 changed files with 754 additions and 31 deletions
+6
View File
@@ -137,6 +137,12 @@ docker compose --profile execution up -d worker watchdog
After source validation, admin queues qualification, reviews actual evidence, and publishes only a passed immutable revision. Missing account credentials/script policy/signing trust, stale Defender or inaccessible QGA remain actionable prerequisites. Never bypass protections or change a queued source revision to make admission pass.
## Windows JSON observation and partial detections
Mutable Windows JSON is read through one bounded, read-only PowerShell snapshot launched over QGA. The filename is encoded data, not shell syntax; the reader opens with `FileShare.ReadWrite | FileShare.Delete`, reads at most 4 MiB from one file generation and closes before returning base64 output. It does not execute file contents, bypass signing policy, or retry an ambiguous launch. Binary/log transport keeps its existing bounded QGA file reads. Holding `guest-file-open` across transport round trips can otherwise block the signed runner's atomic `File.Replace` even after its finite contention retries.
On abnormal completion the worker retains the incomplete outcome and telemetry status, but separately salvages positively correlated Defender detections from the collector. Receipt command/attempt/job, boot, timestamps and exact sample resource must match; old/unrelated detections never imply a verdict, and missing evidence never implies clean. Already captured final reports and positive detections are not overwritten by stale runner status. Historical attempts are not rewritten by this change.
## Grub artifact collection
`settings.grub_paths` is optional immutable Job data, not a worker host path. The interactive runner captures literal local files under the already-selected user/admin token after observation; the SYSTEM telemetry collector exports **only fixed numbered staged files** into the protected control directory after no-reparse/hardlink and hash/size validation. It never opens requested Grub paths as SYSTEM. Requested paths are never resolved on the CT, PVE host or workstation, nor interpolated into shell source.
+73 -2
View File
@@ -1,6 +1,7 @@
package pve
import (
"bytes"
"context"
"encoding/base64"
"encoding/json"
@@ -8,13 +9,19 @@ import (
"io"
"net/url"
"strconv"
"strings"
"time"
"unicode/utf8"
)
// A base64-encoded chunk plus its JSON envelope fits well within responseLimit.
const FileChunk = 1 << 20
const ControlLimit = 40 << 10
// WindowsJSONLimit keeps the base64 snapshot below QGA's 16 MiB output cap and
// the PVE response cap, including both JSON envelopes.
const WindowsJSONLimit = 4 << 20
type Bool bool
func (b *Bool) UnmarshalJSON(v []byte) error {
@@ -78,6 +85,66 @@ func (c *Client) ReadFile(ctx context.Context, node string, id int, path string,
}
}
}
// WindowsJSONReadCommand builds a read-only command; the path is UTF-8 base64
// data, never PowerShell syntax. Exported so the exact production command can be
// exercised locally on Windows without access to a guest or PVE credentials.
func WindowsJSONReadCommand(path string, limit int64) ([]string, error) {
if limit < 1 || limit > WindowsJSONLimit || path == "" || strings.ContainsRune(path, '\x00') || !utf8.ValidString(path) {
return nil, errors.New("invalid Windows JSON read bound or path")
}
script := `$ErrorActionPreference='Stop'; try {
$path=[Text.Encoding]::UTF8.GetString([Convert]::FromBase64String('` + base64.StdEncoding.EncodeToString([]byte(path)) + `'));
$limit=` + strconv.FormatInt(limit, 10) + `;
$stream=[IO.File]::Open($path,[IO.FileMode]::Open,[IO.FileAccess]::Read,([IO.FileShare]::ReadWrite -bor [IO.FileShare]::Delete));
try {
$length=$stream.Length;
if ($length -lt 1 -or $length -gt $limit) { throw 'JSON snapshot exceeds bound or is empty' };
$buffer=New-Object byte[] ([int]$length);
$offset=0;
while ($offset -lt $buffer.Length) {
$count=$stream.Read($buffer,$offset,$buffer.Length-$offset);
if ($count -le 0) { throw 'JSON snapshot shortened while reading' };
$offset+=$count;
};
if ($stream.ReadByte() -ne -1) { throw 'JSON snapshot grew while reading' };
} finally { $stream.Dispose() };
[Console]::Out.Write([Convert]::ToBase64String($buffer));
} catch { [Console]::Error.WriteLine('Windows JSON snapshot read failed'); exit 1 }`
return []string{`C:\Windows\System32\WindowsPowerShell\v1.0\powershell.exe`, "-NoLogo", "-NoProfile", "-NonInteractive", "-Command", script}, nil
}
// ReadWindowsJSON captures one open file generation, closes it before emitting
// output, and never denies the atomic writer's delete/replace sharing. Unlike
// ReadFile, it neither holds a guest file handle across RPCs nor reopens chunks.
func (c *Client) ReadWindowsJSON(ctx context.Context, node string, id int, path string, limit int64) ([]byte, error) {
command, err := WindowsJSONReadCommand(path, limit)
if err != nil {
return nil, err
}
pid, err := c.Exec(ctx, node, id, command)
if err != nil {
// An ambiguous process launch is never retried.
return nil, err
}
out, err := c.execWait(ctx, node, id, pid, base64.StdEncoding.EncodedLen(int(limit)))
if err != nil {
return nil, err
}
return decodeWindowsJSON(out, limit)
}
func decodeWindowsJSON(out []byte, limit int64) ([]byte, error) {
if len(out) == 0 || len(out) > base64.StdEncoding.EncodedLen(int(limit)) || bytes.IndexAny(out, "\r\n") >= 0 {
return nil, errors.New("invalid bounded Windows JSON snapshot")
}
raw := make([]byte, base64.StdEncoding.DecodedLen(len(out)))
n, err := base64.StdEncoding.Strict().Decode(raw, out)
if err != nil || int64(n) > limit || !json.Valid(raw[:n]) {
return nil, errors.New("invalid bounded Windows JSON snapshot")
}
return raw[:n], nil
}
func (c *Client) WriteControl(ctx context.Context, node string, id int, path string, b []byte) error {
if len(b) == 0 || len(b) > ControlLimit || !json.Valid(b) {
return errors.New("invalid bounded control JSON")
@@ -115,6 +182,10 @@ type GuestExitError struct{ Code int }
func (e *GuestExitError) Error() string { return "guest control command failed" }
func (c *Client) ExecWait(ctx context.Context, node string, id, pid int) ([]byte, error) {
return c.execWait(ctx, node, id, pid, ControlLimit)
}
func (c *Client) execWait(ctx context.Context, node string, id, pid, outputLimit int) ([]byte, error) {
t := time.NewTicker(time.Second)
defer t.Stop()
for {
@@ -128,8 +199,8 @@ func (c *Client) ExecWait(ctx context.Context, node string, id, pid int) ([]byte
return nil, err
}
if Text(r.Exited) == "1" || Text(r.Exited) == "true" {
if r.Truncated || len(r.Out) > ControlLimit {
return nil, errors.New("QGA control output exceeded bound")
if r.Truncated || len(r.Out) > outputLimit {
return nil, errors.New("QGA command output exceeded bound")
}
if r.ExitCode != 0 {
return []byte(r.Out), &GuestExitError{Code: r.ExitCode}
+123
View File
@@ -0,0 +1,123 @@
package pve
import (
"bytes"
"context"
"encoding/base64"
"errors"
"net/http"
"strings"
"testing"
)
func TestReadWindowsJSONRejectsUntrustedOutput(t *testing.T) {
for _, tc := range []struct {
name string
out string
limit int64
truncated bool
exit int
}{
{name: "empty", limit: 32},
{name: "malformed base64", out: "!!!!", limit: 32},
{name: "noncanonical base64", out: "e31=", limit: 32},
{name: "unexpected newline", out: "e30=\n", limit: 32},
{name: "malformed JSON", out: base64.StdEncoding.EncodeToString([]byte(`{"unfinished":`)), limit: 32},
{name: "decoded overflow", out: "e30=", limit: 1},
{name: "encoded overflow", out: strings.Repeat("A", 48), limit: 32},
{name: "truncated", out: "e30=", limit: 32, truncated: true},
{name: "guest read error", out: "e30=", limit: 32, exit: 1},
} {
t.Run(tc.name, func(t *testing.T) {
c := fixture(t, func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/api2/json/nodes/node/qemu/9001/agent/exec":
data(w, map[string]any{"pid": 17})
case "/api2/json/nodes/node/qemu/9001/agent/exec-status":
data(w, map[string]any{"exited": true, "exitcode": tc.exit, "out-data": tc.out, "out-truncated": tc.truncated})
default:
t.Errorf("unexpected transport path: %s", r.URL.Path)
w.WriteHeader(http.StatusBadRequest)
}
})
raw, err := c.ReadWindowsJSON(context.Background(), "node", 9001, `C:\status.json`, tc.limit)
if err == nil || raw != nil {
t.Fatalf("untrusted output accepted: raw=%q err=%v", raw, err)
}
if tc.exit != 0 {
var exited *GuestExitError
if !errors.As(err, &exited) || exited.Code != tc.exit {
t.Fatalf("confirmed guest exit lost: %v", err)
}
}
})
}
}
func TestReadWindowsJSONAcceptsMaximumSnapshot(t *testing.T) {
payload := []byte(`"` + strings.Repeat("x", WindowsJSONLimit-2) + `"`)
c := fixture(t, func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/api2/json/nodes/node/qemu/9001/agent/exec":
data(w, map[string]any{"pid": 17})
case "/api2/json/nodes/node/qemu/9001/agent/exec-status":
data(w, map[string]any{"exited": 1, "exitcode": 0, "out-data": base64.StdEncoding.EncodeToString(payload)})
default:
t.Errorf("unexpected transport path: %s", r.URL.Path)
w.WriteHeader(http.StatusBadRequest)
}
})
raw, err := c.ReadWindowsJSON(context.Background(), "node", 9001, `C:\status.json`, WindowsJSONLimit)
if err != nil || !bytes.Equal(raw, payload) {
t.Fatalf("maximum JSON snapshot changed or rejected: size=%d err=%v", len(raw), err)
}
}
func TestReadWindowsJSONNeverRetriesAmbiguousLaunch(t *testing.T) {
calls := 0
c := fixture(t, func(w http.ResponseWriter, r *http.Request) {
calls++
w.WriteHeader(http.StatusBadGateway)
})
raw, err := c.ReadWindowsJSON(context.Background(), "node", 9001, `C:\status.json`, 32)
if err == nil || raw != nil || calls != 1 {
t.Fatalf("ambiguous launch retried or accepted: calls=%d err=%v", calls, err)
}
}
func TestReadWindowsJSONRejectsInvalidRequestsBeforeLaunch(t *testing.T) {
c := fixture(t, func(w http.ResponseWriter, r *http.Request) {
t.Error("invalid snapshot request reached guest")
w.WriteHeader(http.StatusBadRequest)
})
for _, tc := range []struct {
path string
limit int64
}{
{`C:\status.json`, 0},
{`C:\status.json`, WindowsJSONLimit + 1},
{"", 32},
{"C:\\bad\x00.json", 32},
{"C:\\bad\xff.json", 32},
} {
if raw, err := c.ReadWindowsJSON(context.Background(), "node", 9001, tc.path, tc.limit); err == nil || raw != nil {
t.Fatalf("invalid request accepted: path=%q limit=%d", tc.path, tc.limit)
}
}
}
func TestExecWaitRetainsControlOutputBound(t *testing.T) {
for _, n := range []int{ControlLimit, ControlLimit + 1} {
c := fixture(t, func(w http.ResponseWriter, r *http.Request) {
data(w, map[string]any{"exited": true, "exitcode": 0, "out-data": strings.Repeat("x", n)})
})
out, err := c.ExecWait(context.Background(), "node", 9001, 17)
if n == ControlLimit {
if err != nil || len(out) != n {
t.Fatalf("control output at bound rejected: size=%d err=%v", len(out), err)
}
} else if err == nil || out != nil {
t.Fatalf("oversized control output accepted: size=%d err=%v", len(out), err)
}
}
}
+112
View File
@@ -0,0 +1,112 @@
package pve
import (
"bytes"
"context"
"io"
"os"
"os/exec"
"path/filepath"
"strings"
"testing"
"time"
)
func TestWindowsJSONCommandReadsLiteralPath(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "данные';New-Item injected;'$(1+1).json")
payload := []byte(`{"text":"Unicode and PowerShell metacharacters are path data"}`)
if err := os.WriteFile(path, payload, 0600); err != nil {
t.Fatal(err)
}
args, err := WindowsJSONReadCommand(path, int64(len(payload)))
if err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
cmd := exec.CommandContext(ctx, args[0], args[1:]...)
cmd.Dir = dir
out, err := cmd.Output()
if err != nil {
t.Fatal(err)
}
raw, err := decodeWindowsJSON(out, int64(len(payload)))
if err != nil || !bytes.Equal(raw, payload) {
t.Fatalf("literal path not preserved: raw=%q err=%v", raw, err)
}
if _, err := os.Stat(filepath.Join(dir, "injected")); !os.IsNotExist(err) {
t.Fatalf("path was executed instead of read: %v", err)
}
}
func TestWindowsJSONCommandRejectsOversizedFile(t *testing.T) {
path := filepath.Join(t.TempDir(), "status.json")
if err := os.WriteFile(path, []byte(`{"too":"large"}`), 0600); err != nil {
t.Fatal(err)
}
args, err := WindowsJSONReadCommand(path, 2)
if err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
out, err := exec.CommandContext(ctx, args[0], args[1:]...).Output()
if err == nil || len(out) != 0 {
t.Fatalf("oversized file emitted as evidence: stdout=%q err=%v", out, err)
}
}
func TestWindowsJSONCommandClosesSnapshotBeforeOutput(t *testing.T) {
path := filepath.Join(t.TempDir(), "status.json")
payload := []byte(`"` + strings.Repeat("x", WindowsJSONLimit-2) + `"`)
if err := os.WriteFile(path, payload, 0600); err != nil {
t.Fatal(err)
}
args, err := WindowsJSONReadCommand(path, WindowsJSONLimit)
if err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
cmd := exec.CommandContext(ctx, args[0], args[1:]...)
pipe, err := cmd.StdoutPipe()
if err != nil {
t.Fatal(err)
}
defer pipe.Close()
if err := cmd.Start(); err != nil {
t.Fatal(err)
}
defer func() {
_ = cmd.Process.Kill()
_ = cmd.Wait()
}()
prefix := make([]byte, 1)
if _, err := io.ReadFull(pipe, prefix); err != nil {
t.Fatal(err)
}
// The large output is still blocked on its pipe. An exclusive handle is
// stronger than delete sharing: it fails if the reader kept any file open.
check := exec.CommandContext(ctx, args[0], "-NoProfile", "-NonInteractive", "-Command", `$ErrorActionPreference='Stop'; $s=[IO.File]::Open($env:OTCHE_JSON_TEST_PATH,[IO.FileMode]::Open,[IO.FileAccess]::ReadWrite,[IO.FileShare]::None); $s.Dispose()`)
check.Env = append(os.Environ(), "OTCHE_JSON_TEST_PATH="+path)
if out, err := check.CombinedOutput(); err != nil {
t.Fatalf("snapshot file still held during output: %s: %v", out, err)
}
// Replacing the path while output is in flight must not change the captured
// generation or mix newer bytes into the original JSON.
if err := os.WriteFile(path, []byte(`{"generation":"new"}`), 0600); err != nil {
t.Fatal(err)
}
rest, err := io.ReadAll(pipe)
if err != nil {
t.Fatal(err)
}
if err := cmd.Wait(); err != nil {
t.Fatal(err)
}
raw, err := decodeWindowsJSON(append(prefix, rest...), WindowsJSONLimit)
if err != nil || !bytes.Equal(raw, payload) {
t.Fatalf("snapshot mixed file generations: size=%d err=%v", len(raw), err)
}
}
+18 -7
View File
@@ -71,12 +71,11 @@ func guestScript(ctx context.Context, c *pve.Client, b Binding, l allocation, na
return c.ExecWait(ctx, b.Node, l.VMID, pid)
}
func readGuestJSON(ctx context.Context, c *pve.Client, b Binding, l allocation, path string, v any) error {
var raw bytes.Buffer
_, err := c.ReadFile(ctx, b.Node, l.VMID, path, &raw, 4<<20)
raw, err := c.ReadWindowsJSON(ctx, b.Node, l.VMID, path, pve.WindowsJSONLimit)
if err != nil {
return err
}
if json.Unmarshal(raw.Bytes(), v) != nil {
if json.Unmarshal(raw, v) != nil {
return errors.New("guest control file is not valid bounded JSON")
}
return nil
@@ -297,7 +296,7 @@ func (e *engine) execute(ctx context.Context, a attempt) (outcome, findings, tel
}
if l.ID != "" && !collected {
salvage, done := context.WithTimeout(final, 15*time.Second)
partial, detected := e.salvagePartial(salvage, a, l, b, cs, dir)
partial, detected := e.salvagePartial(salvage, a, l, b, cs, dir, report, findings)
if partial != nil {
report = partial
telemetry = "partial"
@@ -546,8 +545,8 @@ func (e *engine) execute(ctx context.Context, a attempt) (outcome, findings, tel
if err = cs.runtime.WriteControl(delivery, b.Node, l.VMID, controlPath, control); err != nil {
return outcome, findings, telemetry, cleanup, nil, err
}
var verify bytes.Buffer
if _, err = cs.runtime.ReadFile(delivery, b.Node, l.VMID, controlPath, &verify, pve.ControlLimit+1); err != nil || !bytes.Equal(control, verify.Bytes()) {
verify, err := cs.runtime.ReadWindowsJSON(delivery, b.Node, l.VMID, controlPath, pve.ControlLimit)
if err != nil || !bytes.Equal(control, verify) {
return outcome, findings, telemetry, cleanup, nil, errors.New("bounded guest manifest verification failed")
}
if err = session.Err(); err != nil {
@@ -770,7 +769,19 @@ func (e *engine) collect(ctx context.Context, a attempt, l allocation, b Binding
if err != nil {
return err
}
n, err := cs.runtime.ReadFile(ctx, b.Node, l.VMID, v.path, f, e.cfg.MaxArtifactBytes-total)
var n int64
if strings.HasSuffix(strings.ToLower(v.path), ".json") || strings.EqualFold(strings.TrimSpace(strings.SplitN(v.ctype, ";", 2)[0]), "application/json") {
limit := min(e.cfg.MaxArtifactBytes-total, int64(pve.WindowsJSONLimit))
var raw []byte
raw, err = cs.runtime.ReadWindowsJSON(ctx, b.Node, l.VMID, v.path, limit)
if err == nil {
var written int
written, err = f.Write(raw)
n = int64(written)
}
} else {
n, err = cs.runtime.ReadFile(ctx, b.Node, l.VMID, v.path, f, e.cfg.MaxArtifactBytes-total)
}
total += n
syncErr := f.Sync()
closeErr := f.Close()
+174 -18
View File
@@ -7,13 +7,17 @@ import (
"errors"
"os"
"path/filepath"
"strings"
"time"
)
// salvagePartial preserves whatever was durably flushed without turning a
// missing final result into clean telemetry. It never schedules guest code.
func (e *engine) salvagePartial(ctx context.Context, a attempt, l allocation, b Binding, cs clients, dir string) (json.RawMessage, string) {
var report json.RawMessage
// salvagePartial preserves durably flushed evidence without treating a missing
// final result as complete telemetry. Guest reads never execute file contents.
func (e *engine) salvagePartial(ctx context.Context, a attempt, l allocation, b Binding, cs clients, dir string, prior json.RawMessage, priorFindings string) (json.RawMessage, string) {
findings := "unknown"
if priorFindings == "detected" {
findings = "detected"
}
if e.leaseAlive(ctx, a) != nil || cs.provisioner.CheckOwned(ctx, e.own(a, l, b)) != nil {
return nil, findings
}
@@ -22,38 +26,190 @@ func (e *engine) salvagePartial(ctx context.Context, a attempt, l allocation, b
return nil, findings
}
var total int64
snapshots := make(map[string]json.RawMessage, 4)
for _, name := range []string{"result.json", "status.json", "receipt.json", "telemetry.json", "defender-events.json"} {
if total >= e.cfg.MaxArtifactBytes {
return report, findings
break
}
var raw bytes.Buffer
remaining := e.cfg.MaxArtifactBytes - total
if remaining > 4<<20 {
remaining = 4 << 20
}
n, err := cs.runtime.ReadFile(ctx, b.Node, l.VMID, guestDir(a)+name, &raw, remaining)
total += n
if err != nil || !json.Valid(raw.Bytes()) {
raw, err := cs.runtime.ReadWindowsJSON(ctx, b.Node, l.VMID, guestDir(a)+name, remaining)
total += int64(len(raw))
if err != nil {
continue
}
path := filepath.Join(salvageDir, name)
if err = os.WriteFile(path, raw.Bytes(), 0600); err != nil {
if err = os.WriteFile(path, raw, 0600); err != nil {
continue
}
_ = e.publishFile(ctx, a, "partial_telemetry", name, "application/json", path)
if name == "result.json" || name == "status.json" {
var r guestResult
if json.Unmarshal(raw.Bytes(), &r) == nil && r.CommandID == a.CommandID && r.AttemptID == a.ID && r.JobID == a.JobID {
if json.Valid(r.Report) {
report = r.Report
}
if r.Findings == "detected" {
if name != "defender-events.json" {
snapshots[name] = raw
}
}
return salvageLiveJSON(a, snapshots, time.Now().UTC(), prior, findings)
}
// The collector has no command IDs of its own. Its command-scoped filename,
// receipt, boot and bounded timestamps must agree before its evidence is used.
// Already observed final reports and on-disk final results take precedence over
// a runner status that stopped refreshing. Positive findings are monotonic.
func salvageLiveJSON(a attempt, snapshots map[string]json.RawMessage, observedAt time.Time, prior json.RawMessage, priorFindings string) (json.RawMessage, string) {
var report map[string]json.RawMessage
if json.Unmarshal(prior, &report) != nil {
report = nil
}
findings := "unknown"
if priorFindings == "detected" {
findings = "detected"
}
for _, name := range []string{"result.json", "status.json"} {
var r guestResult
if json.Unmarshal(snapshots[name], &r) != nil || r.CommandID != a.CommandID || r.AttemptID != a.ID || r.JobID != a.JobID {
continue
}
if name == "result.json" && !validResult(a, r) {
continue
}
var candidate map[string]json.RawMessage
if json.Unmarshal(r.Report, &candidate) != nil || candidate == nil {
continue
}
if report == nil {
report = candidate
}
if r.Findings == "detected" {
findings = "detected"
}
}
var receipt struct {
CommandID string `json:"command_id"`
AttemptID string `json:"attempt_id"`
JobID string `json:"job_id"`
State string `json:"state"`
AcceptedAt time.Time `json:"accepted_at"`
BootID string `json:"boot_id"`
}
var telemetry struct {
UpdatedAt time.Time `json:"updated_at"`
BootID string `json:"boot_id"`
Before json.RawMessage `json:"before"`
After json.RawMessage `json:"after"`
Drift bool `json:"drift"`
Detections []json.RawMessage `json:"detections"`
Environment json.RawMessage `json:"environment"`
Errors []string `json:"errors"`
}
validCollector := json.Unmarshal(snapshots["receipt.json"], &receipt) == nil &&
receipt.CommandID == a.CommandID && receipt.AttemptID == a.ID && receipt.JobID == a.JobID &&
(receipt.State == "accepted" || receipt.State == "running") && !receipt.AcceptedAt.IsZero() &&
json.Unmarshal(snapshots["telemetry.json"], &telemetry) == nil &&
telemetry.BootID == receipt.BootID && !telemetry.UpdatedAt.Before(receipt.AcceptedAt) && !telemetry.UpdatedAt.After(observedAt)
boot, bootErr := time.Parse(time.RFC3339Nano, receipt.BootID)
if validCollector && bootErr == nil && !boot.After(receipt.AcceptedAt) {
if report == nil {
report = make(map[string]json.RawMessage)
}
var defender map[string]json.RawMessage
if json.Unmarshal(report["defender"], &defender) != nil || defender == nil {
defender = make(map[string]json.RawMessage)
}
var detections []json.RawMessage
_ = json.Unmarshal(defender["detections"], &detections)
// Match the collector's two-minute delivery lookback, never an old boot.
since := receipt.AcceptedAt.Add(-2 * time.Minute)
if boot.After(since) {
since = boot
}
for _, raw := range telemetry.Detections {
var detection struct {
Source string `json:"source"`
Resources []string `json:"resources"`
Timestamp time.Time `json:"timestamp"`
}
if json.Unmarshal(raw, &detection) != nil || detection.Source != "defender" || detection.Timestamp.Before(since) || detection.Timestamp.After(telemetry.UpdatedAt) {
continue
}
for _, resource := range detection.Resources {
if matchesSampleResource(a, resource) {
findings = "detected"
detections = appendDistinctJSON(detections, raw)
break
}
}
}
for key, value := range map[string]json.RawMessage{"before": telemetry.Before, "after": telemetry.After} {
if !presentJSON(defender[key]) && presentJSON(value) {
defender[key] = value
}
}
if telemetry.Drift || !presentJSON(defender["drift"]) {
defender["drift"], _ = json.Marshal(telemetry.Drift)
}
if detections == nil {
detections = []json.RawMessage{}
}
defender["detections"], _ = json.Marshal(detections)
report["defender"], _ = json.Marshal(defender)
if !presentJSON(report["environment"]) && presentJSON(telemetry.Environment) {
report["environment"] = telemetry.Environment
}
var collectionErrors []string
_ = json.Unmarshal(report["collection_errors"], &collectionErrors)
for _, message := range telemetry.Errors {
found := false
for _, existing := range collectionErrors {
if message == existing {
found = true
break
}
}
if !found {
collectionErrors = append(collectionErrors, message)
}
}
if len(collectionErrors) > 0 {
report["collection_errors"], _ = json.Marshal(collectionErrors)
}
}
return report, findings
if report == nil {
return nil, findings
}
raw, err := json.Marshal(report)
if err != nil {
return nil, findings
}
return raw, findings
}
func presentJSON(raw json.RawMessage) bool {
return len(raw) > 0 && !bytes.Equal(bytes.TrimSpace(raw), []byte("null"))
}
func appendDistinctJSON(values []json.RawMessage, value json.RawMessage) []json.RawMessage {
for _, existing := range values {
if bytes.Equal(existing, value) {
return values
}
}
return append(values, value)
}
func matchesSampleResource(a attempt, resource string) bool {
if a.Filename == "" || strings.ContainsAny(a.Filename, `\/:`) {
return false
}
resource = strings.ToLower(resource)
resource = strings.TrimPrefix(resource, "file:_")
filename := strings.ToLower(a.Filename)
if resource == `c:\programdata\otche\samples\`+strings.ToLower(a.CommandID)+`\`+filename {
return true
}
// Collect-Otche only adds an ISO candidate after matching its volume label
// against this command's manifest. Drive letters themselves are not stable.
return len(resource) > 2 && resource[0] >= 'a' && resource[0] <= 'z' && resource[1:] == `:\sample\`+filename
}
func appendObservation(path string, v json.RawMessage, limit int64) error {
info, err := os.Stat(path)
+231
View File
@@ -5,6 +5,7 @@ import (
"os"
"path/filepath"
"testing"
"time"
)
func TestExtractedReportRejectsDifferentCommand(t *testing.T) {
@@ -33,3 +34,233 @@ func TestExtractedReportRejectsDifferentCommand(t *testing.T) {
t.Fatal("same command observed detection was lost")
}
}
// This is the collector's real schema, with synthetic identity and environment.
// The EICAR name is metadata only; no executable sample bytes are included.
func partialCollectorFixture(t *testing.T) (attempt, map[string]json.RawMessage, time.Time) {
t.Helper()
a := attempt{ID: newID(), JobID: newID(), CommandID: newID(), Filename: "acceptance-eicar.com"}
status := guestResult{
CommandID: a.CommandID, AttemptID: a.ID, JobID: a.JobID,
Phase: "preparing", Outcome: "pending", Findings: "unknown", Telemetry: "pending",
Report: json.RawMessage(`{"execution":{"exit_code":null,"error":"","actual_duration_seconds":null,"started_at":null,"finished_at":null,"user":"","session_id":null,"pid":null,"path":"","arguments":[],"handler":"","privilege":"admin"},"defender":{"before":null,"after":null,"drift":false,"detections":[]},"environment":null,"collection_errors":[]}`),
}
rawStatus, err := json.Marshal(status)
if err != nil {
t.Fatal(err)
}
rawReceipt, err := json.Marshal(map[string]any{
"command_id": a.CommandID, "attempt_id": a.ID, "job_id": a.JobID,
"state": "accepted", "accepted_at": "2026-09-25T18:14:33.6801866Z",
"started_at": nil, "deadline_at": nil, "session_id": 1,
"user": `TEST-PC\TestUser`, "privilege": "admin", "pid": nil,
"boot_id": "2026-09-25T18:14:10.5000000Z",
})
if err != nil {
t.Fatal(err)
}
telemetry := json.RawMessage(`{"updated_at":"2026-09-25T18:17:35.3563071Z","boot_id":"2026-09-25T18:14:10.5000000Z","before":{"active":true,"fingerprint":"synthetic-baseline"},"after":{"active":true,"fingerprint":"synthetic-baseline"},"drift":false,"detections":[{"name":"Virus:DOS/EICAR_Test_File","id":"{11111111-1111-1111-1111-111111111111}","action":"cleaning_action=9;success=True;status=1","resources":["file:_E:\\sample\\acceptance-eicar.com"],"timestamp":"2026-09-25T18:14:36.9690000Z","stage":"preparation","source":"defender"},{"name":"Microsoft-Windows-Windows Defender","id":"1116","action":"event 1116","resources":["E:\\sample\\acceptance-eicar.com"],"timestamp":"2026-09-25T18:14:36.9896716Z","stage":"preparation","source":"defender"}],"environment":{"os_build":"synthetic","architecture":"x64","powershell_version":"5.1","execution_policy":"AllSigned","runner_version":"1.0.0","security":{}},"errors":[],"initial":false}`)
return a, map[string]json.RawMessage{"status.json": rawStatus, "receipt.json": rawReceipt, "telemetry.json": telemetry}, time.Date(2026, 9, 25, 18, 18, 0, 0, time.UTC)
}
func changeCollectorFixture(t *testing.T, snapshots map[string]json.RawMessage, file string, change func(map[string]any)) {
t.Helper()
var value map[string]any
if err := json.Unmarshal(snapshots[file], &value); err != nil {
t.Fatal(err)
}
change(value)
raw, err := json.Marshal(value)
if err != nil {
t.Fatal(err)
}
snapshots[file] = raw
}
func TestPartialCollectorPreservesPreparingExecutionAndDetection(t *testing.T) {
a, snapshots, now := partialCollectorFixture(t)
changeCollectorFixture(t, snapshots, "telemetry.json", func(v map[string]any) {
v["errors"] = []string{"Defender query bound reached"}
})
raw, findings := salvageLiveJSON(a, snapshots, now, nil, "unknown")
var report struct {
Execution json.RawMessage `json:"execution"`
Defender struct {
Before struct {
Active bool `json:"active"`
} `json:"before"`
Detections []struct {
Name string `json:"name"`
ID string `json:"id"`
} `json:"detections"`
} `json:"defender"`
Environment struct {
Architecture string `json:"architecture"`
} `json:"environment"`
Errors []string `json:"collection_errors"`
}
if err := json.Unmarshal(raw, &report); err != nil {
t.Fatal(err)
}
var status guestResult
if err := json.Unmarshal(snapshots["status.json"], &status); err != nil {
t.Fatal(err)
}
var original map[string]json.RawMessage
if err := json.Unmarshal(status.Report, &original); err != nil {
t.Fatal(err)
}
if findings != "detected" || len(report.Defender.Detections) != 2 || report.Defender.Detections[0].Name != "Virus:DOS/EICAR_Test_File" || report.Defender.Detections[1].ID != "1116" {
t.Fatalf("correlated positive evidence missing: findings=%s report=%s", findings, raw)
}
if string(report.Execution) != string(original["execution"]) || !report.Defender.Before.Active || report.Environment.Architecture != "x64" || len(report.Errors) != 1 || report.Errors[0] != "Defender query bound reached" {
t.Fatalf("execution, environment or incomplete collection evidence lost: %s", raw)
}
}
func TestPartialCollectorRejectsUnboundEvidence(t *testing.T) {
cases := []struct {
name string
file string
change func(map[string]any)
}{
{"foreign command", "receipt.json", func(v map[string]any) { v["command_id"] = newID() }},
{"foreign attempt", "receipt.json", func(v map[string]any) { v["attempt_id"] = newID() }},
{"foreign job", "receipt.json", func(v map[string]any) { v["job_id"] = newID() }},
{"unaccepted command", "receipt.json", func(v map[string]any) { v["state"] = "rejected" }},
{"missing acceptance time", "receipt.json", func(v map[string]any) { delete(v, "accepted_at") }},
{"wrong boot", "telemetry.json", func(v map[string]any) { v["boot_id"] = "2026-09-24T18:14:10.5000000Z" }},
{"stale snapshot", "telemetry.json", func(v map[string]any) { v["updated_at"] = "2026-09-25T18:14:00Z" }},
{"future snapshot", "telemetry.json", func(v map[string]any) { v["updated_at"] = "2026-09-26T18:14:00Z" }},
{"stale detection", "telemetry.json", func(v map[string]any) {
for _, d := range v["detections"].([]any) {
d.(map[string]any)["timestamp"] = "2026-09-25T18:10:00Z"
}
}},
{"post snapshot detection", "telemetry.json", func(v map[string]any) {
for _, d := range v["detections"].([]any) {
d.(map[string]any)["timestamp"] = "2026-09-25T18:17:36Z"
}
}},
{"foreign filename", "telemetry.json", func(v map[string]any) {
for _, d := range v["detections"].([]any) {
d.(map[string]any)["resources"] = []string{`file:_E:\sample\other.com`}
}
}},
{"filename prefix", "telemetry.json", func(v map[string]any) {
for _, d := range v["detections"].([]any) {
d.(map[string]any)["resources"] = []string{`file:_E:\sample\acceptance-eicar.com.exe`}
}
}},
{"foreign command path", "telemetry.json", func(v map[string]any) {
for _, d := range v["detections"].([]any) {
d.(map[string]any)["resources"] = []string{`file:_C:\ProgramData\Otche\samples\other-command\acceptance-eicar.com`}
}
}},
{"empty resources", "telemetry.json", func(v map[string]any) {
for _, d := range v["detections"].([]any) {
d.(map[string]any)["resources"] = []string{}
}
}},
{"policy source", "telemetry.json", func(v map[string]any) {
for _, d := range v["detections"].([]any) {
d.(map[string]any)["source"] = "policy"
}
}},
{"no detections", "telemetry.json", func(v map[string]any) { v["detections"] = []any{} }},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
a, snapshots, now := partialCollectorFixture(t)
changeCollectorFixture(t, snapshots, tc.file, tc.change)
raw, findings := salvageLiveJSON(a, snapshots, now, nil, "unknown")
var report struct {
Defender struct {
Detections []json.RawMessage `json:"detections"`
} `json:"defender"`
}
if err := json.Unmarshal(raw, &report); err != nil {
t.Fatal(err)
}
if findings != "unknown" || len(report.Defender.Detections) != 0 {
t.Fatalf("unbound evidence became detection: %s %s", findings, raw)
}
})
}
}
func TestPartialCollectorMatchesCommandSample(t *testing.T) {
a, snapshots, now := partialCollectorFixture(t)
changeCollectorFixture(t, snapshots, "telemetry.json", func(v map[string]any) {
v["detections"] = v["detections"].([]any)[:1]
v["detections"].([]any)[0].(map[string]any)["resources"] = []string{`file:_C:\ProgramData\Otche\samples\` + a.CommandID + `\ACCEPTANCE-EICAR.COM`}
})
_, findings := salvageLiveJSON(a, snapshots, now, nil, "unknown")
if findings != "detected" {
t.Fatalf("command sample detection lost: %s", findings)
}
}
func TestPartialSalvagePreservesFinalEvidenceAgainstStaleStatus(t *testing.T) {
a, snapshots, now := partialCollectorFixture(t)
final := guestResult{
CommandID: a.CommandID, AttemptID: a.ID, JobID: a.JobID,
Outcome: "executed", Findings: "detected", Telemetry: "partial",
Report: json.RawMessage(`{"execution":{"exit_code":42,"actual_duration_seconds":7},"defender":{"before":{"fingerprint":"final-baseline"},"detections":[{"id":"earlier-detection","source":"defender"}]},"collection_errors":["runner collection incomplete"]}`),
}
rawFinal, err := json.Marshal(final)
if err != nil {
t.Fatal(err)
}
snapshots["result.json"] = rawFinal
for _, invalidCollector := range []bool{false, true} {
if invalidCollector {
snapshots["telemetry.json"] = json.RawMessage(`{"invalid"`)
}
raw, findings := salvageLiveJSON(a, snapshots, now, nil, "unknown")
var report struct {
Execution struct {
ExitCode int `json:"exit_code"`
} `json:"execution"`
Defender struct {
Before struct {
Fingerprint string `json:"fingerprint"`
} `json:"before"`
Detections []struct {
ID string `json:"id"`
} `json:"detections"`
} `json:"defender"`
Errors []string `json:"collection_errors"`
}
if err := json.Unmarshal(raw, &report); err != nil {
t.Fatal(err)
}
if findings != "detected" || report.Execution.ExitCode != 42 || report.Defender.Before.Fingerprint != "final-baseline" || len(report.Defender.Detections) == 0 || report.Defender.Detections[0].ID != "earlier-detection" || len(report.Errors) == 0 || report.Errors[0] != "runner collection incomplete" {
t.Fatalf("stale status downgraded final evidence: %s %s", findings, raw)
}
}
}
func TestPartialSalvagePreservesHeldReportWhenFinalReadFails(t *testing.T) {
a, snapshots, now := partialCollectorFixture(t)
snapshots["result.json"] = json.RawMessage(`{"truncated"`)
prior := json.RawMessage(`{"execution":{"exit_code":42},"defender":{"detections":[{"id":"held-detection","source":"defender"}]},"collection_errors":["final collection failed"]}`)
raw, findings := salvageLiveJSON(a, snapshots, now, prior, "detected")
var report struct {
Execution struct {
ExitCode int `json:"exit_code"`
} `json:"execution"`
Defender struct {
Detections []struct {
ID string `json:"id"`
} `json:"detections"`
} `json:"defender"`
Errors []string `json:"collection_errors"`
}
if err := json.Unmarshal(raw, &report); err != nil {
t.Fatal(err)
}
if findings != "detected" || report.Execution.ExitCode != 42 || len(report.Defender.Detections) != 3 || report.Defender.Detections[0].ID != "held-detection" || len(report.Errors) != 1 || report.Errors[0] != "final collection failed" {
t.Fatalf("stale status or collector erased held final evidence: %s %s", findings, raw)
}
}
+13 -2
View File
@@ -191,9 +191,20 @@ func TestObserveTimeoutLeavesFencedFinalizerHandoff(t *testing.T) {
if err != nil {
t.Fatal(err)
}
executions := 0
c := readinessClient(t, func(w http.ResponseWriter, r *http.Request) {
if strings.HasSuffix(r.URL.Path, "/agent/file-read") && r.URL.Query().Get("file") == guestDir(a)+"receipt.json" {
_ = json.NewEncoder(w).Encode(map[string]any{"data": map[string]any{"content": base64.StdEncoding.EncodeToString(raw), "truncated": false}})
if strings.HasSuffix(r.URL.Path, "/agent/exec") {
executions++
_ = json.NewEncoder(w).Encode(map[string]any{"data": map[string]any{"pid": executions}})
return
}
if strings.HasSuffix(r.URL.Path, "/agent/exec-status") {
response := map[string]any{"exited": true, "exitcode": 1}
if r.URL.Query().Get("pid") == "1" {
response["exitcode"] = 0
response["out-data"] = base64.StdEncoding.EncodeToString(raw)
}
_ = json.NewEncoder(w).Encode(map[string]any{"data": response})
return
}
http.NotFound(w, r)
+4 -2
View File
@@ -81,8 +81,10 @@ Invoke-Otche.ps1 -ManifestPath is the installed interactive runner interface. QG
never invokes the sample or submits a generic shell command. The fixed scheduler
starts prepared InteractiveToken tasks (Limited/Highest), plus a separate SYSTEM
Defender collector. All commands and control JSON must be short/bounded; use ISO
for samples and bounded QGA file-open/read offsets for artifacts, not repeated
file-write overwrites. No credential, password or PVE token belongs in a manifest.
for samples and bounded QGA file-open/read offsets for binary artifacts. Mutable
JSON uses the worker's fixed read-only PowerShell snapshot with ReadWrite|Delete
sharing, closed before output, never file contents as code or a policy bypass.
No repeated file-write overwrites, credential, password or PVE token in a manifest.
Readiness uses a fixed Otche-DesktopReady InteractiveToken/Limited task running the
signed Test-OtcheReady.ps1 -Privilege user -DesktopProbe with no command/input API.
Each request has a new nonce and 60-second deadline; bounded response must match the