1use objectstore_service::id::{ObjectContext, ObjectId};
2use objectstore_service::multipart::{
3 AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse,
4 ListPartsResponse, PartNumber, UploadId, UploadPartResponse,
5};
6use objectstore_service::service::{DeleteResponse, GetResponse, InsertResponse, MetadataResponse};
7
8use objectstore_service::{ClientStream, StorageService};
9use objectstore_types::auth::Permission;
10use objectstore_types::metadata::Metadata;
11use objectstore_types::range::ByteRange;
12use objectstore_types::resumable::{SessionToken, UploadProgress};
13use objectstore_types::time::Timestamp;
14
15use crate::auth::AuthContext;
16use crate::endpoints::common::ApiResult;
17
18#[derive(Debug)]
44pub struct AuthAwareService {
45 service: StorageService,
46 context: AuthContext,
47 enforce: bool,
48}
49
50impl AuthAwareService {
51 pub fn new(service: StorageService, context: AuthContext, enforce: bool) -> Self {
59 Self {
60 service,
61 context,
62 enforce,
63 }
64 }
65
66 pub fn check_permission(&self, perm: Permission, context: &ObjectContext) -> ApiResult<()> {
70 if let Err(error) = self.context.assert_authorized(perm, context) {
71 sentry::with_scope(
72 |s| s.set_tag("perm", perm.to_string()),
73 || error.log(!self.enforce),
74 );
75
76 if self.enforce {
77 return Err(error.into());
78 }
79 }
80
81 Ok(())
82 }
83
84 pub async fn insert_object(
86 &self,
87 context: ObjectContext,
88 key: Option<String>,
89 metadata: Metadata,
90 stream: ClientStream,
91 access_time: Timestamp,
92 ) -> ApiResult<InsertResponse> {
93 self.check_permission(Permission::ObjectWrite, &context)?;
94 Ok(self
95 .service
96 .insert_object(context, key, metadata, stream, access_time)
97 .await?)
98 }
99
100 pub async fn get_metadata(
102 &self,
103 id: ObjectId,
104 access_time: Timestamp,
105 ) -> ApiResult<MetadataResponse> {
106 self.check_permission(Permission::ObjectRead, id.context())?;
107 Ok(self.service.get_metadata(id, access_time).await?)
108 }
109
110 pub async fn get_object(
112 &self,
113 id: ObjectId,
114 access_time: Timestamp,
115 range: Option<ByteRange>,
116 ) -> ApiResult<GetResponse> {
117 self.check_permission(Permission::ObjectRead, id.context())?;
118 Ok(self.service.get_object(id, access_time, range).await?)
119 }
120
121 pub async fn delete_object(
123 &self,
124 id: ObjectId,
125 access_time: Timestamp,
126 ) -> ApiResult<DeleteResponse> {
127 self.check_permission(Permission::ObjectDelete, id.context())?;
128 Ok(self.service.delete_object(id, access_time).await?)
129 }
130
131 pub async fn initiate_multipart(
135 &self,
136 id: ObjectId,
137 metadata: Metadata,
138 ) -> ApiResult<InitiateMultipartResponse> {
139 self.check_permission(Permission::ObjectWrite, id.context())?;
140 Ok(self.service.initiate_multipart(id, metadata).await?)
141 }
142
143 pub async fn upload_part(
145 &self,
146 id: ObjectId,
147 upload_id: UploadId,
148 part_number: PartNumber,
149 content_length: u64,
150 content_md5: Option<String>,
151 body: ClientStream,
152 ) -> ApiResult<UploadPartResponse> {
153 self.check_permission(Permission::ObjectWrite, id.context())?;
154 Ok(self
155 .service
156 .upload_part(
157 id,
158 upload_id,
159 part_number,
160 content_length,
161 content_md5,
162 body,
163 )
164 .await?)
165 }
166
167 pub async fn list_parts(
169 &self,
170 id: ObjectId,
171 upload_id: UploadId,
172 max_parts: Option<u32>,
173 part_number_marker: Option<PartNumber>,
174 ) -> ApiResult<ListPartsResponse> {
175 self.check_permission(Permission::ObjectWrite, id.context())?;
176 Ok(self
177 .service
178 .list_parts(id, upload_id, max_parts, part_number_marker)
179 .await?)
180 }
181
182 pub async fn abort_multipart(
184 &self,
185 id: ObjectId,
186 upload_id: UploadId,
187 ) -> ApiResult<AbortMultipartResponse> {
188 self.check_permission(Permission::ObjectWrite, id.context())?;
189 Ok(self.service.abort_multipart(id, upload_id).await?)
190 }
191
192 pub async fn complete_multipart(
194 &self,
195 id: ObjectId,
196 upload_id: UploadId,
197 parts: Vec<CompletedPart>,
198 access_time: Timestamp,
199 ) -> ApiResult<CompleteMultipartResponse> {
200 self.check_permission(Permission::ObjectWrite, id.context())?;
201 Ok(self
202 .service
203 .complete_multipart(id, upload_id, parts, access_time)
204 .await?)
205 }
206
207 pub async fn create_upload_session(
211 &self,
212 id: ObjectId,
213 metadata: Metadata,
214 total_length: u64,
215 ) -> ApiResult<Option<SessionToken>> {
216 self.check_permission(Permission::ObjectWrite, id.context())?;
217 Ok(self
218 .service
219 .create_upload_session(id, metadata, total_length)
220 .await?)
221 }
222
223 pub async fn put_chunk(
225 &self,
226 id: ObjectId,
227 token: SessionToken,
228 offset: u64,
229 content_length: u64,
230 body: ClientStream,
231 ) -> ApiResult<UploadProgress> {
232 self.check_permission(Permission::ObjectWrite, id.context())?;
233 Ok(self
234 .service
235 .put_chunk(id, token, offset, content_length, body)
236 .await?)
237 }
238
239 pub async fn upload_offset(
241 &self,
242 id: ObjectId,
243 token: SessionToken,
244 ) -> ApiResult<UploadProgress> {
245 self.check_permission(Permission::ObjectWrite, id.context())?;
247 Ok(self.service.upload_offset(id, token).await?)
248 }
249
250 pub async fn cancel_upload(&self, id: ObjectId, token: SessionToken) -> ApiResult<()> {
252 self.check_permission(Permission::ObjectWrite, id.context())?;
254 Ok(self.service.cancel_upload(id, token).await?)
255 }
256}