|
// #2435: carry any lingering failure history so the ordering below |
|
// can deprioritise a recently-failed issue within its bucket. |
|
recentFailures: h.recentFailureCountKey(itemKey), |
|
}) |
|
} |
|
} |
|
|
|
if len(candidates) == 0 { |
|
// #2546: the hub is running and unsuspended but nothing is admissible right |
|
// now (everything is in cooldown, filtered out, disabled, or already held). |
|
// Previously a bare nil — indistinguishable on the wire from "suspended" or |
|
// "hub not ready". Send an explicit no_matching_work negative-ack. |
|
return h.taskUnavailable(taskUnavailableNoMatchingWork) |
|
} |
|
|
|
// Operator priority override (#queue-reorder): the ordered list of issue keys |
|
// the operator dragged to the front of the ready-work queue on the Operations |
|
// tab. It takes precedence over the default ordering below so a prioritised |
|
// issue is OFFERED FIRST. It never bypasses admission: every entry in |
|
// `candidates` already passed the SAME cooldown / failure / disabled-repo / |
|
// filter / in-flight exclusions above, so a pinned-but-no-longer-actionable key |
|
// simply never became a candidate (stale keys are skipped). Rank sentinel: a |
|
// candidate NOT in the override ranks at len(override), so all pinned candidates |
|
// sort ahead of all unpinned ones while their own relative order is the operator's. |
|
var queueOrderIdx map[string]int |
|
if h.server.deps != nil && h.server.deps.Config != nil { |
|
queueOrderIdx = queueOrderIndex(h.server.deps.Config.Hub.ContributeQueueOrder) |
|
} |
|
orderRank := func(c candidate) int { |
|
if len(queueOrderIdx) == 0 { |
|
return 0 // no override → every candidate ties, key is a no-op |
|
} |
|
if r, ok := queueOrderIdx[fmt.Sprintf("%s#%d", c.repoFull, c.number)]; ok { |
|
return r |
|
} |
|
return len(queueOrderIdx) |
|
} |
|
|
|
// Order the admissible set with a STABLE sort so the pick is deterministic |
|
// (easy to reason about and to test — no randomness): |
|
// 0. operator priority override first (#queue-reorder) — pinned issues in the |
|
// operator's dragged order; a no-op when no override is set; |
|
// 1. own-work first (#2390 — preserved unchanged); |
|
// 2. then fewer recent failures first (#2435 remedy 3 backstop) — an issue |
|
// whose short failure cooldown has just elapsed but which still carries |
|
// failure history is deprioritised behind never-failed peers, so a |
|
// flaky issue can no longer monopolise the head of the queue even if the |
|
// ledger is imperfect; |
|
// 3. otherwise the established per-repo / creation scan order is kept. |
|
// When the contributor has no own work AND nothing has failed AND no override is |
|
// set, this is a no-op and behaviour is identical to the previous first-eligible pick. |
|
ownFirst := make([]candidate, len(candidates)) |
|
copy(ownFirst, candidates) |
|
sort.SliceStable(ownFirst, func(i, j int) bool { |
|
if ri, rj := orderRank(ownFirst[i]), orderRank(ownFirst[j]); ri != rj { |
|
return ri < rj // operator-pinned (lower rank) sorts ahead |
|
} |
|
if ownFirst[i].isOwn != ownFirst[j].isOwn { |
|
return ownFirst[i].isOwn // own work sorts ahead of non-own |
|
} |
|
if ownFirst[i].interestMatch != ownFirst[j].interestMatch { |
|
return ownFirst[i].interestMatch // label-affinity matches ahead of non-matches (#2637) |
|
} |
|
if ownFirst[i].recentFailures != ownFirst[j].recentFailures { |
|
return ownFirst[i].recentFailures < ownFirst[j].recentFailures // fewer failures first |
|
} |
|
return false // equal keys → SliceStable preserves original scan order |
|
}) |
|
|
|
chosen := ownFirst[0] |
|
if chosen.isOwn { |
|
h.logger.Info("[contribute-ws] prioritizing contributor's own work (#2390)", |
|
"username", ownUsername, "repo", chosen.repoFull, "number", chosen.number) |
|
} |
|
|
|
// Mint through the shared path so task_assign and the heartbeat token-refresh |
|
// advertise tokens minted the same way (#2393 item 2). tokenMintedAt below |
|
// arms the refresh ticker for the token we hand out here. C4: the token is |
|
// scoped to the chosen issue's REPOSITORY, not the whole installation. |
|
ghToken, err := h.mintScopedToken(c.profile.TrustTier, chosen.repoFull) |
|
if err != nil { |
|
// #2436 finding 1: a mint failure previously returned nil, stranding the |
|
// contributor with no message (the log even said "skipping task" while |
|
// abandoning the whole selection). Send an explicit token_mint_failed |
|
// negative-ack so the failure is diagnosable instead of an indefinite |
|
// hang. We do not fall through to another candidate: the mint is keyed on |
|
// the contributor's tier, not the candidate, so every candidate in this |
|
// pass would fail identically. Preserve the existing Warn log. |
|
h.logger.Warn("[contribute-ws] failed to mint scoped token — task unavailable", |
|
"tier", c.profile.TrustTier, "error", err) |
|
return h.taskUnavailable(taskUnavailableTokenMintFailed) |
|
} |
|
|
|
// The task id carries the item's own identity segment so two zero-numbered |
|
// external items cannot mint the same id within one second |
|
// (kubestellar/hive#4245). GitHub-backed work keeps its historical |
|
// "ct-<repo>-<number>-<unix>" shape byte for byte. |
|
taskID := fmt.Sprintf("ct-%s-%s-%d", chosen.repoFull, taskIDSegment(chosen.ref), time.Now().Unix()) |
|
|
|
// #2539: build the prompt through the shared, credential-free buildTaskPrompt |
|
// so the exact text shipped in task_assign below can also be PREVIEWED |
|
// read-only in the ops tab. The prompt is a pure function of task metadata — |
|
// the minted github_token is attached to the WSMessage separately (never inside |
|
// the prompt), so previewing the prompt can never leak the token. buildTaskPrompt |
|
// itself carries the #2545 workspace-clone instruction (real checkout into |
|
// $HIVE_WORKSPACE_DIR rather than a fork-only --clone=false). |
|
canPush := h.contributorCanPush(chosen.repoFull, ownUsername) |
|
prompt := buildTaskPromptForContributor(chosen.ref, chosen.title, canPush) |
|
if requestedRole != "" { |
|
prompt = buildRoleTaskPromptForContributor(chosen.ref, chosen.title, requestedRole, h.roleKickPrompt(requestedRole), canPush) |
|
} |
|
// #4105: tell the agent up front — from the hub's own handshake-recorded |
|
// invocation values — the exact attribution trailer its PR body must end |
|
// with, so the footer is intentionally produced rather than appended only |
|
// by the post-merge reconciliation safety net (#4088, unchanged). |
|
prompt += attributionPromptInstruction(promptInvocationMeta(c)) |
|
|
|
// #2568: mint a fresh assignment generation for this task. It is stamped on the |
|
// connection, shipped in task_assign below, and echoed back by the relay so a |
|
// later stale-worker completion carrying an older generation is fenced out. |
|
gen := h.nextTaskGen() |
|
|
|
c.mu.Lock() |
|
c.currentTask = &WSTaskAssign{ |
|
TaskID: taskID, |
|
Kind: "issue", |
|
Role: requestedRole, |
|
Repo: chosen.repoFull, |
|
Number: chosen.number, |
|
Title: chosen.title, |
|
Key: chosen.ref.Key(), |
|
SourceType: chosen.ref.SourceType, |
|
ExternalID: chosen.ref.ExternalID, |
|
URL: chosen.url, |
|
} |
|
c.currentTaskGen = gen |
|
// #2568: start the hub-owned lease clock. task_progress renews it; cleanupLoop |
|
// auto-releases the task if it is not renewed within wsTaskTimeout. |
|
c.lastLeaseRenew = time.Now() |
|
// Duration anchor for the run log — lastLeaseRenew moves on every |
|
// progress report, so it cannot serve as the start time. |
|
c.taskAssignedAt = time.Now() |
|
// Store the prompt (never the token) so FleetSnapshot can preview it (#2539), |
|
// and clear any stale idle reason now that this connection has real work. |
|
c.currentPrompt = prompt |
|
c.currentLabels = chosen.labels |
|
c.lastIdleReason = "" |
|
// #2537: hold the minted scoped token as PENDING rather than shipping it in the |
|
// task_assign below. It is delivered only AFTER the acceptance decision — see |
|
// the ready-handler (auto-accept default) and the task_accepted handler |
|
// (explicit-accept mode) — via deliverTaskCredential. A fresh assignment resets |
|
// the delivered flag so the new task's credential is (re)delivered post-accept. |
|
// tokenMintedAt is set here so the #2393 refresh cycle is armed for the token we |
|
// hand out; deliverTaskCredential re-stamps it on actual delivery to anchor the |
|
// 50-minute refresh on when the relay truly received the credential. |
|
c.pendingToken = ghToken |
|
c.credentialDelivered = false |
|
c.tokenMintedAt = time.Now() |
|
c.mu.Unlock() |
|
|
|
// C4: record the SERVER-AUTHORITATIVE lease for this assignment so a later |
|
// reconnect can be validated against what the hub actually issued — the exact |
|
// {task, repo, generation, tier} bound here — instead of reconstructing ownership |
|
// from client-supplied task_progress fields. Revoked on every release path. |
|
// #5681: record the item's canonical key too, so the double-assignment guard can |
|
// recognise the lease after a restart — including for external work, whose |
|
// identity is Key rather than repo#number (#4245). |
|
h.recordLeaseForKey(identityOf(c), taskID, chosen.repoFull, chosen.number, |
|
chosen.ref.Key(), c.profile.TrustTier, gen, time.Now()) |
|
|
|
// #2566: record this assignment against the identity's rolling hourly/daily |
|
// windows so the next selectTask enforces tier_limits.max_per_hour / |
|
// max_per_day. Recorded here — after the task is committed to the connection |
|
// and we are certain a task_assign will ship — so a refused pass (which returns |
|
// early above) never consumes a slot. Uses the same identity key as the |
|
// concurrency gate. |
|
h.recordAssignment(identityOf(c), time.Now()) |
|
|
|
return &WSMessage{ |
|
Type: "task_assign", |
|
Seq: h.nextSeq(), |
|
TaskID: taskID, |
|
TaskGen: gen, |
|
Kind: "issue", |
|
Role: requestedRole, |
|
Repo: chosen.repoFull, |
|
Number: chosen.number, |
|
Title: chosen.title, |
|
URL: chosen.url, |
|
// Source-aware identity (kubestellar/hive#4245). Additive: a GitHub |
|
// assignment carries exactly the fields it always did, and these tell a |
|
// relay working a Linear or Jira item what the item actually IS instead |
|
// of leaving it to infer one from `number: 0`. |
|
TaskKey: chosen.ref.Key(), |
|
SourceType: chosen.ref.SourceType, |
|
ExternalID: chosen.ref.ExternalID, |
|
// #2537: NO github_token / token_expires_at here. The scoped credential is |
|
// split out of task_assign and delivered only after acceptance (see |
|
// pendingToken / deliverTaskCredential). task_assign now carries exactly the |
|
// metadata needed to DECIDE — repo/number/title/url/labels/prompt — plus the |
|
// #2568 TaskGen lease token, and no credential, so nothing an agent could act |
|
// on is authenticated until the task's source has been accepted under the |
|
// operator/contributor policy. |
|
Prompt: prompt, |
|
// The chosen issue's own labels — the Labels envelope field was declared |
Part of #2534.
Downstream observation
projectbluefin/reviewhas removed its separate maintainer dashboard and now runs review batches inside one OMP workbench. We have not merged the Hive contributor path into that workbench because the current public contributor contract does not expose a safe handoff for an already-running OMP session.Verified on 2026-09-14 against Hive
v4commit8cf12bac205912dc615731bd2dccb6dabd89a791:readyhandler clears any prior task, callsselectTask, and sends one server-selectedtask_assign:hive/src/pkg/dashboard/contribute_ws.go
Lines 3860 to 3984 in 8cf12ba
selectTaskchooses one admitted candidate, records the lease/generation, and returns one assignment envelope:hive/src/pkg/dashboard/contribute_ws.go
Lines 6380 to 6584 in 8cf12ba
hive/bin/contributor-relay.js
Lines 140 to 172 in 8cf12ba
hive/bin/contributor-relay.js
Lines 1108 to 1253 in 8cf12ba
The downstream consequence is ours: without another upstream seam, one visual OMP workbench cannot consume contributor assignments without either embedding/reimplementing Hive's relay lifecycle or locally selecting work. We are doing neither.
review-queuenow uses OMP exclusively, while the Hive worker remains a separate foreground path.Question
Does Hive want to support embedding its contributor lifecycle into an already-running agent host such as OMP?
Possible seams, without prescribing one:
Any supported seam needs to preserve Hive's ordering and admission decisions, task generation fencing, credential-after-acceptance rule, lease renewal/release behavior, and output capture. It must not let a client request arbitrary issues or synthesize its own batch.
No downstream protocol shim or client-side assignment logic is being added. Please choose the boundary you want; a
DESIGN-RESPONSEis sufficient.