mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(db): resolve workflow schedule tenant from DB
Co-authored-by: npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf <95cae996907d7cab9f5dbf43c0f53edeac6ab0b032a6feae4abfd784e467b3f5@sprout-oss.stage.blox.sqprod.co> Signed-off-by: npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf <95cae996907d7cab9f5dbf43c0f53edeac6ab0b032a6feae4abfd784e467b3f5@sprout-oss.stage.blox.sqprod.co>
This commit is contained in:
parent
48c53a5d7e
commit
1fa3d837f1
@@ -1122,8 +1122,8 @@ pub async fn release_due_reminder(
|
|||||||
event_id: &[u8],
|
event_id: &[u8],
|
||||||
event_created_at: DateTime<Utc>,
|
event_created_at: DateTime<Utc>,
|
||||||
delivery_stamp: i64,
|
delivery_stamp: i64,
|
||||||
) -> Result<()> {
|
) -> Result<bool> {
|
||||||
sqlx::query(
|
let result = sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
UPDATE events
|
UPDATE events
|
||||||
SET delivered_at = NULL
|
SET delivered_at = NULL
|
||||||
@@ -1138,7 +1138,7 @@ pub async fn release_due_reminder(
|
|||||||
.execute(pool)
|
.execute(pool)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
Ok(())
|
Ok(result.rows_affected() == 1)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
|
|||||||
@@ -664,7 +664,7 @@ impl Db {
|
|||||||
event_id: &[u8],
|
event_id: &[u8],
|
||||||
event_created_at: chrono::DateTime<chrono::Utc>,
|
event_created_at: chrono::DateTime<chrono::Utc>,
|
||||||
delivery_stamp: i64,
|
delivery_stamp: i64,
|
||||||
) -> Result<()> {
|
) -> Result<bool> {
|
||||||
event::release_due_reminder(&self.pool, event_id, event_created_at, delivery_stamp).await
|
event::release_due_reminder(&self.pool, event_id, event_created_at, delivery_stamp).await
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1203,6 +1203,7 @@ impl Db {
|
|||||||
/// Create a new workflow.
|
/// Create a new workflow.
|
||||||
pub async fn create_workflow(
|
pub async fn create_workflow(
|
||||||
&self,
|
&self,
|
||||||
|
community_id: CommunityId,
|
||||||
channel_id: Option<Uuid>,
|
channel_id: Option<Uuid>,
|
||||||
owner_pubkey: &[u8],
|
owner_pubkey: &[u8],
|
||||||
name: &str,
|
name: &str,
|
||||||
@@ -1211,6 +1212,7 @@ impl Db {
|
|||||||
) -> Result<Uuid> {
|
) -> Result<Uuid> {
|
||||||
workflow::create_workflow(
|
workflow::create_workflow(
|
||||||
&self.pool,
|
&self.pool,
|
||||||
|
community_id,
|
||||||
channel_id,
|
channel_id,
|
||||||
owner_pubkey,
|
owner_pubkey,
|
||||||
name,
|
name,
|
||||||
@@ -1250,43 +1252,34 @@ impl Db {
|
|||||||
|
|
||||||
/// Claim a scheduled workflow fire for an authoritative schedule instant.
|
/// Claim a scheduled workflow fire for an authoritative schedule instant.
|
||||||
///
|
///
|
||||||
/// Returns `Some` only for the first pod to claim `(community_id,
|
/// Returns `Some` only for the first pod to claim `(workflow_id,
|
||||||
/// workflow_id, scheduled_for)`; all other pods must skip creating a run.
|
/// scheduled_for)`; all other pods must skip creating a run. The claim SQL
|
||||||
|
/// resolves `community_id` from the workflow row; callers never supply it.
|
||||||
pub async fn claim_scheduled_workflow_fire(
|
pub async fn claim_scheduled_workflow_fire(
|
||||||
&self,
|
&self,
|
||||||
community_id: CommunityId,
|
|
||||||
workflow_id: Uuid,
|
workflow_id: Uuid,
|
||||||
scheduled_for: chrono::DateTime<chrono::Utc>,
|
scheduled_for: chrono::DateTime<chrono::Utc>,
|
||||||
) -> Result<Option<workflow::ScheduledWorkflowFireClaim>> {
|
) -> Result<Option<workflow::ScheduledWorkflowFireClaim>> {
|
||||||
workflow::claim_scheduled_workflow_fire(
|
workflow::claim_scheduled_workflow_fire(&self.pool, workflow_id, scheduled_for).await
|
||||||
&self.pool,
|
|
||||||
community_id,
|
|
||||||
workflow_id,
|
|
||||||
scheduled_for,
|
|
||||||
)
|
|
||||||
.await
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Fetch the latest claimed schedule instant for interval trigger anchoring.
|
/// Fetch the latest claimed schedule instant for interval trigger anchoring.
|
||||||
pub async fn latest_scheduled_workflow_fire(
|
pub async fn latest_scheduled_workflow_fire(
|
||||||
&self,
|
&self,
|
||||||
community_id: CommunityId,
|
|
||||||
workflow_id: Uuid,
|
workflow_id: Uuid,
|
||||||
) -> Result<Option<chrono::DateTime<chrono::Utc>>> {
|
) -> Result<Option<chrono::DateTime<chrono::Utc>>> {
|
||||||
workflow::latest_scheduled_workflow_fire(&self.pool, community_id, workflow_id).await
|
workflow::latest_scheduled_workflow_fire(&self.pool, workflow_id).await
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Attach the workflow run id created from a won scheduled-fire claim.
|
/// Attach the workflow run id created from a won scheduled-fire claim.
|
||||||
pub async fn attach_scheduled_workflow_run(
|
pub async fn attach_scheduled_workflow_run(
|
||||||
&self,
|
&self,
|
||||||
community_id: CommunityId,
|
|
||||||
workflow_id: Uuid,
|
workflow_id: Uuid,
|
||||||
scheduled_for: chrono::DateTime<chrono::Utc>,
|
scheduled_for: chrono::DateTime<chrono::Utc>,
|
||||||
workflow_run_id: Uuid,
|
workflow_run_id: Uuid,
|
||||||
) -> Result<bool> {
|
) -> Result<bool> {
|
||||||
workflow::attach_scheduled_workflow_run(
|
workflow::attach_scheduled_workflow_run(
|
||||||
&self.pool,
|
&self.pool,
|
||||||
community_id,
|
|
||||||
workflow_id,
|
workflow_id,
|
||||||
scheduled_for,
|
scheduled_for,
|
||||||
workflow_run_id,
|
workflow_run_id,
|
||||||
|
|||||||
@@ -145,6 +145,20 @@ mod tests {
|
|||||||
.contains("CREATE TABLE scheduled_workflow_fires"),
|
.contains("CREATE TABLE scheduled_workflow_fires"),
|
||||||
"initial schema migration should include workflow cron claim table"
|
"initial schema migration should include workflow cron claim table"
|
||||||
);
|
);
|
||||||
|
assert!(
|
||||||
|
migrations[0]
|
||||||
|
.sql
|
||||||
|
.as_str()
|
||||||
|
.contains("community_id UUID NOT NULL REFERENCES communities(id)"),
|
||||||
|
"workflows should carry row-owned community_id"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
migrations[0]
|
||||||
|
.sql
|
||||||
|
.as_str()
|
||||||
|
.contains("trg_workflows_community_id_immutable"),
|
||||||
|
"workflow community_id should be immutable after insert"
|
||||||
|
);
|
||||||
assert!(
|
assert!(
|
||||||
migrations[0]
|
migrations[0]
|
||||||
.sql
|
.sql
|
||||||
@@ -156,8 +170,15 @@ mod tests {
|
|||||||
migrations[0]
|
migrations[0]
|
||||||
.sql
|
.sql
|
||||||
.as_str()
|
.as_str()
|
||||||
.contains("PRIMARY KEY (community_id, workflow_id, scheduled_for)"),
|
.contains("PRIMARY KEY (workflow_id, scheduled_for)"),
|
||||||
"workflow cron claim uniqueness must include the community label"
|
"workflow cron claim uniqueness should be the globally unique workflow plus schedule instant"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
migrations[0]
|
||||||
|
.sql
|
||||||
|
.as_str()
|
||||||
|
.contains("FOREIGN KEY (community_id, workflow_id)"),
|
||||||
|
"workflow cron claim table should tie community_id to the workflow row"
|
||||||
);
|
);
|
||||||
assert!(
|
assert!(
|
||||||
migrations[0].sql.as_str().contains("CREATE TABLE channels"),
|
migrations[0].sql.as_str().contains("CREATE TABLE channels"),
|
||||||
|
|||||||
@@ -165,6 +165,8 @@ impl FromStr for ApprovalStatus {
|
|||||||
pub struct WorkflowRecord {
|
pub struct WorkflowRecord {
|
||||||
/// Unique workflow identifier.
|
/// Unique workflow identifier.
|
||||||
pub id: Uuid,
|
pub id: Uuid,
|
||||||
|
/// Server-resolved community that owns this workflow.
|
||||||
|
pub community_id: CommunityId,
|
||||||
/// Human-readable workflow name.
|
/// Human-readable workflow name.
|
||||||
pub name: String,
|
pub name: String,
|
||||||
/// Compressed public key bytes of the workflow owner.
|
/// Compressed public key bytes of the workflow owner.
|
||||||
@@ -215,9 +217,9 @@ pub struct WorkflowRunRecord {
|
|||||||
|
|
||||||
/// A winning scheduled workflow fire claim.
|
/// A winning scheduled workflow fire claim.
|
||||||
///
|
///
|
||||||
/// The primary identity is `(community_id, workflow_id, scheduled_for)`. The
|
/// The primary identity is `(workflow_id, scheduled_for)`. `community_id` is
|
||||||
/// database also returns `claimed_at` so the workflow scheduler can log and audit
|
/// resolved from the workflow row inside the claim SQL and returned for scoped
|
||||||
/// the exact claim row it won without relying on a per-pod clock.
|
/// audit/logging; callers never supply it as a claim.
|
||||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
pub struct ScheduledWorkflowFireClaim {
|
pub struct ScheduledWorkflowFireClaim {
|
||||||
/// Community that owns this scheduled fire.
|
/// Community that owns this scheduled fire.
|
||||||
@@ -263,6 +265,7 @@ pub struct ApprovalRecord {
|
|||||||
/// New workflows start as `active` and `enabled = TRUE`.
|
/// New workflows start as `active` and `enabled = TRUE`.
|
||||||
pub async fn create_workflow(
|
pub async fn create_workflow(
|
||||||
pool: &PgPool,
|
pool: &PgPool,
|
||||||
|
community_id: CommunityId,
|
||||||
channel_id: Option<Uuid>,
|
channel_id: Option<Uuid>,
|
||||||
owner_pubkey: &[u8],
|
owner_pubkey: &[u8],
|
||||||
name: &str,
|
name: &str,
|
||||||
@@ -274,11 +277,12 @@ pub async fn create_workflow(
|
|||||||
sqlx::query(
|
sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
INSERT INTO workflows
|
INSERT INTO workflows
|
||||||
(id, name, owner_pubkey, channel_id, definition, definition_hash, status, enabled)
|
(id, community_id, name, owner_pubkey, channel_id, definition, definition_hash, status, enabled)
|
||||||
VALUES ($1, $2, $3, $4, $5::jsonb, $6, 'active', TRUE)
|
VALUES ($1, $2, $3, $4, $5, $6::jsonb, $7, 'active', TRUE)
|
||||||
"#,
|
"#,
|
||||||
)
|
)
|
||||||
.bind(id)
|
.bind(id)
|
||||||
|
.bind(community_id.as_uuid())
|
||||||
.bind(name)
|
.bind(name)
|
||||||
.bind(owner_pubkey)
|
.bind(owner_pubkey)
|
||||||
.bind(channel_id)
|
.bind(channel_id)
|
||||||
@@ -294,7 +298,7 @@ pub async fn create_workflow(
|
|||||||
pub async fn get_workflow(pool: &PgPool, id: Uuid) -> Result<WorkflowRecord> {
|
pub async fn get_workflow(pool: &PgPool, id: Uuid) -> Result<WorkflowRecord> {
|
||||||
let row = sqlx::query(
|
let row = sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash,
|
SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
|
||||||
status::text AS status, enabled, created_at, updated_at
|
status::text AS status, enabled, created_at, updated_at
|
||||||
FROM workflows
|
FROM workflows
|
||||||
WHERE id = $1
|
WHERE id = $1
|
||||||
@@ -323,7 +327,7 @@ pub async fn list_channel_workflows(
|
|||||||
|
|
||||||
let rows = sqlx::query(
|
let rows = sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash,
|
SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
|
||||||
status::text AS status, enabled, created_at, updated_at
|
status::text AS status, enabled, created_at, updated_at
|
||||||
FROM workflows
|
FROM workflows
|
||||||
WHERE channel_id = $1
|
WHERE channel_id = $1
|
||||||
@@ -352,7 +356,7 @@ pub async fn list_enabled_channel_workflows(
|
|||||||
) -> Result<Vec<WorkflowRecord>> {
|
) -> Result<Vec<WorkflowRecord>> {
|
||||||
let rows = sqlx::query(
|
let rows = sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash,
|
SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
|
||||||
status::text AS status, enabled, created_at, updated_at
|
status::text AS status, enabled, created_at, updated_at
|
||||||
FROM workflows
|
FROM workflows
|
||||||
WHERE channel_id = $1
|
WHERE channel_id = $1
|
||||||
@@ -378,7 +382,7 @@ pub async fn list_enabled_channel_workflows(
|
|||||||
pub async fn list_all_enabled_workflows(pool: &PgPool) -> Result<Vec<WorkflowRecord>> {
|
pub async fn list_all_enabled_workflows(pool: &PgPool) -> Result<Vec<WorkflowRecord>> {
|
||||||
let rows = sqlx::query(
|
let rows = sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash,
|
SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
|
||||||
status::text AS status, enabled, created_at, updated_at
|
status::text AS status, enabled, created_at, updated_at
|
||||||
FROM workflows
|
FROM workflows
|
||||||
WHERE status = 'active'
|
WHERE status = 'active'
|
||||||
@@ -397,27 +401,27 @@ pub async fn list_all_enabled_workflows(pool: &PgPool) -> Result<Vec<WorkflowRec
|
|||||||
|
|
||||||
/// Claim a scheduled workflow fire for an authoritative schedule instant.
|
/// Claim a scheduled workflow fire for an authoritative schedule instant.
|
||||||
///
|
///
|
||||||
/// Returns `Some` only for the first pod that claims `(community_id,
|
/// Returns `Some` only for the first pod that claims `(workflow_id,
|
||||||
/// workflow_id, scheduled_for)`. All other pods receive `None` and must skip
|
/// scheduled_for)`. All other pods receive `None` and must skip creating a
|
||||||
/// creating a workflow run. The `scheduled_for` value must come from an external
|
/// workflow run. The `scheduled_for` value must come from an external
|
||||||
/// schedule anchor (cron expression) or DB-authoritative interval anchor; a
|
/// schedule anchor (cron expression) or DB-authoritative interval anchor; a
|
||||||
/// per-pod in-memory timestamp is not safe because different pods can compute
|
/// per-pod in-memory timestamp is not safe because different pods can compute
|
||||||
/// different claim keys.
|
/// different claim keys.
|
||||||
pub async fn claim_scheduled_workflow_fire(
|
pub async fn claim_scheduled_workflow_fire(
|
||||||
pool: &PgPool,
|
pool: &PgPool,
|
||||||
community_id: CommunityId,
|
|
||||||
workflow_id: Uuid,
|
workflow_id: Uuid,
|
||||||
scheduled_for: DateTime<Utc>,
|
scheduled_for: DateTime<Utc>,
|
||||||
) -> Result<Option<ScheduledWorkflowFireClaim>> {
|
) -> Result<Option<ScheduledWorkflowFireClaim>> {
|
||||||
let row = sqlx::query(
|
let row = sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
INSERT INTO scheduled_workflow_fires (community_id, workflow_id, scheduled_for)
|
INSERT INTO scheduled_workflow_fires (community_id, workflow_id, scheduled_for)
|
||||||
VALUES ($1, $2, $3)
|
SELECT w.community_id, w.id, $2
|
||||||
ON CONFLICT (community_id, workflow_id, scheduled_for) DO NOTHING
|
FROM workflows w
|
||||||
|
WHERE w.id = $1
|
||||||
|
ON CONFLICT (workflow_id, scheduled_for) DO NOTHING
|
||||||
RETURNING community_id, workflow_id, scheduled_for, claimed_at
|
RETURNING community_id, workflow_id, scheduled_for, claimed_at
|
||||||
"#,
|
"#,
|
||||||
)
|
)
|
||||||
.bind(community_id.as_uuid())
|
|
||||||
.bind(workflow_id)
|
.bind(workflow_id)
|
||||||
.bind(scheduled_for)
|
.bind(scheduled_for)
|
||||||
.fetch_optional(pool)
|
.fetch_optional(pool)
|
||||||
@@ -444,18 +448,15 @@ pub async fn claim_scheduled_workflow_fire(
|
|||||||
/// row is the source of truth for schedule deduplication.
|
/// row is the source of truth for schedule deduplication.
|
||||||
pub async fn latest_scheduled_workflow_fire(
|
pub async fn latest_scheduled_workflow_fire(
|
||||||
pool: &PgPool,
|
pool: &PgPool,
|
||||||
community_id: CommunityId,
|
|
||||||
workflow_id: Uuid,
|
workflow_id: Uuid,
|
||||||
) -> Result<Option<DateTime<Utc>>> {
|
) -> Result<Option<DateTime<Utc>>> {
|
||||||
let row = sqlx::query(
|
let row = sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
SELECT MAX(scheduled_for) AS scheduled_for
|
SELECT MAX(scheduled_for) AS scheduled_for
|
||||||
FROM scheduled_workflow_fires
|
FROM scheduled_workflow_fires
|
||||||
WHERE community_id = $1
|
WHERE workflow_id = $1
|
||||||
AND workflow_id = $2
|
|
||||||
"#,
|
"#,
|
||||||
)
|
)
|
||||||
.bind(community_id.as_uuid())
|
|
||||||
.bind(workflow_id)
|
.bind(workflow_id)
|
||||||
.fetch_one(pool)
|
.fetch_one(pool)
|
||||||
.await?;
|
.await?;
|
||||||
@@ -471,7 +472,6 @@ pub async fn latest_scheduled_workflow_fire(
|
|||||||
/// intentional: the schedule instant was claimed and must not duplicate later.
|
/// intentional: the schedule instant was claimed and must not duplicate later.
|
||||||
pub async fn attach_scheduled_workflow_run(
|
pub async fn attach_scheduled_workflow_run(
|
||||||
pool: &PgPool,
|
pool: &PgPool,
|
||||||
community_id: CommunityId,
|
|
||||||
workflow_id: Uuid,
|
workflow_id: Uuid,
|
||||||
scheduled_for: DateTime<Utc>,
|
scheduled_for: DateTime<Utc>,
|
||||||
workflow_run_id: Uuid,
|
workflow_run_id: Uuid,
|
||||||
@@ -479,14 +479,12 @@ pub async fn attach_scheduled_workflow_run(
|
|||||||
let result = sqlx::query(
|
let result = sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
UPDATE scheduled_workflow_fires
|
UPDATE scheduled_workflow_fires
|
||||||
SET workflow_run_id = $4
|
SET workflow_run_id = $3
|
||||||
WHERE community_id = $1
|
WHERE workflow_id = $1
|
||||||
AND workflow_id = $2
|
AND scheduled_for = $2
|
||||||
AND scheduled_for = $3
|
|
||||||
AND workflow_run_id IS NULL
|
AND workflow_run_id IS NULL
|
||||||
"#,
|
"#,
|
||||||
)
|
)
|
||||||
.bind(community_id.as_uuid())
|
|
||||||
.bind(workflow_id)
|
.bind(workflow_id)
|
||||||
.bind(scheduled_for)
|
.bind(scheduled_for)
|
||||||
.bind(workflow_run_id)
|
.bind(workflow_run_id)
|
||||||
@@ -909,8 +907,11 @@ fn row_to_workflow_record(row: sqlx::postgres::PgRow) -> Result<WorkflowRecord>
|
|||||||
|
|
||||||
let enabled: bool = row.try_get("enabled")?;
|
let enabled: bool = row.try_get("enabled")?;
|
||||||
|
|
||||||
|
let community_id: Uuid = row.try_get("community_id")?;
|
||||||
|
|
||||||
Ok(WorkflowRecord {
|
Ok(WorkflowRecord {
|
||||||
id,
|
id,
|
||||||
|
community_id: CommunityId::from_uuid(community_id),
|
||||||
name: row.try_get("name")?,
|
name: row.try_get("name")?,
|
||||||
owner_pubkey: row.try_get("owner_pubkey")?,
|
owner_pubkey: row.try_get("owner_pubkey")?,
|
||||||
channel_id,
|
channel_id,
|
||||||
@@ -975,7 +976,7 @@ pub async fn find_by_owner_and_name(
|
|||||||
) -> Result<Option<WorkflowRecord>> {
|
) -> Result<Option<WorkflowRecord>> {
|
||||||
let row = sqlx::query(
|
let row = sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
SELECT id, name, owner_pubkey, channel_id, definition, definition_hash,
|
SELECT id, community_id, name, owner_pubkey, channel_id, definition, definition_hash,
|
||||||
status::text AS status, enabled, created_at, updated_at
|
status::text AS status, enabled, created_at, updated_at
|
||||||
FROM workflows
|
FROM workflows
|
||||||
WHERE owner_pubkey = $1 AND name = $2
|
WHERE owner_pubkey = $1 AND name = $2
|
||||||
@@ -1099,8 +1100,11 @@ mod tests {
|
|||||||
"steps": [{ "id": "s1", "action": "send_message", "text": "hi" }]
|
"steps": [{ "id": "s1", "action": "send_message", "text": "hi" }]
|
||||||
});
|
});
|
||||||
|
|
||||||
|
let community_id = CommunityId::from_uuid(Uuid::new_v4());
|
||||||
|
|
||||||
let record = WorkflowRecord {
|
let record = WorkflowRecord {
|
||||||
id,
|
id,
|
||||||
|
community_id,
|
||||||
name: "My Workflow".to_owned(),
|
name: "My Workflow".to_owned(),
|
||||||
owner_pubkey: vec![0xab; 32],
|
owner_pubkey: vec![0xab; 32],
|
||||||
channel_id: Some(channel_id),
|
channel_id: Some(channel_id),
|
||||||
@@ -1113,6 +1117,7 @@ mod tests {
|
|||||||
};
|
};
|
||||||
|
|
||||||
assert_eq!(record.id, id);
|
assert_eq!(record.id, id);
|
||||||
|
assert_eq!(record.community_id, community_id);
|
||||||
assert_eq!(record.name, "My Workflow");
|
assert_eq!(record.name, "My Workflow");
|
||||||
assert_eq!(record.owner_pubkey, vec![0xab; 32]);
|
assert_eq!(record.owner_pubkey, vec![0xab; 32]);
|
||||||
assert_eq!(record.channel_id, Some(channel_id));
|
assert_eq!(record.channel_id, Some(channel_id));
|
||||||
@@ -1129,6 +1134,7 @@ mod tests {
|
|||||||
|
|
||||||
let record = WorkflowRecord {
|
let record = WorkflowRecord {
|
||||||
id,
|
id,
|
||||||
|
community_id: CommunityId::from_uuid(Uuid::new_v4()),
|
||||||
name: "Global Workflow".to_owned(),
|
name: "Global Workflow".to_owned(),
|
||||||
owner_pubkey: vec![0x00; 32],
|
owner_pubkey: vec![0x00; 32],
|
||||||
channel_id: None,
|
channel_id: None,
|
||||||
@@ -1150,6 +1156,7 @@ mod tests {
|
|||||||
|
|
||||||
let record = WorkflowRecord {
|
let record = WorkflowRecord {
|
||||||
id,
|
id,
|
||||||
|
community_id: CommunityId::from_uuid(Uuid::new_v4()),
|
||||||
name: "Original".to_owned(),
|
name: "Original".to_owned(),
|
||||||
owner_pubkey: vec![0x01; 32],
|
owner_pubkey: vec![0x01; 32],
|
||||||
channel_id: None,
|
channel_id: None,
|
||||||
@@ -1178,6 +1185,7 @@ mod tests {
|
|||||||
] {
|
] {
|
||||||
let record = WorkflowRecord {
|
let record = WorkflowRecord {
|
||||||
id: Uuid::new_v4(),
|
id: Uuid::new_v4(),
|
||||||
|
community_id: CommunityId::from_uuid(Uuid::new_v4()),
|
||||||
name: "Test".to_owned(),
|
name: "Test".to_owned(),
|
||||||
owner_pubkey: vec![],
|
owner_pubkey: vec![],
|
||||||
channel_id: None,
|
channel_id: None,
|
||||||
@@ -1197,6 +1205,7 @@ mod tests {
|
|||||||
let now = Utc::now();
|
let now = Utc::now();
|
||||||
let record = WorkflowRecord {
|
let record = WorkflowRecord {
|
||||||
id: Uuid::new_v4(),
|
id: Uuid::new_v4(),
|
||||||
|
community_id: CommunityId::from_uuid(Uuid::new_v4()),
|
||||||
name: "Paused".to_owned(),
|
name: "Paused".to_owned(),
|
||||||
owner_pubkey: vec![],
|
owner_pubkey: vec![],
|
||||||
channel_id: None,
|
channel_id: None,
|
||||||
|
|||||||
@@ -612,10 +612,20 @@ async fn handle_workflow_def(
|
|||||||
PersistResult::Inserted(tx) => tx,
|
PersistResult::Inserted(tx) => tx,
|
||||||
};
|
};
|
||||||
|
|
||||||
// 4. Execute: create_workflow
|
// 4. Execute: create_workflow. The workflow's community is resolved from
|
||||||
|
// the server-owned channel row, not from the client-supplied event. The DB
|
||||||
|
// also enforces `(community_id, channel_id)` as a composite FK.
|
||||||
|
let community_id = state
|
||||||
|
.db
|
||||||
|
.community_of_channel(channel_id)
|
||||||
|
.await
|
||||||
|
.map_err(|e| IngestError::Internal(format!("error: db channel community lookup: {e}")))?
|
||||||
|
.ok_or_else(|| IngestError::Rejected("invalid: workflow channel not found".into()))?;
|
||||||
|
|
||||||
let workflow_id = state
|
let workflow_id = state
|
||||||
.db
|
.db
|
||||||
.create_workflow(
|
.create_workflow(
|
||||||
|
community_id,
|
||||||
Some(channel_id),
|
Some(channel_id),
|
||||||
&self_bytes,
|
&self_bytes,
|
||||||
&workflow_name,
|
&workflow_name,
|
||||||
|
|||||||
@@ -58,6 +58,7 @@ CREATE TABLE channels (
|
|||||||
CONSTRAINT chk_channels_id_not_nil CHECK (id <> '00000000-0000-0000-0000-000000000000'::uuid)
|
CONSTRAINT chk_channels_id_not_nil CHECK (id <> '00000000-0000-0000-0000-000000000000'::uuid)
|
||||||
);
|
);
|
||||||
|
|
||||||
|
CREATE UNIQUE INDEX idx_channels_community_id_id ON channels (community_id, id);
|
||||||
CREATE INDEX idx_channels_type ON channels (channel_type);
|
CREATE INDEX idx_channels_type ON channels (channel_type);
|
||||||
CREATE INDEX idx_channels_visibility ON channels (visibility);
|
CREATE INDEX idx_channels_visibility ON channels (visibility);
|
||||||
CREATE INDEX idx_channels_created_by ON channels (created_by);
|
CREATE INDEX idx_channels_created_by ON channels (created_by);
|
||||||
@@ -210,18 +211,37 @@ CREATE TABLE delivery_log_p_future PARTITION OF delivery_log
|
|||||||
|
|
||||||
CREATE TABLE workflows (
|
CREATE TABLE workflows (
|
||||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||||
|
community_id UUID NOT NULL REFERENCES communities(id),
|
||||||
name VARCHAR(255) NOT NULL,
|
name VARCHAR(255) NOT NULL,
|
||||||
owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey),
|
owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey),
|
||||||
channel_id UUID REFERENCES channels(id),
|
channel_id UUID,
|
||||||
definition JSONB NOT NULL,
|
definition JSONB NOT NULL,
|
||||||
definition_hash BYTEA NOT NULL,
|
definition_hash BYTEA NOT NULL,
|
||||||
status workflow_status NOT NULL DEFAULT 'active',
|
status workflow_status NOT NULL DEFAULT 'active',
|
||||||
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
||||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||||
|
FOREIGN KEY (community_id, channel_id)
|
||||||
|
REFERENCES channels(community_id, id)
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE INDEX idx_workflows_channel_active ON workflows (channel_id, status, enabled);
|
CREATE UNIQUE INDEX idx_workflows_community_id_id ON workflows (community_id, id);
|
||||||
|
CREATE INDEX idx_workflows_channel_active ON workflows (community_id, channel_id, status, enabled);
|
||||||
|
|
||||||
|
CREATE OR REPLACE FUNCTION prevent_workflows_community_id_update()
|
||||||
|
RETURNS trigger AS $$
|
||||||
|
BEGIN
|
||||||
|
IF NEW.community_id IS DISTINCT FROM OLD.community_id THEN
|
||||||
|
RAISE EXCEPTION 'workflows.community_id is immutable';
|
||||||
|
END IF;
|
||||||
|
RETURN NEW;
|
||||||
|
END;
|
||||||
|
$$ LANGUAGE plpgsql;
|
||||||
|
|
||||||
|
CREATE TRIGGER trg_workflows_community_id_immutable
|
||||||
|
BEFORE UPDATE OF community_id ON workflows
|
||||||
|
FOR EACH ROW
|
||||||
|
EXECUTE FUNCTION prevent_workflows_community_id_update();
|
||||||
|
|
||||||
-- ── Workflow runs ─────────────────────────────────────────────────────────────
|
-- ── Workflow runs ─────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
@@ -246,12 +266,14 @@ CREATE INDEX idx_workflow_runs_status ON workflow_runs (status);
|
|||||||
-- Each schedule tick must win this insert before creating a workflow run; this
|
-- Each schedule tick must win this insert before creating a workflow run; this
|
||||||
-- replaces per-pod in-memory last-fired state as the deduplication boundary.
|
-- replaces per-pod in-memory last-fired state as the deduplication boundary.
|
||||||
CREATE TABLE scheduled_workflow_fires (
|
CREATE TABLE scheduled_workflow_fires (
|
||||||
community_id UUID NOT NULL REFERENCES communities(id),
|
community_id UUID NOT NULL,
|
||||||
workflow_id UUID NOT NULL REFERENCES workflows(id) ON DELETE CASCADE,
|
workflow_id UUID NOT NULL,
|
||||||
scheduled_for TIMESTAMPTZ NOT NULL,
|
scheduled_for TIMESTAMPTZ NOT NULL,
|
||||||
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||||
workflow_run_id UUID REFERENCES workflow_runs(id) ON DELETE SET NULL,
|
workflow_run_id UUID REFERENCES workflow_runs(id) ON DELETE SET NULL,
|
||||||
PRIMARY KEY (community_id, workflow_id, scheduled_for)
|
PRIMARY KEY (workflow_id, scheduled_for),
|
||||||
|
FOREIGN KEY (community_id, workflow_id)
|
||||||
|
REFERENCES workflows(community_id, id) ON DELETE CASCADE
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE INDEX idx_scheduled_workflow_fires_scheduled_for
|
CREATE INDEX idx_scheduled_workflow_fires_scheduled_for
|
||||||
|
|||||||
+28
-6
@@ -59,6 +59,7 @@ CREATE TABLE channels (
|
|||||||
CONSTRAINT chk_channels_id_not_nil CHECK (id <> '00000000-0000-0000-0000-000000000000'::uuid)
|
CONSTRAINT chk_channels_id_not_nil CHECK (id <> '00000000-0000-0000-0000-000000000000'::uuid)
|
||||||
);
|
);
|
||||||
|
|
||||||
|
CREATE UNIQUE INDEX idx_channels_community_id_id ON channels (community_id, id);
|
||||||
CREATE INDEX idx_channels_type ON channels (channel_type);
|
CREATE INDEX idx_channels_type ON channels (channel_type);
|
||||||
CREATE INDEX idx_channels_visibility ON channels (visibility);
|
CREATE INDEX idx_channels_visibility ON channels (visibility);
|
||||||
CREATE INDEX idx_channels_created_by ON channels (created_by);
|
CREATE INDEX idx_channels_created_by ON channels (created_by);
|
||||||
@@ -215,18 +216,37 @@ CREATE TABLE delivery_log_p_future PARTITION OF delivery_log
|
|||||||
|
|
||||||
CREATE TABLE workflows (
|
CREATE TABLE workflows (
|
||||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||||
|
community_id UUID NOT NULL REFERENCES communities(id),
|
||||||
name VARCHAR(255) NOT NULL,
|
name VARCHAR(255) NOT NULL,
|
||||||
owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey),
|
owner_pubkey BYTEA NOT NULL REFERENCES users(pubkey),
|
||||||
channel_id UUID REFERENCES channels(id),
|
channel_id UUID,
|
||||||
definition JSONB NOT NULL,
|
definition JSONB NOT NULL,
|
||||||
definition_hash BYTEA NOT NULL,
|
definition_hash BYTEA NOT NULL,
|
||||||
status workflow_status NOT NULL DEFAULT 'active',
|
status workflow_status NOT NULL DEFAULT 'active',
|
||||||
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
||||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||||
|
FOREIGN KEY (community_id, channel_id)
|
||||||
|
REFERENCES channels(community_id, id)
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE INDEX idx_workflows_channel_active ON workflows (channel_id, status, enabled);
|
CREATE UNIQUE INDEX idx_workflows_community_id_id ON workflows (community_id, id);
|
||||||
|
CREATE INDEX idx_workflows_channel_active ON workflows (community_id, channel_id, status, enabled);
|
||||||
|
|
||||||
|
CREATE OR REPLACE FUNCTION prevent_workflows_community_id_update()
|
||||||
|
RETURNS trigger AS $$
|
||||||
|
BEGIN
|
||||||
|
IF NEW.community_id IS DISTINCT FROM OLD.community_id THEN
|
||||||
|
RAISE EXCEPTION 'workflows.community_id is immutable';
|
||||||
|
END IF;
|
||||||
|
RETURN NEW;
|
||||||
|
END;
|
||||||
|
$$ LANGUAGE plpgsql;
|
||||||
|
|
||||||
|
CREATE TRIGGER trg_workflows_community_id_immutable
|
||||||
|
BEFORE UPDATE OF community_id ON workflows
|
||||||
|
FOR EACH ROW
|
||||||
|
EXECUTE FUNCTION prevent_workflows_community_id_update();
|
||||||
|
|
||||||
-- ── Workflow runs ─────────────────────────────────────────────────────────────
|
-- ── Workflow runs ─────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
@@ -251,12 +271,14 @@ CREATE INDEX idx_workflow_runs_status ON workflow_runs (status);
|
|||||||
-- Each schedule tick must win this insert before creating a workflow run; this
|
-- Each schedule tick must win this insert before creating a workflow run; this
|
||||||
-- replaces per-pod in-memory last-fired state as the deduplication boundary.
|
-- replaces per-pod in-memory last-fired state as the deduplication boundary.
|
||||||
CREATE TABLE scheduled_workflow_fires (
|
CREATE TABLE scheduled_workflow_fires (
|
||||||
community_id UUID NOT NULL REFERENCES communities(id),
|
community_id UUID NOT NULL,
|
||||||
workflow_id UUID NOT NULL REFERENCES workflows(id) ON DELETE CASCADE,
|
workflow_id UUID NOT NULL,
|
||||||
scheduled_for TIMESTAMPTZ NOT NULL,
|
scheduled_for TIMESTAMPTZ NOT NULL,
|
||||||
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||||
workflow_run_id UUID REFERENCES workflow_runs(id) ON DELETE SET NULL,
|
workflow_run_id UUID REFERENCES workflow_runs(id) ON DELETE SET NULL,
|
||||||
PRIMARY KEY (community_id, workflow_id, scheduled_for)
|
PRIMARY KEY (workflow_id, scheduled_for),
|
||||||
|
FOREIGN KEY (community_id, workflow_id)
|
||||||
|
REFERENCES workflows(community_id, id) ON DELETE CASCADE
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE INDEX idx_scheduled_workflow_fires_scheduled_for
|
CREATE INDEX idx_scheduled_workflow_fires_scheduled_for
|
||||||
|
|||||||
Reference in New Issue
Block a user