fix(scan): Fix-mode scanner + dialog suppression + cancel/timer + importer revive
Build backend images / build content-svc (push) Failing after 1m25s
Build backend images / build file-svc (push) Failing after 1m10s
Build backend images / build gateway (push) Failing after 56s
Build backend images / build identity-svc (push) Failing after 53s
Build backend images / build notification-svc (push) Failing after 57s
Build backend images / build render-svc (push) Failing after 48s
Build backend images / build studio-svc (push) Failing after 1m5s
Build backend images / build content-svc (push) Failing after 1m25s
Build backend images / build file-svc (push) Failing after 1m10s
Build backend images / build gateway (push) Failing after 56s
Build backend images / build identity-svc (push) Failing after 53s
Build backend images / build notification-svc (push) Failing after 57s
Build backend images / build render-svc (push) Failing after 48s
Build backend images / build studio-svc (push) Failing after 1m5s
- scan.jsx: app.beginSuppressDialogs() + clean quit (no AE hang on font/footage dialogs); FIX-mode branch parses frl_c(x)t/m(y) layer names → scenes by c(x); flexible/mockup keep comp-based walk; FR_SCAN_MODE selects. - render-svc: scan job carries project mode; cancel endpoint + node watchdog that kills AE on cancel; parseObjectURL handles minio:// (bucket in host); scan with no template fails cleanly; status guards so late results can't un-cancel. - content importer: revive soft-deleted scenes instead of duplicate-inserting (fixes scenes_project_id_key unique violation); orphan diff ignores deleted. - admin: scan dialog gets project-type picker + elapsed timer + Cancel button. - node-agent: AE-2026 wiring (host port 5010, host-reachable presign endpoint), FR_SCAN_MODE plumbing. docs/aep-template-convention.md: per-type naming + bundles. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -21,8 +21,11 @@ public class AepImportService(ContentDbContext db)
|
||||
|
||||
private async Task<ImportDiff> RunAsync(Guid projectId, ScanResult scan, ScanApplyOptions? apply)
|
||||
{
|
||||
// Load ALL scenes incl. soft-deleted: the UNIQUE(project_id,key) constraint
|
||||
// covers deleted rows too, so a re-scan must REVIVE a soft-deleted match
|
||||
// rather than insert a duplicate (which would violate the constraint).
|
||||
var existing = await db.Scenes
|
||||
.Where(s => s.ProjectId == projectId && s.DeletedAt == null)
|
||||
.Where(s => s.ProjectId == projectId)
|
||||
.Include(s => s.ContentElements)
|
||||
.Include(s => s.ColorElements)
|
||||
.ToListAsync();
|
||||
@@ -74,13 +77,19 @@ public class AepImportService(ContentDbContext db)
|
||||
scene = new Scene { ProjectId = projectId, Key = ss.Key, Sort = ss.Sort ?? sortBase++ };
|
||||
db.Scenes.Add(scene);
|
||||
}
|
||||
else if (scene!.DeletedAt is not null)
|
||||
{
|
||||
scene.DeletedAt = null; // revive a soft-deleted scene the scan re-discovered
|
||||
}
|
||||
ApplySceneHeader(scene!, ss, isNew, apply);
|
||||
ApplyElements(scene!, scElems, exElemByKey, scElemKeys, apply);
|
||||
ApplyColors(scene!, scColors, exColorByKey, scColorKeys, apply);
|
||||
}
|
||||
}
|
||||
|
||||
var orphans = existing.Where(s => !scanKeys.Contains(s.Key)).ToList();
|
||||
// Orphans = currently-active scenes the scan no longer contains (ignore
|
||||
// already soft-deleted rows so they don't inflate the diff).
|
||||
var orphans = existing.Where(s => s.DeletedAt == null && !scanKeys.Contains(s.Key)).ToList();
|
||||
if (apply is not null && apply.RemoveOrphanScenes)
|
||||
foreach (var s in orphans) s.DeletedAt = DateTime.UtcNow;
|
||||
|
||||
|
||||
@@ -31,6 +31,7 @@ import (
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
@@ -255,8 +256,35 @@ func (a *Agent) pollScanOnce(ctx context.Context) {
|
||||
|
||||
outPath := filepath.Join(a.cfg.WorkDir, "scans", claim.ScanJobID, "scan.json")
|
||||
scanCtx, cancel2 := context.WithTimeout(ctx, 10*time.Minute)
|
||||
result, serr := runner.RunScan(scanCtx, a.cfg.AfterFxPath, aepPath, a.cfg.WorkDir, outPath)
|
||||
cancel2()
|
||||
defer cancel2()
|
||||
|
||||
// Watchdog: if the user cancels the scan server-side, cancel the context →
|
||||
// exec.CommandContext kills the AfterFX process.
|
||||
var cancelled atomic.Bool
|
||||
go func() {
|
||||
t := time.NewTicker(3 * time.Second)
|
||||
defer t.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-scanCtx.Done():
|
||||
return
|
||||
case <-t.C:
|
||||
st, e := a.orch.ScanStatus(scanCtx, claim.ScanJobID)
|
||||
if e == nil && (st == "cancelled" || st == "error") {
|
||||
log.Printf("[scan %s] cancelled server-side — killing AE", claim.ScanJobID)
|
||||
cancelled.Store(true)
|
||||
cancel2()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
result, serr := runner.RunScan(scanCtx, a.cfg.AfterFxPath, aepPath, a.cfg.WorkDir, outPath, claim.Mode)
|
||||
if cancelled.Load() {
|
||||
log.Printf("[scan %s] aborted (cancelled)", claim.ScanJobID)
|
||||
return // already cancelled in DB; don't report fail
|
||||
}
|
||||
if serr != nil {
|
||||
a.failScan(claim.ScanJobID, serr.Error())
|
||||
return
|
||||
|
||||
@@ -250,6 +250,7 @@ func (c *Client) ClaimJob(ctx context.Context, nodeID, region string) (*ClaimedJ
|
||||
type ScanClaim struct {
|
||||
ScanJobID string `json:"scan_job_id"`
|
||||
ProjectID string `json:"project_id"`
|
||||
Mode string `json:"mode"` // fix | flexible | mockup | musicvisualizer
|
||||
AEPDownloadURL string `json:"aep_download_url"`
|
||||
IsBundle bool `json:"is_bundle"`
|
||||
BundleMD5 string `json:"bundle_md5"`
|
||||
@@ -306,6 +307,23 @@ func (c *Client) ReportScanFail(ctx context.Context, scanJobID, reason string) e
|
||||
return nil
|
||||
}
|
||||
|
||||
// ScanStatus returns the scan job's current status (for the cancel watchdog).
|
||||
func (c *Client) ScanStatus(ctx context.Context, scanJobID string) (string, error) {
|
||||
resp, err := c.do(ctx, http.MethodGet, fmt.Sprintf("/v1/internal/scan/%s/status", scanJobID), nil)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode >= 300 {
|
||||
return "", fmt.Errorf("scan status: HTTP %d", resp.StatusCode)
|
||||
}
|
||||
var out struct {
|
||||
Status string `json:"status"`
|
||||
}
|
||||
_ = json.NewDecoder(resp.Body).Decode(&out)
|
||||
return out.Status, nil
|
||||
}
|
||||
|
||||
// UpdatePreview sends a base64-encoded preview frame to the orchestrator.
|
||||
// Errors are non-fatal — the UI simply won't update the preview image.
|
||||
func (c *Client) UpdatePreview(ctx context.Context, jobID, imageB64 string) error {
|
||||
|
||||
@@ -0,0 +1,135 @@
|
||||
// bind.jsx — FlatRender data-driven binder (v1 convention).
|
||||
//
|
||||
// Reads a JSON "bind-spec" and writes the user's values INTO a template, then
|
||||
// saves a "bound" .aep that aerender renders. This is the data-driven successor
|
||||
// to the legacy JSXGenerator string-concat approach: render-svc emits the JSON,
|
||||
// this one generic script interprets it.
|
||||
//
|
||||
// Environment:
|
||||
// FR_BIND_AEP — path to the template .aep to open
|
||||
// FR_BIND_SPEC — path to the bind-spec JSON file
|
||||
// FR_BIND_SAVE — path to write the bound .aep (aerender input)
|
||||
// FR_BIND_QUIT — "0" to keep AE open (debug); default quits
|
||||
//
|
||||
// Bind-spec shape (v1):
|
||||
// {
|
||||
// "comp": "frfinal", "duration": 15, "fps": 30,
|
||||
// "colors": [ { "key": "frd_primary", "value": "#3366ff" } ], // frshare text layers
|
||||
// "data": [ { "key": "frd_toggle", "value": "1" } ], // frd_ data layers
|
||||
// "layers": [
|
||||
// { "key":"frl_title","type":"text","value":"Hello",
|
||||
// "font":"Vazirmatn","size":72,"justify":"CENTER_JUSTIFY" },
|
||||
// { "key":"frl_logo","type":"media","value":"C:/work/cdn/logo.png" },
|
||||
// { "key":"frl_music","type":"audio","value":"C:/work/cdn/track.mp3" }
|
||||
// ]
|
||||
// }
|
||||
//
|
||||
// Conventions (see docs/aep-template-convention.md):
|
||||
// - colours live as TEXT layers inside the `frshare` comp; their Source Text IS
|
||||
// the colour value, and template expressions read them → set once, propagate.
|
||||
// - frl_ = editable visible layers (text / media / audio).
|
||||
// - frd_ = data/direction layers (hidden values, RTL companions).
|
||||
|
||||
(function () {
|
||||
function getenv(n) { try { return $.getenv(n); } catch (e) { return null; } }
|
||||
function readFile(p) { var f = new File(p); f.open("r"); var s = f.read(); f.close(); return s; }
|
||||
function marker(p, txt) { try { var f = new File(p); f.open("w"); f.write(txt); f.close(); } catch (e) {} }
|
||||
|
||||
var aep = getenv("FR_BIND_AEP");
|
||||
var specPath = getenv("FR_BIND_SPEC");
|
||||
var savePath = getenv("FR_BIND_SAVE");
|
||||
var donePath = (savePath || "bind") + ".done";
|
||||
var errPath = (savePath || "bind") + ".error";
|
||||
|
||||
try {
|
||||
if (aep) app.open(new File(aep));
|
||||
var proj = app.project;
|
||||
|
||||
var spec = {};
|
||||
if (specPath) spec = eval("(" + readFile(specPath) + ")"); // bind-spec JSON
|
||||
|
||||
// index the spec by layer key for O(1) lookups while walking the project
|
||||
var colorMap = {}, dataMap = {}, layerMap = {}, i;
|
||||
if (spec.colors) for (i = 0; i < spec.colors.length; i++) colorMap[spec.colors[i].key] = spec.colors[i];
|
||||
if (spec.data) for (i = 0; i < spec.data.length; i++) dataMap[spec.data[i].key] = spec.data[i];
|
||||
if (spec.layers) for (i = 0; i < spec.layers.length; i++) layerMap[spec.layers[i].key] = spec.layers[i];
|
||||
|
||||
// ── helpers ───────────────────────────────────────────────────────────
|
||||
function setText(layer, value) {
|
||||
var p = layer.property("Source Text");
|
||||
if (!p) return;
|
||||
var d = p.value;
|
||||
d.text = (value == null) ? "" : ("" + value);
|
||||
p.setValue(d);
|
||||
}
|
||||
function justifyOf(name) {
|
||||
switch (name) {
|
||||
case "LEFT_JUSTIFY": return ParagraphJustification.LEFT_JUSTIFY;
|
||||
case "RIGHT_JUSTIFY": return ParagraphJustification.RIGHT_JUSTIFY;
|
||||
case "FULL_JUSTIFY": return ParagraphJustification.FULL_JUSTIFY;
|
||||
default: return ParagraphJustification.CENTER_JUSTIFY;
|
||||
}
|
||||
}
|
||||
function textBind(layer, item) {
|
||||
var p = layer.property("Source Text");
|
||||
if (!p) return;
|
||||
var d = p.value;
|
||||
d.text = (item.value == null) ? "" : ("" + item.value);
|
||||
if (item.font) d.font = item.font;
|
||||
if (item.size) d.fontSize = item.size;
|
||||
if (item.justify) d.justification = justifyOf(item.justify);
|
||||
p.setValue(d);
|
||||
// NOTE (v1 TODO): box auto-fit + positionMode 0-8 anchoring go here later.
|
||||
}
|
||||
function mediaBind(layer, item) {
|
||||
if (!item.value) return;
|
||||
try {
|
||||
var footage = proj.importFile(new ImportOptions(new File(item.value)));
|
||||
layer.replaceSource(footage, false);
|
||||
// NOTE (v1 TODO): scale-to-fit the box (legacy mediaBind) goes here later.
|
||||
} catch (e) {}
|
||||
}
|
||||
|
||||
function bindLayer(layer) {
|
||||
var name = layer.name;
|
||||
if (layerMap[name]) {
|
||||
var it = layerMap[name];
|
||||
if (it.type === "media" || it.type === "audio") mediaBind(layer, it);
|
||||
else textBind(layer, it);
|
||||
} else if (dataMap[name]) {
|
||||
setText(layer, dataMap[name].value); // frd_ data / direction layers
|
||||
}
|
||||
}
|
||||
|
||||
// ── walk every comp; bind layers by name (frl_/frd_ live anywhere in the tree) ──
|
||||
for (i = 1; i <= proj.numItems; i++) {
|
||||
var item = proj.item(i);
|
||||
if (!(item instanceof CompItem)) continue;
|
||||
if (item.name === "frshare") {
|
||||
for (var c = 1; c <= item.numLayers; c++) {
|
||||
var L = item.layer(c);
|
||||
if (colorMap[L.name]) setText(L, colorMap[L.name].value);
|
||||
}
|
||||
}
|
||||
for (var l = 1; l <= item.numLayers; l++) bindLayer(item.layer(l));
|
||||
}
|
||||
|
||||
// ── comp timing ─────────────────────────────────────────────────────────
|
||||
var comp = null, target = spec.comp || "frfinal";
|
||||
for (i = 1; i <= proj.numItems; i++) {
|
||||
var x = proj.item(i);
|
||||
if (x instanceof CompItem && x.name === target) { comp = x; break; }
|
||||
}
|
||||
if (comp) {
|
||||
if (spec.duration) comp.duration = spec.duration;
|
||||
if (spec.fps) comp.frameRate = spec.fps;
|
||||
}
|
||||
|
||||
// ── save the bound project for aerender ───────────────────────────────────
|
||||
if (savePath) proj.save(new File(savePath));
|
||||
marker(donePath, "ok comp=" + (comp ? comp.name : "?"));
|
||||
} catch (e) {
|
||||
marker(errPath, "bind error: " + e.toString());
|
||||
}
|
||||
try { if (getenv("FR_BIND_QUIT") !== "0") app.quit(); } catch (e) {}
|
||||
})();
|
||||
@@ -32,7 +32,7 @@ func WriteScanScript(workDir string) (string, error) {
|
||||
//
|
||||
// afterfx -r runs the script and the script calls app.quit(); we still poll for the
|
||||
// output file because afterfx can return before the file is flushed.
|
||||
func RunScan(ctx context.Context, afterfxPath, aepPath, workDir, outPath string) ([]byte, error) {
|
||||
func RunScan(ctx context.Context, afterfxPath, aepPath, workDir, outPath, mode string) ([]byte, error) {
|
||||
scriptPath, err := WriteScanScript(workDir)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("write scan script: %w", err)
|
||||
@@ -43,8 +43,11 @@ func RunScan(ctx context.Context, afterfxPath, aepPath, workDir, outPath string)
|
||||
_ = os.Remove(outPath)
|
||||
_ = os.Remove(outPath + ".error")
|
||||
|
||||
if mode == "" {
|
||||
mode = "flexible"
|
||||
}
|
||||
cmd := exec.CommandContext(ctx, afterfxPath, "-r", scriptPath)
|
||||
cmd.Env = append(os.Environ(), "FR_SCAN_AEP="+aepPath, "FR_SCAN_OUT="+outPath)
|
||||
cmd.Env = append(os.Environ(), "FR_SCAN_AEP="+aepPath, "FR_SCAN_OUT="+outPath, "FR_SCAN_MODE="+mode)
|
||||
cmd.Stdout = os.Stdout
|
||||
cmd.Stderr = os.Stderr
|
||||
// afterfx may exit non-zero on app.quit() — don't treat the exit code as fatal;
|
||||
|
||||
@@ -27,6 +27,7 @@
|
||||
|
||||
var aepPath = getenv("FR_SCAN_AEP");
|
||||
var outPath = getenv("FR_SCAN_OUT") || (Folder.temp.fsName + "/fr_scan.json");
|
||||
var mode = String(getenv("FR_SCAN_MODE") || "flexible").toLowerCase(); // fix | flexible | mockup | musicvisualizer
|
||||
|
||||
// ── minimal JSON serializer (older AE has no JSON.stringify) ──────────────
|
||||
function esc(s) {
|
||||
@@ -138,29 +139,32 @@
|
||||
return false;
|
||||
}
|
||||
|
||||
function run() {
|
||||
if (aepPath) app.open(new File(aepPath));
|
||||
var proj = app.project;
|
||||
var result = { source: "ae-jsx", render_comp: null, scenes: [], shared_colors: [] };
|
||||
// ── frshare colours (shared by all modes) ──────────────────────────────────
|
||||
function readSharedColors(proj) {
|
||||
var colors = [];
|
||||
for (var i = 1; i <= proj.numItems; i++) {
|
||||
var item = proj.item(i);
|
||||
if (!(item instanceof CompItem) || item.name !== "frshare") continue;
|
||||
for (var j = 1; j <= item.numLayers; j++) {
|
||||
var cl = item.layer(j), ct = readText(cl);
|
||||
if (ct && isColor(ct.text)) {
|
||||
colors.push({ element_key: cl.name, title: cl.name, attr_value: "fill", default_color: normColor(ct.text), sort: j });
|
||||
}
|
||||
}
|
||||
}
|
||||
return colors;
|
||||
}
|
||||
|
||||
// ── FLEXIBLE / Mockup — each scene IS a comp; layers frl_/frd_ inside ───────
|
||||
function scanFlexible(proj) {
|
||||
var result = { source: "ae-jsx", render_comp: null, scenes: [], shared_colors: readSharedColors(proj) };
|
||||
for (var i = 1; i <= proj.numItems; i++) {
|
||||
var item = proj.item(i);
|
||||
if (!(item instanceof CompItem)) continue;
|
||||
var nm = item.name;
|
||||
|
||||
if (nm === "frfinal") { result.render_comp = "frfinal"; continue; }
|
||||
if (nm === "frshare") {
|
||||
for (var j = 1; j <= item.numLayers; j++) {
|
||||
var cl = item.layer(j), ct = readText(cl);
|
||||
result.shared_colors.push({
|
||||
element_key: cl.name, title: cl.name, attr_value: "fill",
|
||||
default_color: (ct && isColor(ct.text)) ? normColor(ct.text) : "#000000", sort: j
|
||||
});
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (nm === "frfinal" || nm === "flatrender") { result.render_comp = nm; continue; }
|
||||
if (nm === "frshare") continue;
|
||||
if (!compHasEditable(item)) continue;
|
||||
|
||||
var s = scanScene(item);
|
||||
result.scenes.push({
|
||||
key: nm, title: nm, scene_type: "Normal",
|
||||
@@ -168,6 +172,52 @@
|
||||
sort: result.scenes.length, elements: s.elements, colors: s.colors
|
||||
});
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
// ── FIX / MusicVisualizer — scenes encoded in layer names frl_c(x)t/m(y) ────
|
||||
function scanFix(proj) {
|
||||
var result = { source: "ae-jsx", render_comp: "frfinal", scenes: [], shared_colors: readSharedColors(proj) };
|
||||
var sceneMap = {}, order = [];
|
||||
var re = /^frl_c(\d+)([tm])(\d+)$/;
|
||||
for (var i = 1; i <= proj.numItems; i++) {
|
||||
var item = proj.item(i);
|
||||
if (!(item instanceof CompItem) || item.name === "frshare") continue;
|
||||
for (var l = 1; l <= item.numLayers; l++) {
|
||||
var layer = item.layer(l), m = String(layer.name).match(re);
|
||||
if (!m) continue;
|
||||
var sc = parseInt(m[1], 10), typ = m[2], idx = parseInt(m[3], 10);
|
||||
if (!sceneMap[sc]) {
|
||||
sceneMap[sc] = { key: "c" + sc, title: "Scene " + sc, scene_type: "Normal", sort: sc, elements: [], colors: [] };
|
||||
order.push(sc);
|
||||
}
|
||||
var el = { key: layer.name, title: layer.name, sort: idx };
|
||||
if (typ === "t") {
|
||||
var txt = readText(layer);
|
||||
el.type = "Text";
|
||||
if (txt) { el.default_value = txt.text; el.font_face = txt.font; el.font_size = txt.fontSize; el.justify = txt.justify; }
|
||||
} else {
|
||||
el.type = "Media";
|
||||
var src = null; try { src = layer.source; } catch (e) {}
|
||||
if (src) { try { el.width = src.width; el.height = src.height; } catch (e) {} try { el.video_support = !!src.hasVideo; } catch (e) {} }
|
||||
}
|
||||
sceneMap[sc].elements.push(el);
|
||||
}
|
||||
}
|
||||
order.sort(function (a, b) { return a - b; });
|
||||
for (var k = 0; k < order.length; k++) { var s = sceneMap[order[k]]; s.sort = k; result.scenes.push(s); }
|
||||
return result;
|
||||
}
|
||||
|
||||
function run() {
|
||||
// Silence ALL AE dialogs (missing fonts/footage, version prompts, alerts) so a
|
||||
// real project never hangs the headless scan waiting for a click.
|
||||
try { app.beginSuppressDialogs(); } catch (e) {}
|
||||
try { app.preferences.savePrefAsLong("Misc Section", "Play sound when render finishes", 0, PREFType.PREF_Type_MACHINE_INDEPENDENT); } catch (e) {}
|
||||
if (aepPath) app.open(new File(aepPath));
|
||||
var proj = app.project;
|
||||
|
||||
var result = (mode === "fix" || mode === "musicvisualizer") ? scanFix(proj) : scanFlexible(proj);
|
||||
|
||||
var out = new File(outPath);
|
||||
out.encoding = "UTF-8";
|
||||
@@ -181,6 +231,13 @@
|
||||
} catch (e) {
|
||||
try { var ef = new File(outPath + ".error"); ef.open("w"); ef.write("scan error: " + e.toString()); ef.close(); } catch (e2) {}
|
||||
}
|
||||
// Quit AE so the headless afterfx process returns (set FR_SCAN_QUIT=0 to keep open while debugging).
|
||||
try { if (getenv("FR_SCAN_QUIT") !== "0") app.quit(); } catch (e) {}
|
||||
try { app.endSuppressDialogs(false); } catch (e) {}
|
||||
// Quit AE without a save prompt so the headless afterfx process returns
|
||||
// (set FR_SCAN_QUIT=0 to keep it open while debugging).
|
||||
try {
|
||||
if (getenv("FR_SCAN_QUIT") !== "0") {
|
||||
try { app.project.close(CloseOptions.DO_NOT_SAVE_CHANGES); } catch (e2) {}
|
||||
app.quit();
|
||||
}
|
||||
} catch (e) {}
|
||||
})();
|
||||
|
||||
@@ -157,6 +157,7 @@ func main() {
|
||||
v1.POST("/template-scans/:project_id/quick", auth, admin, scanH.QuickScan) // headless Go quick-scan
|
||||
v1.POST("/template-scans/:project_id/jobs", auth, admin, scanH.CreateJob) // queue an AE full scan
|
||||
v1.GET("/template-scan-jobs/:id", auth, admin, scanH.GetJob)
|
||||
v1.POST("/template-scan-jobs/:id/cancel", auth, admin, scanH.Cancel)
|
||||
|
||||
// ── Exports management (admin: all users' rendered videos) ────────────────
|
||||
adminExports := v1.Group("/admin-exports", auth, admin)
|
||||
@@ -185,6 +186,7 @@ func main() {
|
||||
|
||||
// AE scan jobs (node claims, runs scan.jsx, posts the ScanResult back)
|
||||
internal.POST("/scan/claim", scanH.Claim)
|
||||
internal.GET("/scan/:id/status", scanH.Status) // node watchdog (cancel detection)
|
||||
internal.POST("/scan/:id/result", scanH.Result)
|
||||
internal.POST("/scan/:id/fail", scanH.Fail)
|
||||
}
|
||||
|
||||
@@ -22,13 +22,17 @@ type ScanJob struct {
|
||||
type ScanClaim struct {
|
||||
ID uuid.UUID
|
||||
ProjectID uuid.UUID
|
||||
Mode string // fix | flexible | mockup | musicvisualizer → drives scan.jsx parsing
|
||||
}
|
||||
|
||||
func (s *Store) CreateScanJob(ctx context.Context, projectID uuid.UUID, engine string) (uuid.UUID, error) {
|
||||
func (s *Store) CreateScanJob(ctx context.Context, projectID uuid.UUID, engine, mode string) (uuid.UUID, error) {
|
||||
if mode == "" {
|
||||
mode = "flexible"
|
||||
}
|
||||
var id uuid.UUID
|
||||
err := s.pool.QueryRow(ctx,
|
||||
`INSERT INTO render.scan_jobs (project_id, engine, status) VALUES ($1, $2, 'queued') RETURNING id`,
|
||||
projectID, engine).Scan(&id)
|
||||
`INSERT INTO render.scan_jobs (project_id, engine, status, mode) VALUES ($1, $2, 'queued', $3) RETURNING id`,
|
||||
projectID, engine, mode).Scan(&id)
|
||||
return id, err
|
||||
}
|
||||
|
||||
@@ -44,7 +48,7 @@ func (s *Store) ClaimScanJob(ctx context.Context, nodeID uuid.UUID) (*ScanClaim,
|
||||
ORDER BY created_at
|
||||
LIMIT 1 FOR UPDATE SKIP LOCKED
|
||||
)
|
||||
RETURNING id, project_id`, nodeID).Scan(&c.ID, &c.ProjectID)
|
||||
RETURNING id, project_id, mode`, nodeID).Scan(&c.ID, &c.ProjectID, &c.Mode)
|
||||
if err != nil {
|
||||
if err == pgx.ErrNoRows {
|
||||
return nil, nil
|
||||
@@ -54,19 +58,42 @@ func (s *Store) ClaimScanJob(ctx context.Context, nodeID uuid.UUID) (*ScanClaim,
|
||||
return &c, nil
|
||||
}
|
||||
|
||||
// SetScanResult / SetScanError only act on a 'running' job, so a result that
|
||||
// arrives after the user cancelled doesn't un-cancel it.
|
||||
func (s *Store) SetScanResult(ctx context.Context, id uuid.UUID, resultJSON string) error {
|
||||
_, err := s.pool.Exec(ctx,
|
||||
`UPDATE render.scan_jobs SET status = 'done', result = $2::jsonb, error = NULL, updated_at = NOW() WHERE id = $1`,
|
||||
`UPDATE render.scan_jobs SET status = 'done', result = $2::jsonb, error = NULL, updated_at = NOW() WHERE id = $1 AND status = 'running'`,
|
||||
id, resultJSON)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *Store) SetScanError(ctx context.Context, id uuid.UUID, msg string) error {
|
||||
_, err := s.pool.Exec(ctx,
|
||||
`UPDATE render.scan_jobs SET status = 'error', error = $2, updated_at = NOW() WHERE id = $1`, id, msg)
|
||||
`UPDATE render.scan_jobs SET status = 'error', error = $2, updated_at = NOW() WHERE id = $1 AND status = 'running'`, id, msg)
|
||||
return err
|
||||
}
|
||||
|
||||
// CancelScanJob marks a queued/running scan as cancelled (user-requested). The
|
||||
// node's watchdog sees this and kills the AE process.
|
||||
func (s *Store) CancelScanJob(ctx context.Context, id uuid.UUID) error {
|
||||
_, err := s.pool.Exec(ctx,
|
||||
`UPDATE render.scan_jobs SET status = 'cancelled', error = 'cancelled by user', updated_at = NOW() WHERE id = $1 AND status IN ('queued','running')`, id)
|
||||
return err
|
||||
}
|
||||
|
||||
// GetScanStatus returns just the status string (lightweight, for the node watchdog).
|
||||
func (s *Store) GetScanStatus(ctx context.Context, id uuid.UUID) (string, error) {
|
||||
var st string
|
||||
err := s.pool.QueryRow(ctx, `SELECT status FROM render.scan_jobs WHERE id = $1`, id).Scan(&st)
|
||||
if err != nil {
|
||||
if err == pgx.ErrNoRows {
|
||||
return "", nil
|
||||
}
|
||||
return "", err
|
||||
}
|
||||
return st, nil
|
||||
}
|
||||
|
||||
func (s *Store) GetScanJob(ctx context.Context, id uuid.UUID) (*ScanJob, error) {
|
||||
var j ScanJob
|
||||
err := s.pool.QueryRow(ctx,
|
||||
|
||||
@@ -173,19 +173,27 @@ func extractAepFromZip(zb []byte) ([]byte, error) {
|
||||
}
|
||||
|
||||
// ── AE scan jobs (async, full fidelity) ───────────────────────────────────────
|
||||
// POST /v1/template-scans/:project_id/jobs (admin)
|
||||
// POST /v1/template-scans/:project_id/jobs (admin) body: {"mode":"fix|flexible|mockup|musicvisualizer"}
|
||||
func (h *ScanHandler) CreateJob(c *gin.Context) {
|
||||
pid, err := uuid.Parse(c.Param("project_id"))
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadRequest, models.APIError{Code: "bad_request", Message: "invalid project_id"})
|
||||
return
|
||||
}
|
||||
id, err := h.store.CreateScanJob(c.Request.Context(), pid, "ae-jsx")
|
||||
var req struct {
|
||||
Mode string `json:"mode"`
|
||||
}
|
||||
_ = c.ShouldBindJSON(&req)
|
||||
mode := strings.ToLower(strings.TrimSpace(req.Mode))
|
||||
if mode == "" {
|
||||
mode = "flexible"
|
||||
}
|
||||
id, err := h.store.CreateScanJob(c.Request.Context(), pid, "ae-jsx", mode)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, models.APIError{Code: "internal_error", Message: err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"id": id, "status": "queued"})
|
||||
c.JSON(http.StatusOK, gin.H{"id": id, "status": "queued", "mode": mode})
|
||||
}
|
||||
|
||||
// GET /v1/template-scan-jobs/:id (admin)
|
||||
@@ -223,9 +231,18 @@ func (h *ScanHandler) Claim(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
url, isBundle, md5 := resolveTemplateObject(h.minio, h.templatesBucket, claim.ProjectID)
|
||||
if url == "" {
|
||||
// No template object stored for this project — fail the job with a clear
|
||||
// message instead of handing the node an empty URL.
|
||||
_ = h.store.SetScanError(c.Request.Context(), claim.ID,
|
||||
"no template stored for this project — upload the .aep from «فایلها» first")
|
||||
c.Status(http.StatusNoContent)
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{
|
||||
"scan_job_id": claim.ID,
|
||||
"project_id": claim.ProjectID,
|
||||
"mode": claim.Mode,
|
||||
"aep_download_url": url,
|
||||
"is_bundle": isBundle,
|
||||
"bundle_md5": md5,
|
||||
@@ -272,6 +289,35 @@ func (h *ScanHandler) Fail(c *gin.Context) {
|
||||
c.Status(http.StatusNoContent)
|
||||
}
|
||||
|
||||
// POST /v1/template-scan-jobs/:id/cancel (admin)
|
||||
func (h *ScanHandler) Cancel(c *gin.Context) {
|
||||
id, err := uuid.Parse(c.Param("id"))
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadRequest, models.APIError{Code: "bad_request", Message: "invalid id"})
|
||||
return
|
||||
}
|
||||
if err := h.store.CancelScanJob(c.Request.Context(), id); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, models.APIError{Code: "internal_error", Message: err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"status": "cancelled"})
|
||||
}
|
||||
|
||||
// GET /v1/internal/scan/:id/status (node watchdog, HMAC) → {"status": "..."}
|
||||
func (h *ScanHandler) Status(c *gin.Context) {
|
||||
id, err := uuid.Parse(c.Param("id"))
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadRequest, models.APIError{Code: "bad_request", Message: "invalid id"})
|
||||
return
|
||||
}
|
||||
st, err := h.store.GetScanStatus(c.Request.Context(), id)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, models.APIError{Code: "internal_error", Message: err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"status": st})
|
||||
}
|
||||
|
||||
// resolveTemplateObject presigns the canonical template object for a project,
|
||||
// probing bundle.zip → template.aep → template.aepx (same order as render claim).
|
||||
func resolveTemplateObject(mc *minio.Client, bucket string, projectID uuid.UUID) (url string, isBundle bool, md5 string) {
|
||||
|
||||
@@ -128,13 +128,30 @@ func (h *TemplateBundleHandler) Set(c *gin.Context) {
|
||||
})
|
||||
}
|
||||
|
||||
// parseObjectURL extracts the bucket and object key from a path-style storage URL
|
||||
// such as http://host:9000/<bucket>/<key...>. Query/fragment are ignored.
|
||||
// parseObjectURL extracts the bucket and object key from a storage URL. Handles
|
||||
// two forms:
|
||||
// - minio://<bucket>/<key...> (file-svc FileAddress — bucket is the HOST)
|
||||
// - http(s)://host[:port]/<bucket>/<key> (public path-style — bucket is the 1st path seg)
|
||||
// Query/fragment are ignored.
|
||||
func parseObjectURL(raw string) (bucket, key string, err error) {
|
||||
u, perr := url.Parse(strings.TrimSpace(raw))
|
||||
if perr != nil {
|
||||
return "", "", fmt.Errorf("parse url: %w", perr)
|
||||
}
|
||||
|
||||
// minio://<bucket>/<key> — the bucket lives in the host component.
|
||||
if u.Scheme == "minio" {
|
||||
k := strings.TrimPrefix(u.Path, "/")
|
||||
if u.Host == "" || k == "" {
|
||||
return "", "", fmt.Errorf("cannot derive bucket/key from %q", raw)
|
||||
}
|
||||
if dec, derr := url.PathUnescape(k); derr == nil {
|
||||
k = dec
|
||||
}
|
||||
return u.Host, k, nil
|
||||
}
|
||||
|
||||
// http(s)://host[:port]/<bucket>/<key...>
|
||||
p := strings.TrimPrefix(u.Path, "/")
|
||||
if p == "" {
|
||||
return "", "", fmt.Errorf("url has no path: %q", raw)
|
||||
|
||||
Reference in New Issue
Block a user