Files
ombi-mcp/internal/tools/jobs.go
gronod eab16991ac Implement the runnable MCP server: 31 tools, auth, projections (Phase 07)
Lands the executable half of the Phase 01-06 contracts: stdio MCP
wiring with bundle policy enforced at call time, JWT/api_key upstream
auth with single-flight renewal and 401 retry, strict argument
decoding, allowlisted result projections with enum label twins, TV
season expansion, settings read-modify-write under a revision lock,
and a sanitized ToolError envelope that never forwards raw upstream
bodies.
2026-09-18 19:14:10 +01:00

178 lines
5.0 KiB
Go

package tools
import (
"context"
"encoding/json"
"fmt"
"ombi-mcp/internal/ombi"
)
// permittedJobs is the closed enum of POST Job routes — names match
// the ledger exactly, including case-sensitive arrAvailability.
var permittedJobs = map[string]bool{
"update": true, "plexuserimporter": true, "plexwatchlist": true,
"embyuserimporter": true, "jellyfinuserimporter": true,
"plexcontentcacher": true, "clearmediaserverdata": true,
"plexrecentlyadded": true, "embycontentcacher": true,
"embyrecentlyadded": true, "jellyfincontentcacher": true,
"arrAvailability": true, "autodeleterequests": true,
"newsletter": true,
}
// write_job_run — trigger a permitted job or revalidate watchlist
// users. Accepted does not mean completed.
func handleJobRun(ctx context.Context, env *Env, raw json.RawMessage) *ToolResult {
o := newOp(ctx, env, "write_job_run", true)
var a JobRunArgs
if fail := o.args(raw, &a); fail != nil {
return fail
}
switch a.Action {
case "run":
if !permittedJobs[a.Job] {
return o.invalid("job", "unsupported job %q", a.Job)
}
raw, fail := o.call("POST", "/api/v1/Job/"+a.Job, nil, nil)
if fail != nil {
return fail
}
return o.jobResult(raw, a.Job)
case "revalidate_watchlist":
raw, fail := o.call("POST",
"/api/v1/Plex/WatchlistUsers/revalidate", nil, nil)
if fail != nil {
return fail
}
return o.jobResult(raw, "revalidate_watchlist")
default:
return o.invalid("action", "unsupported action %q", a.Action)
}
}
func (o *op) jobResult(raw []byte, name string) *ToolResult {
m := &Mutation{Kind: "mutation", Outcome: "accepted"}
if b, fail := o.decodeBool(raw); fail == nil {
m.UpstreamResult = &b
if !b {
return o.fail("UPSTREAM_REJECTED",
fmt.Sprintf("job %q was not accepted", name), false)
}
} else if obj, fail2 := o.decodeObject(raw); fail2 == nil {
if s := jstr(obj, "message"); s != "" {
m.Message = sanitizeText(s, maxSanitizedMsg)
}
}
return o.ok(m)
}
// write_notification_send — mass email to explicit resolved user
// IDs. Recipients are resolved privately; partial resolution failure
// is reported per item rather than silently dropped.
func handleNotificationSend(ctx context.Context, env *Env, raw json.RawMessage) *ToolResult {
o := newOp(ctx, env, "write_notification_send", true)
var a NotificationSendArgs
if fail := o.args(raw, &a); fail != nil {
return fail
}
if !nonempty(a.Subject) {
return o.invalid("subject", "subject must contain non-whitespace text")
}
if !nonempty(a.Body) {
return o.invalid("body", "body must contain non-whitespace text")
}
if len(a.UserIDs) == 0 {
return o.invalid("user_ids", "user_ids must be a nonempty list")
}
if len(a.UserIDs) > 100 {
return o.invalid("user_ids", "user_ids exceeds the 100-recipient cap")
}
seen := map[string]bool{}
users := []ombi.MassEmailUser{}
var itemResults []ItemResult
failed := 0
for _, id := range a.UserIDs {
if !nonempty(id) {
return o.invalid("user_ids", "user ids must be nonempty")
}
if seen[id] {
return o.invalid("user_ids", "duplicate user id %q", id)
}
seen[id] = true
raw, fail := o.call("GET", "/api/v1/Identity/User/"+seg(id), nil, nil)
if fail != nil {
failed++
itemResults = append(itemResults,
ItemResult{Identifier: id, Outcome: "rejected",
Message: "user lookup failed"})
continue
}
u, fail := o.decodeObject(raw)
if fail != nil {
failed++
itemResults = append(itemResults,
ItemResult{Identifier: id, Outcome: "rejected",
Message: "user record unreadable"})
continue
}
users = append(users, ombi.MassEmailUser{
ID: id,
UserName: jstr(u, "userName", "username"),
Email: jstr(u, "emailAddress", "email"),
})
itemResults = append(itemResults,
ItemResult{Identifier: id, Outcome: "accepted"})
}
if len(users) == 0 {
res := o.fail("PARTIAL_FAILURE",
"no recipients could be resolved", false)
return res
}
if failed > 0 {
o.warnf("%d of %d recipients could not be resolved", failed, len(a.UserIDs))
}
body := ombi.MassEmailModel{
Subject: a.Subject,
Body: a.Body,
Bcc: a.Bcc,
Users: users,
}
raw, fail := o.call("POST", "/api/v1/Notifications/massemail", nil, body)
if fail != nil {
return fail
}
m, fail := o.boolMutation(raw, nil)
if fail != nil {
return fail
}
if failed > 0 {
m.Outcome = "partial"
m.ItemResults = itemResults
}
n := len(users)
m.AffectedCount = &n
return o.ok(m)
}
// write_retry_remove — DELETE the queue ID only; never calls a
// nonexistent POST RequestRetry and does not touch the underlying
// request.
func handleRetryRemove(ctx context.Context, env *Env, raw json.RawMessage) *ToolResult {
o := newOp(ctx, env, "write_retry_remove", true)
var a RetryRemoveArgs
if fail := o.args(raw, &a); fail != nil {
return fail
}
if a.QueueID < 1 {
return o.invalid("queue_id", "queue_id must be a positive integer")
}
raw, fail := o.call("DELETE",
"/api/v1/RequestRetry/"+segInt(a.QueueID), nil, nil)
if fail != nil {
return fail
}
m := emptyMutation("completed", nil)
m.QueueID = &a.QueueID
return o.ok(m)
}