Wire SCUM schema probe durable run path

This commit is contained in:
npc0-hue
2026-08-12 16:49:05 +08:00
parent ca3fb39df8
commit d854d6fda3
21 changed files with 309 additions and 100 deletions
+2 -2
View File
@@ -637,7 +637,7 @@ func firstEligibleSupportedJob(jobs []domain.Job, capabilities []string, stamp t
func assignmentFromJob(job domain.Job, leaseToken string) domain.RunJobAssignment {
fencingToken := uint64(0)
if isProtectedRequestCapability(job.Capability) {
if isProtectedRequestCapability(job.Capability) || job.Capability == domain.JobCapabilityRemoteRunDBSQLiteProbe {
fencingToken = uint64(job.Attempt)
}
return domain.RunJobAssignment{
@@ -651,7 +651,7 @@ func assignmentFromJob(job domain.Job, leaseToken string) domain.RunJobAssignmen
State: job.State,
Progress: domain.RunJobProgressReport{Percent: job.Progress.Percent, Phase: job.Progress.Phase, Message: job.Progress.Message},
ResultRef: job.ResultRef,
ExecutionInput: domain.JobExecutionInput{WorkspaceScope: job.ExecutionInput.WorkspaceScope, Content: job.ExecutionInput.Content, ExpectedVersion: job.ExecutionInput.ExpectedVersion, ExpectedChecksum: job.ExecutionInput.ExpectedChecksum, MaxReadBytes: job.ExecutionInput.MaxReadBytes, RemoteAdapterKey: job.ExecutionInput.RemoteAdapterKey, RemoteAdapterKind: job.ExecutionInput.RemoteAdapterKind, TimeoutSeconds: job.ExecutionInput.TimeoutSeconds, PluginID: job.ExecutionInput.PluginID, LifecycleOperation: job.ExecutionInput.LifecycleOperation, TargetVersion: job.ExecutionInput.TargetVersion, Inputs: domain.CopyStringMap(job.ExecutionInput.Inputs), LogSource: domain.CopyRuntimeLogSourcePtr(job.ExecutionInput.LogSource), LogSources: domain.CopyRuntimeLogSources(job.ExecutionInput.LogSources), DLLExtensions: append([]domain.RuntimeDLLExtensionPlan(nil), job.ExecutionInput.DLLExtensions...), SourceRCON: domain.CopyRuntimeSourceRCONPlan(job.ExecutionInput.SourceRCON), Deployment: deploymentPlanForDispatchValue(job.ExecutionInput.Deployment), ServerDeploymentPlan: domain.CopyServerDeploymentPlan(job.ExecutionInput.ServerDeploymentPlan)},
ExecutionInput: domain.JobExecutionInput{WorkspaceScope: job.ExecutionInput.WorkspaceScope, Content: job.ExecutionInput.Content, ExpectedVersion: job.ExecutionInput.ExpectedVersion, ExpectedChecksum: job.ExecutionInput.ExpectedChecksum, MaxReadBytes: job.ExecutionInput.MaxReadBytes, RemoteAdapterKey: job.ExecutionInput.RemoteAdapterKey, RemoteAdapterKind: job.ExecutionInput.RemoteAdapterKind, TimeoutSeconds: job.ExecutionInput.TimeoutSeconds, PluginID: job.ExecutionInput.PluginID, LifecycleOperation: job.ExecutionInput.LifecycleOperation, TargetVersion: job.ExecutionInput.TargetVersion, Inputs: domain.CopyStringMap(job.ExecutionInput.Inputs), LogSource: domain.CopyRuntimeLogSourcePtr(job.ExecutionInput.LogSource), LogSources: domain.CopyRuntimeLogSources(job.ExecutionInput.LogSources), DLLExtensions: append([]domain.RuntimeDLLExtensionPlan(nil), job.ExecutionInput.DLLExtensions...), SourceRCON: domain.CopyRuntimeSourceRCONPlan(job.ExecutionInput.SourceRCON), Deployment: deploymentPlanForDispatchValue(job.ExecutionInput.Deployment), ServerDeploymentPlan: domain.CopyServerDeploymentPlan(job.ExecutionInput.ServerDeploymentPlan), SQLiteSchemaProbe: domain.CopySCUMSchemaProbeRequestPtr(job.ExecutionInput.SQLiteSchemaProbe)},
LeaseToken: leaseToken,
Attempt: job.Attempt,
FencingToken: fencingToken,
+21 -3
View File
@@ -86,20 +86,30 @@ func (svc *CoreService) RequestRemoteAdapterForSession(sessionID string, request
return domain.RemoteAdapterResult{}, validationError("remote adapter timeout or retry exceeds declaration")
}
inputRef := request.InputRef
if inputRef == "" {
isSchemaProbe := request.Capability == domain.JobCapabilityRemoteRunDBSQLiteProbe && request.PlatformScheduled && request.SQLiteSchemaProbe != nil
if inputRef == "" && !isSchemaProbe {
inputRef = fmt.Sprintf("input://remote-adapters/%s/%s", instance.ID, request.DeclarationKey)
}
targetKey := request.TargetKey
executionInput := domain.JobExecutionInput{WorkspaceScope: svc.runtimeProfileScope(instance.ID), RemoteAdapterKey: selected.Key, RemoteAdapterKind: string(selected.Kind), TimeoutSeconds: timeout, Inputs: domain.CopyStringMap(request.Inputs), SQLiteSchemaProbe: domain.CopySCUMSchemaProbeRequestPtr(request.SQLiteSchemaProbe)}
if isSchemaProbe {
targetKey = sqliteSchemaProbeRunTargetKey(request.TargetKey)
inputRef = ""
executionInput.RemoteAdapterKey = ""
executionInput.RemoteAdapterKind = ""
executionInput.Inputs = nil
}
job := domain.Job{
ID: jobIDFromParts("job-remote-adapter", instance.ID, request.IdempotencyKey),
ServerInstanceID: instance.ID,
RunEndpointID: instance.RunEndpointID,
Capability: request.Capability,
TargetKey: request.TargetKey,
TargetKey: targetKey,
InputRef: inputRef,
IdempotencyKey: request.IdempotencyKey,
Progress: domain.JobProgress{Percent: 0, Message: "scoped remote adapter queued"},
RetryPolicy: domain.JobRetryPolicy{MaxAttempts: attempts, InitialBackoffSeconds: 2, MaxBackoffSeconds: 30},
ExecutionInput: domain.JobExecutionInput{WorkspaceScope: svc.runtimeProfileScope(instance.ID), RemoteAdapterKey: selected.Key, RemoteAdapterKind: string(selected.Kind), TimeoutSeconds: timeout, Inputs: domain.CopyStringMap(request.Inputs)},
ExecutionInput: executionInput,
}
created, err := svc.CreateJob(job)
if err != nil {
@@ -116,6 +126,14 @@ func (svc *CoreService) RequestRemoteAdapterForSession(sessionID string, request
return domain.RemoteAdapterResult{RequestID: created.ID, ServerInstanceID: instance.ID, DeclarationKey: selected.Key, TargetKey: request.TargetKey, Kind: selected.Kind, Status: string(created.State), Retryable: attempts > 1, Message: "scoped remote adapter queued", ResultRef: "job://" + created.ID, AuditEventID: auditID}, nil
}
func sqliteSchemaProbeRunTargetKey(targetKey string) string {
trimmed := strings.TrimSpace(targetKey)
if strings.HasPrefix(trimmed, "databases/") {
return trimmed
}
return "databases/" + trimmed
}
func intersectRemoteCapabilities(profile []string, declared []string, endpoint []string) []string {
result := make([]string, 0, len(profile))
for _, capability := range profile {
+29 -5
View File
@@ -140,11 +140,11 @@ func TestSCUMSchemaProbeDispatchIsPlatformScheduledAndFenced(t *testing.T) {
if err != nil {
t.Fatalf("get probe job: %v", err)
}
if job.Capability != domain.JobCapabilityRemoteRunDBSQLiteProbe || job.TargetKey != "scum-database" || job.ExecutionInput.RemoteAdapterKind != string(domain.RemoteAdapterDatabase) || job.ExecutionInput.Inputs["databaseIdentity"] != "logical:scum-database" {
if job.Capability != domain.JobCapabilityRemoteRunDBSQLiteProbe || job.TargetKey != "databases/scum-database" || job.InputRef != "" || job.ExecutionInput.RemoteAdapterKey != "" || job.ExecutionInput.RemoteAdapterKind != "" || len(job.ExecutionInput.Inputs) != 0 {
t.Fatalf("unexpected probe job envelope: %+v", job)
}
if job.ExecutionInput.Inputs["jobId"] != job.ID || job.ExecutionInput.Inputs["requestId"] != probeRequest.RequestID || job.ExecutionInput.Inputs["maxResultBytes"] != "524288" {
t.Fatalf("probe inputs are not fenced and bounded: %+v", job.ExecutionInput.Inputs)
if job.ExecutionInput.SQLiteSchemaProbe == nil || job.ExecutionInput.SQLiteSchemaProbe.RequestID != probeRequest.RequestID || job.ExecutionInput.SQLiteSchemaProbe.Binding.DatabaseIdentity != "scum-database" || job.ExecutionInput.SQLiteSchemaProbe.Bounds.MaxResultBytes != 524288 {
t.Fatalf("probe job did not include typed SQLite schema probe request: %+v", job.ExecutionInput.SQLiteSchemaProbe)
}
helloRequest := validRunControlHello()
@@ -158,8 +158,32 @@ func TestSCUMSchemaProbeDispatchIsPlatformScheduledAndFenced(t *testing.T) {
if err != nil {
t.Fatalf("claim probe job: %v", err)
}
if !claim.HasJob || claim.Job == nil || claim.Job.ExecutionInput.Inputs["adapterVersion"] != "scum-live-data-v0" || claim.Job.ExecutionInput.Inputs["targetKey"] != "scum-database" {
t.Fatalf("claimed probe job lost typed inputs: %+v", claim.Job)
if !claim.HasJob || claim.Job == nil || claim.Job.TargetKey != "databases/scum-database" || claim.Job.InputRef != "" || claim.Job.FencingToken == 0 || claim.Job.MaxAttempts != 1 || claim.Job.ExecutionInput.RemoteAdapterKey != "" || len(claim.Job.ExecutionInput.Inputs) != 0 {
t.Fatalf("claimed probe job lost fenced typed envelope: %+v", claim.Job)
}
if claim.Job.ExecutionInput.SQLiteSchemaProbe == nil || claim.Job.ExecutionInput.SQLiteSchemaProbe.RequestID != probeRequest.RequestID || claim.Job.ExecutionInput.SQLiteSchemaProbe.Binding.RunBindingID != probeRequest.Binding.RunBindingID {
t.Fatalf("claimed probe job lost typed schema probe request: %+v", claim.Job.ExecutionInput.SQLiteSchemaProbe)
}
assignmentBody := dto.RunJobAssignmentFromDomain(*claim.Job)
payload, err := json.Marshal(assignmentBody)
if err != nil {
t.Fatalf("marshal probe Run assignment: %v", err)
}
var runWire struct {
ExecutionInput struct {
SQLiteSchemaProbe struct {
RequestID string `json:"requestId"`
JobID string `json:"jobId"`
Bounds *dto.SCUMSchemaProbeBoundsDTO `json:"bounds"`
Limits dto.SCUMSchemaProbeBoundsDTO `json:"limits"`
} `json:"sqliteSchemaProbe"`
} `json:"executionInput"`
}
if err := json.Unmarshal(payload, &runWire); err != nil {
t.Fatalf("unmarshal probe Run assignment: %v", err)
}
if runWire.ExecutionInput.SQLiteSchemaProbe.RequestID != probeRequest.RequestID || runWire.ExecutionInput.SQLiteSchemaProbe.JobID != "" || runWire.ExecutionInput.SQLiteSchemaProbe.Bounds != nil || runWire.ExecutionInput.SQLiteSchemaProbe.Limits.MaxResultBytes != probeRequest.Bounds.MaxResultBytes {
t.Fatalf("probe Run assignment JSON does not match Run contract: %s", payload)
}
badBinding := probeRequest.Binding
badBinding.RunBindingID = "runtime-binding-other"
+14 -44
View File
@@ -2,7 +2,6 @@ package service
import (
"errors"
"strconv"
"strings"
"browser.local/platform/domain"
@@ -10,7 +9,7 @@ import (
"browser.local/platform/validator"
)
const scumSchemaProbeExecutionKind = "sqlite.schema.probe"
const scumSchemaProbeExecutionKind = "sqlite.schema-probe"
func (svc *CoreService) RequestSCUMSchemaProbeForSession(sessionID, serverInstanceID, idempotencyKey string) (domain.SCUMSchemaProbeRequest, domain.RemoteAdapterResult, error) {
idempotencyKey = strings.TrimSpace(idempotencyKey)
@@ -65,7 +64,7 @@ func (svc *CoreService) RequestSCUMSchemaProbeForSession(sessionID, serverInstan
PluginID: plugin.ID,
PluginVersion: plugin.Version,
AdapterVersion: adapterVersion,
DatabaseIdentity: "logical:" + probe.TargetKey,
DatabaseIdentity: scumSchemaProbeDatabaseIdentity(probe.TargetKey),
},
Bounds: bounds,
RequestedAt: svc.now(),
@@ -81,9 +80,8 @@ func (svc *CoreService) RequestSCUMSchemaProbeForSession(sessionID, serverInstan
TimeoutSeconds: scumSchemaProbeTimeoutSeconds(bounds),
MaxAttempts: 1,
IdempotencyKey: idempotencyKey,
InputRef: "input://scum-schema-probe/" + instance.ID + "/" + idempotencyKey,
Inputs: scumSchemaProbeInputs(request, probe.TargetKey),
PlatformScheduled: true,
SQLiteSchemaProbe: &request,
})
if err != nil {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, err
@@ -91,6 +89,11 @@ func (svc *CoreService) RequestSCUMSchemaProbeForSession(sessionID, serverInstan
return request, result, nil
}
func scumSchemaProbeDatabaseIdentity(targetKey string) string {
trimmed := strings.TrimSpace(targetKey)
return strings.TrimPrefix(trimmed, "databases/")
}
func scumSchemaProbeAdapterVersion(manifest domain.SCUMLiveDataManifest) string {
for _, gate := range manifest.CapabilityGates {
if gate.Capability == domain.SCUMDataCapabilitySchemaProbe {
@@ -116,51 +119,18 @@ func scumSchemaProbeTimeoutSeconds(bounds domain.SCUMSchemaProbeBounds) int {
return seconds
}
func scumSchemaProbeInputs(request domain.SCUMSchemaProbeRequest, targetKey string) map[string]string {
return map[string]string{
"requestId": request.RequestID,
"jobId": request.JobID,
"targetKey": targetKey,
"serverInstanceId": request.Binding.ServerInstanceID,
"runBindingId": request.Binding.RunBindingID,
"runEndpointId": request.Binding.RunEndpointID,
"pluginId": request.Binding.PluginID,
"pluginVersion": request.Binding.PluginVersion,
"adapterVersion": request.Binding.AdapterVersion,
"gameVersion": request.Binding.GameVersion,
"databaseIdentity": request.Binding.DatabaseIdentity,
"maxObjects": strconv.Itoa(request.Bounds.MaxObjects),
"maxColumnsPerObject": strconv.Itoa(request.Bounds.MaxColumnsPerObject),
"maxIndexesPerObject": strconv.Itoa(request.Bounds.MaxIndexesPerObject),
"maxForeignKeys": strconv.Itoa(request.Bounds.MaxForeignKeys),
"maxCardinalityReads": strconv.Itoa(request.Bounds.MaxCardinalityReads),
"maxSampleRows": strconv.Itoa(request.Bounds.MaxSampleRows),
"timeoutMs": strconv.Itoa(request.Bounds.TimeoutMS),
"maxResultBytes": strconv.Itoa(request.Bounds.MaxResultBytes),
}
}
func scumSchemaProbeBindingFromInputs(inputs map[string]string) domain.SCUMBindingIdentity {
return domain.SCUMBindingIdentity{
ServerInstanceID: inputs["serverInstanceId"],
RunBindingID: inputs["runBindingId"],
RunEndpointID: inputs["runEndpointId"],
PluginID: inputs["pluginId"],
PluginVersion: inputs["pluginVersion"],
AdapterVersion: inputs["adapterVersion"],
GameVersion: inputs["gameVersion"],
DatabaseIdentity: inputs["databaseIdentity"],
}
}
func validateSCUMSchemaProbeResultForJob(job domain.Job, result domain.SCUMSchemaProbeResult) error {
if err := validator.ValidateSCUMSchemaProbeResult(result); err != nil {
return err
}
if result.JobID != job.ID || result.JobID != job.ExecutionInput.Inputs["jobId"] || result.RequestID != job.ExecutionInput.Inputs["requestId"] {
expected := job.ExecutionInput.SQLiteSchemaProbe
if expected == nil {
return validationError("SQLite schema probe request is missing from leased job")
}
if result.JobID != job.ID || result.JobID != expected.JobID || result.RequestID != expected.RequestID {
return validationError("SQLite schema probe result does not match leased job identity")
}
if !sameSCUMSchemaProbeBinding(result.Binding, scumSchemaProbeBindingFromInputs(job.ExecutionInput.Inputs)) {
if !sameSCUMSchemaProbeBinding(result.Binding, expected.Binding) {
return validationError("SQLite schema probe result does not match leased binding identity")
}
return nil