diff --git a/backends/factos_cf_workers/src/factos/factos_cf_workers.gleam b/backends/factos_cf_workers/src/factos/factos_cf_workers.gleam index 85ea1e0..99ccb43 100644 --- a/backends/factos_cf_workers/src/factos/factos_cf_workers.gleam +++ b/backends/factos_cf_workers/src/factos/factos_cf_workers.gleam @@ -31,6 +31,7 @@ pub type Proposed(event) { id: String, event: event, type_: factos.EventType, + version: Int, tags: List(factos.Tag), metadata: factos.Metadata, data: String, @@ -45,6 +46,7 @@ pub type StoredEvent { stream: String, revision: Int, type_: factos.EventType, + version: Int, tags: List(factos.Tag), metadata: factos.Metadata, data: String, @@ -102,6 +104,7 @@ pub fn migrate(database: d1.Database) -> Promise(Result(Nil, Error(_, _))) { stream text not null, revision integer not null, type text not null, + version integer not null, tags text not null, metadata text not null, data text not null, @@ -329,7 +332,7 @@ fn read_matching_events( ) { d1.prepare( database, - "select position, id, stream, revision, type, tags, metadata, data from factos_events order by position", + "select position, id, stream, revision, type, version, tags, metadata, data from factos_events order by position", ) |> d1.raw |> promise.map(fn(result) { @@ -348,7 +351,7 @@ fn read_stream_events( ) { d1.prepare( database, - "select position, id, stream, revision, type, tags, metadata, data from factos_events where stream = ? order by revision", + "select position, id, stream, revision, type, version, tags, metadata, data from factos_events where stream = ? order by revision", ) |> d1.bind([stream_name]) |> d1.raw @@ -395,8 +398,8 @@ fn decode_row( use decoded <- result.try( decode_event(stored) |> result.map_error(DecodeError), ) - let factos.Decoded(event, type_, tags, metadata) = decoded - let StoredEvent(position, id, stream, revision, _, _, _, _) = stored + let factos.Decoded(event, type_, version, tags, metadata) = decoded + let StoredEvent(position, id, stream, revision, _, _, _, _, _) = stored Ok(factos.Recorded( id: id, @@ -404,6 +407,7 @@ fn decode_row( revision: revision, position: factos.SequencePosition(position), type_: type_, + version: version, tags: tags, metadata: metadata, event: event, @@ -418,9 +422,10 @@ fn decode_stored_event( use stream <- result.try(decode_string_field(row, 2)) use revision <- result.try(decode_int_field(row, 3)) use type_name <- result.try(decode_string_field(row, 4)) - use tags <- result.try(decode_string_field(row, 5)) - use metadata <- result.try(decode_string_field(row, 6)) - use data <- result.try(decode_string_field(row, 7)) + use version <- result.try(decode_int_field(row, 5)) + use tags <- result.try(decode_string_field(row, 6)) + use metadata <- result.try(decode_string_field(row, 7)) + use data <- result.try(decode_string_field(row, 8)) Ok(StoredEvent( position: position, @@ -428,6 +433,7 @@ fn decode_stored_event( stream: stream, revision: revision, type_: factos.event_type(type_name), + version: version, tags: tags_from_text(tags), metadata: metadata_from_text(metadata), data: data, @@ -483,7 +489,7 @@ fn append_sql( }) let sql = - "insert into factos_events (id, stream, revision, type, tags, metadata, data) " + "insert into factos_events (id, stream, revision, type, version, tags, metadata, data) " <> string.join(list.map(rows, fn(row) { row.0 }), with: " union all ") <> " returning position, revision" @@ -499,13 +505,14 @@ fn append_select_sql( index: Int, ) -> #(String, List(String)) { let EventCodec(encode, _) = codec - let Proposed(id, _, type_, tags, metadata, data) = encode(event) + let Proposed(id, _, type_, version, tags, metadata, data) = encode(event) let base_values = [ id, stream_name, stream_name, int.to_string(index + 1), factos.event_type_name(type_), + int.to_string(version), tags_to_text(tags), metadata_to_text(metadata), data, @@ -513,7 +520,7 @@ fn append_select_sql( let #(condition, condition_values) = append_condition_sql(stream_name, mode) #( - "select ?, ?, (select coalesce(max(revision), -1) from factos_events where stream = ?) + cast(? as integer), ?, ?, ?, ? where " + "select ?, ?, (select coalesce(max(revision), -1) from factos_events where stream = ?) + cast(? as integer), ?, cast(? as integer), ?, ?, ? where " <> condition, list.append(base_values, condition_values), ) diff --git a/backends/factos_cf_workers/test/factos_cf_workers_test.gleam b/backends/factos_cf_workers/test/factos_cf_workers_test.gleam index 9611b14..32b2301 100644 --- a/backends/factos_cf_workers/test/factos_cf_workers_test.gleam +++ b/backends/factos_cf_workers/test/factos_cf_workers_test.gleam @@ -286,6 +286,7 @@ fn encode(event: Event) -> backend.Proposed(Event) { id: "event-" <> name, event: event, type_: factos.event_type("reserved"), + version: 1, tags: [factos.tag("name:" <> name)], metadata: factos.empty_metadata(), data: name, @@ -301,6 +302,7 @@ fn decode( Ok(factos.Decoded( event: Reserved(stored.data), type_: stored.type_, + version: stored.version, tags: stored.tags, metadata: stored.metadata, )) diff --git a/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam b/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam index 1b19c46..3e636af 100644 --- a/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam +++ b/backends/factos_kurrentdb_erlang/src/factos/factos_kurrentdb_erlang.gleam @@ -33,6 +33,7 @@ pub type Proposed(event) { Proposed( event: event, type_: factos.EventType, + version: Int, tags: List(factos.Tag), metadata: factos.Metadata, message: append_to_stream.Event, @@ -541,13 +542,14 @@ fn decode_recorded( |> result.map_error(DecodeError), ) - let factos.Decoded(event, type_, tags, metadata) = decoded + let factos.Decoded(event, type_, version, tags, metadata) = decoded Ok(factos.Recorded( id: uuid.to_string(recorded_event.id), stream: recorded_event.stream, revision: recorded_event.revision, position: factos.SequencePosition(recorded_event.commit_position), type_: type_, + version: version, tags: tags, metadata: metadata, event: event, diff --git a/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam b/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam index b0de283..26cd640 100644 --- a/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam +++ b/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam @@ -205,6 +205,7 @@ fn encode_counter_event( factos_kurrentdb_erlang.Proposed( event: event, type_: factos.event_type(event_type), + version: 1, tags: [factos.tag("counter:load")], metadata: factos.empty_metadata(), message: append_to_stream.binary_event( @@ -233,6 +234,7 @@ fn decode_counter_event( Ok(factos.Decoded( event: Incremented(value), type_: factos.event_type(event_type), + version: 1, tags: [factos.tag("counter:load")], metadata: factos.empty_metadata(), )) diff --git a/backends/factos_pog/src/factos/factos_pog.gleam b/backends/factos_pog/src/factos/factos_pog.gleam index e3c357a..ead590c 100644 --- a/backends/factos_pog/src/factos/factos_pog.gleam +++ b/backends/factos_pog/src/factos/factos_pog.gleam @@ -31,6 +31,7 @@ pub type Proposed(event) { id: String, event: event, type_: factos.EventType, + version: Int, tags: List(factos.Tag), metadata: factos.Metadata, data: BitArray, @@ -48,6 +49,7 @@ pub type StoredEvent { stream: String, revision: Int, type_: factos.EventType, + version: Int, tags: List(factos.Tag), metadata: factos.Metadata, data: BitArray, @@ -108,6 +110,7 @@ pub fn migrate(connection: pog.Connection) -> Result(Nil, Error(_, _)) { stream text not null, revision integer not null, type text not null, + version integer not null, tags text not null, metadata text not null, data bytea not null, @@ -383,12 +386,12 @@ fn insert_events( [] -> Ok(Append(current_revision: revision - 1, position: position)) [event, ..rest] -> { let EventCodec(encode, _) = codec - let Proposed(id, _, type_, tags, metadata, data) = encode(event) + let Proposed(id, _, type_, version, tags, metadata, data) = encode(event) use returned <- result.try( pog.query( " - insert into factos_events (id, stream, revision, type, tags, metadata, data) - values ($1, $2, $3, $4, $5, $6, $7) + insert into factos_events (id, stream, revision, type, version, tags, metadata, data) + values ($1, $2, $3, $4, $5, $6, $7, $8) returning position ", ) @@ -396,6 +399,7 @@ fn insert_events( |> pog.parameter(pog.text(stream_name)) |> pog.parameter(pog.int(revision)) |> pog.parameter(pog.text(factos.event_type_name(type_))) + |> pog.parameter(pog.int(version)) |> pog.parameter(pog.text(tags_to_text(tags))) |> pog.parameter(pog.text(metadata_to_text(metadata))) |> pog.parameter(pog.bytea(data)) @@ -426,7 +430,7 @@ fn read_matching_events( ) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { use rows <- result.try( pog.query( - "select position, id, stream, revision, type, tags, metadata, data from factos_events order by position", + "select position, id, stream, revision, type, version, tags, metadata, data from factos_events order by position", ) |> pog.returning(stored_event_decoder()) |> pog.execute(on: connection) @@ -444,7 +448,7 @@ fn read_stream_events( ) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { use rows <- result.try( pog.query( - "select position, id, stream, revision, type, tags, metadata, data from factos_events where stream = $1 order by revision", + "select position, id, stream, revision, type, version, tags, metadata, data from factos_events where stream = $1 order by revision", ) |> pog.parameter(pog.text(stream_name)) |> pog.returning(stored_event_decoder()) @@ -475,8 +479,8 @@ fn decode_row( ) -> Result(factos.Recorded(event), Error(domain_error, decode_error)) { let EventCodec(_, decode_event) = codec use decoded <- result.try(decode_event(row) |> result.map_error(DecodeError)) - let factos.Decoded(event, type_, tags, metadata) = decoded - let StoredEvent(position, id, stream, revision, _, _, _, _) = row + let factos.Decoded(event, type_, version, tags, metadata) = decoded + let StoredEvent(position, id, stream, revision, _, _, _, _, _) = row Ok(factos.Recorded( id: id, @@ -484,6 +488,7 @@ fn decode_row( revision: revision, position: factos.SequencePosition(position), type_: type_, + version: version, tags: tags, metadata: metadata, event: event, @@ -496,15 +501,17 @@ fn stored_event_decoder() -> decode.Decoder(StoredEvent) { use stream <- decode.field(2, decode.string) use revision <- decode.field(3, decode.int) use type_name <- decode.field(4, decode.string) - use tags <- decode.field(5, decode.string) - use metadata <- decode.field(6, decode.string) - use data <- decode.field(7, decode.bit_array) + use version <- decode.field(5, decode.int) + use tags <- decode.field(6, decode.string) + use metadata <- decode.field(7, decode.string) + use data <- decode.field(8, decode.bit_array) decode.success(StoredEvent( position: position, id: id, stream: stream, revision: revision, type_: factos.event_type(type_name), + version: version, tags: tags_from_text(tags), metadata: metadata_from_text(metadata), data: data, @@ -575,6 +582,7 @@ fn matches_query_parts( revision: 0, position: factos.NoPosition, type_: type_, + version: 1, tags: tags, metadata: factos.empty_metadata(), event: Nil, diff --git a/backends/factos_pog/test/factos_pog_test.gleam b/backends/factos_pog/test/factos_pog_test.gleam index b07a811..bcf280c 100644 --- a/backends/factos_pog/test/factos_pog_test.gleam +++ b/backends/factos_pog/test/factos_pog_test.gleam @@ -199,6 +199,7 @@ fn encode(event: Event) -> factos_pog.Proposed(Event) { id: "event-" <> event.username, event: event, type_: factos.event_type("UserRegistered"), + version: 1, tags: [factos.tag("username:" <> event.username)], metadata: factos.empty_metadata(), data: bit_array.from_string(event.username), @@ -217,6 +218,7 @@ fn decode( Ok(factos.Decoded( event: UserRegistered(username), type_: stored.type_, + version: stored.version, tags: stored.tags, metadata: stored.metadata, )) @@ -307,6 +309,7 @@ fn encode_counter_event( id: "counter-event-" <> int.to_string(value), event: event, type_: factos.event_type("Incremented"), + version: 1, tags: [factos.tag("counter:load")], metadata: factos.empty_metadata(), data: bit_array.from_string(int.to_string(value)), @@ -330,6 +333,7 @@ fn decode_counter_event( Ok(factos.Decoded( event: Incremented(value), type_: stored.type_, + version: stored.version, tags: stored.tags, metadata: stored.metadata, )) diff --git a/backends/factos_sqlight/src/factos/factos_sqlight.gleam b/backends/factos_sqlight/src/factos/factos_sqlight.gleam index 0ae97fb..b2c5f61 100644 --- a/backends/factos_sqlight/src/factos/factos_sqlight.gleam +++ b/backends/factos_sqlight/src/factos/factos_sqlight.gleam @@ -34,6 +34,7 @@ pub type Proposed(event) { id: String, event: event, type_: factos.EventType, + version: Int, tags: List(factos.Tag), metadata: factos.Metadata, data: BitArray, @@ -52,6 +53,7 @@ pub type StoredEvent { stream: String, revision: Int, type_: factos.EventType, + version: Int, tags: List(factos.Tag), metadata: factos.Metadata, data: BitArray, @@ -108,6 +110,7 @@ pub fn migrate(connection: sqlight.Connection) -> Result(Nil, Error(_, _)) { stream text not null, revision integer not null, type text not null, + version integer not null, tags text not null, metadata text not null default '', data blob not null, @@ -373,12 +376,12 @@ fn insert_events( [] -> Ok(Append(current_revision: revision - 1, position: position)) [event, ..rest] -> { let EventCodec(encode, _) = codec - let Proposed(id, _, type_, tags, metadata, data) = encode(event) + let Proposed(id, _, type_, version, tags, metadata, data) = encode(event) use positions <- result.try( sqlight.query( " - insert into factos_events (id, stream, revision, type, tags, metadata, data) - values (?, ?, ?, ?, ?, ?, ?) + insert into factos_events (id, stream, revision, type, version, tags, metadata, data) + values (?, ?, ?, ?, ?, ?, ?, ?) returning position ", on: connection, @@ -387,6 +390,7 @@ fn insert_events( sqlight.text(stream_name), sqlight.int(revision), sqlight.text(factos.event_type_name(type_)), + sqlight.int(version), sqlight.text(tags_to_text(tags)), sqlight.text(metadata_to_text(metadata)), sqlight.blob(data), @@ -418,7 +422,7 @@ fn read_matching_events( ) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { use rows <- result.try( sqlight.query( - "select position, id, stream, revision, type, tags, metadata, data from factos_events order by position", + "select position, id, stream, revision, type, version, tags, metadata, data from factos_events order by position", on: connection, with: [], expecting: stored_event_decoder(), @@ -436,7 +440,7 @@ fn read_stream_events( ) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { use rows <- result.try( sqlight.query( - "select position, id, stream, revision, type, tags, metadata, data from factos_events where stream = ? order by revision", + "select position, id, stream, revision, type, version, tags, metadata, data from factos_events where stream = ? order by revision", on: connection, with: [sqlight.text(stream_name)], expecting: stored_event_decoder(), @@ -466,8 +470,8 @@ fn decode_row( ) -> Result(factos.Recorded(event), Error(domain_error, decode_error)) { let EventCodec(_, decode_event) = codec use decoded <- result.try(decode_event(row) |> result.map_error(DecodeError)) - let factos.Decoded(event, type_, tags, metadata) = decoded - let StoredEvent(position, id, stream, revision, _, _, _, _) = row + let factos.Decoded(event, type_, version, tags, metadata) = decoded + let StoredEvent(position, id, stream, revision, _, _, _, _, _) = row Ok(factos.Recorded( id: id, @@ -475,6 +479,7 @@ fn decode_row( revision: revision, position: factos.SequencePosition(position), type_: type_, + version: version, tags: tags, metadata: metadata, event: event, @@ -487,15 +492,17 @@ fn stored_event_decoder() -> decode.Decoder(StoredEvent) { use stream <- decode.field(2, decode.string) use revision <- decode.field(3, decode.int) use type_name <- decode.field(4, decode.string) - use tags <- decode.field(5, decode.string) - use metadata <- decode.field(6, decode.string) - use data <- decode.field(7, decode.bit_array) + use version <- decode.field(5, decode.int) + use tags <- decode.field(6, decode.string) + use metadata <- decode.field(7, decode.string) + use data <- decode.field(8, decode.bit_array) decode.success(StoredEvent( position: position, id: id, stream: stream, revision: revision, type_: factos.event_type(type_name), + version: version, tags: tags_from_text(tags), metadata: metadata_from_text(metadata), data: data, @@ -567,6 +574,7 @@ fn matches_query_parts( revision: 0, position: factos.NoPosition, type_: type_, + version: 1, tags: tags, metadata: factos.empty_metadata(), event: Nil, diff --git a/backends/factos_sqlight/test/factos_sqlight_test.gleam b/backends/factos_sqlight/test/factos_sqlight_test.gleam index 3a9034d..2caa161 100644 --- a/backends/factos_sqlight/test/factos_sqlight_test.gleam +++ b/backends/factos_sqlight/test/factos_sqlight_test.gleam @@ -145,6 +145,7 @@ fn encode(event: Event) -> factos_sqlight.Proposed(Event) { id: "event-" <> event.username, event: event, type_: factos.event_type("UserRegistered"), + version: 1, tags: [factos.tag("username:" <> event.username)], metadata: factos.empty_metadata(), data: bit_array.from_string(event.username), @@ -163,6 +164,7 @@ fn decode( Ok(factos.Decoded( event: UserRegistered(username), type_: stored.type_, + version: stored.version, tags: stored.tags, metadata: stored.metadata, )) @@ -284,6 +286,7 @@ fn encode_counter_event( id: "counter-event-" <> int.to_string(value), event: event, type_: factos.event_type("Incremented"), + version: 1, tags: [factos.tag("counter:load")], metadata: factos.empty_metadata(), data: bit_array.from_string(int.to_string(value)), @@ -307,6 +310,7 @@ fn decode_counter_event( Ok(factos.Decoded( event: Incremented(value), type_: stored.type_, + version: stored.version, tags: stored.tags, metadata: stored.metadata, )) diff --git a/examples/orders_sqlight/src/order_workflow.gleam b/examples/orders_sqlight/src/order_workflow.gleam index 35403e9..010776b 100644 --- a/examples/orders_sqlight/src/order_workflow.gleam +++ b/examples/orders_sqlight/src/order_workflow.gleam @@ -701,6 +701,7 @@ fn proposed( id: "example-" <> type_name <> "-" <> fields_to_payload(fields), event: event, type_: factos.event_type(type_name), + version: 1, tags: [factos.tag("restaurant"), ..tags], metadata: factos.empty_metadata(), data: bit_array.from_string(fields_to_payload(fields)), @@ -721,6 +722,7 @@ fn decode_event( Ok(factos.Decoded( event: event, type_: stored.type_, + version: stored.version, tags: stored.tags, metadata: stored.metadata, )) diff --git a/examples/tickets_pog/src/ticket_sale.gleam b/examples/tickets_pog/src/ticket_sale.gleam index 1b139ad..87cf8ea 100644 --- a/examples/tickets_pog/src/ticket_sale.gleam +++ b/examples/tickets_pog/src/ticket_sale.gleam @@ -282,6 +282,7 @@ fn encode_event(event: Event) -> factos_pog.Proposed(Event) { id: "ticket-sold-" <> buyer, event: event, type_: factos.event_type("TicketSold"), + version: 1, tags: [factos.tag("event:" <> event_id)], metadata: factos.empty_metadata(), data: bit_array.from_string(buyer), @@ -301,6 +302,7 @@ fn decode_event( Ok(factos.Decoded( event: TicketSold(buyer), type_: stored.type_, + version: stored.version, tags: stored.tags, metadata: stored.metadata, )) diff --git a/src/factos.gleam b/src/factos.gleam index e377ab2..121663a 100644 --- a/src/factos.gleam +++ b/src/factos.gleam @@ -137,7 +137,13 @@ pub type Decoded(event) { /// Backends use codecs supplied by the application. The decoded value includes /// the domain event plus the event type and tags that should participate in /// query matching, along with non-query event metadata. - Decoded(event: event, type_: EventType, tags: List(Tag), metadata: Metadata) + Decoded( + event: event, + type_: EventType, + version: Int, + tags: List(Tag), + metadata: Metadata, + ) } pub type Recorded(event) { @@ -152,6 +158,7 @@ pub type Recorded(event) { revision: Int, position: SequencePosition, type_: EventType, + version: Int, tags: List(Tag), metadata: Metadata, event: event, diff --git a/test/factos_test.gleam b/test/factos_test.gleam index 196ef57..f8cfcbe 100644 --- a/test/factos_test.gleam +++ b/test/factos_test.gleam @@ -238,6 +238,7 @@ fn recorded( revision: revision, position: factos.SequencePosition(revision), type_: type_, + version: 1, tags: tags, metadata: factos.empty_metadata(), event: event,