Review: one local name is still a local name; a local name is unique; recovery knows it; recipes read as instructions; a tag before a digest; ask fails at once when nothing serves
A secrets object with one local name delivered no file. Two requirements could share a local name. secret recover and the export could not tell two locals apart. The recipe check missed continued lines and read heredoc bodies as bases. repo:tag@digest kept the tag in the repository. ask now publishes mandatory, so a tool nothing serves is said at once rather than after the wait.
This commit is contained in:
@@ -831,11 +831,7 @@ func undeclaredFetches(recipe string, declared map[string]bool) (bases, copies [
|
||||
*out = append(*out, ref)
|
||||
}
|
||||
}
|
||||
for _, raw := range strings.Split(recipe, "\n") {
|
||||
line := strings.TrimSpace(raw)
|
||||
if line == "" || strings.HasPrefix(line, "#") {
|
||||
continue
|
||||
}
|
||||
for _, line := range instructions(recipe) {
|
||||
fields := strings.Fields(line)
|
||||
switch strings.ToUpper(fields[0]) {
|
||||
case "FROM":
|
||||
@@ -860,7 +856,64 @@ func undeclaredFetches(recipe string, declared map[string]bool) (bases, copies [
|
||||
note(strings.TrimPrefix(f, "--from="))
|
||||
}
|
||||
}
|
||||
case "RUN":
|
||||
// RUN --mount=type=bind,from=<image>,… reaches for an image exactly as COPY --from does.
|
||||
out = &copies
|
||||
for _, f := range fields[1:] {
|
||||
if !strings.HasPrefix(f, "--mount=") {
|
||||
continue
|
||||
}
|
||||
for _, opt := range strings.Split(strings.TrimPrefix(f, "--mount="), ",") {
|
||||
if from, found := strings.CutPrefix(opt, "from="); found {
|
||||
note(from)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return bases, copies
|
||||
}
|
||||
|
||||
// instructions is a recipe as its instructions, one per line: continuations joined, comments and
|
||||
// blank lines dropped, and heredoc bodies (`COPY <<EOF … EOF`) skipped — a Python file written into
|
||||
// an image is not a list of images to fetch. The review found a `COPY \` continued onto the next
|
||||
// line slip past the check, and a stage named on a continuation line refused as a fetch.
|
||||
func instructions(recipe string) []string {
|
||||
var out []string
|
||||
var current strings.Builder
|
||||
var heredoc string
|
||||
flush := func() {
|
||||
if line := strings.TrimSpace(current.String()); line != "" && !strings.HasPrefix(line, "#") {
|
||||
out = append(out, line)
|
||||
}
|
||||
current.Reset()
|
||||
}
|
||||
for _, raw := range strings.Split(recipe, "\n") {
|
||||
if heredoc != "" {
|
||||
if strings.TrimSpace(raw) == heredoc {
|
||||
heredoc = ""
|
||||
}
|
||||
continue
|
||||
}
|
||||
line := strings.TrimRight(raw, " \t")
|
||||
if strings.HasPrefix(strings.TrimSpace(line), "#") && current.Len() == 0 {
|
||||
continue
|
||||
}
|
||||
if strings.HasSuffix(line, "\\") {
|
||||
current.WriteString(strings.TrimSuffix(line, "\\"))
|
||||
current.WriteString(" ")
|
||||
continue
|
||||
}
|
||||
current.WriteString(line)
|
||||
if at := strings.Index(current.String(), "<<"); at >= 0 {
|
||||
// `<<EOF`, `<<-EOF`, `<<'EOF'`, `<<"EOF"`: the body runs to a line that is the word.
|
||||
word := strings.Fields(current.String()[at+2:])
|
||||
if len(word) > 0 {
|
||||
heredoc = strings.Trim(strings.TrimPrefix(word[0], "-"), `'"`)
|
||||
}
|
||||
}
|
||||
flush()
|
||||
}
|
||||
flush()
|
||||
return out
|
||||
}
|
||||
|
||||
@@ -52,6 +52,11 @@ func parseReference(ref string) (upstream, error) {
|
||||
name, reference := ref, "latest"
|
||||
if at := strings.Index(ref, "@"); at >= 0 {
|
||||
name, reference = ref[:at], ref[at+1:]
|
||||
// `repo:tag@digest` is what a runtime prints; the digest names the image and the tag is
|
||||
// only what it was called. The tag is not part of the repository.
|
||||
if colon := strings.LastIndex(name, ":"); colon > strings.LastIndex(name, "/") {
|
||||
name = name[:colon]
|
||||
}
|
||||
} else if colon := strings.LastIndex(ref, ":"); colon > strings.LastIndex(ref, "/") {
|
||||
name, reference = ref[:colon], ref[colon+1:]
|
||||
}
|
||||
|
||||
@@ -208,3 +208,14 @@ func TestAReferenceIsReadTheWayARuntimeReadsIt(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// `repo:tag@digest` is what a runtime prints; the tag is not part of the repository (review C6).
|
||||
func TestATagBeforeTheDigestIsNotPartOfTheRepository(t *testing.T) {
|
||||
got, err := parseReference("quay.io/minio/mc:RELEASE.2025@sha256:abc")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got.repository != "minio/mc" || got.reference != "sha256:abc" {
|
||||
t.Fatalf("got %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -133,3 +133,22 @@ FROM golang:1.25-alpine AS go
|
||||
t.Fatalf("declared arguments are not fetches: %v", copies)
|
||||
}
|
||||
}
|
||||
|
||||
// Continued lines are one instruction, heredoc bodies are not instructions, and a RUN --mount reaches
|
||||
// for an image as a COPY --from does (review C4, C5).
|
||||
func TestARecipeIsReadAsInstructions(t *testing.T) {
|
||||
recipe := "ARG RUNTIME_BASE\n" +
|
||||
"FROM ${RUNTIME_BASE} \\\n AS build\n" +
|
||||
"COPY \\\n --from=docker.io/vendor/one:latest /a /a\n" +
|
||||
"COPY --from=build /out /out\n" +
|
||||
"COPY <<EOF /app/x.py\nfrom os import path\nEOF\n" +
|
||||
"RUN --mount=type=bind,from=docker.io/vendor/two:1,target=/t cp /t/x /x\n" +
|
||||
"FROM scratch\n"
|
||||
bases, copies := undeclaredFetches(recipe, map[string]bool{"RUNTIME_BASE": true})
|
||||
if strings.Join(copies, "|") != "docker.io/vendor/one:latest|docker.io/vendor/two:1" {
|
||||
t.Fatalf("copies: %v", copies)
|
||||
}
|
||||
if len(bases) != 0 {
|
||||
t.Fatalf("a heredoc line or a continued stage was read as a base: %v", bases)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -651,7 +651,10 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant,
|
||||
if sorted[i].Consumer != sorted[j].Consumer {
|
||||
return sorted[i].Consumer < sorted[j].Consumer
|
||||
}
|
||||
return sorted[i].From < sorted[j].From
|
||||
if sorted[i].From != sorted[j].From {
|
||||
return sorted[i].From < sorted[j].From
|
||||
}
|
||||
return sorted[i].Local < sorted[j].Local
|
||||
})
|
||||
// A consumer already carried by the grants loop, keyed (provision, module). When provider and
|
||||
// consumer are co-located, `grantsFor` enumerates the same-node consumer too, so without this the
|
||||
@@ -788,6 +791,8 @@ type Kept struct {
|
||||
Sealed string `json:"sealed"`
|
||||
Key string `json:"key"`
|
||||
MadeAt time.Time `json:"made-at"`
|
||||
// Local is the credential's name inside the consumer where it holds several (ADR 0094).
|
||||
Local string `json:"local,omitempty"`
|
||||
}
|
||||
|
||||
// KeptExport is what a person keeps beside the operator key, and what a vault keeps on its disk:
|
||||
|
||||
@@ -1050,6 +1050,7 @@ func ParseManifest(raw []byte) (Manifest, error) {
|
||||
problems = append(problems, m.Module+" needs a secret with no name")
|
||||
}
|
||||
}
|
||||
localOf := map[string]string{}
|
||||
for _, to := range m.SecretRequirements() {
|
||||
if _, plain := m.Secrets[to]; plain {
|
||||
if _, also := m.SecretsMany[to]; also {
|
||||
@@ -1069,9 +1070,15 @@ func ParseManifest(raw []byte) (Manifest, error) {
|
||||
m.Module, to, f.Local))
|
||||
}
|
||||
// A local name is what `${secret:<name>}` says, so it may not be another requirement's
|
||||
// name or one of the module's own secrets — the file would hold the wrong credential
|
||||
// while every check passed.
|
||||
// name, another requirement's local name, or one of the module's own secrets — the
|
||||
// file would hold the wrong credential while every check passed.
|
||||
if f.Local != "" {
|
||||
if other, taken := localOf[f.Local]; taken && other != to {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%s keeps credentials for %q and %q both under %q — a local name names one",
|
||||
m.Module, other, to, f.Local))
|
||||
}
|
||||
localOf[f.Local] = to
|
||||
if _, own := m.OwnSecrets[f.Local]; own {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%s keeps a credential for %q under %q, which is also one of its own secrets",
|
||||
|
||||
@@ -881,7 +881,9 @@ func perConsumer(needs []Needed, order []string, catalogue map[string]Manifest)
|
||||
// from one provider (ADR 0094). Each is its own pair credential downstream.
|
||||
func eachLocal(needs []Needed, catalogue map[string]Manifest, n Needed) []Needed {
|
||||
files := catalogue[n.For].SecretFiles(n.Name)
|
||||
if len(files) <= 1 {
|
||||
// One file under a local name is still a local name: the review found a module keeping ONE
|
||||
// named secret given a need with no local, and so no file, while everything reported success.
|
||||
if len(files) == 0 || (len(files) == 1 && files[0].Local == "") {
|
||||
return append(needs, n)
|
||||
}
|
||||
for _, f := range files {
|
||||
|
||||
@@ -159,3 +159,47 @@ func TestAProviderKeepsOneFilePerHolder(t *testing.T) {
|
||||
t.Fatalf("two holders are two grant files: %v", ids)
|
||||
}
|
||||
}
|
||||
|
||||
// One file under a local name is still a local name (review C1): the need carries it, the file
|
||||
// is written, and ${secret:<name>} is filled.
|
||||
func TestOneLocalNameIsStillALocalName(t *testing.T) {
|
||||
only, _ := ParseManifest([]byte(`{"module":"one","version":"1","requires":["secret"],
|
||||
"secrets":{"secret":{"only":"/var/lib/one/only"}}}`))
|
||||
vault := vaultAndCA()["mesh-vault"]
|
||||
got, err := Resolve(shelf(vault, only), []string{"mesh-vault", "one"}, workstation(), World{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var found *Needed
|
||||
for i, n := range got.Needs {
|
||||
if n.For == "one" && n.Name == "secret" {
|
||||
found = &got.Needs[i]
|
||||
}
|
||||
}
|
||||
if found == nil || found.Local != "only" {
|
||||
t.Fatalf("the one named file did not become a need under its name: %v", got.Needs)
|
||||
}
|
||||
found.Sealed = "sealed-only"
|
||||
out, err := got.Declaration(Rendering{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var written bool
|
||||
for _, r := range out {
|
||||
if r["path"] == "/var/lib/one/only" && r["sealed"] == "sealed-only" {
|
||||
written = true
|
||||
}
|
||||
}
|
||||
if !written {
|
||||
t.Fatal("the file under the one local name was not written")
|
||||
}
|
||||
}
|
||||
|
||||
// A local name names one credential: two requirements may not share it (review C2).
|
||||
func TestALocalNameIsUniqueAcrossRequirements(t *testing.T) {
|
||||
_, err := ParseManifest([]byte(`{"module":"x","version":"1","requires":["secret","postgres-database"],
|
||||
"secrets":{"secret":{"x":"/var/lib/x/a"},"postgres-database":{"x":"/var/lib/x/b"}}}`))
|
||||
if err == nil || !strings.Contains(err.Error(), "both under") {
|
||||
t.Fatalf("two requirements under one local name were accepted: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -124,20 +124,20 @@ func (i *Inventory) KeptForOperator(ctx context.Context) (kept, earlier, unrecov
|
||||
}
|
||||
rows, err := i.store.Pool().Query(ctx,
|
||||
`select 'own', n.name, s.module, s.name, '', s.origin, coalesce(s.operator_sealed, ''),
|
||||
coalesce(s.operator_key, ''), s.made_at
|
||||
coalesce(s.operator_key, ''), s.made_at, ''
|
||||
from module_secret s join node n on n.id = s.node
|
||||
union all
|
||||
select 'pair', c.name, s.consumer_module, s.name, p.name, 'made', coalesce(s.operator_sealed, ''),
|
||||
coalesce(s.operator_key, ''), s.created_at
|
||||
select 'pair', c.name, s.consumer_module, s.name, p.name, s.origin, coalesce(s.operator_sealed, ''),
|
||||
coalesce(s.operator_key, ''), s.created_at, s.local
|
||||
from secret s join node c on c.id = s.consumer join node p on p.id = s.provider
|
||||
order by 1, 2, 3, 4`)
|
||||
order by 1, 2, 3, 4, 10`)
|
||||
if err != nil {
|
||||
return nil, nil, nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
for rows.Next() {
|
||||
var k Kept
|
||||
if err := rows.Scan(&k.Kind, &k.Node, &k.Module, &k.Name, &k.Provider, &k.Origin, &k.Sealed, &k.Key, &k.MadeAt); err != nil {
|
||||
if err := rows.Scan(&k.Kind, &k.Node, &k.Module, &k.Name, &k.Provider, &k.Origin, &k.Sealed, &k.Key, &k.MadeAt, &k.Local); err != nil {
|
||||
return nil, nil, nil, err
|
||||
}
|
||||
switch {
|
||||
@@ -159,7 +159,7 @@ func (i *Inventory) KeptForOperator(ctx context.Context) (kept, earlier, unrecov
|
||||
// credential is keyed by provider as well, and a consumer whose provision moved leaves the old
|
||||
// provider's row behind: two rows is refused with both providers named, never answered with
|
||||
// whichever came first, unless `provider` says which.
|
||||
func (i *Inventory) KeptSecret(ctx context.Context, node, module, name, provider string) (Kept, error) {
|
||||
func (i *Inventory) KeptSecret(ctx context.Context, node, module, name, provider, local string) (Kept, error) {
|
||||
var k Kept
|
||||
err := i.store.Pool().QueryRow(ctx,
|
||||
`select 'own', n.name, s.module, s.name, '', s.origin, coalesce(s.operator_sealed, ''),
|
||||
@@ -169,11 +169,12 @@ func (i *Inventory) KeptSecret(ctx context.Context, node, module, name, provider
|
||||
Scan(&k.Kind, &k.Node, &k.Module, &k.Name, &k.Provider, &k.Origin, &k.Sealed, &k.Key, &k.MadeAt)
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
rows, qerr := i.store.Pool().Query(ctx,
|
||||
`select 'pair', c.name, s.consumer_module, s.name, p.name, 'made', coalesce(s.operator_sealed, ''),
|
||||
coalesce(s.operator_key, ''), s.created_at
|
||||
`select 'pair', c.name, s.consumer_module, s.name, p.name, s.origin, coalesce(s.operator_sealed, ''),
|
||||
coalesce(s.operator_key, ''), s.created_at, s.local
|
||||
from secret s join node c on c.id = s.consumer join node p on p.id = s.provider
|
||||
where c.name = $1 and s.consumer_module = $2 and s.name = $3 and ($4 = '' or p.name = $4)
|
||||
order by p.name`, node, module, name, provider)
|
||||
and s.local = $5
|
||||
order by p.name`, node, module, name, provider, local)
|
||||
if qerr != nil {
|
||||
return Kept{}, qerr
|
||||
}
|
||||
@@ -181,7 +182,7 @@ func (i *Inventory) KeptSecret(ctx context.Context, node, module, name, provider
|
||||
var found []Kept
|
||||
for rows.Next() {
|
||||
var row Kept
|
||||
if err := rows.Scan(&row.Kind, &row.Node, &row.Module, &row.Name, &row.Provider, &row.Origin, &row.Sealed, &row.Key, &row.MadeAt); err != nil {
|
||||
if err := rows.Scan(&row.Kind, &row.Node, &row.Module, &row.Name, &row.Provider, &row.Origin, &row.Sealed, &row.Key, &row.MadeAt, &row.Local); err != nil {
|
||||
return Kept{}, err
|
||||
}
|
||||
found = append(found, row)
|
||||
|
||||
@@ -21,7 +21,7 @@ func TestAnOwnSecretIsSealedToTheOperatorToo(t *testing.T) {
|
||||
if _, err := inv.SecretForModule(ctx, "consumer", "postgres", "superuser"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := inv.KeptSecret(ctx, "consumer", "postgres", "superuser", ""); err == nil {
|
||||
if _, err := inv.KeptSecret(ctx, "consumer", "postgres", "superuser", "", ""); err == nil {
|
||||
t.Fatal("a secret made before the operator key was reported recoverable")
|
||||
}
|
||||
kept, _, unrecoverable, err := inv.KeptForOperator(ctx)
|
||||
@@ -54,7 +54,7 @@ func TestAnOwnSecretIsSealedToTheOperatorToo(t *testing.T) {
|
||||
if len(kept) != 2 || len(unrecoverable) != 1 {
|
||||
t.Fatalf("after a key: %d kept, %d unrecoverable", len(kept), len(unrecoverable))
|
||||
}
|
||||
got, err := inv.KeptSecret(ctx, "provider", "postgres", "replication", "")
|
||||
got, err := inv.KeptSecret(ctx, "provider", "postgres", "replication", "", "")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -68,7 +68,7 @@ func TestAnOwnSecretIsSealedToTheOperatorToo(t *testing.T) {
|
||||
if got.Origin != "accepted" || got.Key != pub {
|
||||
t.Fatalf("kept as %+v", got)
|
||||
}
|
||||
minted, err := inv.KeptSecret(ctx, "provider", "postgres", "superuser", "")
|
||||
minted, err := inv.KeptSecret(ctx, "provider", "postgres", "superuser", "", "")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -81,7 +81,7 @@ func TestAnOwnSecretIsSealedToTheOperatorToo(t *testing.T) {
|
||||
if _, err := inv.SecretForModule(ctx, "consumer", "postgres", "superuser"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := inv.KeptSecret(ctx, "consumer", "postgres", "superuser", ""); err == nil {
|
||||
if _, err := inv.KeptSecret(ctx, "consumer", "postgres", "superuser", "", ""); err == nil {
|
||||
t.Fatal("asking again did not remake, yet it became recoverable")
|
||||
}
|
||||
// Until the node rejoins with a new sealing key: then the secret is remade, and the remake is
|
||||
@@ -97,7 +97,7 @@ func TestAnOwnSecretIsSealedToTheOperatorToo(t *testing.T) {
|
||||
if _, err := inv.SecretForModule(ctx, "consumer", "postgres", "superuser"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
remade, err := inv.KeptSecret(ctx, "consumer", "postgres", "superuser", "")
|
||||
remade, err := inv.KeptSecret(ctx, "consumer", "postgres", "superuser", "", "")
|
||||
if err != nil {
|
||||
t.Fatalf("the remade secret is not recoverable: %v", err)
|
||||
}
|
||||
@@ -168,7 +168,7 @@ func TestAPairCredentialIsSealedToTheOperatorToo(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
kept, err := inv.KeptSecret(ctx, "consumer", "gitea", "secret", "")
|
||||
kept, err := inv.KeptSecret(ctx, "consumer", "gitea", "secret", "", "")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -194,10 +194,10 @@ func TestAPairCredentialIsSealedToTheOperatorToo(t *testing.T) {
|
||||
if _, err := inv.SecretFor(ctx, "secret", "consumer", "gitea", "consumer", ""); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := inv.KeptSecret(ctx, "consumer", "gitea", "secret", ""); err == nil || !strings.Contains(err.Error(), "--provider") {
|
||||
if _, err := inv.KeptSecret(ctx, "consumer", "gitea", "secret", "", ""); err == nil || !strings.Contains(err.Error(), "--provider") {
|
||||
t.Fatalf("two providers were not refused: %v", err)
|
||||
}
|
||||
if byName, err := inv.KeptSecret(ctx, "consumer", "gitea", "secret", "provider"); err != nil || byName.Provider != "provider" {
|
||||
if byName, err := inv.KeptSecret(ctx, "consumer", "gitea", "secret", "provider", ""); err != nil || byName.Provider != "provider" {
|
||||
t.Fatalf("naming the provider did not select it: %+v %v", byName, err)
|
||||
}
|
||||
pub2, _, _ := secrets.Keypair()
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"crypto/rand"
|
||||
"encoding/base64"
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/secrets"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
@@ -673,3 +674,44 @@ func TestTwoLocalNamesAreTwoCredentials(t *testing.T) {
|
||||
t.Fatalf("the provider is told two credentials to create: %+v", from)
|
||||
}
|
||||
}
|
||||
|
||||
// The operator can recover either of two local names apart (review C3).
|
||||
func TestTheOperatorRecoversEachLocalNameApart(t *testing.T) {
|
||||
inv, ctx := twoNodesWithKeys(t)
|
||||
pub, _, err := secrets.Keypair()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := inv.SetOperatorKey(ctx, pub); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, local := range []string{"root-key", "root-pass"} {
|
||||
if _, err := inv.SecretFor(ctx, "secret", "consumer", "gitea", "provider", local); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
key, err := inv.KeptSecret(ctx, "consumer", "gitea", "secret", "provider", "root-key")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
pass, err := inv.KeptSecret(ctx, "consumer", "gitea", "secret", "provider", "root-pass")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if key.Local != "root-key" || pass.Local != "root-pass" || key.Kind != "pair" {
|
||||
t.Fatalf("recovery does not tell the two apart: %+v %+v", key, pass)
|
||||
}
|
||||
kept, _, _, err := inv.KeptForOperator(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var locals []string
|
||||
for _, k := range kept {
|
||||
if k.Kind == "pair" {
|
||||
locals = append(locals, k.Local)
|
||||
}
|
||||
}
|
||||
if strings.Join(locals, ",") != "root-key,root-pass" {
|
||||
t.Fatalf("the export does not name the local names: %v", kept)
|
||||
}
|
||||
}
|
||||
|
||||
+11
-1
@@ -49,9 +49,12 @@ func Ask(ctx context.Context, channel *amqp.Channel, module, tool string, args j
|
||||
return Answer{}, err
|
||||
}
|
||||
|
||||
// Mandatory, so a request nothing consumes comes straight back: a module that is down, or a
|
||||
// tool that does not exist, is said at once rather than after the whole wait.
|
||||
returned := channel.NotifyReturn(make(chan amqp.Return, 1))
|
||||
id := fmt.Sprintf("ask-%d", time.Now().UnixNano())
|
||||
key := module + "." + tool
|
||||
if err := channel.PublishWithContext(ctx, RPCExchange, key, false, false, amqp.Publishing{
|
||||
if err := channel.PublishWithContext(ctx, RPCExchange, key, true, false, amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
CorrelationId: id,
|
||||
ReplyTo: replies.Name,
|
||||
@@ -64,6 +67,13 @@ func Ask(ctx context.Context, channel *amqp.Channel, module, tool string, args j
|
||||
defer cancel()
|
||||
for {
|
||||
select {
|
||||
case back := <-returned:
|
||||
if back.CorrelationId == id {
|
||||
return Answer{}, fmt.Errorf(
|
||||
"nothing serves %s: no runtime has bound %q on the broker. The module is not "+
|
||||
"assigned, its runtime is not up, or it serves no such tool — `status` "+
|
||||
"says whether the machine carrying it has applied", module, key)
|
||||
}
|
||||
case <-waiting.Done():
|
||||
return Answer{}, fmt.Errorf(
|
||||
"%s did not answer within %s. Its runtime serves %q when it is up and has bound "+
|
||||
|
||||
Reference in New Issue
Block a user