g1t/services/work/src/lib.rs

714 lines25,622 bytesCodeBlame
1//! The work service: intents, attempts and sessions.
2//!
3//! Other services reach it over `POST /rpc/<method>`; see
4//! `g1t_contracts::work` for the methods and their arguments. It also
5//! consumes its queue of events from the bus.
6
7mod rows;
8
9use g1t_contracts::events::{
10 AttemptEvent, Delivered, IntentClosed, IntentOpened, NewEvent, SessionAppended,
11};
12use g1t_contracts::repos::{ForkArgs, GetArgs, GetByIdArgs, LandArgs, Landed, Repo};
13use g1t_contracts::time::rfc3339;
14use g1t_contracts::work::*;
15use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
16use g1t_kit::{args, js, now_ms, reply, rpc_method};
17use serde::Serialize;
18use worker::wasm_bindgen::JsValue;
19use worker::{
20 Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
21};
22
23use rows::{AttemptRow, IntentRow, SessionRow};
24
25const SOURCE: &str = "work";
26const MAX_ENTRY_BATCH: usize = 200;
27const MAX_ENTRY_CHARS: usize = 64_000;
28const SESSION_PAGE: u32 = 500;
29const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
30
31const INTENT_COLUMNS: &str = "intents.*,
32 (SELECT count(*) FROM attempts WHERE attempts.intent_id = intents.id) AS attempt_count";
33
34fn no_intent<T>() -> Outcome<T> {
35 Outcome::fail(FailureCode::NotFound, "Intent not found.")
36}
37
38fn no_attempt<T>() -> Outcome<T> {
39 Outcome::fail(FailureCode::NotFound, "Attempt not found.")
40}
41
42fn optional(value: &Option<String>) -> JsValue {
43 value.as_deref().map_or(JsValue::NULL, JsValue::from)
44}
45
46struct Work {
47 db: D1Database,
48 repos: Fetcher,
49 /// The events service, an RPC stub.
50 events: JsValue,
51}
52
53impl Work {
54 async fn publish<T: Serialize>(&self, event: NewEvent<T>) -> Result<()> {
55 js::call(&self.events, "publish", &[js::to_js(&[event])?]).await?;
56 Ok(())
57 }
58
59 async fn repo_by_path(
60 &self,
61 path: &g1t_contracts::repos::RepoPath,
62 viewer: &Viewer,
63 ) -> Result<Outcome<Repo>> {
64 g1t_kit::call(
65 &self.repos,
66 "get",
67 &GetArgs {
68 path: path.clone(),
69 viewer: viewer.clone(),
70 },
71 )
72 .await
73 }
74
75 async fn repo_by_id(&self, id: &str, viewer: &Viewer) -> Result<Outcome<Repo>> {
76 g1t_kit::call(
77 &self.repos,
78 "get_by_id",
79 &GetByIdArgs {
80 id: id.to_owned(),
81 viewer: viewer.clone(),
82 },
83 )
84 .await
85 }
86
87 async fn intent_by_id(&self, id: &str) -> Result<Option<Intent>> {
88 Ok(self
89 .db
90 .prepare(format!("SELECT {INTENT_COLUMNS} FROM intents WHERE id = ?"))
91 .bind(&[id.into()])?
92 .first::<IntentRow>(None)
93 .await?
94 .map(Intent::from))
95 }
96
97 async fn attempt_by_id(&self, id: &str) -> Result<Option<Attempt>> {
98 Ok(self
99 .db
100 .prepare("SELECT * FROM attempts WHERE id = ?")
101 .bind(&[id.into()])?
102 .first::<AttemptRow>(None)
103 .await?
104 .map(Attempt::from))
105 }
106
107 /// The attempt, if `actor` is the one running it.
108 async fn own_attempt(&self, actor: &User, id: &str) -> Result<Outcome<Attempt>> {
109 let Some(attempt) = self.attempt_by_id(id).await? else {
110 return Ok(no_attempt());
111 };
112 // Reading it must be allowed before "forbidden" may reveal it exists.
113 let viewer = Some(actor.clone());
114 if let Outcome::Fail(_) = self.repo_by_id(&attempt.repo_id, &viewer).await? {
115 return Ok(no_attempt());
116 }
117 if attempt.started_by.id != actor.id {
118 return Ok(Outcome::fail(
119 FailureCode::Forbidden,
120 "Only the person who started an attempt can change it.",
121 ));
122 }
123 Ok(Outcome::Ok(attempt))
124 }
125
126 fn attempt_event(attempt: &Attempt) -> AttemptEvent {
127 AttemptEvent {
128 attempt_id: attempt.id.clone(),
129 intent_id: attempt.intent_id.clone(),
130 repo_id: attempt.repo_id.clone(),
131 ..AttemptEvent::default()
132 }
133 }
134
135 async fn open_intent(&self, a: OpenIntentArgs) -> Result<Outcome<Intent>> {
136 if !a.actor.verified {
137 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
138 }
139 let title = a.title.trim();
140 if title.is_empty() {
141 return Ok(Outcome::fail(
142 FailureCode::Invalid,
143 "An intent needs a title.",
144 ));
145 }
146 let repo = match self.repo_by_path(&a.repo, &Some(a.actor.clone())).await? {
147 Outcome::Ok(repo) => repo,
148 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
149 };
150 let checks: Vec<&str> = a
151 .checks
152 .iter()
153 .map(|check| check.trim())
154 .filter(|check| !check.is_empty())
155 .collect();
156
157 let now = now_ms();
158 let id = new_id("int", now);
159 // Numbering and insert are one statement, so concurrent opens on the
160 // same repo cannot take the same number.
161 self.db
162 .prepare(
163 "INSERT INTO intents
164 (id, repo_id, number, title, brief, checks, author_id, author_name, created_at)
165 SELECT ?, ?, COALESCE(MAX(number), 0) + 1, ?, ?, ?, ?, ?, ?
166 FROM intents WHERE repo_id = ?",
167 )
168 .bind(&[
169 id.as_str().into(),
170 repo.id.as_str().into(),
171 title.into(),
172 a.brief.trim().into(),
173 serde_json::to_string(&checks)?.into(),
174 a.actor.id.as_str().into(),
175 a.actor.username.as_str().into(),
176 rfc3339(now).into(),
177 repo.id.as_str().into(),
178 ])?
179 .run()
180 .await?;
181 let Some(intent) = self.intent_by_id(&id).await? else {
182 return Ok(no_intent());
183 };
184 self.publish(NewEvent {
185 kind: "intent.opened",
186 source: SOURCE,
187 repo_id: Some(intent.repo_id.clone()),
188 actor: Some(a.actor.id),
189 data: IntentOpened {
190 intent_id: intent.id.clone(),
191 repo_id: intent.repo_id.clone(),
192 number: intent.number,
193 title: intent.title.clone(),
194 },
195 })
196 .await?;
197 Ok(Outcome::Ok(intent))
198 }
199
200 async fn list_intents(&self, a: ListIntentsArgs) -> Result<Outcome<Vec<Intent>>> {
201 let repo = match self.repo_by_path(&a.repo, &a.viewer).await? {
202 Outcome::Ok(repo) => repo,
203 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
204 };
205 let status = a
206 .status
207 .and_then(|status| serde_json::to_value(status).ok())
208 .and_then(|value| value.as_str().map(str::to_owned));
209 let rows = self
210 .db
211 .prepare(format!(
212 "SELECT {INTENT_COLUMNS} FROM intents
213 WHERE repo_id = ? AND (? IS NULL OR status = ?)
214 ORDER BY number DESC LIMIT 100"
215 ))
216 .bind(&[repo.id.into(), optional(&status), optional(&status)])?
217 .all()
218 .await?
219 .results::<IntentRow>()?;
220 Ok(Outcome::Ok(rows.into_iter().map(Intent::from).collect()))
221 }
222
223 async fn get_intent(&self, a: GetIntentArgs) -> Result<Outcome<IntentDetail>> {
224 let repo = match self.repo_by_path(&a.repo, &a.viewer).await? {
225 Outcome::Ok(repo) => repo,
226 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
227 };
228 let row = self
229 .db
230 .prepare(format!(
231 "SELECT {INTENT_COLUMNS} FROM intents WHERE repo_id = ? AND number = ?"
232 ))
233 .bind(&[repo.id.into(), a.number.into()])?
234 .first::<IntentRow>(None)
235 .await?;
236 let Some(intent) = row.map(Intent::from) else {
237 return Ok(no_intent());
238 };
239 let attempts = self
240 .db
241 .prepare("SELECT * FROM attempts WHERE intent_id = ? ORDER BY number")
242 .bind(&[intent.id.as_str().into()])?
243 .all()
244 .await?
245 .results::<AttemptRow>()?;
246 Ok(Outcome::Ok(IntentDetail {
247 intent,
248 attempts: attempts.into_iter().map(Attempt::from).collect(),
249 }))
250 }
251
252 async fn withdraw_intent(&self, a: IntentActionArgs) -> Result<Outcome<Intent>> {
253 let Some(mut intent) = self.intent_by_id(&a.intent_id).await? else {
254 return Ok(no_intent());
255 };
256 let repo = match self
257 .repo_by_id(&intent.repo_id, &Some(a.actor.clone()))
258 .await?
259 {
260 Outcome::Ok(repo) => repo,
261 Outcome::Fail(_) => return Ok(no_intent()),
262 };
263 if intent.author.id != a.actor.id && repo.owner_id != a.actor.id {
264 return Ok(Outcome::fail(
265 FailureCode::Forbidden,
266 "Only the author or the repo owner can withdraw an intent.",
267 ));
268 }
269 if intent.status != IntentStatus::Open {
270 return Ok(Outcome::fail(
271 FailureCode::Conflict,
272 "This intent is already closed.",
273 ));
274 }
275 self.db
276 .prepare("UPDATE intents SET status = 'withdrawn' WHERE id = ?")
277 .bind(&[intent.id.as_str().into()])?
278 .run()
279 .await?;
280 self.publish(NewEvent {
281 kind: "intent.closed",
282 source: SOURCE,
283 repo_id: Some(intent.repo_id.clone()),
284 actor: Some(a.actor.id),
285 data: IntentClosed {
286 intent_id: intent.id.clone(),
287 repo_id: intent.repo_id.clone(),
288 reason: "withdrawn",
289 },
290 })
291 .await?;
292 intent.status = IntentStatus::Withdrawn;
293 Ok(Outcome::Ok(intent))
294 }
295
296 async fn start_attempt(&self, a: StartAttemptArgs) -> Result<Outcome<Attempt>> {
297 if !a.actor.verified {
298 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
299 }
300 let Some(intent) = self.intent_by_id(&a.intent_id).await? else {
301 return Ok(no_intent());
302 };
303 if intent.status != IntentStatus::Open {
304 return Ok(Outcome::fail(
305 FailureCode::Conflict,
306 "This intent is already closed.",
307 ));
308 }
309 let agent = match a.agent.trim() {
310 "" => "agent",
311 agent => agent,
312 };
313 let runtime = match a.runtime {
314 AttemptRuntime::Hosted => "hosted",
315 AttemptRuntime::External => "external",
316 };
317
318 let now = now_ms();
319 let id = new_id("att", now);
320 let fork: Outcome<Repo> = g1t_kit::call(
321 &self.repos,
322 "fork_for_attempt",
323 &ForkArgs {
324 source_id: intent.repo_id.clone(),
325 attempt_id: id.clone(),
326 actor: a.actor.clone(),
327 },
328 )
329 .await?;
330 let fork = match fork {
331 Outcome::Ok(fork) => fork,
332 // The repo is not visible to this actor, so neither is the intent.
333 Outcome::Fail(failure) if failure.code == FailureCode::NotFound => {
334 return Ok(no_intent());
335 }
336 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
337 };
338
339 let timestamp = rfc3339(now);
340 self.db
341 .prepare(
342 "INSERT INTO attempts
343 (id, intent_id, repo_id, number, agent, runtime, fork_repo_id,
344 fork_namespace, fork_name, started_by_id, started_by_name,
345 created_at, updated_at)
346 SELECT ?, ?, ?, COALESCE(MAX(number), 0) + 1, ?, ?, ?, ?, ?, ?, ?, ?, ?
347 FROM attempts WHERE intent_id = ?",
348 )
349 .bind(&[
350 id.as_str().into(),
351 intent.id.as_str().into(),
352 intent.repo_id.as_str().into(),
353 agent.into(),
354 runtime.into(),
355 fork.id.into(),
356 fork.namespace.into(),
357 fork.name.into(),
358 a.actor.id.as_str().into(),
359 a.actor.username.as_str().into(),
360 timestamp.as_str().into(),
361 timestamp.as_str().into(),
362 intent.id.as_str().into(),
363 ])?
364 .run()
365 .await?;
366 let Some(attempt) = self.attempt_by_id(&id).await? else {
367 return Ok(no_attempt());
368 };
369 self.publish(NewEvent {
370 kind: "attempt.started",
371 source: SOURCE,
372 repo_id: Some(attempt.repo_id.clone()),
373 actor: Some(a.actor.id),
374 data: AttemptEvent {
375 agent: Some(attempt.agent.clone()),
376 ..Self::attempt_event(&attempt)
377 },
378 })
379 .await?;
380 Ok(Outcome::Ok(attempt))
381 }
382
383 async fn get_attempt(&self, a: AttemptViewArgs) -> Result<Outcome<AttemptDetail>> {
384 let Some(attempt) = self.attempt_by_id(&a.attempt_id).await? else {
385 return Ok(no_attempt());
386 };
387 if let Outcome::Fail(_) = self.repo_by_id(&attempt.repo_id, &a.viewer).await? {
388 return Ok(no_attempt());
389 }
390 let Some(intent) = self.intent_by_id(&attempt.intent_id).await? else {
391 return Ok(no_attempt());
392 };
393 Ok(Outcome::Ok(AttemptDetail { attempt, intent }))
394 }
395
396 /// Submits or abandons an attempt on behalf of the one running it.
397 async fn close_attempt(
398 &self,
399 a: AttemptActionArgs,
400 status: AttemptStatus,
401 ) -> Result<Outcome<Attempt>> {
402 let mut attempt = match self.own_attempt(&a.actor, &a.attempt_id).await? {
403 Outcome::Ok(attempt) => attempt,
404 failed => return Ok(failed),
405 };
406 if !attempt.status.is_active() {
407 return Ok(Outcome::fail(
408 FailureCode::Conflict,
409 format!("This attempt is already {}.", attempt.status.as_str()),
410 ));
411 }
412 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
413 let now = rfc3339(now_ms());
414 self.db
415 .prepare(
416 "UPDATE attempts SET status = ?, summary = COALESCE(?, summary), updated_at = ?
417 WHERE id = ?",
418 )
419 .bind(&[
420 status.as_str().into(),
421 optional(&summary),
422 now.as_str().into(),
423 attempt.id.as_str().into(),
424 ])?
425 .run()
426 .await?;
427 let submitted = status == AttemptStatus::Submitted;
428 self.publish(NewEvent {
429 kind: if submitted {
430 "attempt.submitted"
431 } else {
432 "attempt.updated"
433 },
434 source: SOURCE,
435 repo_id: Some(attempt.repo_id.clone()),
436 actor: Some(a.actor.id),
437 data: AttemptEvent {
438 status: (!submitted).then_some(status.as_str()),
439 ..Self::attempt_event(&attempt)
440 },
441 })
442 .await?;
443 attempt.status = status;
444 attempt.summary = summary.or(attempt.summary);
445 attempt.updated_at = now;
446 Ok(Outcome::Ok(attempt))
447 }
448
449 async fn ship_attempt(&self, a: AttemptActionArgs) -> Result<Outcome<Attempt>> {
450 let Some(mut attempt) = self.attempt_by_id(&a.attempt_id).await? else {
451 return Ok(no_attempt());
452 };
453 // Whether the actor may see and write the repo is decided by repos.
454 let viewer = Some(a.actor.clone());
455 if let Outcome::Fail(_) = self.repo_by_id(&attempt.repo_id, &viewer).await? {
456 return Ok(no_attempt());
457 }
458 if !attempt.status.is_active() {
459 return Ok(Outcome::fail(
460 FailureCode::Conflict,
461 format!("This attempt is already {}.", attempt.status.as_str()),
462 ));
463 }
464 let Some(intent) = self.intent_by_id(&attempt.intent_id).await? else {
465 return Ok(no_intent());
466 };
467 if intent.status != IntentStatus::Open {
468 return Ok(Outcome::fail(
469 FailureCode::Conflict,
470 "This intent is already closed.",
471 ));
472 }
473
474 let landed: Outcome<Landed> = g1t_kit::call(
475 &self.repos,
476 "land",
477 &LandArgs {
478 fork_id: attempt.fork_repo_id.clone(),
479 actor: a.actor.clone(),
480 },
481 )
482 .await?;
483 let landed = match landed {
484 Outcome::Ok(landed) => landed,
485 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
486 };
487
488 let now = rfc3339(now_ms());
489 self.db
490 .batch(vec![
491 self.db
492 .prepare(
493 "UPDATE attempts
494 SET status = 'shipped', head_commit = ?, landed_base = ?, updated_at = ?
495 WHERE id = ?",
496 )
497 .bind(&[
498 landed.commit.as_str().into(),
499 optional(&landed.previous),
500 now.as_str().into(),
501 attempt.id.as_str().into(),
502 ])?,
503 self.db
504 .prepare("UPDATE intents SET status = 'shipped' WHERE id = ?")
505 .bind(&[intent.id.as_str().into()])?,
506 ])
507 .await?;
508 self.publish(NewEvent {
509 kind: "attempt.shipped",
510 source: SOURCE,
511 repo_id: Some(attempt.repo_id.clone()),
512 actor: Some(a.actor.id.clone()),
513 data: AttemptEvent {
514 commit: Some(landed.commit.clone()),
515 ..Self::attempt_event(&attempt)
516 },
517 })
518 .await?;
519 self.publish(NewEvent {
520 kind: "intent.closed",
521 source: SOURCE,
522 repo_id: Some(attempt.repo_id.clone()),
523 actor: Some(a.actor.id),
524 data: IntentClosed {
525 intent_id: intent.id,
526 repo_id: attempt.repo_id.clone(),
527 reason: "shipped",
528 },
529 })
530 .await?;
531
532 attempt.status = AttemptStatus::Shipped;
533 attempt.head_commit = Some(landed.commit);
534 attempt.landed_base = landed.previous;
535 attempt.updated_at = now;
536 Ok(Outcome::Ok(attempt))
537 }
538
539 async fn list_active_attempts(&self, a: ViewerArgs) -> Result<Vec<AttemptDetail>> {
540 let Some(viewer) = a.viewer else {
541 return Ok(Vec::new());
542 };
543 let rows = self
544 .db
545 .prepare(
546 "SELECT * FROM attempts
547 WHERE started_by_id = ? AND status IN ('working', 'submitted')
548 ORDER BY updated_at DESC LIMIT 50",
549 )
550 .bind(&[viewer.id.into()])?
551 .all()
552 .await?
553 .results::<AttemptRow>()?;
554 let mut active = Vec::with_capacity(rows.len());
555 for attempt in rows.into_iter().map(Attempt::from) {
556 if let Some(intent) = self.intent_by_id(&attempt.intent_id).await? {
557 active.push(AttemptDetail { attempt, intent });
558 }
559 }
560 Ok(active)
561 }
562
563 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
564 if a.entries.is_empty() {
565 return Ok(Outcome::Ok(Appended { count: 0 }));
566 }
567 if a.entries.len() > MAX_ENTRY_BATCH {
568 return Ok(Outcome::fail(
569 FailureCode::Invalid,
570 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
571 ));
572 }
573 let attempt = match self.own_attempt(&a.actor, &a.attempt_id).await? {
574 Outcome::Ok(attempt) => attempt,
575 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
576 };
577
578 let now = rfc3339(now_ms());
579 let count = a.entries.len() as u32;
580 let mut statements = Vec::with_capacity(a.entries.len() + 1);
581 for entry in a.entries {
582 let kind = serde_json::to_value(entry.kind)?;
583 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
584 // Each insert takes the next sequence number itself, so two
585 // writers appending at once cannot collide.
586 statements.push(
587 self.db
588 .prepare(
589 "INSERT INTO session_entries (attempt_id, seq, kind, text, tool, \"commit\", at)
590 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
591 FROM session_entries WHERE attempt_id = ?",
592 )
593 .bind(&[
594 attempt.id.as_str().into(),
595 kind.as_str().unwrap_or("note").into(),
596 text.into(),
597 optional(&entry.tool),
598 optional(&entry.commit.or_else(|| attempt.head_commit.clone())),
599 now.as_str().into(),
600 attempt.id.as_str().into(),
601 ])?,
602 );
603 }
604 statements.push(
605 self.db
606 .prepare("UPDATE attempts SET updated_at = ? WHERE id = ?")
607 .bind(&[now.as_str().into(), attempt.id.as_str().into()])?,
608 );
609 self.db.batch(statements).await?;
610 self.publish(NewEvent {
611 kind: "session.appended",
612 source: SOURCE,
613 repo_id: Some(attempt.repo_id.clone()),
614 actor: Some(a.actor.id),
615 data: SessionAppended {
616 attempt_id: attempt.id.clone(),
617 session_id: attempt.id,
618 count,
619 },
620 })
621 .await?;
622 Ok(Outcome::Ok(Appended { count }))
623 }
624
625 async fn read_session(&self, a: AttemptViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
626 let after_seq = a.after_seq;
627 let attempt_id = a.attempt_id.clone();
628 if let Outcome::Fail(failure) = self.get_attempt(a).await? {
629 return Ok(Outcome::Fail(failure));
630 }
631 let rows = self
632 .db
633 .prepare(
634 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
635 WHERE attempt_id = ? AND seq > ? ORDER BY seq LIMIT ?",
636 )
637 .bind(&[attempt_id.into(), after_seq.into(), SESSION_PAGE.into()])?
638 .all()
639 .await?
640 .results::<SessionRow>()?;
641 Ok(Outcome::Ok(
642 rows.into_iter().map(SessionEntry::from).collect(),
643 ))
644 }
645
646 /// A push to an attempt's fork moves that attempt's head.
647 async fn on_event(&self, event: &Delivered) -> Result<()> {
648 if event.kind != "git.push" {
649 return Ok(());
650 }
651 let (Some(repo_id), Some(after)) = (event.repo_id.as_deref(), event.data["after"].as_str())
652 else {
653 return Ok(());
654 };
655 self.db
656 .prepare("UPDATE attempts SET head_commit = ?, updated_at = ? WHERE fork_repo_id = ?")
657 .bind(&[after.into(), rfc3339(now_ms()).into(), repo_id.into()])?
658 .run()
659 .await?;
660 Ok(())
661 }
662}
663
664fn service(env: &Env) -> Result<Work> {
665 Ok(Work {
666 db: env.d1("DB")?,
667 repos: env.service("REPOS")?,
668 events: js::binding(env, "EVENTS")?,
669 })
670}
671
672#[event(fetch)]
673async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
674 let Some(method) = rpc_method(&request) else {
675 return Response::error("Not found", 404);
676 };
677 let body: serde_json::Value = request.json().await?;
678 let work = service(&env)?;
679
680 match method.as_str() {
681 "open_intent" => reply(&work.open_intent(args(body)?).await?),
682 "list_intents" => reply(&work.list_intents(args(body)?).await?),
683 "get_intent" => reply(&work.get_intent(args(body)?).await?),
684 "withdraw_intent" => reply(&work.withdraw_intent(args(body)?).await?),
685 "start_attempt" => reply(&work.start_attempt(args(body)?).await?),
686 "get_attempt" => reply(&work.get_attempt(args(body)?).await?),
687 "submit_attempt" => reply(
688 &work
689 .close_attempt(args(body)?, AttemptStatus::Submitted)
690 .await?,
691 ),
692 "abandon_attempt" => reply(
693 &work
694 .close_attempt(args(body)?, AttemptStatus::Abandoned)
695 .await?,
696 ),
697 "ship_attempt" => reply(&work.ship_attempt(args(body)?).await?),
698 "list_active_attempts" => reply(&work.list_active_attempts(args(body)?).await?),
699 "append_session" => reply(&work.append_session(args(body)?).await?),
700 "read_session" => reply(&work.read_session(args(body)?).await?),
701 _ => Response::error("Unknown method", 404),
702 }
703}
704
705/// Events from the bus, delivered on this service's own queue.
706#[event(queue)]
707async fn queue(batch: MessageBatch<Delivered>, env: Env, _ctx: Context) -> Result<()> {
708 let work = service(&env)?;
709 for message in batch.messages()? {
710 work.on_event(message.body()).await?;
711 message.ack();
712 }
713 Ok(())
714}